src/transport/cache.ts
v0.3.0 · 7.7 KB
// Weight-saving cache for heavy info reads.
//
// The IP budget is burned by concurrent bursts, not by total volume. Failure patterns this helper
// guards against:
// - `meta` stampede: getMeta() on every WS message, a fixed TTL, no in-flight dedup. At the TTL
// boundary N concurrent misses fire N x 20 weight at once.
// - Per-user caches created at one restart expire in the same second: N users x maxBuilderFee
// (weight 20 each) drain the bucket at once, and every caller waits behind it.
// - During an HL degradation a stale cache without a failure timestamp re-requests `meta`
// from the hot path on every call, each call with its own retries.
// - Event-driven invalidation (every own fill) bypasses a short position cache, and a periodic
// sweep then hits REST on every tick, feeding a 429 storm.
/** Options for {@link createCachedLoader}. */
export interface CachedLoaderOptions {
/** Time-to-live of a successful value, ms. */
ttlMs: number;
/**
* Relative TTL jitter applied once per write: 0.2 means `ttl * (0.8 .. 1.2)`. Default 0.2.
* Spreads the expiry of many per-key entries created at the same moment.
*/
ttlJitter?: number;
/**
* Minimum interval between two loads of the same key, ms, that `invalidate()` does NOT reset.
* Default 0. Use it for caches invalidated by events (fills, WS staleness).
*/
minRefreshIntervalMs?: number;
/**
* Return the stale value immediately and refresh in the background. Default false.
*
* Pair it with `failCooldownMs`: without a cooldown every call made while HL is degraded starts a
* new background refresh (each one retried by the info client), and a stale `meta` cache keeps
* hammering the API from the hot path.
*/
staleWhileRevalidate?: boolean;
/**
* After a failed load, keep serving the last good value (if any) for this long without new
* requests, ms. Default 0: a failure is rethrown and the next call retries.
*/
failCooldownMs?: number;
/** Never serve a value older than this, ms, even stale. Default: no limit. */
maxStaleMs?: number;
/** Clock in ms. Default `Date.now`. */
now?: () => number;
/** Random source in [0, 1). Default `Math.random`. */
random?: () => number;
/** Receives background refresh failures (stale-while-revalidate). */
onBackgroundError?: (err: unknown, key: string) => void;
}
/** Cached value with its age. */
export interface CachedValue<T> {
value: T;
loadedAt: number;
/** True while within TTL and not invalidated. */
fresh: boolean;
}
/** Cache handle returned by {@link createCachedLoader}. */
export interface CachedLoader<T, K extends string = string> {
/** Returns the cached value or loads it; concurrent calls for one key share a single request. */
get(key?: K): Promise<T>;
/** Current value without loading. */
peek(key?: K): CachedValue<T> | undefined;
/** Marks the value stale (subject to `minRefreshIntervalMs`). */
invalidate(key?: K): void;
/**
* Forgets the key entirely, including the refresh floor. An in-flight load for it will not write
* its result back - use this when the subject changes (another account on the same slot).
*/
delete(key?: K): void;
clear(): void;
/** True while a load for the key is in flight. */
isLoading(key?: K): boolean;
}
interface Entry<T> {
hasValue: boolean;
value: T | undefined;
loadedAt: number;
expiresAt: number;
invalidated: boolean;
lastFailAt: number | undefined;
inflight: Promise<T> | undefined;
}
/**
* Keyed TTL cache with single-flight, TTL jitter, a refresh floor that survives invalidation,
* optional stale-while-revalidate and a failure cooldown.
*
* Errors are never cached as values: "could not read" is not "empty". An empty array from HL is a
* valid value and is cached; a failure either rethrows or (with `failCooldownMs`) serves the last
* good value - check `peek(key).fresh` before authorizing anything that grows exposure.
*
* If a key includes a priority, include it in the key yourself (`wallet:high`): a high-priority
* read should not inherit the queue position of a background read for the same wallet.
*/
export function createCachedLoader<T, K extends string = string>(
load: (key: K) => Promise<T>,
options: CachedLoaderOptions,
): CachedLoader<T, K> {
if (!Number.isFinite(options.ttlMs) || options.ttlMs < 0) {
throw new RangeError(`ttlMs must be a non-negative number, got ${options.ttlMs}`);
}
const now = options.now ?? Date.now;
const random = options.random ?? Math.random;
const jitter = Math.min(1, Math.max(0, options.ttlJitter ?? 0.2));
const minRefresh = Math.max(0, options.minRefreshIntervalMs ?? 0);
const failCooldown = Math.max(0, options.failCooldownMs ?? 0);
const maxStale = options.maxStaleMs ?? Number.POSITIVE_INFINITY;
const entries = new Map<string, Entry<T>>();
const keyOf = (key: K | undefined): string => (key === undefined ? '' : key);
function entryFor(k: string): Entry<T> {
let e = entries.get(k);
if (!e) {
e = {
hasValue: false,
value: undefined,
loadedAt: 0,
expiresAt: 0,
invalidated: false,
lastFailAt: undefined,
inflight: undefined,
};
entries.set(k, e);
}
return e;
}
function jitteredTtl(): number {
return Math.round(options.ttlMs * (1 - jitter + random() * 2 * jitter));
}
function refresh(k: string, e: Entry<T>): Promise<T> {
if (e.inflight) return e.inflight;
// `.then` defers `load` to a microtask, so even a synchronously throwing loader cannot clear
// `inflight` before it is assigned below.
const p: Promise<T> = Promise.resolve()
.then(() => load(k as K))
.then(
(value) => {
if (entries.get(k) === e) {
const t = now();
e.hasValue = true;
e.value = value;
e.loadedAt = t;
e.expiresAt = t + jitteredTtl();
e.invalidated = false;
e.lastFailAt = undefined;
}
return value;
},
(err: unknown) => {
const t = now();
// A deleted/replaced entry (subject changed) must not receive this result.
if (entries.get(k) === e) {
e.lastFailAt = t;
if (failCooldown > 0 && e.hasValue && t - e.loadedAt <= maxStale) return e.value as T;
}
throw err;
},
)
.finally(() => {
if (e.inflight === p) e.inflight = undefined;
});
e.inflight = p;
return p;
}
function get(key?: K): Promise<T> {
const k = keyOf(key);
const e = entryFor(k);
if (e.hasValue) {
const t = now();
const age = t - e.loadedAt;
if (!e.invalidated && t < e.expiresAt) return Promise.resolve(e.value as T);
if (age <= maxStale) {
if (minRefresh > 0 && age < minRefresh) return Promise.resolve(e.value as T);
if (failCooldown > 0 && e.lastFailAt !== undefined && t - e.lastFailAt < failCooldown) {
return Promise.resolve(e.value as T);
}
if (options.staleWhileRevalidate) {
refresh(k, e).catch((err: unknown) => options.onBackgroundError?.(err, k));
return Promise.resolve(e.value as T);
}
}
}
return refresh(k, e);
}
function peek(key?: K): CachedValue<T> | undefined {
const e = entries.get(keyOf(key));
if (!e || !e.hasValue) return undefined;
return { value: e.value as T, loadedAt: e.loadedAt, fresh: !e.invalidated && now() < e.expiresAt };
}
return {
get,
peek,
invalidate(key?: K) {
const e = entries.get(keyOf(key));
if (e) e.invalidated = true;
},
delete(key?: K) {
entries.delete(keyOf(key));
},
clear() {
entries.clear();
},
isLoading(key?: K) {
return entries.get(keyOf(key))?.inflight !== undefined;
},
};
}