src/experimental/ledger.ts
v0.1.0 · 4.3 KB
import type { QfexFill } from '../frames/index.js';
import { type Dec, decAdd, decCmp, decNeg, decSign, decToString, parseDec } from '../numbers/index.js';
const ZERO: Dec = { mant: 0n, exp: 0 };
/**
* @experimental REST has no watermark and a fill stream can miss external activity.
* An agreeing value is evidence of consistency, never proof of exchange freshness.
* This module is intentionally absent from the root barrel and account reader.
*/
export function createFillLedger(options: {
now?: () => number;
dedupTtlMs: number;
maxTradeIds: number;
maxHistory: number;
gapSettleMs: number;
}) {
for (const field of ['dedupTtlMs', 'maxTradeIds', 'maxHistory', 'gapSettleMs'] as const)
if (!Number.isFinite(options?.[field]) || options[field] <= 0)
throw new TypeError('Required ledger policy: ' + field);
const now = options.now ?? Date.now;
const seen = new Map<string, number>();
const states = new Map<string, { anchor: Dec | null; expected: Dec | null; history: Dec[]; incomplete: boolean }>();
let sequence = 0;
let gapAt: number | null = null;
function state(symbol: string) {
let value = states.get(symbol);
if (!value) {
value = { anchor: null, expected: null, history: [], incomplete: false };
states.set(symbol, value);
}
return value;
}
function noteFill(fill: QfexFill): boolean {
const at = now();
for (const [id, time] of seen) if (at - time >= options.dedupTtlMs) seen.delete(id);
if (seen.has(fill.tradeId)) return false;
if (seen.size >= options.maxTradeIds) {
for (const value of states.values()) value.incomplete = true;
return false;
}
const quantity = parseDec(fill.qty);
if (!quantity || decSign(quantity) <= 0) throw new TypeError('Invalid fill');
seen.set(fill.tradeId, at);
sequence++;
const selected = state(fill.symbol);
if (selected.expected) {
const previous = selected.expected;
selected.expected = decAdd(previous, fill.side === 'BUY' ? quantity : decNeg(quantity));
selected.history.push(previous);
if (selected.history.length > options.maxHistory) {
selected.incomplete = true;
selected.history.shift();
}
} else selected.incomplete = true;
return true;
}
function mark() {
return { sequence, at: now() };
}
function noteGap() {
gapAt = now();
for (const value of states.values()) value.incomplete = true;
}
function judge(
positions: ReadonlyMap<string, Dec>,
readMark: ReturnType<typeof mark>,
): { kind: 'consistent' | 'lagging' | 'uncertain'; symbols: string[] } {
if (sequence !== readMark.sequence || gapAt !== null) return { kind: 'uncertain', symbols: [...states.keys()] };
const lagging: string[] = [],
uncertain: string[] = [];
for (const symbol of new Set([...states.keys(), ...positions.keys()])) {
const selected = state(symbol),
actual = positions.get(symbol) ?? ZERO;
if (selected.incomplete || !selected.expected) {
uncertain.push(symbol);
continue;
}
if (decCmp(actual, selected.expected) !== 0) {
if (selected.history.some((value) => decCmp(value, actual) === 0)) lagging.push(symbol);
else uncertain.push(symbol);
}
}
return uncertain.length
? { kind: 'uncertain', symbols: uncertain }
: lagging.length
? { kind: 'lagging', symbols: lagging }
: { kind: 'consistent', symbols: [] };
}
/** Caller explicitly establishes a baseline after independently validating account state. */
function anchor(positions: ReadonlyMap<string, Dec>, readMark: ReturnType<typeof mark>): boolean {
if (sequence !== readMark.sequence || (gapAt !== null && now() - gapAt < options.gapSettleMs)) return false;
states.clear();
for (const [symbol, quantity] of positions)
states.set(symbol, { anchor: quantity, expected: quantity, history: [], incomplete: false });
gapAt = null;
return true;
}
return {
noteFill,
noteGap,
mark,
judge,
anchor,
expected: (symbol: string) => states.get(symbol)?.expected ?? null,
report: () => ({
sequence,
gapAt,
symbols: [...states].map(([symbol, value]) => ({
symbol,
expected: value.expected ? decToString(value.expected) : null,
incomplete: value.incomplete,
})),
}),
};
}