Skip to content
markpaper

src/transport/cache.ts

v0.3.0 · 7.7 KB

Download file
// 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;
    },
  };
}
All files