Skip to content
markpaper

src/pending/memory.ts

v0.1.0 · 11.8 KB

Download file
import { isSignatureString } from '../address/base58.js';
import { orderKey } from '../ids/order.js';
import type { PhxWriteAnswer } from '../signer/decode.js';
export interface EchoPlace {
  symbol: string;
  side: 'bid' | 'ask';
  priceTicks: string;
  seq: string;
  lots: bigint;
  reduceOnly: boolean;
  slot: number;
  at: number;
}
export interface SlotVerdict {
  trusted: boolean;
  reason: string | null;
}
export interface EchoRow {
  symbol: string;
  side: 'bid' | 'ask';
  priceTicks: string;
  seq: string;
  lots: bigint;
  reduceOnly: boolean;
  echoed?: boolean;
}
export interface UnknownIntent {
  intentId: string;
  kind: 'ioc' | 'gtc';
  symbol: string;
  side: 'bid' | 'ask';
  placementKey: string;
  reduceOnly: boolean;
  signature: string | null;
  lastValidSlot: number | null;
  sentSlot: number;
  sentAt: number;
}
export interface WriteMemoryOptions {
  maxReadLagSlots: number;
  uncertainFillGateMs: number;
  echoTtlMs: number;
  iocSlotTtl: number;
  signWindowMs: number;
  slotMsCeil: number;
  noResponseLockSlots: number;
  noResponseLockMs: number;
  noResponseLockGtcMs: number;
  bandRejectTtlMs: number;
  bandAfterHoursMaxMs: number;
  now?: () => number;
}
/** All timing and lag policy is explicit and scoped to one account. Band memory is experimental. */
export function createWriteMemory(opts: WriteMemoryOptions) {
  for (const k of [
    'maxReadLagSlots',
    'uncertainFillGateMs',
    'echoTtlMs',
    'iocSlotTtl',
    'signWindowMs',
    'slotMsCeil',
    'noResponseLockSlots',
    'noResponseLockMs',
    'noResponseLockGtcMs',
    'bandRejectTtlMs',
    'bandAfterHoursMaxMs',
  ] as const)
    if (!Number.isSafeInteger(opts[k]) || opts[k] < 0) throw new TypeError('Invalid memory policy: ' + k);
  const clock = opts.now ?? Date.now;
  let lastStateSlot = 0;
  let lastStateReadStartedAt = 0;
  let lastFillSlot = 0;
  let lastChangeSlot = 0;
  let uncertainFill: {
    slot: number;
    until: number;
  } | null = null;
  const UNCERTAIN_FILL_GATE_MS = opts.uncertainFillGateMs;
  const placeEcho = new Map<string, EchoPlace>();
  const cancelEcho = new Map<
    string,
    {
      slot: number;
      at: number;
    }
  >();
  const bandRejected = new Map<
    string,
    {
      at: number;
      afterHours: boolean;
    }
  >();
  const BAND_AFTER_HOURS_MAX_MS = opts.bandAfterHoursMaxMs;
  const placementKeyFor = (symbol: string, side: 'bid' | 'ask', priceTicks: string): string =>
    `${symbol}|${side}|${priceTicks}`;
  function noteStateSlot(slot: number, startedAt = clock()): void {
    if (slot > lastStateSlot) lastStateSlot = slot;
    if (startedAt > lastStateReadStartedAt) lastStateReadStartedAt = startedAt;
  }
  function getLastStateSlot(): number {
    return lastStateSlot;
  }
  let lastOwnWriteAt = 0;
  function ownWriteWithin(withinMs: number, now = clock()): boolean {
    return withinMs > 0 && lastOwnWriteAt > 0 && now - lastOwnWriteAt <= withinMs;
  }
  function lastOwnWriteTime(): number {
    return lastOwnWriteAt;
  }
  function noteStateChangeSlot(slot: number, now = clock()): void {
    if (slot > lastChangeSlot) lastChangeSlot = slot;
    if (now > lastOwnWriteAt) lastOwnWriteAt = now;
  }
  function noteFillSlot(slot: number, now = clock()): void {
    noteStateChangeSlot(slot, now);
    if (slot > lastFillSlot) lastFillSlot = slot;
  }
  function noteUncertainFillSlot(slot: number, now = clock()): void {
    if (!uncertainFill || slot >= uncertainFill.slot) uncertainFill = { slot, until: now + UNCERTAIN_FILL_GATE_MS };
    if (now > lastOwnWriteAt) lastOwnWriteAt = now;
  }
  function judgeStateSlot(slot: number, now = clock()): SlotVerdict {
    if (slot < lastStateSlot)
      return { trusted: false, reason: `Snapshot slot ${slot} is older than accepted slot ${lastStateSlot}` };
    if (slot < lastFillSlot)
      return { trusted: false, reason: `Snapshot slot ${slot} precedes own fill slot ${lastFillSlot}` };
    if (uncertainFill && slot < uncertainFill.slot && now < uncertainFill.until) {
      return {
        trusted: false,
        reason: `Snapshot slot ${slot} precedes an unresolved confirmed write at slot ${uncertainFill.slot}`,
      };
    }
    if (lastChangeSlot - slot > opts.maxReadLagSlots) {
      return {
        trusted: false,
        reason: `Snapshot slot ${slot} lags own write slot ${lastChangeSlot} by more than ${opts.maxReadLagSlots} slots`,
      };
    }
    return { trusted: true, reason: null };
  }
  function notePlacedEcho(e: Omit<EchoPlace, 'at'>, now = clock()): void {
    noteStateChangeSlot(e.slot, now);
    placeEcho.set(orderKey(e.symbol, e.priceTicks, e.seq), { ...e, at: now });
  }
  function noteCancelledEcho(symbol: string, priceTicks: string, seq: string, slot: number, now = clock()): void {
    noteStateChangeSlot(slot, now);
    const key = orderKey(symbol, priceTicks, seq);
    placeEcho.delete(key);
    cancelEcho.set(key, { slot, at: now });
  }
  function applyEcho<T extends EchoRow>(rows: T[], stateSlot: number, now = clock()): Array<T | EchoRow> {
    const ttl = opts.echoTtlMs;
    for (const [k, c] of cancelEcho) if (c.slot <= stateSlot || now - c.at > ttl) cancelEcho.delete(k);
    for (const [k, p] of placeEcho) if (p.slot <= stateSlot || now - p.at > ttl) placeEcho.delete(k);
    const out: Array<T | EchoRow> = rows.filter((r) => !cancelEcho.has(orderKey(r.symbol, r.priceTicks, r.seq)));
    const present = new Set(rows.map((r) => orderKey(r.symbol, r.priceTicks, r.seq)));
    for (const [k, p] of placeEcho) {
      if (present.has(k)) continue;
      out.push({
        symbol: p.symbol,
        side: p.side,
        priceTicks: p.priceTicks,
        seq: p.seq,
        lots: p.lots,
        reduceOnly: p.reduceOnly,
        echoed: true,
      });
    }
    return out;
  }
  /** @experimental Band rejection memory is an opt-in retry policy; band-edge behavior is unverified. */
  function noteBandReject(key: string, now = clock(), afterHours = false): void {
    bandRejected.set(key, { at: now, afterHours });
  }
  /** @experimental This policy must not be treated as a venue guarantee. */
  function isBandRejected(key: string, now = clock(), afterHoursNow: boolean | null = null): boolean {
    const e = bandRejected.get(key);
    if (e === undefined) return false;
    if (now - e.at <= opts.bandRejectTtlMs) return true;
    if (e.afterHours && afterHoursNow === true && now - e.at <= BAND_AFTER_HOURS_MAX_MS) return true;
    bandRejected.delete(key);
    return false;
  }
  const unknowns = new Map<string, UnknownIntent>();
  const SIGN_WINDOW_MS = opts.signWindowMs;
  const SLOT_MS_CEIL = opts.slotMsCeil;
  const iocPhysicalEnd = (sentAt: number): number => sentAt + SIGN_WINDOW_MS + opts.iocSlotTtl * SLOT_MS_CEIL;
  function noteUnknown(u: UnknownIntent): void {
    if (u.signature !== null && !isSignatureString(u.signature))
      throw new TypeError('Invalid unknown-outcome signature');
    const previous = unknowns.get(u.intentId);
    if (previous) {
      if (JSON.stringify(previous) !== JSON.stringify(u))
        throw new TypeError('Unknown intentId reused with conflicting evidence');
      return;
    }
    unknowns.set(u.intentId, { ...u });
  }
  function iocLockOf(symbol: string): UnknownIntent | undefined {
    for (const u of unknowns.values()) if (u.kind === 'ioc' && u.symbol === symbol) return u;
    return undefined;
  }
  function placementLockOf(placementKey: string): UnknownIntent | undefined {
    for (const u of unknowns.values()) if (u.kind === 'gtc' && u.placementKey === placementKey) return u;
    return undefined;
  }
  function unknownCount(): number {
    return unknowns.size;
  }
  function applyConfirmedOrder(
    ctx: {
      symbol: string;
      side: 'bid' | 'ask';
      reduceOnly: boolean;
    },
    a: PhxWriteAnswer,
    now = clock(),
  ): void {
    if (a.status !== 'confirmed' || a.slot === null) return;
    if (!a.cpi) {
      noteUncertainFillSlot(a.slot, now);
      return;
    }
    if (a.cpi.filledBaseLots > 0n) noteFillSlot(a.slot, now);
    if (a.cpi.basePosted > 0n) {
      notePlacedEcho(
        {
          symbol: ctx.symbol,
          side: ctx.side,
          priceTicks: a.cpi.priceInTicks.toString(),
          seq: a.cpi.seq.toString(),
          lots: a.cpi.basePosted,
          reduceOnly: ctx.reduceOnly,
          slot: a.slot,
        },
        now,
      );
    }
  }
  async function resolveUnknowns(
    fetchTx: (sig: string, side: 'bid' | 'ask') => Promise<PhxWriteAnswer | null>,
    now = clock(),
  ): Promise<
    Array<{
      u: UnknownIntent;
      how: string;
    }>
  > {
    const released: Array<{
      u: UnknownIntent;
      how: string;
    }> = [];
    for (const u of [...unknowns.values()]) {
      if (u.signature) {
        let a: PhxWriteAnswer | null = null;
        try {
          a = await fetchTx(u.signature, u.side);
        } catch {
          a = null;
        }
        if (a && a.signature === u.signature && a.status === 'confirmed') {
          applyConfirmedOrder(u, a, now);
          unknowns.delete(u.intentId);
          const filled = a.cpi ? a.cpi.filledBaseLots.toString() : '?';
          released.push({
            u,
            how: `Confirmed at slot ${a.slot}: filled ${filled} lots, posted ${a.cpi ? a.cpi.basePosted.toString() : '?'} lots`,
          });
          continue;
        }
        if (a && a.signature === u.signature && a.status === 'failed') {
          unknowns.delete(u.intentId);
          released.push({ u, how: `Transaction failed: ${a.rejectReason ?? 'failed'}` });
          continue;
        }
      }
      if (u.kind === 'ioc' && u.lastValidSlot !== null) {
        if (lastStateSlot > u.lastValidSlot + 2) {
          unknowns.delete(u.intentId);
          released.push({
            u,
            how: `lastValidSlot ${u.lastValidSlot} expired; accepted snapshot slot ${lastStateSlot}`,
          });
        }
        // A known absolute lifetime is stronger evidence than the estimated no-response window.
        continue;
      }
      if (!u.signature || (u.kind === 'ioc' && u.lastValidSlot === null)) {
        const bySlot =
          u.kind === 'ioc' &&
          u.sentSlot > 0 &&
          lastStateSlot > u.sentSlot + opts.iocSlotTtl + opts.noResponseLockSlots &&
          lastStateReadStartedAt >= iocPhysicalEnd(u.sentAt);
        const windowEnd = u.sentAt + (u.kind === 'ioc' ? opts.noResponseLockMs : opts.noResponseLockGtcMs);
        const byTime = now >= windowEnd && lastStateReadStartedAt >= windowEnd;
        if (bySlot || byTime) {
          unknowns.delete(u.intentId);
          released.push({
            u,
            how: bySlot
              ? `Execution window closed: accepted snapshot slot ${lastStateSlot}, sent at slot ${u.sentSlot}`
              : `Execution window closed after ${Math.round((now - u.sentAt) / 1000)} seconds and a read started after the window`,
          });
        }
      }
      // A signed resting order has no proven expiry in this memory's evidence.
      // Elapsed wall time cannot establish that it was never posted; retain until /tx resolves it.
    }
    return released;
  }
  function counts(): {
    placeEcho: number;
    cancelEcho: number;
    band: number;
    unknown: number;
    lastStateSlot: number;
    lastFillSlot: number;
    lastChangeSlot: number;
  } {
    return {
      placeEcho: placeEcho.size,
      cancelEcho: cancelEcho.size,
      band: bandRejected.size,
      unknown: unknowns.size,
      lastStateSlot,
      lastFillSlot,
      lastChangeSlot,
    };
  }
  return {
    placementKeyFor,
    noteStateSlot,
    getLastStateSlot,
    ownWriteWithin,
    lastOwnWriteTime,
    noteStateChangeSlot,
    noteFillSlot,
    noteUncertainFillSlot,
    judgeStateSlot,
    notePlacedEcho,
    noteCancelledEcho,
    applyEcho,
    noteBandReject,
    isBandRejected,
    iocPhysicalEnd,
    noteUnknown,
    iocLockOf,
    placementLockOf,
    unknownCount,
    applyConfirmedOrder,
    resolveUnknowns,
    counts,
  };
}
export type WriteMemory = ReturnType<typeof createWriteMemory>;
All files