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