Skip to content
markpaper

src/history/paginate.ts

v0.3.0 · 4.6 KB

Download file
// Forward time pagination shared by fills and funding feeds.

/** Result of a forward time pagination. */
export interface ForwardPageResult<T> {
  /** Unique items inside [startTime, endTime], in the order received (not sorted). */
  items: T[];
  /** Number of requests made. */
  pages: number;
  /** true when the feed was exhausted for the window; false when stopped by a cap or a stall. */
  complete: boolean;
  /** true when `maxItems` stopped the loop. */
  aborted: boolean;
  /**
   * Milliseconds in which a full page consisted of a single timestamp. Items of
   * that millisecond beyond the page cap cannot be reached through time
   * pagination and may be missing.
   */
  denseMillis: number[];
}

export interface ForwardPaginationParams<T> {
  startTime: number;
  endTime?: number;
  /** Performs one request starting at `cursor` (inclusive). */
  fetchPage: (cursor: number) => Promise<unknown>;
  timeOf: (item: T) => number;
  keyOf: (item: T) => string;
  /** Known server cap per response. When unknown, the loop stops on a page with nothing new. */
  pageLimit?: number;
  maxPages?: number;
  maxItems?: number;
  pageDelayMs?: number;
  signal?: AbortSignal;
  /** Used in error messages. */
  label: string;
}

/** Sleeps `ms`, rejecting early when the signal aborts. */
export function sleep(ms: number, signal?: AbortSignal): Promise<void> {
  if (ms <= 0) return Promise.resolve();
  return new Promise((resolve, reject) => {
    const onAbort = (): void => {
      clearTimeout(timer);
      reject(signal?.reason ?? new Error('aborted'));
    };
    const timer = setTimeout(() => {
      signal?.removeEventListener('abort', onAbort);
      resolve();
    }, ms);
    signal?.addEventListener('abort', onAbort, { once: true });
  });
}

/**
 * Pages a chronological feed forward in time.
 *
 * Cursor rule: after a page that spans several milliseconds the next request
 * starts AT the newest millisecond (not newest + 1) and duplicates are removed
 * by key. The common `newest + 1` cursor silently drops the tail of a
 * millisecond that the page cap cut in half; re-reading that millisecond costs
 * no extra request. Only when a whole capped page is one millisecond does the
 * cursor move to `ms + 1`, and that millisecond is reported in `denseMillis`.
 */
export async function paginateForward<T>(p: ForwardPaginationParams<T>): Promise<ForwardPageResult<T>> {
  const seen = new Set<string>();
  const items: T[] = [];
  const denseMillis: number[] = [];
  const endTime = p.endTime;
  let cursor = p.startTime;
  let pages = 0;
  let complete = false;
  let aborted = false;

  for (;;) {
    if (endTime !== undefined && cursor > endTime) {
      complete = true;
      break;
    }
    if (p.maxItems !== undefined && items.length >= p.maxItems) {
      aborted = true;
      break;
    }
    if (p.maxPages !== undefined && pages >= p.maxPages) break;
    p.signal?.throwIfAborted();
    if (pages > 0) await sleep(p.pageDelayMs ?? 0, p.signal);

    const raw = await p.fetchPage(cursor);
    pages += 1;
    // Never turn a bad response into []: an empty array means "no data", and
    // that silently hid real failures before.
    if (!Array.isArray(raw)) {
      throw new TypeError(`history: ${p.label} returned a non-array response`);
    }
    const page = raw as T[];
    if (page.length === 0) {
      complete = true;
      break;
    }

    let fresh = 0;
    let minT = Number.POSITIVE_INFINITY;
    let maxT = Number.NEGATIVE_INFINITY;
    for (const item of page) {
      const t = p.timeOf(item);
      if (!Number.isFinite(t)) continue;
      if (t < minT) minT = t;
      if (t > maxT) maxT = t;
      const key = p.keyOf(item);
      if (seen.has(key)) continue;
      seen.add(key);
      if (t < p.startTime || (endTime !== undefined && t > endTime)) continue;
      items.push(item);
      fresh += 1;
    }
    if (!Number.isFinite(maxT)) {
      throw new TypeError(`history: ${p.label} page has no item with a numeric time`);
    }

    const knownCap = p.pageLimit !== undefined;
    if (knownCap && page.length < (p.pageLimit as number)) {
      complete = true;
      break;
    }
    if (maxT < cursor) {
      // The server ignored the cursor; stop instead of looping forever.
      break;
    }
    if (minT === maxT) {
      if (knownCap) denseMillis.push(maxT);
      else if (fresh === 0) {
        complete = true;
        break;
      }
      cursor = maxT + 1;
      continue;
    }
    if (fresh === 0) {
      // Several milliseconds yet nothing new: the feed does not honour the cursor.
      break;
    }
    cursor = maxT;
  }

  return { items, pages, complete, aborted, denseMillis };
}
All files