Skip to content
markpaper

src/pending/index.ts

v0.1.0 · 10.3 KB

Download file
import {
  canonicalStatus,
  isFinalOrderEvent,
  type QfexOrderEvent,
  remainingIsZero,
  statusClass,
} from '../frames/index.js';
import { decCmp, decSign, decSub, decToString, parseDec } from '../numbers/index.js';
export interface UnknownIntent {
  cloid: string;
  symbol: string;
  side: 'BUY' | 'SELL';
  kind: 'ioc' | 'gtc' | 'close';
  reduceOnly?: boolean;
  placementKey?: string;
  quantity: string;
  sentAt: number;
  orderId?: string;
}
export interface HistoricOutcome {
  terminalStatus: string;
  filledQty: string;
  avgPrice: string | null;
  orderId: string | null;
  symbol: string;
  side: 'BUY' | 'SELL';
}
export interface UnknownRelease {
  intent: UnknownIntent;
  outcome: 'alive' | 'terminal' | 'expired';
  filledQty?: string;
  how: string;
}
export const placementKey = (symbol: string, side: 'BUY' | 'SELL', price: string): string => {
  const value = parseDec(price);
  if (!value || decSign(value) <= 0) throw new TypeError('Invalid placement price');
  return symbol + '|' + side + '|' + decToString(value);
};
export function cancelTruth(
  event: Pick<QfexOrderEvent, 'status' | 'qty' | 'remaining'>,
  options: { fillsSeen: boolean },
):
  | { verdict: 'cancelled' }
  | { verdict: 'not_cancelled'; gone: boolean; why: string }
  | { verdict: 'not_final'; why: string } {
  const status = canonicalStatus(event.status),
    kind = statusClass(status);
  if (kind === 'fill')
    return remainingIsZero(event)
      ? { verdict: 'not_cancelled', gone: true, why: 'Order filled' }
      : { verdict: 'not_final', why: 'Partial fill is not a cancellation response' };
  if (kind === 'cancel_reply')
    return { verdict: 'not_cancelled', gone: false, why: 'Order absence is not proven by ' + status };
  if (kind === 'reject') return { verdict: 'not_cancelled', gone: false, why: status };
  if (kind === 'terminal') {
    const quantity = parseDec(event.qty),
      remaining = parseDec(event.remaining);
    if (status !== 'CANCELLED')
      return { verdict: 'not_cancelled', gone: true, why: 'Terminal order was not cancelled by this request' };
    if (!quantity || !remaining || decCmp(remaining, quantity) < 0 || options.fillsSeen)
      return { verdict: 'not_cancelled', gone: true, why: 'Cancellation does not exclude a prior fill' };
    return { verdict: 'cancelled' };
  }
  return { verdict: 'not_final', why: status };
}
export function createWriteMemory(options: {
  echoTtlMs: number;
  readAfterWriteMs: number;
  unknownLockMs: number;
  unknownLockGtcMs: number;
  failedHistoryLockMs: number;
  canExpire: (intent: UnknownIntent) => boolean | Promise<boolean>;
  now?: () => number;
  knownFilled?: (intent: UnknownIntent) => string | null;
  onTerminalFill?: (intent: UnknownIntent, filledQty: string) => void;
}) {
  for (const key of [
    'echoTtlMs',
    'readAfterWriteMs',
    'unknownLockMs',
    'unknownLockGtcMs',
    'failedHistoryLockMs',
  ] as const)
    if (!Number.isFinite(options?.[key]) || options[key] <= 0)
      throw new TypeError('Required write memory policy: ' + key);
  if (typeof options.canExpire !== 'function')
    throw new TypeError('Unknown expiry requires explicit independent evidence');
  const now = options.now ?? Date.now;
  const placed = new Map<string, { row: QfexOrderEvent; at: number; key: string }>();
  const gone = new Map<string, number>();
  const keys = new Map<string, { at: number; orderId?: string }>();
  const unknown = new Map<string, UnknownIntent>();
  let resolving: Promise<UnknownRelease[]> | null = null;
  function sweep(at: number) {
    for (const [id, entry] of placed) if (at - entry.at > options.echoTtlMs) placed.delete(id);
    for (const [id, since] of gone) if (at - since > options.echoTtlMs) gone.delete(id);
    for (const [key, entry] of keys) if (at - entry.at > options.echoTtlMs) keys.delete(key);
  }
  function forget(orderId: string) {
    placed.delete(orderId);
    for (const [key, value] of keys) if (value.orderId === orderId) keys.delete(key);
  }
  function notePlaced(row: QfexOrderEvent, at = now()) {
    if (!row.orderId || !row.cloid || statusClass(row.status) !== 'live')
      throw new TypeError('Placement echo requires a live acknowledgment');
    const key = placementKey(row.symbol, row.side, row.price);
    placed.set(row.orderId, { row: { ...row }, at, key });
    keys.set(key, { at, orderId: row.orderId });
  }
  function noteTerminal(orderId: string, at = now()) {
    forget(orderId);
    gone.set(orderId, Math.max(gone.get(orderId) ?? 0, at));
  }
  function applyEcho(rows: readonly QfexOrderEvent[], readStartedAt: number, at = now()): QfexOrderEvent[] {
    sweep(at);
    const present = new Set(rows.map((row) => row.orderId));
    for (const [id, since] of gone) if (!present.has(id) && readStartedAt >= since) gone.delete(id);
    const result = rows.filter((row) => !gone.has(row.orderId)).map((row) => ({ ...row }));
    for (const [id, entry] of placed) {
      if (present.has(id)) {
        forget(id);
        continue;
      }
      if (readStartedAt >= entry.at + options.readAfterWriteMs) {
        forget(id);
        continue;
      }
      result.push({ ...entry.row });
    }
    return result;
  }
  function release(intent: UnknownIntent, outcome: UnknownRelease['outcome'], filledQty?: string): UnknownRelease {
    unknown.delete(intent.cloid);
    if (intent.orderId && outcome === 'terminal') noteTerminal(intent.orderId);
    if (filledQty && decSign(parseDec(filledQty)!) > 0) options.onTerminalFill?.({ ...intent }, filledQty);
    return { intent: { ...intent }, outcome, how: outcome, ...(filledQty === undefined ? {} : { filledQty }) };
  }
  async function resolveOnce(
    io: {
      openOrders: () => Promise<{ orders: QfexOrderEvent[]; complete: boolean }>;
      historicByCloid: (cloid: string) => Promise<HistoricOutcome | null>;
    },
    at: number,
  ) {
    const releases: UnknownRelease[] = [];
    let open: { orders: QfexOrderEvent[]; complete: boolean } | null = null;
    try {
      open = await io.openOrders();
    } catch {}
    for (const intent of [...unknown.values()]) {
      const hits = (open?.orders ?? []).filter((row) => row.cloid === intent.cloid);
      if (hits.length > 1) continue;
      const hit = hits[0];
      if (hit) {
        if (
          hit.symbol !== intent.symbol ||
          hit.side !== intent.side ||
          (intent.orderId && hit.orderId.toLowerCase() !== intent.orderId.toLowerCase())
        )
          continue;
        if (intent.kind === 'gtc' && statusClass(hit.status) === 'live') releases.push(release(intent, 'alive'));
        continue;
      }
      let historic: HistoricOutcome | null | undefined;
      try {
        historic = await io.historicByCloid(intent.cloid);
      } catch {
        historic = undefined;
      }
      if (historic) {
        const fill = parseDec(historic.filledQty),
          sent = parseDec(intent.quantity);
        if (
          !fill ||
          !sent ||
          decSign(fill) < 0 ||
          decCmp(fill, sent) > 0 ||
          historic.symbol !== intent.symbol ||
          historic.side !== intent.side ||
          (intent.orderId && historic.orderId?.toLowerCase() !== intent.orderId.toLowerCase()) ||
          !isFinalOrderEvent({ status: historic.terminalStatus, remaining: '0' })
        )
          continue;
        releases.push(release(intent, 'terminal', decToString(fill)));
        continue;
      }
      const lock = intent.kind === 'gtc' ? options.unknownLockGtcMs : options.unknownLockMs;
      const hold = historic === undefined ? options.failedHistoryLockMs : lock;
      if (open?.complete && at >= intent.sentAt + hold && (await options.canExpire({ ...intent })))
        releases.push(release(intent, 'expired'));
    }
    return releases;
  }
  function resolveUnknowns(io: Parameters<typeof resolveOnce>[0], at = now()) {
    if (resolving) return resolving;
    resolving = resolveOnce(io, at).finally(() => {
      resolving = null;
    });
    return resolving;
  }
  function resolveUnknownByEvent(event: QfexOrderEvent): UnknownRelease | null {
    const intent = event.cloid ? unknown.get(event.cloid) : undefined;
    if (
      !intent ||
      event.symbol !== intent.symbol ||
      event.side !== intent.side ||
      (intent.orderId && event.orderId.toLowerCase() !== intent.orderId.toLowerCase())
    )
      return null;
    intent.orderId = event.orderId;
    const kind = statusClass(event.status);
    if (kind === 'live' && intent.kind === 'gtc') return release(intent, 'alive');
    if (!isFinalOrderEvent(event) || ((intent.kind === 'ioc' || intent.kind === 'close') && kind === 'terminal'))
      return null;
    if (kind === 'reject') {
      const known = options.knownFilled?.(intent);
      if (known && decSign(parseDec(known) ?? { mant: 1n, exp: 0 }) > 0) return null;
      return release(intent, 'terminal', '0');
    }
    const sent = parseDec(intent.quantity),
      remaining = parseDec(event.remaining);
    if (!sent || !remaining || decCmp(remaining, sent) > 0) return null;
    return release(intent, 'terminal', decToString(decSub(sent, remaining)));
  }
  return {
    notePlaced,
    noteCancelled: noteTerminal,
    noteTerminal,
    applyEcho,
    placementKey,
    noteNoSuchOrder(orderId: string, at = now()) {
      const entry = placed.get(orderId);
      if (entry && at - entry.at < options.readAfterWriteMs) return false;
      forget(orderId);
      return true;
    },
    notePlacedKey(key: string, at = now()) {
      keys.set(key, { at });
    },
    isPlacedKey(key: string, at = now()) {
      sweep(at);
      return keys.has(key);
    },
    forgetKey: (key: string) => keys.delete(key),
    noteUnknown(intent: UnknownIntent) {
      if (unknown.has(intent.cloid)) throw new TypeError('Unknown intent already registered');
      unknown.set(intent.cloid, { ...intent });
    },
    iocLock: (symbol: string) =>
      [...unknown.values()].find((intent) => intent.symbol === symbol && intent.kind !== 'gtc'),
    reduceLock: (symbol: string, side: 'BUY' | 'SELL') =>
      [...unknown.values()].find(
        (intent) => intent.symbol === symbol && intent.side === side && (intent.kind === 'close' || intent.reduceOnly),
      ),
    placementLock: (key: string) =>
      [...unknown.values()].find((intent) => intent.kind === 'gtc' && intent.placementKey === key),
    unknownIntents: () => [...unknown.values()].map((intent) => ({ ...intent })),
    resolveUnknowns,
    resolveUnknownByEvent,
  };
}
export type WriteMemory = ReturnType<typeof createWriteMemory>;
All files