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