src/ws/selfHeal.ts
v0.3.0 · 5.5 KB
// 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 };
},
};
}