Skip to content
markpaper

src/ws/selfHeal.ts

v0.3.0 · 5.5 KB

Download file
// Self-healing of silently lost user subscriptions: freshness tracking with position-aware TTL and
// resubscribe backoff. Pure logic, driven by a timer the caller owns.

export interface StalenessTrackerOptions {
  /**
   * Stale threshold for addresses WITH positions and for addresses whose first snapshot has not
   * arrived yet. Default 90 000 ms.
   */
  staleMs?: number;
  /**
   * Stale threshold for addresses confirmed empty (snapshot seen, no positions). Default 15 min.
   * An empty account legitimately stays silent after its snapshot; a single 90 s TTL marks those
   * addresses "stale" every few minutes and causes endless resubscribe churn.
   */
  emptyStaleMs?: number;
  /** Stale ticks before the first resubscribe. Default 2. */
  resubAfterStaleTicks?: number;
  /** Cap of the required stale ticks (with 30 s ticks: ~10 min). Default 20. */
  maxResubBackoffTicks?: number;
  /** Resubscribes per tick at most: big subscribe bursts appear to get refused. Default 5. */
  maxResubsPerTick?: number;
}

export interface StalenessStats {
  tracked: number;
  /** Keys that sent anything in the last 90 s. `subscriptionCount` proves nothing; this does. */
  fresh90s: number;
  /** Keys that sent anything in the last 15 min. */
  fresh15m: number;
  /** Keys confirmed to have positions. */
  withPositions: number;
}

export interface StalenessTracker {
  /** Starts tracking; the grace period starts now (first snapshot arrives 0.5-2 s after subscribe). */
  track(key: string, now: number): void;
  /** Stops tracking and clears EVERY structure for the key. */
  untrack(key: string): void;
  /**
   * Records a message for `key`. Call it BEFORE dedupe/processing, otherwise repeated snapshots count
   * as silence. Pass `hasPositions` only after a successfully parsed snapshot. A fresh message resets
   * both the stale-tick and the resubscribe-attempt counters.
   */
  markSeen(key: string, now: number, hasPositions?: boolean): void;
  /**
   * One reconcile tick (run it from a ~30 s timer, never from the event path: dozens of calls per
   * millisecond mark addresses "stale" before their first snapshot and start a ghost-slot storm).
   * Returns the keys to resubscribe now - on the SAME connection. `isEligible(key)` can skip keys
   * whose connection is not healthy.
   *
   * A resubscribe does NOT refresh `lastSeen`: otherwise the backoff would reset every tick.
   */
  tick(now: number, isEligible?: (key: string) => boolean): string[];
  stats(now: number): StalenessStats;
}

interface KeyState { lastSeen: number; lastMessageAt: number | undefined; hasPositions: boolean | undefined; staleTicks: number; resubAttempts: number }

/**
 * Creates a staleness tracker. Required stale ticks before a resubscribe grow as
 * `min(resubAfterStaleTicks * 2^attempts, maxResubBackoffTicks)`; without backoff silent addresses are
 * resubscribed forever, and each resubscribe produces ghost slots that knock out live subscriptions.
 *
 * @example
 * const tracker = createStalenessTracker();
 * setInterval(() => {
 *   for (const user of tracker.tick(Date.now(), () => client.isHealthy)) {
 *     client.resubscribe({ type: 'allDexsClearinghouseState', user });
 *   }
 * }, 30_000).unref();
 */
export function createStalenessTracker(options: StalenessTrackerOptions = {}): StalenessTracker {
  const staleMs = options.staleMs ?? 90_000;
  const emptyStaleMs = options.emptyStaleMs ?? 15 * 60_000;
  const resubAfterStaleTicks = options.resubAfterStaleTicks ?? 2;
  const maxResubBackoffTicks = options.maxResubBackoffTicks ?? 20;
  const maxResubsPerTick = options.maxResubsPerTick ?? 5;
  const keys = new Map<string, KeyState>();

  return {
    track(key, now) {
      if (!keys.has(key)) keys.set(key, { lastSeen: now, lastMessageAt: undefined, hasPositions: undefined, staleTicks: 0, resubAttempts: 0 });
    },

    untrack(key) {
      keys.delete(key);
    },

    markSeen(key, now, hasPositions) {
      const state = keys.get(key);
      if (!state) return;
      state.lastSeen = Math.max(state.lastSeen, now);
      state.lastMessageAt = Math.max(state.lastMessageAt ?? now, now);
      if (hasPositions !== undefined) state.hasPositions = hasPositions;
      state.staleTicks = 0;
      state.resubAttempts = 0;
    },

    tick(now, isEligible) {
      const out: string[] = [];
      for (const [key, state] of keys) {
        if (isEligible && !isEligible(key)) continue;
        // hasPositions === undefined -> no snapshot yet -> short threshold.
        const threshold = state.hasPositions === false ? emptyStaleMs : staleMs;
        if (now - state.lastSeen <= threshold) {
          state.staleTicks = 0;
          state.resubAttempts = 0;
          continue;
        }
        state.staleTicks++;
        const required = Math.min(resubAfterStaleTicks * 2 ** state.resubAttempts, maxResubBackoffTicks);
        if (state.staleTicks >= required && out.length < maxResubsPerTick) {
          out.push(key);
          state.staleTicks = 0;
          state.resubAttempts++;
        }
      }
      return out;
    },

    stats(now) {
      let fresh90s = 0;
      let fresh15m = 0;
      let withPositions = 0;
      for (const state of keys.values()) {
        // Only real messages count: the grace period after track() is not evidence of a live subscription.
        if (state.lastMessageAt !== undefined && now - state.lastMessageAt <= 90_000) fresh90s++;
        if (state.lastMessageAt !== undefined && now - state.lastMessageAt <= 15 * 60_000) fresh15m++;
        if (state.hasPositions === true) withPositions++;
      }
      return { tracked: keys.size, fresh90s, fresh15m, withPositions };
    },
  };
}
All files