src/mds/stream.ts
v0.1.0 · 17.2 KB
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>;