Skip to content
markpaper

src/ws/dedupe.ts

v0.3.0 · 3.1 KB

Download file
// Snapshot handling and id-based dedupe for fill/trade streams.

import type { WsFill, WsUserFills } from './types.js';

export interface TidDeduper {
  /** Returns `true` when `id` was already seen; otherwise records it and returns `false`. */
  seen(id: number | string): boolean;
  has(id: number | string): boolean;
  readonly size: number;
}

/**
 * Bounded "seen ids" set (oldest ids evicted first). Frames repeat: `userFills` sends its history
 * again after every (re)subscribe and `trades` can repeat after a reconnect - dedupe by `tid`.
 */
export function createTidDeduper(maxSize = 10_000): TidDeduper {
  if (!(maxSize >= 1)) throw new RangeError('maxSize must be >= 1');
  const ids = new Set<string>();
  return {
    seen(id) {
      const k = String(id);
      if (ids.has(k)) return true;
      ids.add(k);
      if (ids.size > maxSize) {
        const oldest = ids.values().next();
        if (!oldest.done) ids.delete(oldest.value);
      }
      return false;
    },
    has(id) {
      return ids.has(String(id));
    },
    get size() {
      return ids.size;
    },
  };
}

/** Returns true for a history snapshot frame (`isSnapshot: true`). */
export function isSnapshotFrame(data: unknown): boolean {
  return !!data && typeof data === 'object' && (data as { isSnapshot?: unknown }).isSnapshot === true;
}

export interface UserFillsSplit {
  isSnapshot: boolean;
  /** Fills to apply to a ledger: not seen before and, for snapshots, not older than the session start. */
  fresh: WsFill[];
  /** Fills skipped as duplicates or pre-session history. */
  skipped: WsFill[];
  /**
   * `true` for snapshots: positions and open orders must be re-read from REST (`clearinghouseState`,
   * `openOrders`) instead of being derived from the snapshot.
   */
  reconcileFromRest: boolean;
}

/**
 * Splits a `userFills` frame into fills to account for and fills to skip.
 *
 * The first frame after every (re)subscribe is `isSnapshot: true` HISTORY, including fills that
 * happened while the socket was down. Treating it as new fills doubles positions and PnL after a
 * reconnect. Rule: from a snapshot take only fills with `time >= sessionStartTime` not
 * yet seen by `tid` into the PnL ledger, never change positions from it; reconcile from REST.
 *
 * Pass `user` to drop frames of another address (compared in lowercase).
 */
export function splitUserFills(
  data: WsUserFills,
  opts: { sessionStartTime: number; deduper: TidDeduper; user?: string },
): UserFillsSplit {
  const isSnapshot = isSnapshotFrame(data);
  const fresh: WsFill[] = [];
  const skipped: WsFill[] = [];
  const fills = Array.isArray(data?.fills) ? data.fills : [];
  if (opts.user !== undefined && String(data?.user ?? '').toLowerCase() !== opts.user.toLowerCase()) {
    return { isSnapshot, fresh, skipped: [...fills], reconcileFromRest: false };
  }
  for (const fill of fills) {
    if (isSnapshot && !(fill.time >= opts.sessionStartTime)) { skipped.push(fill); continue; }
    if (opts.deduper.seen(fill.tid)) { skipped.push(fill); continue; }
    fresh.push(fill);
  }
  return { isSnapshot, fresh, skipped, reconcileFromRest: isSnapshot };
}
All files