src/markets/cache.ts
v0.2.0 · 7.1 KB
// `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;
},
};
}