Skip to content
markpaper

src/account/index.ts

v0.1.0 · 10.9 KB

Download file
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;
}
All files