src/session/session.ts
v0.1.0 · 74 KB
import type { ClientRequest, IncomingMessage } from 'node:http';
import WebSocket from 'ws';
import {
buildRestAuthHeaders,
buildWsAuthFrame,
classifyQfexAuthFailure,
createNonceSource,
isQfexAccountId,
maskSecrets,
type QfexKeyPair,
qfexWsFrameIsWrite,
wsQueryAuthUrl,
} from '../auth/auth.js';
import {
decodeTradeFrame,
type QfexFill,
type QfexLeverageRow,
type QfexOrderEvent,
remainingIsZero,
statusClass,
type TradeFrame,
} from '../frames/index.js';
import { type Dec, decAdd, decCmp, decSign, decSub, decToString, parseDec } from '../numbers/index.js';
import { QfexHttpError } from '../rest/client.js';
import { type AddSlot, QFEX_WEIGHTS, type QfexAddSlot, type QfexBucket, type RateBudget } from '../throttle/index.js';
type Obj = Record<string, unknown>;
export interface TradeSession {
ensureStarted(): void;
close(): Promise<void>;
isClosed(): boolean;
whenReady(timeoutMs: number): Promise<boolean>;
submitOrder(
wire: AddOrderWire,
waitTerminal: boolean,
ackTimeoutMs: number,
terminalTimeoutMs: number,
options?: QfexSubmitOptions,
): Promise<SubmitResult>;
submitCancel(symbol: string, orderId: string, timeoutMs: number): Promise<CancelResult>;
getUserOrders(limit: number, offset: number, timeoutMs: number): Promise<QfexOrdersRead>;
getUserLeverage(timeoutMs: number): Promise<Map<string, number>>;
getAvailableLeverageLevels(timeoutMs: number): Promise<QfexLeverageRow[]>;
onOrderEvent(fn: (event: QfexOrderEvent) => void): void;
onFill(fn: (fill: QfexFill) => void): void;
onDrop(fn: (epoch: number) => void): void;
onReady(fn: (epoch: number) => void): void;
onError(fn: (event: QfexSessionErrorEvent) => void): void;
writeBlock(): string | null;
cancelBlock(): string | null;
health(): QfexSessionHealth;
}
export class QfexSessionError extends Error {
override name = 'QfexSessionError';
constructor(
readonly kind: QfexSessionErrorKind,
message: string,
readonly code: string | null = null,
) {
super(maskSecrets(message));
}
}
export type SessionState =
| 'idle'
| 'connecting'
| 'authenticating'
| 'subscribing'
| 'ready'
| 'backoff'
| 'auth_failed'
| 'closed';
export interface QfexSessionHealth {
state: SessionState;
since: number;
epoch: number;
reconnects: number;
lastFrameAt: number;
lastPongAt: number;
authFailure: {
at: number;
httpStatus: number | null;
text: string;
} | null;
subscribed: string[];
ready: boolean;
backoffUntil: number | null;
inflight: {
adds: number;
cancels: number;
reads: number;
};
accountFlags: Record<string, number>;
recycle: string | null;
stats: {
framesIn: number;
framesOut: number;
unknownFrames: number;
errFrames: number;
unsolicitedReads: number;
unsolicitedOrders: number;
foreignFills: number;
dryBlocked: number;
rateLimitedUncorrelated: number;
readTimeouts: number;
};
}
export interface AddOrderWire {
symbol: string;
side: 'BUY' | 'SELL';
tif: 'GTC' | 'IOC';
quantity: string;
price: string;
reduceOnly: boolean;
cloid: string;
}
export type QfexNotSentWhy =
| 'not_ready'
| 'rate_limited'
| 'backpressure'
| 'closed'
| 'dry_run'
| 'duplicate_cloid'
| 'invalid';
export interface QfexRateLimitedSuspect {
at: number;
retryAfterUntil: number;
message: string | null;
}
export type SubmitResult =
| {
kind: 'not_sent';
why: QfexNotSentWhy;
detail?: string;
}
| {
kind: 'rejected';
status?: string;
code?: string;
message?: string;
orderId?: string;
afterAck?: boolean;
events?: QfexOrderEvent[];
fills?: QfexFill[];
}
| {
kind: 'acked';
orderId: string;
events: QfexOrderEvent[];
fills: QfexFill[];
terminal: QfexOrderEvent | null;
}
| {
kind: 'unknown';
why: 'timeout' | 'dropped_after_send' | 'error_reply';
orderId?: string;
events: QfexOrderEvent[];
fills: QfexFill[];
rateLimitedSuspect?: QfexRateLimitedSuspect;
};
export type CancelResult =
| {
kind: 'cancelled';
event: QfexOrderEvent | null;
fillsSeen: boolean;
}
| {
kind: 'not_found';
status: string;
event: QfexOrderEvent | null;
}
| {
kind: 'rejected';
status: string;
event?: QfexOrderEvent | null;
message?: string;
}
| {
kind: 'not_sent';
why?: QfexNotSentWhy;
}
| {
kind: 'unknown';
why?: 'timeout' | 'dropped_after_send';
fillsSeen?: boolean;
};
export interface QfexOrdersRead {
epoch: number;
orders: QfexOrderEvent[];
twaps: unknown[];
stopOrders: unknown[];
malformed: string[];
}
export type QfexSessionErrorKind =
| 'closed'
| 'not_ready'
| 'timeout'
| 'dropped'
| 'rejected'
| 'rate_limited'
| 'malformed'
| 'correlation'
| 'dry_run'
| 'send_failed'
| 'invalid';
export interface QfexSessionErrorEvent {
code: string;
message: string | null;
requestType: string | null;
at: number;
epoch: number;
}
export interface QfexSessionClock {
now(): number;
setTimer(
ms: number,
fn: () => void,
): {
cancel(): void;
};
}
export interface QfexSessionOptions {
url: string;
authMode?: 'query' | 'headers';
credentials: () => QfexKeyPair;
accountId: () => string | null | Promise<string | null>;
dryRun: () => boolean;
clock?: QfexSessionClock;
wallNow?: () => number;
random?: () => number;
authTimeoutMs?: number;
subscribeTimeoutMs?: number;
handshakeTimeoutMs?: number;
pingIntervalMs?: number;
deadAfterMs?: number;
backoffMinMs?: number;
backoffMaxMs?: number;
stableAfterMs?: number;
authFailBackoffMs?: number;
maxBufferedBytes?: number;
replyTimeoutMs?: number;
fillGraceMs?: number;
ordersPageLimit: number;
logWindowMs: number;
strayFillMs: number;
strayFillMax: number;
seenTradesMax: number;
sentCloidsMax: number;
closeWaitMs: number;
httpBodyMax: number;
maxPayloadBytes: number;
}
export interface QfexSubmitOptions {
settleAfterAckMs?: number;
slotWaitMs?: number;
}
interface TimerH {
c: {
cancel(): void;
} | null;
dead: boolean;
}
type Stage = 'handshake' | 'auth' | 'subscribe' | 'ready';
interface Conn {
epoch: number;
ws: WebSocket;
account: string | null;
stage: Stage;
openedAt: number | null;
authedAt: number | null;
readyAt: number | null;
subs: Set<string>;
subWaiting: string | null;
stageTimer: TimerH | null;
pingTimer: TimerH | null;
lastFrameAt: number;
lastPongAt: number;
dropped: boolean;
closedP: Promise<void>;
closedResolve: () => void;
recycle: string | null;
fillsSubscribed: boolean;
authFailed: 'hard' | 'soft' | null;
readsSent: number;
}
interface AddIntent {
wire: AddOrderWire;
epoch: number;
sentAt: number;
waitTerminal: boolean;
settleMs: number;
termMs: number;
orderId: string | null;
acked: boolean;
events: QfexOrderEvent[];
fills: QfexFill[];
fillIds: Set<string>;
terminal: QfexOrderEvent | null;
slot: QfexAddSlot | null;
ackTimer: TimerH | null;
termTimer: TimerH | null;
graceTimer: TimerH | null;
settleTimer: TimerH | null;
rateLimitedSuspect: QfexRateLimitedSuspect | null;
done: boolean;
resolve: (r: SubmitResult) => void;
}
interface CancelWait {
orderId: string;
symbol: string;
epoch: number;
fillsSeen: boolean;
timer: TimerH | null;
done: boolean;
promise: Promise<CancelResult>;
resolve: (r: CancelResult) => void;
}
type ReadKey = 'orders' | 'leverage' | 'levels';
interface ReadWait {
key: ReadKey;
epoch: number;
timer: TimerH | null;
resolve: (v: { epoch: number; frame: TradeFrame }) => void;
reject: (e: Error) => void;
}
type SendOutcome = 'sent' | 'not_open' | 'dry_run' | 'rate_limited' | 'backpressure' | 'forbidden' | 'error';
const FORBIDDEN_TYPES: ReadonlySet<string> = new Set([
'cancel_all_orders',
'cancel_on_disconnect',
'modify_order',
'close_position',
]);
const CHANNELS = ['order_responses', 'fills'] as const;
const FLAG_CODES: ReadonlySet<string> = new Set(['KycRequired', 'TncRequired', 'PermissionDenied']);
const DEFINITE_ADD_ERR_CODES: ReadonlySet<string> = new Set([
'InvalidJSONFormat',
'InvalidParameter',
'PermissionDenied',
'InvalidOrder',
'InvalidOrderId',
'KycRequired',
'TncRequired',
'Unauthenticated',
]);
const SYMBOL_RE = /^[A-Za-z0-9][A-Za-z0-9._]{0,31}-[A-Za-z0-9]{2,8}$/;
const UUID_RE = /^[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$/;
const CLOID_RE = /^[0-9a-f]{32}$/;
const LITERAL_RE = /^\d+(?:\.\d+)?$/;
export const realSessionClock: QfexSessionClock = {
now: () => Date.now(),
setTimer: (ms, fn) => {
const t = setTimeout(fn, Math.max(0, ms));
t.unref();
return { cancel: () => clearTimeout(t) };
},
};
const isObj = (v: unknown): v is Obj => typeof v === 'object' && v !== null && !Array.isArray(v);
export function decodeRestOpenOrders(payload: unknown): Map<string, number> {
if (!isObj(payload)) throw new Error('REST positions response must be an object');
const list = payload.positions;
const out = new Map<string, number>();
if (list === null || list === undefined) return out;
if (!Array.isArray(list)) throw new Error('REST positions must be an array');
for (const row of list) {
if (!isObj(row) || typeof row.symbol !== 'string' || row.symbol === '')
throw new Error('Position row must contain a symbol');
if (out.has(row.symbol)) throw new Error(`Duplicate position symbol: ${row.symbol}`);
const v = row.open_orders;
const d = typeof v === 'string' ? parseDec(v) : null;
if (d === null || d.exp !== 0 || decSign(d) < 0 || d.mant > BigInt(Number.MAX_SAFE_INTEGER)) continue;
out.set(row.symbol, Number(d.mant));
}
return out;
}
const READ_KEY_BY_REQUEST: Readonly<Record<string, ReadKey>> = {
get_user_orders: 'orders',
get_user_leverage: 'leverage',
get_available_leverage_levels: 'levels',
};
function errText(e: unknown): string {
return maskSecrets(String((e as Error)?.message ?? e)).slice(0, 300);
}
function notSentWhy(outcome: SendOutcome): QfexNotSentWhy {
if (outcome === 'dry_run') return 'dry_run';
if (outcome === 'rate_limited') return 'rate_limited';
if (outcome === 'backpressure') return 'backpressure';
return 'not_ready';
}
function canonLiteral(value: unknown): string | null {
if (typeof value !== 'string' || !LITERAL_RE.test(value)) return null;
const parsed = parseDec(value);
return parsed && decSign(parsed) > 0 ? decToString(parsed) : null;
}
export function buildAddOrderFrame(wire: AddOrderWire): string {
const quantity = canonLiteral(wire.quantity),
price = canonLiteral(wire.price);
if (
!SYMBOL_RE.test(wire.symbol) ||
!['BUY', 'SELL'].includes(wire.side) ||
!['GTC', 'IOC'].includes(wire.tif) ||
!CLOID_RE.test(wire.cloid) ||
typeof wire.reduceOnly !== 'boolean' ||
!quantity ||
!price
)
throw new TypeError('Invalid add order frame');
return `{"type":"add_order","params":{"symbol":${JSON.stringify(wire.symbol)},"side":"${wire.side}","order_type":"LIMIT","order_time_in_force":"${wire.tif}","quantity":${quantity},"price":${price},"reduce_only":${wire.reduceOnly},"client_order_id":"${wire.cloid}"}}`;
}
export interface TradeSessionOptions extends QfexSessionOptions {
budget: RateBudget;
addSlot: AddSlot;
webSocketFactory?: (url: string, options: WebSocket.ClientOptions) => WebSocket;
isOwnOrder?: (cloid: string | null) => boolean;
isKnownOrder?: (orderId: string) => boolean;
onConfirmedFill?: (fill: QfexFill) => void;
onGap?: (epoch: number) => void;
onLog?: (message: string) => void;
}
export function createTradeSessionInternal(options: TradeSessionOptions): TradeSession {
if (
!options?.url ||
!options.credentials ||
!options.accountId ||
!options.dryRun ||
!options.budget ||
!options.addSlot
)
throw new TypeError('Required session transport and policy');
const webSocketFactory = options.webSocketFactory ?? ((url, opts) => new WebSocket(url, opts));
const budget = options.budget;
const nonce = createNonceSource();
const newNonce = (at: number) => nonce.next(at);
const acquireAddSlot = options.addSlot.acquire;
const noteWeight = budget.noteWeight;
const noteRateLimited = budget.noteRateLimited;
const rateHold = budget.hold;
function redact(message: string): string {
let secret: string | undefined;
try {
secret = options.credentials().secret;
} catch {}
return maskSecrets(message, { secret });
}
const log = {
info: (message: string) => {
try {
options.onLog?.(redact(message));
} catch {}
},
warn: (message: string) => {
try {
options.onLog?.(redact(message));
} catch {}
},
error: (message: string) => {
try {
options.onLog?.(redact(message));
} catch {}
},
};
function defaultWireAccount(): never {
throw new TypeError('Explicit account selector required');
}
function resolveOptions(o: QfexSessionOptions) {
for (const name of [
'authTimeoutMs',
'subscribeTimeoutMs',
'handshakeTimeoutMs',
'pingIntervalMs',
'deadAfterMs',
'backoffMinMs',
'backoffMaxMs',
'stableAfterMs',
'authFailBackoffMs',
'maxBufferedBytes',
'replyTimeoutMs',
'fillGraceMs',
'ordersPageLimit',
'logWindowMs',
'strayFillMs',
'strayFillMax',
'seenTradesMax',
'sentCloidsMax',
'closeWaitMs',
'httpBodyMax',
'maxPayloadBytes',
] as const)
if (!Number.isFinite(o[name]) || !(o[name]! > 0)) throw new TypeError('Required session policy: ' + name);
return {
...o,
url: o.url!,
credentials: o.credentials!,
accountId: o.accountId!,
dryRun: o.dryRun!,
authMode: o.authMode ?? 'query',
clock: o.clock ?? realSessionClock,
wallNow: o.wallNow ?? Date.now,
random: o.random ?? Math.random,
authTimeoutMs: o.authTimeoutMs!,
subscribeTimeoutMs: o.subscribeTimeoutMs!,
handshakeTimeoutMs: o.handshakeTimeoutMs!,
pingIntervalMs: o.pingIntervalMs!,
deadAfterMs: o.deadAfterMs!,
backoffMinMs: o.backoffMinMs!,
backoffMaxMs: o.backoffMaxMs!,
stableAfterMs: o.stableAfterMs!,
authFailBackoffMs: o.authFailBackoffMs!,
maxBufferedBytes: o.maxBufferedBytes!,
replyTimeoutMs: o.replyTimeoutMs!,
fillGraceMs: o.fillGraceMs!,
ordersPageLimit: o.ordersPageLimit!,
};
}
type Resolved = ReturnType<typeof resolveOptions>;
class QfexSession {
private readonly optsSrc: () => QfexSessionOptions;
private st: SessionState = 'idle';
private stSince: number;
private epoch = 0;
private reconnects = 0;
private attempts = 0;
private conn: Conn | null = null;
private connecting: Promise<void> | null = null;
private reconnectTimer: TimerH | null = null;
private backoffUntil: number | null = null;
private closed = false;
private closeP: Promise<void> | null = null;
private authFailure: {
at: number;
httpStatus: number | null;
text: string;
} | null = null;
private authFailStreak = 0;
private lastFrameAt = 0;
private lastPongAt = 0;
private readonly timers = new Set<TimerH>();
private readonly sleepers = new Set<() => void>();
private readonly readyWaiters = new Set<(ok: boolean) => void>();
private readonly adds = new Map<string, AddIntent>();
private readonly addsByOid = new Map<string, AddIntent>();
private readonly cancels = new Map<string, CancelWait>();
private readonly reads = new Map<ReadKey, ReadWait>();
private readonly readChains = new Map<ReadKey, Promise<void>>();
private readonly sentCloids = new Set<string>();
private readonly strayFills = new Map<
string,
Array<{
f: QfexFill;
at: number;
}>
>();
private readonly seenTrades = new Map<string, number>();
private readonly logGate = new Map<string, number>();
private readonly flags = new Map<string, number>();
private readonly listeners = {
order: [] as Array<(o: QfexOrderEvent) => void>,
fill: [] as Array<(f: QfexFill) => void>,
drop: [] as Array<(epoch: number) => void>,
error: [] as Array<(e: QfexSessionErrorEvent) => void>,
ready: [] as Array<(epoch: number) => void>,
};
private readonly stats = {
framesIn: 0,
framesOut: 0,
unknownFrames: 0,
errFrames: 0,
unsolicitedReads: 0,
unsolicitedOrders: 0,
foreignFills: 0,
dryBlocked: 0,
rateLimitedUncorrelated: 0,
readTimeouts: 0,
};
constructor(opts: QfexSessionOptions | (() => QfexSessionOptions)) {
this.optsSrc = typeof opts === 'function' ? opts : () => opts;
this.stSince = this.o().clock.now();
}
private cachedFor: QfexSessionOptions | null = null;
private cached: Resolved | null = null;
private o(): Resolved {
const src = this.optsSrc();
if (this.cached === null || src !== this.cachedFor) {
this.cached = resolveOptions(src);
this.cachedFor = src;
}
return this.cached;
}
private now(): number {
return this.o().clock.now();
}
private timer(ms: number, fn: () => void): TimerH {
const h: TimerH = { c: null, dead: false };
this.timers.add(h);
h.c = this.o().clock.setTimer(ms, () => {
if (h.dead) return;
h.dead = true;
this.timers.delete(h);
try {
fn();
} catch (e) {
log.error(`Session timer callback failed: ${errText(e)}`);
}
});
return h;
}
private cancelTimer(h: TimerH | null): void {
if (!h || h.dead) return;
h.dead = true;
this.timers.delete(h);
try {
h.c?.cancel();
} catch {}
}
private sleep(ms: number): Promise<void> {
if (this.closed) return Promise.resolve();
return new Promise<void>((resolve) => {
let t: TimerH | null = null;
const wake = (): void => {
this.sleepers.delete(wake);
this.cancelTimer(t);
resolve();
};
this.sleepers.add(wake);
t = this.timer(ms, wake);
});
}
private note(key: string, level: 'info' | 'warn' | 'error', msg: string): void {
const now = this.now();
const last = this.logGate.get(key);
if (last !== undefined && now - last < this.o().logWindowMs) return;
this.logGate.set(key, now);
if (this.logGate.size > 500) {
for (const [k, at] of this.logGate) if (now - at >= this.o().logWindowMs) this.logGate.delete(k);
}
log[level](maskSecrets(msg));
}
private setState(s: SessionState): void {
if (this.st === s) return;
this.st = s;
this.stSince = this.now();
}
ensureStarted(): void {
if (this.closed || this.conn || this.connecting || this.reconnectTimer) return;
this.startConnect();
}
private startConnect(): void {
if (this.closed || this.conn || this.connecting) return;
this.setState('connecting');
const p = this.openConn()
.catch((e) => {
log.error(`Session connection task failed: ${errText(e)}`);
if (!this.closed && !this.conn && !this.reconnectTimer)
this.scheduleReconnect(this.backoffDelay(), 'backoff');
})
.finally(() => {
if (this.connecting === p) this.connecting = null;
});
this.connecting = p;
}
private async resolveAccount(): Promise<string | null> {
const f = this.o().accountId;
const v = f ? await f() : await defaultWireAccount();
if (v !== null && !isQfexAccountId(v))
throw new QfexSessionError('invalid', 'Account selector returned an invalid UUID');
return v;
}
private async openConn(): Promise<void> {
const o = this.o();
let keys: QfexKeyPair;
let account: string | null;
try {
keys = o.credentials();
account = await this.resolveAccount();
} catch (e) {
if (this.closed) return;
const http = e instanceof QfexHttpError ? e : null;
const transient = http !== null && http.kind !== 'auth' && http.kind !== 'client';
const text = `Account selection failed: ${errText(e)}`;
if (transient) {
this.note('acct_transient', 'warn', `Session will retry after transient account-selection failure: ${text}`);
this.scheduleReconnect(this.backoffDelay(), 'backoff');
} else {
this.authFailure = { at: this.now(), httpStatus: http?.status ?? null, text };
this.authFailStreak++;
log.error(
`Session account-selection failure: ${text}; retry backoff ${Math.round(o.authFailBackoffMs / 1000)} seconds`,
);
this.scheduleReconnect(o.authFailBackoffMs, 'auth_failed');
}
return;
}
if (this.closed) return;
const epoch = ++this.epoch;
const wall = o.wallNow();
let url: string;
let headers: Record<string, string> | undefined;
try {
url = o.authMode === 'query' ? wsQueryAuthUrl(o.url, keys.publicKey) : o.url;
headers = o.authMode === 'headers' ? buildRestAuthHeaders(keys, account, wall, newNonce(wall)) : undefined;
} catch (e) {
this.authFailure = { at: this.now(), httpStatus: null, text: `Credential loading failed: ${errText(e)}` };
log.error(`Session authentication unavailable: ${this.authFailure.text}`);
this.scheduleReconnect(o.authFailBackoffMs, 'auth_failed');
return;
}
let ws: WebSocket;
try {
ws = webSocketFactory(url, {
headers,
handshakeTimeout: o.handshakeTimeoutMs,
maxPayload: o.maxPayloadBytes,
perMessageDeflate: false,
followRedirects: false,
});
} catch (e) {
log.error(`WebSocket construction failed: ${errText(e)}`);
this.scheduleReconnect(this.backoffDelay(), 'backoff');
return;
}
let closedResolve: () => void = () => undefined;
const closedP = new Promise<void>((r) => {
closedResolve = r;
});
const now = this.now();
const conn: Conn = {
epoch,
ws,
account,
stage: 'handshake',
openedAt: null,
authedAt: null,
readyAt: null,
subs: new Set(),
subWaiting: null,
stageTimer: null,
pingTimer: null,
lastFrameAt: now,
lastPongAt: 0,
dropped: false,
closedP,
closedResolve,
recycle: null,
fillsSubscribed: false,
authFailed: null,
readsSent: 0,
};
this.conn = conn;
ws.on('error', (err) => this.onWsError(conn, err));
ws.on('unexpected-response', (req, res) => this.onUnexpected(conn, req, res));
ws.on('open', () => this.onOpen(conn, keys));
ws.on('message', (data, isBinary) => {
try {
this.onMessage(conn, data, isBinary);
} catch (e) {
log.error(`Session frame handler failed: ${errText(e)}`);
}
});
ws.on('pong', () => {
const t = this.now();
conn.lastPongAt = t;
this.lastPongAt = t;
});
ws.on('close', (code) => {
conn.closedResolve();
this.dropConn(conn, `WebSocket closed with code ${code}`);
});
}
private onWsError(conn: Conn, err: Error): void {
if (conn.dropped) return;
this.note(`ws_error:${conn.epoch}`, 'warn', `WebSocket error at epoch ${conn.epoch}: ${errText(err)}`);
this.dropConn(conn, `WebSocket error: ${errText(err)}`);
}
private onUnexpected(conn: Conn, req: ClientRequest, res: IncomingMessage): void {
const status = res.statusCode ?? 0;
let body = '';
let finished = false;
let cap: ReturnType<typeof setTimeout> | null = null;
const finish = (): void => {
if (finished) return;
finished = true;
if (cap) clearTimeout(cap);
const auth = classifyQfexAuthFailure(status, body);
const text = `WebSocket handshake returned HTTP ${status}${body ? `: ${maskSecrets(body.slice(0, 200))}` : ''}`;
this.authFailure = { at: this.now(), httpStatus: status, text };
if (auth) {
this.authFailStreak++;
conn.authFailed = auth.retryFreshNonce && this.authFailStreak < 2 ? 'soft' : 'hard';
}
log.error(`WebSocket handshake rejected: ${text}${auth ? ` (${auth.kind})` : ''}`);
try {
conn.ws.terminate();
} catch {}
this.dropConn(conn, text);
};
cap = setTimeout(() => finish(), this.o().closeWaitMs);
cap.unref();
res.on('error', () => finish());
req.on('error', () => undefined);
res.setEncoding('utf8');
res.on('data', (d: string) => {
if (body.length < this.o().httpBodyMax) body += d;
});
res.on('end', () => finish());
res.on('close', () => finish());
}
private onOpen(conn: Conn, keys: QfexKeyPair): void {
if (conn.dropped || this.closed) {
try {
conn.ws.terminate();
} catch {}
return;
}
const o = this.o();
const now = this.now();
conn.openedAt = now;
conn.lastFrameAt = now;
conn.stage = 'auth';
this.setState('authenticating');
let frame: string;
try {
const wall = o.wallNow();
frame = buildWsAuthFrame(keys, conn.account, wall, newNonce(wall));
} catch (e) {
this.authFailure = { at: now, httpStatus: null, text: `Auth frame construction failed: ${errText(e)}` };
conn.authFailed = 'hard';
this.dropConn(conn, this.authFailure.text);
return;
}
const r = this.sendRaw(conn, frame, 'auth', null, 0);
if (r !== 'sent') {
this.dropConn(conn, `Auth frame was not sent: ${r}`);
return;
}
conn.stageTimer = this.timer(o.authTimeoutMs, () =>
this.dropConn(conn, `Authentication reply timed out after ${o.authTimeoutMs} ms`),
);
this.startLiveness(conn);
}
private startLiveness(conn: Conn): void {
const o = this.o();
const tick = (): void => {
if (conn.dropped) return;
const now = this.now();
const lastSign = Math.max(conn.lastFrameAt, conn.lastPongAt, conn.openedAt ?? 0);
if (now - lastSign >= o.deadAfterMs) {
this.dropConn(conn, `Heartbeat missing for ${Math.round((now - lastSign) / 1000)} seconds`);
return;
}
try {
if (conn.ws.readyState === WebSocket.OPEN) conn.ws.ping();
} catch (e) {
this.note('ping_fail', 'warn', `WebSocket ping failed: ${errText(e)}`);
}
conn.pingTimer = this.timer(o.pingIntervalMs, tick);
};
conn.pingTimer = this.timer(o.pingIntervalMs, tick);
}
private onAuthed(conn: Conn): void {
this.cancelTimer(conn.stageTimer);
conn.stageTimer = null;
conn.authedAt = this.now();
conn.stage = 'subscribe';
this.authFailStreak = 0;
this.authFailure = null;
this.setState('subscribing');
this.subscribeNext(conn);
}
private authRejected(conn: Conn, text: string, hard: boolean): void {
this.authFailStreak++;
conn.authFailed = hard || this.authFailStreak >= 2 ? 'hard' : 'soft';
this.authFailure = { at: this.now(), httpStatus: null, text };
log.error(`Session authentication rejected: ${text}`);
this.dropConn(conn, text);
}
private subscribeNext(conn: Conn): void {
if (conn.dropped || conn.stage !== 'subscribe') return;
const o = this.o();
const next = CHANNELS.find((c) => !conn.subs.has(c));
if (next === undefined) {
this.becomeReady(conn);
return;
}
const left = this.retryAfterLeft();
if (left > 0) {
conn.stageTimer = this.timer(left + 1, () => this.subscribeNext(conn));
return;
}
conn.subWaiting = next;
const r = this.sendRaw(
conn,
`{"type":"subscribe","params":{"channels":["${next}"]}}`,
'subscribe',
'general',
0.1,
);
if (r !== 'sent') {
this.dropConn(conn, `Subscription ${next} was not sent: ${r}`);
return;
}
conn.stageTimer = this.timer(o.subscribeTimeoutMs, () =>
this.dropConn(conn, `Subscription ${next} reply timed out after ${o.subscribeTimeoutMs} ms`),
);
}
private onSubscribed(conn: Conn, channel: string): void {
if (conn.stage !== 'subscribe' || channel !== conn.subWaiting) {
this.note(`sub_unexpected:${channel}`, 'warn', `Unexpected subscription acknowledgement: ${channel}`);
return;
}
this.cancelTimer(conn.stageTimer);
conn.stageTimer = null;
conn.subWaiting = null;
conn.subs.add(channel);
if (channel === 'fills') {
conn.fillsSubscribed = true;
options.onGap?.(conn.epoch);
}
this.subscribeNext(conn);
}
private becomeReady(conn: Conn): void {
const now = this.now();
conn.stage = 'ready';
conn.readyAt = now;
this.setState('ready');
log.info(`Trade session ready at epoch ${conn.epoch}; subscribed to ${[...conn.subs].join(', ')}`);
for (const w of [...this.readyWaiters]) w(true);
this.readyWaiters.clear();
for (const fn of this.listeners.ready) this.safe(() => fn(conn.epoch));
}
private backoffDelay(): number {
const o = this.o();
const base = Math.min(o.backoffMaxMs, o.backoffMinMs * 2 ** Math.min(this.attempts, 16));
this.attempts++;
const r = Math.min(1, Math.max(0, o.random()));
return Math.max(o.backoffMinMs, Math.round(base * (0.5 + 0.5 * r)));
}
private scheduleReconnect(delay: number, state: 'backoff' | 'auth_failed'): void {
if (this.closed) return;
this.cancelTimer(this.reconnectTimer);
this.setState(state);
this.backoffUntil = this.now() + delay;
if (state === 'auth_failed') {
for (const w of [...this.readyWaiters]) w(false);
this.readyWaiters.clear();
}
this.reconnectTimer = this.timer(delay, () => {
this.reconnectTimer = null;
this.backoffUntil = null;
this.reconnects++;
this.startConnect();
});
}
private dropConn(conn: Conn, reason: string): void {
if (conn.dropped) return;
conn.dropped = true;
this.cancelTimer(conn.stageTimer);
this.cancelTimer(conn.pingTimer);
conn.stageTimer = null;
conn.pingTimer = null;
if (this.conn === conn) this.conn = null;
const now = this.now();
const wasAuthed = conn.authedAt !== null;
if (conn.fillsSubscribed) options.onGap?.(conn.epoch);
for (const it of [...this.adds.values()]) {
if (it.epoch !== conn.epoch) continue;
this.finishAdd(it, {
kind: 'unknown',
why: 'dropped_after_send',
...(it.orderId ? { orderId: it.orderId } : {}),
events: [...it.events],
fills: [...it.fills],
...(it.rateLimitedSuspect ? { rateLimitedSuspect: it.rateLimitedSuspect } : {}),
});
}
for (const cw of [...this.cancels.values()]) {
if (cw.epoch === conn.epoch)
this.finishCancel(cw, { kind: 'unknown', why: 'dropped_after_send', fillsSeen: cw.fillsSeen });
}
for (const w of [...this.reads.values()]) {
if (w.epoch !== conn.epoch) continue;
this.reads.delete(w.key);
this.cancelTimer(w.timer);
w.reject(new QfexSessionError('dropped', `Session epoch ${conn.epoch} dropped: ${reason}`));
}
try {
conn.ws.terminate();
} catch {}
if (wasAuthed) {
log.warn(`Session epoch ${conn.epoch} dropped: ${maskSecrets(reason)}`);
for (const fn of this.listeners.drop) this.safe(() => fn(conn.epoch));
}
if (this.closed) {
this.setState('closed');
return;
}
const o = this.o();
if (conn.authFailed === 'hard') {
log.error(`Authentication backoff ${Math.round(o.authFailBackoffMs / 1000)} seconds`);
this.scheduleReconnect(o.authFailBackoffMs, 'auth_failed');
return;
}
if (conn.authFailed === 'soft') {
this.scheduleReconnect(o.backoffMinMs, 'backoff');
return;
}
if (conn.readyAt !== null && now - conn.readyAt >= o.stableAfterMs) this.attempts = 0;
this.scheduleReconnect(this.backoffDelay(), 'backoff');
}
private requestRecycle(conn: Conn, reason: string): void {
if (conn.dropped) return;
if (!conn.recycle) {
conn.recycle = reason;
log.warn(`Session epoch ${conn.epoch} will recycle after active writes settle: ${reason}`);
}
this.maybeRecycleNow();
}
private maybeRecycleNow(): void {
const conn = this.conn;
if (!conn || !conn.recycle || conn.dropped) return;
for (const it of this.adds.values()) if (it.epoch === conn.epoch) return;
for (const cw of this.cancels.values()) if (cw.epoch === conn.epoch) return;
this.dropConn(conn, `Session recycled: ${conn.recycle}`);
}
close(): Promise<void> {
if (this.closeP) return this.closeP;
this.closed = true;
this.closeP = (async () => {
this.cancelTimer(this.reconnectTimer);
this.reconnectTimer = null;
this.backoffUntil = null;
const pending = this.connecting;
if (pending) {
await Promise.race([
pending.catch(() => undefined),
new Promise<void>((r) => {
const t = setTimeout(r, this.o().closeWaitMs);
t.unref();
void pending.finally(() => clearTimeout(t)).catch(() => undefined);
}),
]);
}
const conn = this.conn;
if (conn) {
this.dropConn(conn, 'Session closed by caller');
await Promise.race([
conn.closedP,
new Promise<void>((r) => {
const t = setTimeout(r, this.o().closeWaitMs);
t.unref();
void conn.closedP.then(() => clearTimeout(t));
}),
]);
}
for (const it of [...this.adds.values()]) {
this.finishAdd(it, {
kind: 'unknown',
why: 'dropped_after_send',
...(it.orderId ? { orderId: it.orderId } : {}),
events: [...it.events],
fills: [...it.fills],
});
}
for (const cw of [...this.cancels.values()])
this.finishCancel(cw, { kind: 'unknown', why: 'dropped_after_send', fillsSeen: cw.fillsSeen });
for (const w of [...this.reads.values()]) {
this.reads.delete(w.key);
this.cancelTimer(w.timer);
w.reject(new QfexSessionError('closed', 'Session is closed'));
}
for (const w of [...this.readyWaiters]) w(false);
this.readyWaiters.clear();
for (const wake of [...this.sleepers]) wake();
for (const h of [...this.timers]) this.cancelTimer(h);
this.setState('closed');
})();
return this.closeP;
}
isClosed(): boolean {
return this.closed;
}
whenReady(timeoutMs: number): Promise<boolean> {
this.ensureStarted();
if (this.conn?.stage === 'ready' && !this.conn.dropped) return Promise.resolve(true);
if (this.closed || this.st === 'auth_failed') return Promise.resolve(false);
return new Promise<boolean>((resolve) => {
let t: TimerH | null = null;
const done = (ok: boolean): void => {
this.readyWaiters.delete(done);
this.cancelTimer(t);
resolve(ok);
};
this.readyWaiters.add(done);
t = this.timer(timeoutMs, () => done(false));
});
}
private retryAfterLeft(): number {
return budget.hold('general') || budget.hold('cancel') ? 1 : 0;
}
private sendRaw(conn: Conn, text: string, type: string, bucket: QfexBucket | null, weight: number): SendOutcome {
if (FORBIDDEN_TYPES.has(type)) {
log.error(`Forbidden WebSocket write type: ${type}`);
return 'forbidden';
}
if (conn.dropped || conn.ws.readyState !== WebSocket.OPEN) return 'not_open';
if (qfexWsFrameIsWrite(text) && this.o().dryRun()) {
this.stats.dryBlocked++;
this.note(`dry:${type}`, 'warn', `Read-only mode blocks WebSocket frame: ${type}`);
return 'dry_run';
}
if (type !== 'auth' && this.retryAfterLeft() > 0) return 'rate_limited';
if (conn.ws.bufferedAmount > this.o().maxBufferedBytes) return 'backpressure';
try {
conn.ws.send(text, (err) => {
if (err) {
this.note(`send_err:${conn.epoch}`, 'warn', `WebSocket ${type} send failed: ${errText(err)}`);
this.dropConn(conn, `WebSocket send failed: ${errText(err)}`);
}
});
} catch (e) {
this.note(`send_throw:${conn.epoch}`, 'warn', `WebSocket ${type} send threw: ${errText(e)}`);
return 'error';
}
this.stats.framesOut++;
if (bucket) noteWeight(bucket, weight);
return 'sent';
}
private readyConn(): Conn | null {
const c = this.conn;
return c && c.stage === 'ready' && !c.dropped ? c : null;
}
private onMessage(conn: Conn, data: WebSocket.RawData, isBinary: boolean): void {
if (conn.dropped) return;
const now = this.now();
conn.lastFrameAt = now;
this.lastFrameAt = now;
this.stats.framesIn++;
if (isBinary) {
this.note('binary', 'warn', 'Binary trade frame ignored');
return;
}
const text = Buffer.isBuffer(data)
? data.toString('utf8')
: Array.isArray(data)
? Buffer.concat(data).toString('utf8')
: Buffer.from(data).toString('utf8');
const f = decodeTradeFrame(text);
if (f.k === 'err' && f.message !== null) f.message = redact(f.message);
if (f.k === 'unknown' && f.detail !== null) f.detail = redact(f.detail);
switch (f.k) {
case 'auth':
if (conn.stage !== 'auth') {
this.note('auth_late', 'warn', `Late auth response ignored: ok=${f.ok}`);
return;
}
if (f.ok) this.onAuthed(conn);
else this.authRejected(conn, `Authentication rejected in ${f.form}${f.detail ? `: ${f.detail}` : ''}`, false);
return;
case 'subscribed':
this.onSubscribed(conn, f.channel);
return;
case 'unsubscribed':
this.note(`unsub:${f.channel}`, 'warn', `Server unsubscribed from ${f.channel}`);
return;
case 'order':
this.routeOrder(conn, f.o);
return;
case 'order_partial':
this.routePartial(conn, f);
return;
case 'fill':
this.routeFill(conn, f.f);
return;
case 'orders':
this.routeRead(conn, 'orders', f);
return;
case 'leverage':
this.routeRead(conn, 'leverage', f);
return;
case 'leverage_levels':
this.routeRead(conn, 'levels', f);
return;
case 'err':
this.routeErr(conn, f);
return;
case 'ping':
return;
case 'ack':
this.note('ack', 'info', 'Uncorrelated acknowledgement ignored');
return;
case 'positions':
case 'balances':
case 'twap':
case 'stop_order':
case 'user_trades':
this.note(`push:${f.k}`, 'info', `Unsolicited ${f.k} read ignored`);
return;
case 'unknown':
this.stats.unknownFrames++;
this.note(
`unknown:${f.why}`,
'warn',
`Unknown trade frame: ${f.why}${f.detail ? `: ${String(f.detail).slice(0, 120)}` : ''}`,
);
return;
}
}
private safe(fn: () => void): void {
try {
fn();
} catch (e) {
log.error(`Session listener failed: ${errText(e)}`);
}
}
private isOurTag(cloid: string | null): boolean {
return options.isOwnOrder?.(cloid) ?? false;
}
private routeOrder(conn: Conn, o: QfexOrderEvent): void {
const oid = o.orderId.toLowerCase();
const cls = statusClass(o.status);
const cw = this.cancels.get(oid);
if (cw && cw.symbol !== o.symbol) {
this.finishCancel(cw, { kind: 'unknown', fillsSeen: cw.fillsSeen });
this.requestRecycle(conn, 'Contradictory cancellation symbol');
return;
}
if (cw && cw.epoch === conn.epoch && !cw.done) {
if (cls === 'fill' && !remainingIsZero(o)) {
cw.fillsSeen = true;
} else if (cls === 'fill') {
this.finishCancel(cw, { kind: 'rejected', status: o.status, event: o });
return;
} else if (cls === 'terminal') {
this.finishCancel(cw, { kind: 'cancelled', event: o, fillsSeen: cw.fillsSeen });
return;
} else if (cls === 'cancel_reply') {
this.finishCancel(cw, { kind: 'not_found', status: o.status, event: o });
return;
} else if (cls === 'reject') {
this.finishCancel(cw, { kind: 'rejected', status: o.status, event: o });
return;
}
}
const it = (o.cloid !== null ? this.adds.get(o.cloid) : undefined) ?? this.addsByOid.get(oid);
if (it && it.epoch === conn.epoch && !it.done) {
if (!this.identityMatches(it, o)) {
this.quarantine(it);
this.requestRecycle(conn, 'Contradictory order identity');
return;
} else {
this.feedAdd(it, o);
return;
}
}
this.unsolicitedOrder(o);
}
private bindOrderId(it: AddIntent, orderId: string): void {
if (it.orderId !== null) return;
it.orderId = orderId;
const oid = orderId.toLowerCase();
this.addsByOid.set(oid, it);
const stray = this.strayFills.get(oid);
if (stray) {
this.strayFills.delete(oid);
for (const s of stray) this.attachFill(it, s.f);
}
}
private identityMatches(
it: AddIntent,
event: { symbol: string; side: string; orderId: string; cloid: string | null },
): boolean {
return (
event.symbol === it.wire.symbol &&
event.side === it.wire.side &&
(event.cloid === null || event.cloid === it.wire.cloid) &&
(it.orderId === null || event.orderId.toLowerCase() === it.orderId.toLowerCase())
);
}
private quarantine(it: AddIntent): void {
this.note(
'contradictory_identity',
'error',
'Correlated frame contradicts the submitted symbol, side or order identity; outcome is unknown',
);
this.finishAdd(it, {
kind: 'unknown',
why: 'error_reply',
...(it.orderId ? { orderId: it.orderId } : {}),
events: [...it.events],
fills: [...it.fills],
});
}
private firstResponse(it: AddIntent): void {
this.cancelTimer(it.ackTimer);
it.ackTimer = null;
this.releaseSlot(it);
}
private feedAdd(it: AddIntent, o: QfexOrderEvent): void {
this.bindOrderId(it, o.orderId);
if (it.done) return;
it.events.push(o);
const cls = statusClass(o.status);
switch (cls) {
case 'live':
this.firstResponse(it);
it.acked = true;
if (!it.waitTerminal) this.onGtcAcked(it);
return;
case 'fill':
this.firstResponse(it);
it.acked = true;
if (remainingIsZero(o)) {
it.terminal = o;
this.onTerminal(it);
} else if (!it.waitTerminal) {
this.onGtcAcked(it);
}
return;
case 'terminal':
this.firstResponse(it);
it.terminal = o;
this.onTerminal(it);
return;
case 'reject':
this.firstResponse(it);
if (it.fills.length > 0 || decSign(this.impliedFilled(it)) > 0) {
it.terminal = o;
this.onTerminal(it);
return;
}
this.finishAdd(it, {
kind: 'rejected',
status: o.status,
orderId: o.orderId,
afterAck: it.acked,
events: [...it.events],
fills: [],
});
return;
default:
this.note(
`add_status:${o.status}`,
'warn',
`Unclassified order status for ${it.wire.cloid}: ${o.status} (${cls}); waiting for evidence`,
);
}
}
private onGtcAcked(it: AddIntent): void {
if (it.settleMs > 0) {
if (it.settleTimer) return;
it.settleTimer = this.timer(it.settleMs, () => this.finishAcked(it));
return;
}
this.finishAcked(it);
}
private finishAcked(it: AddIntent): void {
if (it.done || it.orderId === null) return;
this.finishAdd(it, {
kind: 'acked',
orderId: it.orderId,
events: [...it.events],
fills: [...it.fills],
terminal: it.terminal,
});
}
private impliedFilled(it: AddIntent): Dec {
let best: Dec = { mant: 0n, exp: 0 };
const sent = parseDec(it.wire.quantity);
if (!sent) return best;
for (const e of it.events) {
if (statusClass(e.status) !== 'fill') continue;
const r = parseDec(e.remaining);
if (!r) continue;
const f = decSub(sent, r);
if (decCmp(f, best) > 0) best = f;
}
return best;
}
private fillsSum(it: AddIntent): Dec {
let s: Dec = { mant: 0n, exp: 0 };
for (const f of it.fills) {
const q = parseDec(f.qty);
if (q) s = decAdd(s, q);
}
return s;
}
private fillsBehind(it: AddIntent): boolean {
const sum = this.fillsSum(it);
if (decCmp(sum, this.impliedFilled(it)) < 0) return true;
return it.terminal !== null && it.terminal.status === 'IOC_PARTIALLY_FILLED' && decSign(sum) === 0;
}
private errCannotProveReject(it: AddIntent, code: string): boolean {
if (it.acked || it.fills.length > 0 || decSign(this.impliedFilled(it)) > 0) return true;
return code !== 'RateLimited' && !DEFINITE_ADD_ERR_CODES.has(code);
}
private onTerminal(it: AddIntent): void {
if (it.done) return;
if (!it.waitTerminal) {
this.finishAcked(it);
return;
}
if (!this.fillsBehind(it)) {
this.finishAcked(it);
return;
}
if (it.graceTimer) return;
const left = it.sentAt + it.termMs - this.now();
const grace = Math.max(0, Math.min(this.o().fillGraceMs, left));
if (grace <= 0) {
this.finishAcked(it);
return;
}
it.graceTimer = this.timer(grace, () => this.finishAcked(it));
}
private attachFill(it: AddIntent, f: QfexFill): void {
if (!this.identityMatches(it, f)) {
this.quarantine(it);
return;
}
if (it.done) return;
if (it.fillIds.has(f.tradeId)) return;
it.fillIds.add(f.tradeId);
it.fills.push(f);
if (it.graceTimer && !this.fillsBehind(it)) this.finishAcked(it);
}
private routePartial(
conn: Conn,
f: Extract<
TradeFrame,
{
k: 'order_partial';
}
>,
): void {
const cls = statusClass(f.status);
if (f.orderId) {
const cw = this.cancels.get(f.orderId.toLowerCase());
if (cw && cw.epoch === conn.epoch && !cw.done) {
if (cls === 'cancel_reply') {
this.finishCancel(cw, { kind: 'not_found', status: f.status, event: null });
return;
}
if (cls === 'terminal') {
this.finishCancel(cw, { kind: 'cancelled', event: null, fillsSeen: cw.fillsSeen });
return;
}
if (cls === 'reject' || cls === 'fill') {
this.finishCancel(cw, { kind: 'rejected', status: f.status, event: null });
return;
}
}
}
if (f.cloid) {
const it = this.adds.get(f.cloid);
if (it && it.epoch === conn.epoch && !it.done && cls === 'reject') {
if (
(f.symbol !== null && f.symbol !== it.wire.symbol) ||
(f.orderId !== null && it.orderId !== null && f.orderId.toLowerCase() !== it.orderId.toLowerCase())
) {
this.quarantine(it);
this.requestRecycle(conn, 'Contradictory partial order identity');
return;
}
this.firstResponse(it);
if (it.fills.length > 0 || decSign(this.impliedFilled(it)) > 0) {
this.finishAdd(it, {
kind: 'unknown',
why: 'error_reply',
...(it.orderId ? { orderId: it.orderId } : {}),
events: [...it.events],
fills: [...it.fills],
...(it.rateLimitedSuspect ? { rateLimitedSuspect: it.rateLimitedSuspect } : {}),
});
return;
}
this.finishAdd(it, {
kind: 'rejected',
status: f.status,
...(f.orderId ? { orderId: f.orderId } : {}),
afterAck: it.acked,
events: [...it.events],
});
return;
}
}
this.note(
`partial:${f.status}`,
'warn',
`Partial order response (${f.status}) cannot confirm terminal outcome: ${f.why.slice(0, 120)}`,
);
}
private routeFill(conn: Conn, f: QfexFill): void {
const now = this.now();
const oid = f.orderId.toLowerCase();
const it = (f.cloid ? this.adds.get(f.cloid) : undefined) ?? this.addsByOid.get(oid);
const live = it && it.epoch === conn.epoch && !it.done ? it : undefined;
if (live && !this.identityMatches(live, f)) {
this.quarantine(live);
this.requestRecycle(conn, 'Contradictory fill identity');
return;
}
const noCloid = f.cloid === null || f.cloid === '';
if (live || this.isOurTag(f.cloid) || (noCloid && options.isKnownOrder?.(f.orderId) === true)) {
try {
options.onConfirmedFill?.(f);
} catch (e) {
log.error(`Confirmed-fill callback failed: ${errText(e)}`);
}
} else {
this.stats.foreignFills++;
this.note('foreign_fill', 'info', `Fill for an unowned order ${f.orderId} (${f.symbol})`);
}
const cw = this.cancels.get(oid);
if (cw && !cw.done) cw.fillsSeen = true;
if (live) {
if (live.orderId === null) this.bindOrderId(live, f.orderId);
this.releaseSlot(live);
this.attachFill(live, f);
} else if (!it) {
let arr = this.strayFills.get(oid);
if (!arr) {
arr = [];
this.strayFills.set(oid, arr);
}
arr.push({ f, at: now });
this.pruneStray(now);
}
const dup = this.seenTrades.has(f.tradeId);
if (!dup) {
this.seenTrades.set(f.tradeId, now);
if (this.seenTrades.size > this.o().seenTradesMax) {
const drop = this.seenTrades.size - this.o().seenTradesMax;
let i = 0;
for (const k of this.seenTrades.keys()) {
if (i++ >= drop) break;
this.seenTrades.delete(k);
}
}
for (const fn of this.listeners.fill) this.safe(() => fn(f));
}
}
private pruneStray(now: number): void {
let total = 0;
for (const [k, arr] of this.strayFills) {
const keep = arr.filter((s) => now - s.at < this.o().strayFillMs);
if (keep.length === 0) this.strayFills.delete(k);
else {
this.strayFills.set(k, keep);
total += keep.length;
}
}
if (total > this.o().strayFillMax) this.strayFills.clear();
}
private unsolicitedOrder(o: QfexOrderEvent): void {
this.stats.unsolicitedOrders++;
const now = this.now();
if (o.status === 'CANCELLED_STP') {
this.note(`stp:${o.symbol}`, 'info', `Self-trade prevention cancelled order ${o.orderId} on ${o.symbol}`);
}
for (const fn of this.listeners.order) this.safe(() => fn(o));
}
private routeRead(conn: Conn, key: ReadKey, frame: TradeFrame): void {
const w = this.reads.get(key);
if (w && w.epoch === conn.epoch) {
this.reads.delete(key);
this.cancelTimer(w.timer);
w.resolve({ epoch: conn.epoch, frame });
return;
}
this.stats.unsolicitedReads++;
if (conn.readsSent === 0) {
this.note(
`read_push:${key}`,
'warn',
`Unexpected ${key} read at epoch ${conn.epoch}; retaining active writes before recycle`,
);
return;
}
log.warn(`Unexpected ${key} read at epoch ${conn.epoch}; correlation is uncertain`);
this.requestRecycle(conn, `Unexpected ${key} read`);
}
private routeErr(
conn: Conn,
f: Extract<
TradeFrame,
{
k: 'err';
}
>,
): void {
const now = this.now();
this.stats.errFrames++;
if (FLAG_CODES.has(f.code)) this.flags.set(f.code, now);
const inc = f.incoming;
const requestType = inc.kind === 'request' ? inc.type : null;
for (const fn of this.listeners.error)
this.safe(() => fn({ code: f.code, message: f.message, requestType, at: now, epoch: conn.epoch }));
const msg = f.message ?? undefined;
if (conn.stage === 'auth') {
if (f.code === 'AlreadyAuthenticated' && this.o().authMode === 'headers') {
log.info('Uncorrelated rate-limit response; holding writes');
this.onAuthed(conn);
return;
}
this.authRejected(conn, `${f.code}${f.message ? `: ${f.message}` : ''}`, FLAG_CODES.has(f.code));
return;
}
if (inc.kind === 'request') {
if (f.code === 'RateLimited') noteRateLimited('all', f.message);
switch (inc.type) {
case 'add_order': {
const it = inc.cloid ? this.adds.get(inc.cloid) : undefined;
if (!it || it.epoch !== conn.epoch || it.done) break;
this.firstResponse(it);
if (this.errCannotProveReject(it, f.code)) {
this.note(
`err_unknown:${f.code}`,
'warn',
`Error ${f.code} cannot establish rejection of ${it.wire.cloid}; acknowledged=${it.acked}, observed fills=${it.fills.length}; outcome is unknown`,
);
this.finishAdd(it, {
kind: 'unknown',
why: 'error_reply',
...(it.orderId ? { orderId: it.orderId } : {}),
events: [...it.events],
fills: [...it.fills],
...(it.rateLimitedSuspect ? { rateLimitedSuspect: it.rateLimitedSuspect } : {}),
});
} else if (f.code === 'RateLimited') {
this.finishAdd(it, { kind: 'not_sent', why: 'rate_limited', detail: f.message ?? 'RateLimited' });
} else {
this.finishAdd(it, {
kind: 'rejected',
code: f.code,
...(msg ? { message: msg } : {}),
...(it.orderId ? { orderId: it.orderId } : {}),
afterAck: it.acked,
events: [...it.events],
});
}
return;
}
case 'cancel_order': {
const cw = inc.orderId ? this.cancels.get(inc.orderId.toLowerCase()) : undefined;
if (!cw || cw.epoch !== conn.epoch || cw.done) break;
if (f.code === 'RateLimited') this.finishCancel(cw, { kind: 'not_sent', why: 'rate_limited' });
else this.finishCancel(cw, { kind: 'rejected', status: f.code, ...(msg ? { message: msg } : {}) });
return;
}
case 'subscribe':
if (conn.stage === 'subscribe' && conn.subWaiting !== null) {
this.cancelTimer(conn.stageTimer);
conn.stageTimer = null;
conn.subWaiting = null;
if (f.code === 'RateLimited') {
this.subscribeNext(conn);
return;
}
this.dropConn(conn, `Uncorrelated permission error ${f.code}${f.message ? ` (${f.message})` : ''}`);
return;
}
break;
default: {
const key = READ_KEY_BY_REQUEST[inc.type];
if (key) {
const w = this.reads.get(key);
if (w && w.epoch === conn.epoch) {
this.reads.delete(key);
this.cancelTimer(w.timer);
w.reject(
new QfexSessionError(
f.code === 'RateLimited' ? 'rate_limited' : 'rejected',
`${inc.type}: ${f.code}${f.message ? ` (${f.message})` : ''}`,
f.code,
),
);
return;
}
this.stats.unsolicitedReads++;
if (conn.readsSent === 0) {
this.note(
`read_err_push:${key}`,
'warn',
`Unsolicited ${key} read error (${f.code}); waiting for active writes before recycle`,
);
return;
}
this.requestRecycle(conn, `Unsolicited ${key} read error (${f.code})`);
return;
}
}
}
this.note(
`err_req:${inc.type}:${f.code}`,
'warn',
`Error ${f.code} for unsupported request type ${inc.type}; no write rejected by inference`,
);
return;
}
if (f.code === 'RateLimited') {
const ms = noteRateLimited('all', f.message);
this.stats.rateLimitedUncorrelated++;
const suspect: QfexRateLimitedSuspect = { at: now, retryAfterUntil: now + ms, message: f.message };
for (const it of this.adds.values())
if (!it.done && !it.acked && it.fills.length === 0 && it.epoch === conn.epoch)
it.rateLimitedSuspect = suspect;
if (conn.stage === 'subscribe' && conn.subWaiting !== null) {
this.cancelTimer(conn.stageTimer);
conn.stageTimer = null;
conn.subWaiting = null;
this.subscribeNext(conn);
}
return;
}
this.note(
`err_bare:${f.code}`,
'warn',
`Uncorrelated error ${f.code} (${f.form})${f.message ? `: ${f.message.slice(0, 120)}` : ''}; no write rejected by inference`,
);
}
private finishAdd(it: AddIntent, r: SubmitResult): void {
if (it.done) return;
it.done = true;
this.cancelTimer(it.ackTimer);
this.cancelTimer(it.termTimer);
this.cancelTimer(it.graceTimer);
this.cancelTimer(it.settleTimer);
this.releaseSlot(it);
if (this.adds.get(it.wire.cloid) === it) this.adds.delete(it.wire.cloid);
if (it.orderId !== null && this.addsByOid.get(it.orderId.toLowerCase()) === it)
this.addsByOid.delete(it.orderId.toLowerCase());
it.resolve(r);
this.maybeRecycleNow();
}
private releaseSlot(it: AddIntent): void {
const s = it.slot;
it.slot = null;
if (s) s.release();
}
private validateWire(o: AddOrderWire):
| {
ok: true;
frame: string;
}
| {
ok: false;
why: string;
} {
if (!o || typeof o !== 'object') return { ok: false, why: 'Order input must be an object' };
if (typeof o.symbol !== 'string' || !SYMBOL_RE.test(o.symbol))
return { ok: false, why: `Invalid symbol: ${String(o.symbol).slice(0, 40)}` };
if (o.side !== 'BUY' && o.side !== 'SELL') return { ok: false, why: `Invalid side: ${String(o.side)}` };
if (o.tif !== 'GTC' && o.tif !== 'IOC') return { ok: false, why: `tif ${String(o.tif)}` };
if (typeof o.reduceOnly !== 'boolean') return { ok: false, why: 'reduceOnly must be boolean' };
if (typeof o.cloid !== 'string' || !CLOID_RE.test(o.cloid))
return { ok: false, why: 'Client order ID must contain exactly 32 lowercase hexadecimal characters' };
const q = canonLiteral(o.quantity);
const p = canonLiteral(o.price);
if (q === null)
return { ok: false, why: `Invalid positive quantity literal: ${String(o.quantity).slice(0, 40)}` };
if (p === null) return { ok: false, why: `Invalid positive price literal: ${String(o.price).slice(0, 40)}` };
const frame =
`{"type":"add_order","params":{"symbol":${JSON.stringify(o.symbol)},"side":"${o.side}","order_type":"LIMIT",` +
`"order_time_in_force":"${o.tif}","quantity":${q},"price":${p},"reduce_only":${o.reduceOnly ? 'true' : 'false'},` +
`"client_order_id":"${o.cloid}"}}`;
return { ok: true, frame };
}
private notReadyWhy(): QfexNotSentWhy {
return this.closed ? 'closed' : 'not_ready';
}
async submitOrder(
o: AddOrderWire,
waitTerminal: boolean,
ackTimeoutMs: number,
terminalTimeoutMs: number,
opts: QfexSubmitOptions = {},
): Promise<SubmitResult> {
const v = this.validateWire(o);
if (!v.ok) {
log.error(`Order was not sent: ${v.why}`);
return { kind: 'not_sent', why: 'invalid', detail: v.why };
}
if (this.closed) return { kind: 'not_sent', why: 'closed' };
if (this.o().dryRun()) {
this.stats.dryBlocked++;
return { kind: 'not_sent', why: 'dry_run' };
}
if (this.sentCloids.has(o.cloid)) {
log.error(`Client order ID ${o.cloid} was already sent; refusing reuse`);
return { kind: 'not_sent', why: 'duplicate_cloid' };
}
const pre = this.preWriteBlock();
if (pre) return pre;
const rh = rateHold('general');
if (rh) return { kind: 'not_sent', why: 'rate_limited', detail: rh };
if (this.sentCloids.size >= this.o().sentCloidsMax)
return {
kind: 'not_sent',
why: 'backpressure',
detail:
'Client order-ID memory is full; use a new session only with persisted ownership and unknown-intent state',
};
const slot = await acquireAddSlot(o.cloid, opts.slotWaitMs ?? this.o().replyTimeoutMs);
if (!slot) return { kind: 'not_sent', why: 'backpressure', detail: 'Another add order is awaiting correlation' };
const again =
this.preWriteBlock() ?? (this.closed ? { kind: 'not_sent' as const, why: 'closed' as const } : null);
const rh2 = again ? null : rateHold('general');
if (again || rh2) {
slot.release();
return again ?? { kind: 'not_sent', why: 'rate_limited', detail: rh2 ?? undefined };
}
if (this.sentCloids.has(o.cloid)) {
slot.release();
return { kind: 'not_sent', why: 'duplicate_cloid' };
}
const conn = this.readyConn();
if (!conn) {
slot.release();
return { kind: 'not_sent', why: this.notReadyWhy() };
}
this.rememberCloid(o.cloid);
const out = this.sendRaw(conn, v.frame, 'add_order', 'general', QFEX_WEIGHTS.addOrder);
if (out !== 'sent') {
slot.release();
return { kind: 'not_sent', why: notSentWhy(out), detail: out };
}
return new Promise<SubmitResult>((resolve) => {
const now = this.now();
const it: AddIntent = {
wire: { ...o },
epoch: conn.epoch,
sentAt: now,
waitTerminal,
settleMs: waitTerminal ? 0 : Math.max(0, opts.settleAfterAckMs ?? 0),
termMs: Math.max(ackTimeoutMs, terminalTimeoutMs),
orderId: null,
acked: false,
events: [],
fills: [],
fillIds: new Set(),
terminal: null,
slot,
ackTimer: null,
termTimer: null,
graceTimer: null,
settleTimer: null,
rateLimitedSuspect: null,
done: false,
resolve,
};
this.adds.set(o.cloid, it);
it.ackTimer = this.timer(ackTimeoutMs, () => {
if (it.done || it.acked) return;
if (it.waitTerminal && it.orderId !== null) return;
this.finishAdd(it, {
kind: 'unknown',
why: 'timeout',
...(it.orderId ? { orderId: it.orderId } : {}),
events: [...it.events],
fills: [...it.fills],
...(it.rateLimitedSuspect ? { rateLimitedSuspect: it.rateLimitedSuspect } : {}),
});
});
if (waitTerminal) {
it.termTimer = this.timer(it.termMs, () => {
if (it.done) return;
if (it.terminal) return this.finishAcked(it);
this.finishAdd(it, {
kind: 'unknown',
why: 'timeout',
...(it.orderId ? { orderId: it.orderId } : {}),
events: [...it.events],
fills: [...it.fills],
...(it.rateLimitedSuspect ? { rateLimitedSuspect: it.rateLimitedSuspect } : {}),
});
});
}
});
}
private rememberCloid(cloid: string): void {
this.sentCloids.add(cloid);
}
private preWriteBlock(): {
kind: 'not_sent';
why: QfexNotSentWhy;
detail?: string;
} | null {
if (this.closed) return { kind: 'not_sent', why: 'closed' };
const conn = this.readyConn();
if (!conn) return { kind: 'not_sent', why: 'not_ready', detail: `Session is not ready: ${this.st}` };
if (conn.recycle)
return { kind: 'not_sent', why: 'not_ready', detail: `Session is waiting to recycle: ${conn.recycle}` };
return null;
}
private finishCancel(cw: CancelWait, r: CancelResult): void {
if (cw.done) return;
cw.done = true;
this.cancelTimer(cw.timer);
if (this.cancels.get(cw.orderId.toLowerCase()) === cw) this.cancels.delete(cw.orderId.toLowerCase());
cw.resolve(r);
this.maybeRecycleNow();
}
async submitCancel(symbol: string, orderId: string, timeoutMs: number): Promise<CancelResult> {
if (
typeof symbol !== 'string' ||
!SYMBOL_RE.test(symbol) ||
typeof orderId !== 'string' ||
!UUID_RE.test(orderId)
) {
log.error('Cancel order ID must be a UUID');
return { kind: 'not_sent', why: 'invalid' };
}
if (this.closed) return { kind: 'not_sent', why: 'closed' };
if (this.o().dryRun()) {
this.stats.dryBlocked++;
return { kind: 'not_sent', why: 'dry_run' };
}
const conn = this.readyConn();
if (!conn || conn.recycle) return { kind: 'not_sent', why: 'not_ready' };
const oid = orderId.toLowerCase();
const existing = this.cancels.get(oid);
if (existing && !existing.done) return existing.promise;
if (rateHold('cancel')) return { kind: 'not_sent', why: 'rate_limited' };
const text = `{"type":"cancel_order","params":{"symbol":${JSON.stringify(symbol)},"order_id":"${orderId}","cancel_order_id_type":"order_id"}}`;
const out = this.sendRaw(conn, text, 'cancel_order', 'cancel', QFEX_WEIGHTS.cancelOrder);
if (out !== 'sent') return { kind: 'not_sent', why: notSentWhy(out) };
const now = this.now();
let resolve: (r: CancelResult) => void = () => undefined;
const promise = new Promise<CancelResult>((r) => {
resolve = r;
});
const cw: CancelWait = {
orderId,
symbol,
epoch: conn.epoch,
fillsSeen: false,
timer: null,
done: false,
promise,
resolve,
};
this.cancels.set(oid, cw);
cw.timer = this.timer(timeoutMs, () =>
this.finishCancel(cw, { kind: 'unknown', why: 'timeout', fillsSeen: cw.fillsSeen }),
);
return promise;
}
private async withReadLock<T>(key: ReadKey, fn: () => Promise<T>): Promise<T> {
const prev = this.readChains.get(key) ?? Promise.resolve();
let release: () => void = () => undefined;
const mine = new Promise<void>((r) => {
release = r;
});
const chain = prev.then(() => mine);
this.readChains.set(key, chain);
try {
await prev;
return await fn();
} finally {
release();
if (this.readChains.get(key) === chain) this.readChains.delete(key);
}
}
private readOnce(
key: ReadKey,
text: string,
type: string,
weight: number,
timeoutMs: number,
): Promise<{
epoch: number;
frame: TradeFrame;
}> {
if (this.closed) return Promise.reject(new QfexSessionError('closed', 'Session is closed'));
const conn = this.readyConn();
if (!conn) return Promise.reject(new QfexSessionError('not_ready', `Session is not ready: ${this.st}`));
if (conn.recycle)
return Promise.reject(
new QfexSessionError('correlation', `Read correlation is uncertain; session will recycle: ${conn.recycle}`),
);
if (this.reads.has(key))
return Promise.reject(new QfexSessionError('correlation', `A ${key} read is already pending`));
const out = this.sendRaw(conn, text, type, 'general', weight);
if (out !== 'sent') {
const kind: QfexSessionErrorKind =
out === 'rate_limited'
? 'rate_limited'
: out === 'dry_run'
? 'dry_run'
: out === 'not_open'
? 'not_ready'
: 'send_failed';
return Promise.reject(new QfexSessionError(kind, `${type} request was not sent: ${out}`));
}
conn.readsSent++;
return new Promise((resolve, reject) => {
const w: ReadWait = { key, epoch: conn.epoch, timer: null, resolve, reject };
this.reads.set(key, w);
w.timer = this.timer(timeoutMs, () => {
if (this.reads.get(key) !== w) return;
this.reads.delete(key);
this.stats.readTimeouts++;
reject(new QfexSessionError('timeout', `${type} reply timed out after ${timeoutMs} ms`));
this.dropConn(conn, `${type} reply timed out after ${timeoutMs} ms`);
});
});
}
async getUserOrders(limit: number, offset: number, timeoutMs: number): Promise<QfexOrdersRead> {
if (!Number.isSafeInteger(limit) || limit < 1 || !Number.isSafeInteger(offset) || offset < 0) {
throw new QfexSessionError('invalid', `Invalid order page limit=${limit}, offset=${offset}`);
}
return this.withReadLock('orders', async () => {
const { epoch, frame } = await this.readOnce(
'orders',
`{"type":"get_user_orders","params":{"limit":${limit},"offset":${offset}}}`,
'get_user_orders',
QFEX_WEIGHTS.getUserOrders,
timeoutMs,
);
if (frame.k !== 'orders') throw new QfexSessionError('malformed', `Expected order page; received ${frame.k}`);
if (frame.malformed.length > 0) {
throw new QfexSessionError('malformed', `Malformed order page: ${frame.malformed.slice(0, 3).join('; ')}`);
}
return { epoch, orders: frame.orders, twaps: frame.twaps, stopOrders: frame.stopOrders, malformed: [] };
});
}
async getUserLeverage(timeoutMs: number): Promise<Map<string, number>> {
return this.withReadLock('leverage', async () => {
const { frame } = await this.readOnce(
'leverage',
'{"type":"get_user_leverage","params":{"limit":1000,"offset":0}}',
'get_user_leverage',
0.1,
timeoutMs,
);
if (frame.k !== 'leverage')
throw new QfexSessionError('malformed', `Expected leverage response; received ${frame.k}`);
if (frame.malformed.length > 0)
throw new QfexSessionError(
'malformed',
`Malformed leverage response: ${frame.malformed.slice(0, 3).join('; ')}`,
);
const out = new Map<string, number>();
for (const r of frame.rows) {
const n = Number(r.leverage);
if (!Number.isFinite(n) || n <= 0)
throw new QfexSessionError('malformed', `Invalid leverage ${r.leverage} for ${r.symbol}`);
out.set(r.symbol, n);
}
return out;
});
}
async getAvailableLeverageLevels(timeoutMs: number): Promise<QfexLeverageRow[]> {
return this.withReadLock('levels', async () => {
const { frame } = await this.readOnce(
'levels',
'{"type":"get_available_leverage_levels","params":{}}',
'get_available_leverage_levels',
0.1,
timeoutMs,
);
if (frame.k !== 'leverage_levels')
throw new QfexSessionError('malformed', `Expected leverage levels; received ${frame.k}`);
if (frame.malformed.length > 0)
throw new QfexSessionError(
'malformed',
`Malformed leverage levels: ${frame.malformed.slice(0, 3).join('; ')}`,
);
return frame.rows;
});
}
onOrderEvent(fn: (o: QfexOrderEvent) => void): void {
this.listeners.order.push(fn);
}
onFill(fn: (f: QfexFill) => void): void {
this.listeners.fill.push(fn);
}
onDrop(fn: (epoch: number) => void): void {
this.listeners.drop.push(fn);
}
onReady(fn: (epoch: number) => void): void {
this.listeners.ready.push(fn);
}
onError(fn: (e: QfexSessionErrorEvent) => void): void {
this.listeners.error.push(fn);
}
writeBlock(): string | null {
if (this.closed) return 'Session is closed';
const c = this.readyConn();
if (!c) {
const af = this.authFailure ? ` — ${this.authFailure.text}` : '';
return `Session is not ready: ${this.st}${af}`;
}
if (c.recycle) return `Session will recycle: ${c.recycle}`;
return null;
}
cancelBlock(): string | null {
if (this.closed) return 'Session is closed';
const c = this.readyConn();
if (!c) {
const af = this.authFailure ? ` — ${this.authFailure.text}` : '';
return `Session is not ready: ${this.st}${af}`;
}
if (c.recycle) return `Session will recycle: ${c.recycle}`;
return null;
}
healthReasons(): string[] {
const r: string[] = [];
if (this.st === 'auth_failed') r.push('qfex_auth_failed');
else if (this.st !== 'idle' && this.st !== 'ready' && this.st !== 'closed') r.push('qfex_ws_down');
return r;
}
health(): QfexSessionHealth {
const c = this.conn;
const flags: Record<string, number> = {};
for (const [k, v] of this.flags) flags[k] = v;
return {
state: this.st,
since: this.stSince,
epoch: this.epoch,
reconnects: this.reconnects,
lastFrameAt: this.lastFrameAt,
lastPongAt: this.lastPongAt,
authFailure: this.authFailure ? { ...this.authFailure } : null,
subscribed: c ? [...c.subs] : [],
ready: this.readyConn() !== null,
backoffUntil: this.backoffUntil,
inflight: { adds: this.adds.size, cancels: this.cancels.size, reads: this.reads.size },
accountFlags: flags,
recycle: c?.recycle ?? null,
stats: { ...this.stats },
};
}
}
return new QfexSession(options);
}