src/account/index.ts
v0.1.0 · 10.9 KB
import { type QfexOrderEvent, statusClass } from '../frames/index.js';
import { type QfexMarket, qfexQuant } from '../markets/index.js';
import { type Dec, decCmp, decSign, decToString, parseDec, snapToGrid, snapToStep } from '../numbers/index.js';
import type { HistoricOutcome } from '../pending/index.js';
import type { RestClient } from '../rest/index.js';
import type { TradeSession } from '../session/index.js';
import {
decodeSubaccountEquity,
deriveEquity,
judgeEquity,
judgeOrderPages,
type QfexBalance,
type QfexEquityJudgeOpts,
type QfexOrdersPage,
} from './decoders.js';
export * from './decoders.js';
export class QfexReadError extends Error {
override name = 'QfexReadError';
}
const object = (value: unknown): value is Record<string, unknown> =>
!!value && typeof value === 'object' && !Array.isArray(value);
const required = (record: Record<string, unknown>, key: string): Dec => {
const value = parseDec(record[key]);
if (!value) throw new QfexReadError('Missing or invalid ' + key);
return value;
};
const optional = (record: Record<string, unknown>, key: string): Dec | null =>
record[key] === undefined || record[key] === null ? null : required(record, key);
export interface Position {
symbol: string;
quantity: Dec;
rawQuantity: Dec;
onGrid: boolean;
averagePrice: Dec | null;
unrealisedPnl: Dec | null;
leverage: Dec | null;
openOrders: number | null;
}
export function decodePositions(
payload: unknown,
markets: ReadonlyMap<string, Pick<QfexMarket, 'lot'>>,
): { positions: Map<string, Position>; balance: QfexBalance; listNull: boolean } {
if (
!object(payload) ||
!Object.hasOwn(payload, 'positions') ||
!object(payload.balance) ||
(payload.positions !== null && !Array.isArray(payload.positions))
)
throw new QfexReadError('Invalid positions response');
const balance: QfexBalance = {
available: required(payload.balance, 'available_balance'),
orderMargin: required(payload.balance, 'order_margin'),
positionMargin: required(payload.balance, 'position_margin'),
deposit: optional(payload.balance, 'deposit'),
realised: optional(payload.balance, 'realised_pnl'),
unrealised: optional(payload.balance, 'unrealised_pnl'),
netFunding: optional(payload.balance, 'net_funding'),
builderRewards: optional(payload.balance, 'builder_rewards'),
};
const positions = new Map<string, Position>();
for (const row of (payload.positions as unknown[] | null) ?? []) {
if (!object(row) || typeof row.symbol !== 'string' || row.symbol === '' || positions.has(row.symbol))
throw new QfexReadError('Malformed or duplicate position');
const rawQuantity = required(row, 'position');
const lot = parseDec(markets.get(row.symbol)?.lot);
const snapped = lot && decSign(lot) > 0 ? snapToGrid(rawQuantity, lot) : null;
const open = optional(row, 'open_orders');
const openOrders =
open && open.exp === 0 && decSign(open) >= 0 && open.mant <= BigInt(Number.MAX_SAFE_INTEGER)
? Number(open.mant)
: null;
positions.set(row.symbol, {
symbol: row.symbol,
quantity: snapped?.value ?? rawQuantity,
rawQuantity,
onGrid: snapped?.noise ?? false,
averagePrice: optional(row, 'average_price'),
unrealisedPnl: optional(row, 'unrealised_pnl'),
leverage: optional(row, 'leverage'),
openOrders,
});
}
if (positions.size === 0 && decSign(balance.positionMargin) > 0)
throw new QfexReadError('Empty positions contradict positive position margin');
return { positions, balance, listNull: payload.positions === null };
}
export function normalizeOpenOrders(
rows: readonly QfexOrderEvent[],
markets: ReadonlyMap<string, QfexMarket>,
isOwnOrder: (row: QfexOrderEvent) => boolean,
): { managed: QfexOrderEvent[]; foreign: QfexOrderEvent[]; invalid: QfexOrderEvent[] } {
const managed: QfexOrderEvent[] = [],
foreign: QfexOrderEvent[] = [],
invalid: QfexOrderEvent[] = [];
for (const row of rows) {
if (statusClass(row.status) !== 'live' && statusClass(row.status) !== 'fill') continue;
const market = markets.get(row.symbol);
const q = market ? qfexQuant(market) : null;
const remaining = parseDec(row.remaining);
if (!q || !remaining || decSign(remaining) <= 0 || !q.pxStrOnTick(row.price)) {
invalid.push(row);
continue;
}
(isOwnOrder(row) ? managed : foreign).push({ ...row });
}
return { managed, foreign, invalid };
}
/** Exact cloid equality is required: the exchange echoes it byte for byte. */
export function decodeHistoricByCloid(
payload: unknown,
cloid: string,
lotOf: (symbol: string) => string | undefined,
): HistoricOutcome | null {
if (!object(payload) || (payload.data !== null && !Array.isArray(payload.data)))
throw new QfexReadError('Invalid historic order response');
const matches = ((payload.data as unknown[] | null) ?? []).filter(
(row) => object(row) && row.client_order_id === cloid,
) as Record<string, unknown>[];
if (matches.length === 0) return null;
if (matches.length !== 1) throw new QfexReadError('Duplicate historic cloid');
const row = matches[0]!;
if (
typeof row.symbol !== 'string' ||
typeof row.terminal_status !== 'string' ||
(row.side !== 'BUY' && row.side !== 'SELL')
)
throw new QfexReadError('Historic order lacks symbol, side or status');
const filled = required(row, 'filled_qty');
if (decSign(filled) < 0) throw new QfexReadError('Negative historic fill');
const lot = parseDec(lotOf(row.symbol));
return {
symbol: row.symbol,
side: row.side,
terminalStatus: row.terminal_status,
filledQty: decToString(lot && decSign(lot) > 0 ? snapToStep(filled, lot, 'nearest') : filled),
avgPrice: typeof row.avg_price === 'string' ? row.avg_price : null,
orderId: typeof row.order_id === 'string' ? row.order_id : null,
};
}
export async function readCompleteOpenOrders(
session: Pick<TradeSession, 'getUserOrders' | 'health'>,
options: { pageLimit: number; maxPages: number; timeoutMs: number; allowMultiPage: boolean },
): Promise<QfexOrderEvent[]> {
for (const key of ['pageLimit', 'maxPages', 'timeoutMs'] as const)
if (!Number.isSafeInteger(options?.[key]) || options[key] <= 0)
throw new TypeError('Required order read policy: ' + key);
const epoch = session.health().epoch;
const pages: QfexOrdersPage[] = [];
for (let offset = 0; pages.length < options.maxPages; offset += options.pageLimit) {
const page = await session.getUserOrders(options.pageLimit, offset, options.timeoutMs);
if (page.epoch !== epoch || session.health().epoch !== epoch)
throw new QfexReadError('Order pages span socket epochs');
pages.push(page);
if (page.orders.length < options.pageLimit) break;
}
const judgment = judgeOrderPages({ pages, limit: options.pageLimit }, { allowMultiPage: options.allowMultiPage });
if (!judgment.ok || !judgment.complete) throw new QfexReadError(judgment.why ?? 'Incomplete order overview');
if (judgment.twaps || judgment.stopOrders)
throw new QfexReadError('Account includes unsupported conditional or TWAP orders');
return judgment.orders;
}
export function createAccountReader(options: {
rest: RestClient;
session: Pick<TradeSession, 'getUserOrders' | 'health'>;
accountId: string;
markets: () => ReadonlyMap<string, QfexMarket>;
equity: QfexEquityJudgeOpts;
pageLimit: number;
maxPages: number;
timeoutMs: number;
allowMultiPage: boolean;
now?: () => number;
sleep?: (ms: number) => Promise<void>;
}) {
if (!options?.rest || !options.session || !options.accountId || typeof options.markets !== 'function')
throw new TypeError('Required account sources');
const now = options.now ?? Date.now;
const sleep = options.sleep ?? ((ms: number) => new Promise<void>((resolve) => setTimeout(resolve, ms)));
async function positions() {
return decodePositions(await options.rest.userGet('/user/positions', options.accountId), options.markets());
}
async function orders() {
return readCompleteOpenOrders(options.session, options);
}
async function readAccountSnapshot() {
const startedAt = now();
const before = await positions();
const open = await orders();
const after = await positions();
if (
before.positions.size !== after.positions.size ||
[...before.positions].some(
([symbol, position]) =>
!after.positions.has(symbol) || decCmp(position.quantity, after.positions.get(symbol)!.quantity) !== 0,
)
)
throw new QfexReadError('Positions changed across snapshot');
const endpoint = decodeSubaccountEquity(
await options.rest.userGet('/user/subaccounts/equity', options.accountId),
options.accountId,
);
if (!endpoint.found) throw new QfexReadError('Requested account absent from equity response');
return {
startedAt,
receivedAt: now(),
positions: after.positions,
orders: open,
ordersComplete: true as const,
equity: judgeEquity(deriveEquity(after.balance), endpoint.equity, options.equity),
canonicalAccountId: endpoint.accountId,
consistency: 'positions_agree' as const,
};
}
/** REST has no fill watermark. Caller evidence is mandatory; agreeing reads alone prove no freshness. */
async function readFreshPosition(
symbol: string,
policy: {
timeoutMs: number;
pollMs: number;
evidence: (position: Position | null, readStartedAt: number) => boolean;
},
): Promise<{ kind: 'confirmed'; position: Position | null } | { kind: 'uncertain'; position: Position | null }> {
if (!(policy?.timeoutMs > 0) || !(policy.pollMs > 0) || typeof policy.evidence !== 'function')
throw new TypeError('Fresh position read requires explicit evidence and polling policy');
const deadline = now() + policy.timeoutMs;
let last: Position | null = null;
do {
const started = now();
try {
const read = await positions();
last = read.positions.get(symbol) ?? null;
if (policy.evidence(last, started)) return { kind: 'confirmed', position: last };
} catch {}
if (now() >= deadline) break;
await sleep(Math.min(policy.pollMs, deadline - now()));
} while (now() < deadline);
return { kind: 'uncertain', position: last };
}
return { readAccountSnapshot, readFreshPosition, readCompleteOpenOrders: orders, readPositions: positions };
}
export function mids(
contracts: readonly { tickerId: string; lastPrice: string | null; indexPrice: string | null }[],
bbo?: (symbol: string) => { bid: number; ask: number } | null,
): Record<string, string> {
const result: Record<string, string> = {};
for (const contract of contracts) {
const top = bbo?.(contract.tickerId);
const candidate =
top && Number.isFinite(top.bid) && Number.isFinite(top.ask) && top.bid > 0 && top.ask >= top.bid
? String(top.bid / 2 + top.ask / 2)
: (contract.indexPrice ?? contract.lastPrice);
const parsed = parseDec(candidate);
if (parsed && decSign(parsed) > 0) result[contract.tickerId] = decToString(parsed);
}
return result;
}