Skip to content
markpaper

src/mds/stream.ts

v0.1.0 · 17.2 KB

Download file
import WebSocket from 'ws';
import { type Dec, decCmp, decSign, decToNumber, parseDec, parseLossless } from '../numbers/index.js';
export interface QfexMdsClock {
  now(): number;
  setTimer(
    ms: number,
    fn: () => void,
  ): {
    cancel(): void;
  };
}
export interface QfexMdsOptions {
  url?: string;
  clock?: QfexMdsClock;
  random?: () => number;
  pingIntervalMs?: number;
  deadAfterMs?: number;
  backoffMinMs?: number;
  backoffMaxMs?: number;
  stableAfterMs?: number;
  handshakeTimeoutMs?: number;
  symbolsPerFrame?: number;
  logWindowMs: number;
  closeWaitMs: number;
  maxPayloadBytes: number;
  bandAgeFrom?: 'frame' | 'stream';
}
export type QfexMdsState = 'idle' | 'connecting' | 'open' | 'backoff' | 'closed';
export interface QfexMdsHealth {
  state: QfexMdsState;
  since: number;
  epoch: number;
  reconnects: number;
  lastFrameAt: number;
  watched: number;
  subscribed: number;
  bands: number;
  bbos: number;
  errors: number;
  unknownFrames: number;
  lastError: string | null;
}
interface Resolved {
  url: string;
  clock: QfexMdsClock;
  random: () => number;
  pingIntervalMs: number;
  deadAfterMs: number;
  backoffMinMs: number;
  backoffMaxMs: number;
  stableAfterMs: number;
  handshakeTimeoutMs: number;
  symbolsPerFrame: number;
  bandAgeFrom: 'frame' | 'stream';
}
interface Timer {
  c: {
    cancel(): void;
  } | null;
  dead: boolean;
}
interface MdsConn {
  epoch: number;
  ws: WebSocket;
  openedAt: number | null;
  lastFrameAt: number;
  lastPongAt: number;
  requested: Set<string>;
  confirmed: Set<string>;
  pingTimer: Timer | null;
  dropped: boolean;
  closedP: Promise<void>;
  closedResolve: () => void;
}
interface BandRec {
  min: number;
  max: number;
  minDec: Dec;
  maxDec: Dec;
  at: number;
  epoch: number;
}
interface BboRec {
  bid: number;
  ask: number;
  bidSz: number | null;
  askSz: number | null;
  at: number;
  epoch: number;
}
type Obj = Record<string, unknown>;
export interface MarketDataOptions extends QfexMdsOptions {
  webSocketFactory?: (url: string, opts: WebSocket.ClientOptions) => WebSocket;
  onLog?: (message: string) => void;
}
/** @experimental Stream-based band freshness is not confirmed; frame age is the conservative mode. */
export function createMarketDataStreamInternal(options: MarketDataOptions) {
  if (!options?.url) throw new TypeError('Required MDS URL');
  for (const name of [
    'pingIntervalMs',
    'deadAfterMs',
    'backoffMinMs',
    'backoffMaxMs',
    'stableAfterMs',
    'handshakeTimeoutMs',
    'symbolsPerFrame',
    'logWindowMs',
    'closeWaitMs',
    'maxPayloadBytes',
  ] as const)
    if (!Number.isFinite(options[name]) || !(options[name]! > 0)) throw new TypeError('Required MDS policy: ' + name);
  const webSocketFactory = options.webSocketFactory ?? ((url, opts) => new WebSocket(url, opts));
  const log = {
    info: (message: string) => {
      try {
        options.onLog?.(message);
      } catch {}
    },
    warn: (message: string) => {
      try {
        options.onLog?.(message);
      } catch {}
    },
    error: (message: string) => {
      try {
        options.onLog?.(message);
      } catch {}
    },
  };
  const opts = (): Resolved => ({
    url: options.url!,
    clock: options.clock ?? realClock,
    random: options.random ?? Math.random,
    pingIntervalMs: options.pingIntervalMs!,
    deadAfterMs: options.deadAfterMs!,
    backoffMinMs: options.backoffMinMs!,
    backoffMaxMs: options.backoffMaxMs!,
    stableAfterMs: options.stableAfterMs!,
    handshakeTimeoutMs: options.handshakeTimeoutMs!,
    symbolsPerFrame: options.symbolsPerFrame!,
    bandAgeFrom: options.bandAgeFrom ?? 'frame',
  });
  const realClock: QfexMdsClock = {
    now: () => Date.now(),
    setTimer: (ms, fn) => {
      const t = setTimeout(fn, Math.max(0, ms));
      t.unref();
      return { cancel: () => clearTimeout(t) };
    },
  };
  const SYMBOL_RE = /^[A-Za-z0-9][A-Za-z0-9._]{0,31}-[A-Za-z0-9]{2,8}$/;
  const watch = new Set<string>();
  const bands = new Map<string, BandRec>();
  const bbos = new Map<string, BboRec>();
  const symFrameAt = new Map<
    string,
    {
      at: number;
      epoch: number;
    }
  >();
  const timers = new Set<Timer>();
  const logGate = new Map<string, number>();
  let conn: MdsConn | null = null;
  let state: QfexMdsState = 'idle';
  let stateSince = 0;
  let epoch = 0;
  let reconnects = 0;
  let attempts = 0;
  let reconnectTimer: Timer | null = null;
  let closed = false;
  const wired = false;
  let lastFrameAt = 0;
  let errors = 0;
  let unknownFrames = 0;
  let lastError: string | null = null;
  function now(): number {
    return opts().clock.now();
  }
  function setState(s: QfexMdsState): void {
    if (state === s) return;
    state = s;
    stateSince = now();
  }
  function timer(ms: number, fn: () => void): Timer {
    const h: Timer = { c: null, dead: false };
    timers.add(h);
    h.c = opts().clock.setTimer(ms, () => {
      if (h.dead) return;
      h.dead = true;
      timers.delete(h);
      try {
        fn();
      } catch (e) {
        log.error(`MDS timer callback failed: ${errText(e)}`);
      }
    });
    return h;
  }
  function cancelTimer(h: Timer | null): void {
    if (!h || h.dead) return;
    h.dead = true;
    timers.delete(h);
    try {
      h.c?.cancel();
    } catch {}
  }
  function errText(e: unknown): string {
    return String((e as Error)?.message ?? e).slice(0, 300);
  }
  function note(key: string, level: 'info' | 'warn' | 'error', msg: string): void {
    const t = now();
    const last = logGate.get(key);
    if (last !== undefined && t - last < options.logWindowMs) return;
    logGate.set(key, t);
    if (logGate.size > 200) for (const [k, at] of logGate) if (t - at >= options.logWindowMs) logGate.delete(k);
    log[level](msg);
  }
  function ensureQfexMds(): void {
    if (closed || conn || reconnectTimer || watch.size === 0) return;
    connect();
  }
  function connect(): void {
    if (closed || conn) return;
    const o = opts();
    const ep = ++epoch;
    setState('connecting');
    let ws: WebSocket;
    try {
      ws = webSocketFactory(o.url, {
        handshakeTimeout: o.handshakeTimeoutMs,
        maxPayload: options.maxPayloadBytes,
        perMessageDeflate: false,
        followRedirects: false,
      });
    } catch (e) {
      lastError = errText(e);
      log.error(`MDS connection failed: ${lastError}`);
      scheduleReconnect(backoffDelay());
      return;
    }
    let closedResolve: () => void = () => undefined;
    const closedP = new Promise<void>((r) => {
      closedResolve = r;
    });
    const c: MdsConn = {
      epoch: ep,
      ws,
      openedAt: null,
      lastFrameAt: now(),
      lastPongAt: 0,
      requested: new Set(),
      confirmed: new Set(),
      pingTimer: null,
      dropped: false,
      closedP,
      closedResolve,
    };
    conn = c;
    ws.on('error', (err) => {
      if (c.dropped) return;
      errors++;
      lastError = errText(err);
      note(`err:${c.epoch}`, 'warn', `[qfex] MDS (epoch ${c.epoch}): ${lastError}`);
      drop(c, `MDS socket error: ${lastError}`);
    });
    ws.on('unexpected-response', (req, res) => {
      const status = res.statusCode ?? 0;
      req.on('error', () => undefined);
      res.on('error', () => undefined);
      res.resume();
      lastError = `MDS handshake returned HTTP ${status}`;
      log.error(`[qfex] MDS: ${lastError}`);
      try {
        ws.terminate();
      } catch {}
      drop(c, lastError);
    });
    ws.on('open', () => onOpen(c));
    ws.on('message', (data, isBinary) => {
      try {
        onMessage(c, data, isBinary);
      } catch (e) {
        note('handler', 'error', `MDS frame callback failed: ${errText(e)}`);
      }
    });
    ws.on('pong', () => {
      c.lastPongAt = now();
    });
    ws.on('close', (code) => {
      c.closedResolve();
      drop(c, `MDS socket closed with code ${code}`);
    });
  }
  function onOpen(c: MdsConn): void {
    if (c.dropped || closed) {
      try {
        c.ws.terminate();
      } catch {}
      return;
    }
    const t = now();
    c.openedAt = t;
    c.lastFrameAt = t;
    setState('open');
    log.info(`MDS ready at epoch ${c.epoch}; watching ${watch.size} symbols`);
    subscribe(c, [...watch]);
    startLiveness(c);
  }
  function startLiveness(c: MdsConn): void {
    const o = opts();
    const tick = (): void => {
      if (c.dropped) return;
      const t = now();
      const lastSign = Math.max(c.lastFrameAt, c.lastPongAt, c.openedAt ?? 0);
      if (t - lastSign >= o.deadAfterMs) {
        drop(c, `MDS heartbeat missing for ${Math.round((t - lastSign) / 1000)} seconds`);
        return;
      }
      try {
        if (c.ws.readyState === WebSocket.OPEN) c.ws.ping();
      } catch (e) {
        note('ping', 'warn', `MDS ping failed: ${errText(e)}`);
      }
      c.pingTimer = timer(o.pingIntervalMs, tick);
    };
    c.pingTimer = timer(o.pingIntervalMs, tick);
  }
  function subscribe(c: MdsConn, symbols: string[]): void {
    if (c.dropped || c.ws.readyState !== WebSocket.OPEN) return;
    const fresh = symbols.filter((s) => SYMBOL_RE.test(s) && !c.requested.has(s)).sort();
    const per = opts().symbolsPerFrame;
    for (let i = 0; i < fresh.length; i += per) {
      const chunk = fresh.slice(i, i + per);
      const text = JSON.stringify({ type: 'subscribe', channels: ['bbo', 'minmax_price'], symbols: chunk });
      try {
        c.ws.send(text, (err) => {
          if (err) drop(c, `MDS subscription send failed: ${errText(err)}`);
        });
      } catch (e) {
        drop(c, `MDS subscription send failed: ${errText(e)}`);
        return;
      }
      for (const s of chunk) c.requested.add(s);
    }
  }
  function backoffDelay(): number {
    const o = opts();
    const base = Math.min(o.backoffMaxMs, o.backoffMinMs * 2 ** Math.min(attempts, 16));
    attempts++;
    const r = Math.min(1, Math.max(0, o.random()));
    return Math.max(o.backoffMinMs, Math.round(base * (0.5 + 0.5 * r)));
  }
  function scheduleReconnect(delay: number): void {
    if (closed) return;
    cancelTimer(reconnectTimer);
    setState('backoff');
    reconnectTimer = timer(delay, () => {
      reconnectTimer = null;
      reconnects++;
      if (watch.size > 0) connect();
      else setState('idle');
    });
  }
  function drop(c: MdsConn, reason: string): void {
    if (c.dropped) return;
    c.dropped = true;
    cancelTimer(c.pingTimer);
    c.pingTimer = null;
    if (conn === c) conn = null;
    try {
      c.ws.terminate();
    } catch {}
    if (c.openedAt !== null) log.warn(`MDS epoch ${c.epoch} dropped: ${reason}`);
    if (closed) {
      setState('closed');
      return;
    }
    const o = opts();
    if (c.openedAt !== null && now() - c.openedAt >= o.stableAfterMs) attempts = 0;
    scheduleReconnect(backoffDelay());
  }
  const isObj = (v: unknown): v is Obj => typeof v === 'object' && v !== null && !Array.isArray(v);
  function topLevel(v: unknown): {
    px: Dec;
    sz: Dec | null;
  } | null {
    if (!Array.isArray(v) || v.length === 0) return null;
    const row = v[0];
    if (!Array.isArray(row) || row.length < 1) return null;
    const px = typeof row[0] === 'string' ? parseDec(row[0]) : null;
    const sz = row.length > 1 && typeof row[1] === 'string' ? parseDec(row[1]) : null;
    if (!px || decSign(px) <= 0 || !Number.isFinite(decToNumber(px))) return null;
    if (row.length > 1 && (!sz || decSign(sz) <= 0 || !Number.isFinite(decToNumber(sz)))) return null;
    return { px, sz };
  }
  function onMessage(c: MdsConn, data: WebSocket.RawData, isBinary: boolean): void {
    if (c.dropped) return;
    const t = now();
    c.lastFrameAt = t;
    lastFrameAt = t;
    if (isBinary) {
      note('binary', 'warn', 'Binary MDS frame ignored');
      return;
    }
    const text = Buffer.isBuffer(data)
      ? data.toString('utf8')
      : Array.isArray(data)
        ? Buffer.concat(data).toString('utf8')
        : Buffer.from(data).toString('utf8');
    let v: unknown;
    try {
      v = parseLossless(text);
    } catch {
      unknownFrames++;
      note('not_json', 'warn', 'Invalid JSON in MDS frame');
      return;
    }
    if (!isObj(v)) {
      unknownFrames++;
      return;
    }
    if (typeof v.error_code === 'string') {
      errors++;
      lastError = `${v.error_code}${typeof v.message === 'string' ? `: ${v.message.slice(0, 160)}` : ''}`;
      note(`mds_err:${v.error_code}`, 'warn', `MDS exchange error: ${lastError}`);
      return;
    }
    const type = typeof v.type === 'string' ? v.type : null;
    const sym = typeof v.symbol === 'string' ? v.symbol : null;
    if (type === 'minmax_price') {
      const min = typeof v.min_price === 'string' ? parseDec(v.min_price) : null;
      const max = typeof v.max_price === 'string' ? parseDec(v.max_price) : null;
      if (
        !sym ||
        !min ||
        !max ||
        decSign(min) < 0 ||
        decCmp(max, min) <= 0 ||
        !Number.isFinite(decToNumber(min)) ||
        !Number.isFinite(decToNumber(max))
      ) {
        unknownFrames++;
        if (sym) bands.delete(sym);
        note(`bad_band:${sym}`, 'warn', `Invalid price band for ${sym ?? '?'}`);
        return;
      }
      bands.set(sym, { min: decToNumber(min), max: decToNumber(max), minDec: min, maxDec: max, at: t, epoch: c.epoch });
      symFrameAt.set(sym, { at: t, epoch: c.epoch });
      return;
    }
    if (type === 'bbo') {
      const bid = topLevel(v.bid);
      const ask = topLevel(v.ask);
      if (!sym) {
        unknownFrames++;
        return;
      }
      symFrameAt.set(sym, { at: t, epoch: c.epoch });
      if (!bid || !ask || decCmp(ask.px, bid.px) < 0) {
        bbos.delete(sym);
        return;
      }
      bbos.set(sym, {
        bid: decToNumber(bid.px),
        ask: decToNumber(ask.px),
        bidSz: bid.sz ? decToNumber(bid.sz) : null,
        askSz: ask.sz ? decToNumber(ask.sz) : null,
        at: t,
        epoch: c.epoch,
      });
      return;
    }
    if (type === 'subscribed' || type === 'unsubscribed') {
      const syms = Array.isArray(v.symbols) ? v.symbols.filter((s): s is string => typeof s === 'string') : [];
      for (const s of syms) {
        if (type === 'subscribed') c.confirmed.add(s);
        else {
          c.confirmed.delete(s);
          bands.delete(s);
          bbos.delete(s);
          symFrameAt.delete(s);
        }
      }
      return;
    }
    unknownFrames++;
    note(`unknown:${type}`, 'info', `Unknown MDS frame type: ${type ?? 'missing'}`);
  }
  function mdsWatch(symbols: string[]): void {
    if (closed) return;
    const add: string[] = [];
    for (const s of symbols) {
      if (typeof s !== 'string' || !SYMBOL_RE.test(s)) {
        note(`bad_sym:${String(s).slice(0, 20)}`, 'warn', `Invalid MDS symbol: ${String(s).slice(0, 40)}`);
        continue;
      }
      if (!watch.has(s)) {
        watch.add(s);
        add.push(s);
      }
    }
    if (add.length === 0) return;
    if (conn && !conn.dropped && conn.openedAt !== null) subscribe(conn, add);
    else ensureQfexMds();
  }
  function liveConn(): MdsConn | null {
    return conn && !conn.dropped && conn.openedAt !== null ? conn : null;
  }
  function mdsBbo(symbol: string): {
    bid: number;
    ask: number;
    at: number;
  } | null {
    const c = liveConn();
    const e = bbos.get(symbol);
    if (!c || !e || e.epoch !== c.epoch) return null;
    if (!(e.bid > 0) || !(e.ask > 0) || e.bid >= e.ask) return null;
    return { bid: e.bid, ask: e.ask, at: e.at };
  }
  function mdsBand(symbol: string): {
    min: number;
    max: number;
    at: number;
  } | null {
    const c = liveConn();
    const e = bands.get(symbol);
    if (!c || !e || e.epoch !== c.epoch) return null;
    let at = e.at;
    if (opts().bandAgeFrom === 'stream') {
      const f = symFrameAt.get(symbol);
      if (f && f.epoch === c.epoch && f.at > at) at = f.at;
    }
    return { min: e.min, max: e.max, at };
  }
  function mdsFreshBand(
    symbol: string,
    maxAgeMs: number,
    at: number = now(),
  ): {
    min: number;
    max: number;
    at: number;
  } | null {
    const b = mdsBand(symbol);
    if (!(maxAgeMs > 0) || !Number.isFinite(maxAgeMs)) throw new TypeError('Required band freshness policy');
    if (!b || at < b.at || at - b.at > maxAgeMs) return null;
    return b;
  }
  function qfexMdsHealth(): QfexMdsHealth {
    const c = liveConn();
    return {
      state,
      since: stateSince,
      epoch,
      reconnects,
      lastFrameAt,
      watched: watch.size,
      subscribed: c ? c.requested.size : 0,
      bands: c ? [...bands.values()].filter((b) => b.epoch === c.epoch).length : 0,
      bbos: c ? [...bbos.values()].filter((b) => b.epoch === c.epoch).length : 0,
      errors,
      unknownFrames,
      lastError,
    };
  }
  async function closeQfexMds(): Promise<void> {
    closed = true;
    cancelTimer(reconnectTimer);
    reconnectTimer = null;
    const c = conn;
    if (c) {
      drop(c, 'MDS closed by caller');
      await Promise.race([
        c.closedP,
        new Promise<void>((r) => {
          const t = setTimeout(r, options.closeWaitMs);
          t.unref();
          void c.closedP.then(() => clearTimeout(t));
        }),
      ]);
    }
    for (const h of [...timers]) cancelTimer(h);
    setState('closed');
  }
  return {
    watch: mdsWatch,
    bbo: mdsBbo,
    band: mdsBand,
    freshBand: mdsFreshBand,
    health: qfexMdsHealth,
    close: closeQfexMds,
  };
}
export type MarketDataStream = ReturnType<typeof createMarketDataStreamInternal>;
All files