Skip to content
markpaper

src/markets/cache.ts

v0.2.1 · 5.7 KB

Download file
// Meta cache with TTL, single-flight refresh, stale-on-error and a POINTWISE sanity check
// (knowledge base: instances-and-api.md §2.2, "rules around the meta cache").
//
// Why stale-on-error: a meta failure aborts every order, closes included. Serving the old map with a
// loud event and a short retry backoff is safer than failing.
// Why the sanity check is pointwise: a rule "no previously seen coin may disappear" freezes a cache
// after a market is renamed (every refresh is rejected, new listings stay invisible). Here a refresh is
// rejected only if a market from the WORKING SET (positions or orders on it) disappeared or changed
// market_id / decimals.

import { type ParseOrderBookDetailsOptions, parseOrderBookDetails } from './parse.js';
import { type LighterMarket, LighterMetaError } from './types.js';

/** Default cache TTL: 5 minutes. */
export const DEFAULT_META_TTL_MS = 5 * 60_000;
/** After a failed refresh, do not retry sooner than this while serving the stale map. */
export const DEFAULT_META_FAILED_BACKOFF_MS = 30_000;

export type MarketCacheEvent =
  | { type: 'refreshed'; markets: number }
  | { type: 'stale'; ageMs: number; error: string }
  | { type: 'rejected'; error: string };

export interface MarketCacheOptions {
  /** Transport: returns the raw `orderBookDetails` JSON (any client, any fetch; stub it in tests). */
  fetchDetails: () => Promise<unknown>;
  ttlMs?: number;
  failedBackoffMs?: number;
  parse?: ParseOrderBookDetailsOptions;
  /** Clock, for tests. */
  now?: () => number;
  /** Refresh / stale / rejected notifications. Exceptions thrown by the callback are swallowed. */
  onEvent?: (event: MarketCacheEvent) => void;
}

export interface MarketCache {
  /** Cached map, refreshed when older than `ttlMs`; throws only when there is no map at all. */
  get(): Promise<Map<string, LighterMarket>>;
  /** One market by exact symbol. */
  market(symbol: string): Promise<LighterMarket | undefined>;
  /** Marks a symbol as "we have a position or an order here": its disappearance rejects a refresh. */
  noteInterest(symbol: string): void;
  forgetInterest(symbol: string): void;
  workingSet(): string[];
  /** Current map and its age without triggering a refresh; `null` before the first successful load. */
  peek(): { markets: Map<string, LighterMarket>; ageMs: number } | null;
  /** Forces the next `get()` to refresh. */
  invalidate(): void;
}

/**
 * Pointwise sanity check between two decoded metas: for every symbol in `workingSet` that existed
 * before, it must still exist with the same `market_id` and decimals.
 *
 * @throws LighterMetaError describing the first violation.
 */
export function assertMarketsSane(
  previous: Map<string, LighterMarket> | null,
  next: Map<string, LighterMarket>,
  workingSet: Iterable<string>,
): void {
  if (!previous) return;
  for (const symbol of workingSet) {
    const before = previous.get(symbol);
    if (!before) continue;
    const after = next.get(symbol);
    if (!after) throw new LighterMetaError(`meta refresh lost ${symbol}, which is in the working set`);
    if (after.marketId !== before.marketId) {
      throw new LighterMetaError(`market_id of ${symbol} changed (${before.marketId} -> ${after.marketId})`);
    }
    if (after.sizeDecimals !== before.sizeDecimals || after.priceDecimals !== before.priceDecimals) {
      throw new LighterMetaError(`size/price decimals of ${symbol} changed under resting orders`);
    }
  }
}

export function createMarketCache(opts: MarketCacheOptions): MarketCache {
  const ttlMs = opts.ttlMs ?? DEFAULT_META_TTL_MS;
  const backoffMs = opts.failedBackoffMs ?? DEFAULT_META_FAILED_BACKOFF_MS;
  const now = opts.now ?? Date.now;
  const working = new Set<string>();
  let cache: { ts: number; map: Map<string, LighterMarket> } | null = null;
  let inflight: Promise<Map<string, LighterMarket>> | null = null;
  let retryAfter = 0;
  let forceRefresh = false;

  const emit = (e: MarketCacheEvent) => {
    try {
      opts.onEvent?.(e);
    } catch {
      // telemetry must never break the read path
    }
  };

  async function refresh(): Promise<Map<string, LighterMarket>> {
    try {
      const raw = await opts.fetchDetails();
      const map = parseOrderBookDetails(raw, opts.parse);
      assertMarketsSane(cache?.map ?? null, map, working);
      cache = { ts: now(), map };
      retryAfter = 0;
      emit({ type: 'refreshed', markets: map.size });
      return map;
    } catch (err) {
      const message = err instanceof Error ? err.message : String(err);
      if (cache) {
        retryAfter = now() + backoffMs;
        // A LighterMetaError is a rejected payload (corrupt or failing the working-set check);
        // anything else is transport. Both serve the old map, the event tells them apart.
        emit(
          err instanceof LighterMetaError
            ? { type: 'rejected', error: message }
            : { type: 'stale', ageMs: now() - cache.ts, error: message },
        );
        return cache.map;
      }
      throw err;
    }
  }

  return {
    get() {
      const t = now();
      if (cache && !forceRefresh && (t - cache.ts < ttlMs || t < retryAfter)) return Promise.resolve(cache.map);
      if (inflight) return inflight;
      forceRefresh = false;
      inflight = refresh().finally(() => {
        inflight = null;
      });
      return inflight;
    },
    async market(symbol) {
      return (await this.get()).get(symbol);
    },
    noteInterest(symbol) {
      working.add(symbol);
    },
    forgetInterest(symbol) {
      working.delete(symbol);
    },
    workingSet() {
      return [...working];
    },
    peek() {
      return cache ? { markets: cache.map, ageMs: now() - cache.ts } : null;
    },
    invalidate() {
      forceRefresh = true;
      retryAfter = 0;
    },
  };
}
All files