Skip to content
markpaper

src/markets/cache.ts

v0.2.0 · 7.1 KB

Download file
// `symbols` cache: TTL, single-flight, stale-serve on failure, stability check across refreshes and an
// explicit freshness rule for `trading_status`.
//
// Rules:
//   - TTL 5 min; on a failed refresh serve STALE (product ids and lots essentially never change, and a
//     thrown meta error would abort every order, including closes) and retry no sooner than 30 s later;
//   - fresh = the last SUCCESSFUL refresh is younger than two TTLs (one failure is forgiven). Only a
//     fresh cache publishes `trading_status` as a fact (see `resolveMarketMode`);
//   - stability: a known coin must not vanish, change product id, lot or tick. `trading_status` changes
//     always pass. With `pinnedCoins` the check tolerates a dropped coin that has no positions/orders -
//     the strict rule freezes the cache after a public market rename (CIRCLE -> CRCL, 2026-08-10).

import { checkProductsStable, parseSymbols, type StabilityCheck } from './parse.js';
import type { PerpProduct, PerpUniverse } from './types.js';

/** Default TTL: how late a listing or a status change can show up. */
export const DEFAULT_SYMBOLS_TTL_MS = 5 * 60_000;
/** Default pause after a failed refresh before the gateway is asked again (stale data served meanwhile). */
export const DEFAULT_SYMBOLS_FAIL_COOLDOWN_MS = 30_000;
/** Freshness window as a multiple of the TTL: one failed refresh is forgiven, two are not. */
export const FRESH_TTL_MULTIPLIER = 2;

/** Diagnostic events of {@link createSymbolsCache}. */
export type SymbolsCacheEvent =
  | {
      readonly type: 'refreshed';
      readonly added: readonly string[];
      readonly dropped: readonly string[];
      readonly count: number;
    }
  | {
      readonly type: 'refresh-failed';
      readonly error: unknown;
      readonly servingStale: boolean;
      readonly ageMs: number | null;
    }
  | { readonly type: 'unstable'; readonly check: Extract<StabilityCheck, { ok: false }>; readonly consecutive: number };

/** Options of {@link createSymbolsCache}. */
export interface SymbolsCacheOptions {
  /** Loads the raw `data` of `{type:'symbols', product_type:'perp'}` (e.g. `() => client.query(...)`). */
  load: () => Promise<unknown>;
  /** Cache TTL. Default {@link DEFAULT_SYMBOLS_TTL_MS}. */
  ttlMs?: number;
  /** Cooldown after a failed refresh. Default {@link DEFAULT_SYMBOLS_FAIL_COOLDOWN_MS}. */
  failCooldownMs?: number;
  /**
   * Coins that must survive a refresh (positions, resting orders, traded markets). Evaluated on every
   * refresh so it can be a live view. Omitted = strict mode (any dropped coin rejects the refresh).
   */
  pinnedCoins?: () => Iterable<string>;
  /** Monotonic clock in ms. Default `performance.now()`. */
  now?: () => number;
  /** Diagnostic events; exceptions from the handler are swallowed. */
  onEvent?: (event: SymbolsCacheEvent) => void;
}

/** Health of the cache for logs and `/health`. */
export interface SymbolsCacheStatus {
  readonly loaded: boolean;
  /** Age of the last successful refresh, ms. */
  readonly ageMs: number | null;
  /** Younger than {@link FRESH_TTL_MULTIPLIER} TTLs. */
  readonly fresh: boolean;
  readonly count: number;
  /** Consecutive refreshes rejected by the stability check; >= 3 means the cache is effectively frozen. */
  readonly consecutiveUnstable: number;
  readonly lastError: unknown;
}

export interface SymbolsCache {
  /** Universe from cache (refreshes when the TTL expired; serves stale on failure once loaded). */
  get(): Promise<PerpUniverse>;
  /** Product by coin name (`'BTC'`), or `undefined` when not listed in the current universe. */
  resolve(coin: string): Promise<PerpProduct | undefined>;
  /** Cached universe without triggering a refresh. */
  peek(): PerpUniverse | null;
  /** True while the last successful refresh is younger than two TTLs. */
  isFresh(): boolean;
  status(): SymbolsCacheStatus;
  /** Forces a refresh on the next `get()`. */
  invalidate(): void;
}

/** Thrown when a refresh is rejected by the stability check. */
export class SymbolsUnstableError extends Error {
  override readonly name = 'SymbolsUnstableError';
  readonly check: Extract<StabilityCheck, { ok: false }>;
  constructor(check: Extract<StabilityCheck, { ok: false }>) {
    super(check.message);
    this.check = check;
  }
}

/** Creates a `symbols` cache (see the module header for the rules). */
export function createSymbolsCache(options: SymbolsCacheOptions): SymbolsCache {
  const ttlMs = options.ttlMs ?? DEFAULT_SYMBOLS_TTL_MS;
  const cooldownMs = options.failCooldownMs ?? DEFAULT_SYMBOLS_FAIL_COOLDOWN_MS;
  if (!(ttlMs > 0) || !(cooldownMs >= 0))
    throw new RangeError('ttlMs must be positive and failCooldownMs non-negative');
  const now = options.now ?? (() => performance.now());
  const emit = (event: SymbolsCacheEvent): void => {
    try {
      options.onEvent?.(event);
    } catch {
      // Diagnostics must never break a refresh.
    }
  };

  let cache: { ts: number; universe: PerpUniverse } | null = null;
  let inflight: Promise<PerpUniverse> | null = null;
  let retryAfter = Number.NEGATIVE_INFINITY;
  let forceRefresh = false;
  let consecutiveUnstable = 0;
  let lastError: unknown;

  async function build(): Promise<PerpUniverse> {
    const raw = await options.load();
    const universe = parseSymbols(raw);
    const check = checkProductsStable(cache?.universe ?? null, universe, {
      pinnedCoins: options.pinnedCoins?.(),
    });
    if (!check.ok) {
      consecutiveUnstable++;
      emit({ type: 'unstable', check, consecutive: consecutiveUnstable });
      throw new SymbolsUnstableError(check);
    }
    consecutiveUnstable = 0;
    emit({ type: 'refreshed', added: check.added, dropped: check.dropped, count: universe.byCoin.size });
    return universe;
  }

  function get(): Promise<PerpUniverse> {
    const t = now();
    if (cache && !forceRefresh && (t - cache.ts < ttlMs || t < retryAfter)) return Promise.resolve(cache.universe);
    if (inflight) return inflight;
    inflight = build()
      .then((universe) => {
        cache = { ts: now(), universe };
        retryAfter = Number.NEGATIVE_INFINITY;
        forceRefresh = false;
        lastError = undefined;
        return universe;
      })
      .catch((err: unknown) => {
        lastError = err;
        forceRefresh = false;
        if (cache) {
          retryAfter = now() + cooldownMs;
          emit({ type: 'refresh-failed', error: err, servingStale: true, ageMs: now() - cache.ts });
          return cache.universe;
        }
        emit({ type: 'refresh-failed', error: err, servingStale: false, ageMs: null });
        throw err;
      })
      .finally(() => {
        inflight = null;
      });
    return inflight;
  }

  function isFresh(): boolean {
    return cache !== null && now() - cache.ts < FRESH_TTL_MULTIPLIER * ttlMs;
  }

  return {
    get,
    resolve: async (coin) => (await get()).byCoin.get(coin),
    peek: () => cache?.universe ?? null,
    isFresh,
    status: () => ({
      loaded: cache !== null,
      ageMs: cache ? now() - cache.ts : null,
      fresh: isFresh(),
      count: cache?.universe.byCoin.size ?? 0,
      consecutiveUnstable,
      lastError,
    }),
    invalidate: () => {
      retryAfter = Number.NEGATIVE_INFINITY;
      forceRefresh = true;
    },
  };
}
All files