From 5b22bee0f49659a31394d7bf5a0d7ccf44f73a13 Mon Sep 17 00:00:00 2001 From: 4gray Date: Mon, 27 Jul 2026 10:27:16 +0200 Subject: [PATCH] feat(stalker): add session watchdog scheduler --- .../stalker-session/stalker-watchdog.spec.ts | 422 ++++++++++++++++++ .../stalker-session/stalker-watchdog.ts | 285 ++++++++++++ 2 files changed, 707 insertions(+) create mode 100644 apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.spec.ts create mode 100644 apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.ts diff --git a/apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.spec.ts b/apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.spec.ts new file mode 100644 index 000000000..faf567970 --- /dev/null +++ b/apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.spec.ts @@ -0,0 +1,422 @@ +import { + STALKER_WATCHDOG_INTERVALS, + StalkerWatchdog, + resolveStalkerWatchdogIntervalMs, +} from './stalker-watchdog'; +import type { + StalkerWatchdogActivation, + StalkerWatchdogClock, + StalkerWatchdogClockCallback, + StalkerWatchdogTarget, +} from './stalker-watchdog'; + +interface ScheduledTask { + readonly callback: StalkerWatchdogClockCallback; + readonly delayMs: number; + cancelled: boolean; +} + +class ManualWatchdogClock implements StalkerWatchdogClock { + readonly scheduledDelays: number[] = []; + readonly #tasks: ScheduledTask[] = []; + + get pendingCount(): number { + return this.#tasks.filter((task) => !task.cancelled).length; + } + + schedule( + callback: StalkerWatchdogClockCallback, + delayMs: number + ): () => void { + const task: ScheduledTask = { + callback, + cancelled: false, + delayMs, + }; + this.#tasks.push(task); + this.scheduledDelays.push(delayMs); + return () => { + task.cancelled = true; + }; + } + + async runNext(): Promise { + const index = this.#tasks.findIndex((task) => !task.cancelled); + if (index < 0) { + return false; + } + const [task] = this.#tasks.splice(index, 1); + task.cancelled = true; + await task.callback(); + return true; + } +} + +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, reject, resolve }; +} + +async function flushMicrotasks(): Promise { + for (let index = 0; index < 4; index += 1) { + await Promise.resolve(); + } +} + +function activation( + overrides: Partial = {} +): StalkerWatchdogActivation { + return { + leaseRef: 'lease-a', + principal: 'principal-a', + profile: { timeslot: 30 }, + sessionKey: 'session-a', + ...overrides, + }; +} + +function createHarness( + options: { + readonly isWireActive?: ( + target: Readonly + ) => boolean | Promise; + readonly jitter?: () => number; + readonly joinRefresh?: ( + target: Readonly + ) => Promise | undefined; + readonly ping?: ( + target: Readonly + ) => Promise; + } = {} +) { + const clock = new ManualWatchdogClock(); + const isWireActive = jest.fn( + options.isWireActive ?? + ((_target: Readonly) => true) + ); + const jitter = jest.fn(options.jitter ?? (() => 0)); + const joinRefresh = jest.fn( + options.joinRefresh ?? + ((_target: Readonly) => undefined) + ); + const ping = jest.fn( + options.ping ?? + (async (_target: Readonly) => undefined) + ); + const watchdog = new StalkerWatchdog({ + clock, + isWireActive, + jitter, + joinRefresh, + ping, + }); + return { + clock, + isWireActive, + jitter, + joinRefresh, + ping, + watchdog, + }; +} + +describe('resolveStalkerWatchdogIntervalMs', () => { + it('prefers a valid profile timeslot and accepts decimal strings', () => { + expect( + resolveStalkerWatchdogIntervalMs({ + timeslot: '42.5', + watchdog_timeout: '80', + }) + ).toBe(42_500); + }); + + it('falls back to watchdog_timeout when timeslot is invalid', () => { + expect( + resolveStalkerWatchdogIntervalMs({ + timeslot: 'not-a-number', + watchdog_timeout: '40', + }) + ).toBe(40_000); + }); + + it.each([ + undefined, + {}, + { timeslot: 0 }, + { timeslot: -2 }, + { timeslot: Number.NaN }, + { timeslot: Number.POSITIVE_INFINITY }, + { timeslot: '20 seconds' }, + ])('uses the conservative default for invalid profile %p', (profile) => { + expect(resolveStalkerWatchdogIntervalMs(profile)).toBe( + STALKER_WATCHDOG_INTERVALS.defaultMs + ); + }); + + it('clamps finite profile values to safe minimum and maximum delays', () => { + expect(resolveStalkerWatchdogIntervalMs({ timeslot: 0.5 })).toBe( + STALKER_WATCHDOG_INTERVALS.minimumMs + ); + expect( + resolveStalkerWatchdogIntervalMs({ + timeslot: Number.MAX_VALUE, + }) + ).toBe(STALKER_WATCHDOG_INTERVALS.maximumMs); + }); + + it('applies bounded deterministic jitter without escaping the clamps', () => { + expect(resolveStalkerWatchdogIntervalMs({ timeslot: 100 }, -1)).toBe( + 90_000 + ); + expect(resolveStalkerWatchdogIntervalMs({ timeslot: 100 }, 1)).toBe( + 110_000 + ); + expect(resolveStalkerWatchdogIntervalMs({ timeslot: 100 }, 99)).toBe( + 110_000 + ); + expect( + resolveStalkerWatchdogIntervalMs({ timeslot: 100 }, Number.NaN) + ).toBe(100_000); + expect(resolveStalkerWatchdogIntervalMs({ timeslot: 1 }, -1)).toBe( + STALKER_WATCHDOG_INTERVALS.minimumMs + ); + expect(resolveStalkerWatchdogIntervalMs({ timeslot: 1_000 }, 1)).toBe( + STALKER_WATCHDOG_INTERVALS.maximumMs + ); + }); +}); + +describe('StalkerWatchdog', () => { + it('starts a profile-derived timer and retains one watchdog for multiple leases', () => { + const { clock, watchdog } = createHarness(); + + watchdog.activate(activation()); + watchdog.activate( + activation({ + leaseRef: 'lease-b', + profile: { timeslot: 90 }, + }) + ); + watchdog.activate(activation()); + + expect(clock.pendingCount).toBe(1); + expect(clock.scheduledDelays).toEqual([30_000]); + + watchdog.deactivate(activation()); + expect(clock.pendingCount).toBe(1); + + watchdog.deactivate(activation({ leaseRef: 'lease-b' })); + expect(clock.pendingCount).toBe(0); + }); + + it('uses the injected jitter source independently for each scheduled tick', async () => { + const jitterValues = [-1, 1]; + const { clock, jitter, watchdog } = createHarness({ + jitter: () => jitterValues.shift() ?? 0, + }); + + watchdog.activate(activation({ profile: { timeslot: 100 } })); + expect(clock.scheduledDelays).toEqual([90_000]); + + await clock.runNext(); + + expect(clock.scheduledDelays).toEqual([90_000, 110_000]); + expect(jitter).toHaveBeenCalledTimes(2); + }); + + it('pings only the wire-active principal and keeps inactive sessions scheduled', async () => { + let wireActivePrincipal = 'principal-b'; + const { clock, ping, watchdog } = createHarness({ + isWireActive: (target) => target.principal === wireActivePrincipal, + }); + watchdog.activate(activation()); + watchdog.activate( + activation({ + leaseRef: 'lease-b', + principal: 'principal-b', + }) + ); + + await clock.runNext(); + await clock.runNext(); + + expect(ping).toHaveBeenCalledTimes(1); + expect(ping).toHaveBeenCalledWith({ + principal: 'principal-b', + sessionKey: 'session-a', + }); + expect(clock.pendingCount).toBe(2); + + wireActivePrincipal = 'principal-a'; + await clock.runNext(); + expect(ping).toHaveBeenLastCalledWith({ + principal: 'principal-a', + sessionKey: 'session-a', + }); + }); + + it('joins an in-flight refresh, rechecks wire ownership, and skips a stale ping', async () => { + const refresh = deferred(); + const ownership = [true, false]; + const { clock, isWireActive, joinRefresh, ping, watchdog } = + createHarness({ + isWireActive: () => ownership.shift() ?? false, + joinRefresh: () => refresh.promise, + }); + watchdog.activate(activation()); + + const tick = clock.runNext(); + await flushMicrotasks(); + + expect(joinRefresh).toHaveBeenCalledWith({ + principal: 'principal-a', + sessionKey: 'session-a', + }); + expect(ping).not.toHaveBeenCalled(); + + refresh.resolve(); + await tick; + + expect(isWireActive).toHaveBeenCalledTimes(2); + expect(ping).not.toHaveBeenCalled(); + expect(clock.pendingCount).toBe(1); + }); + + it('waits for an in-flight refresh before pinging the still-active principal', async () => { + const refresh = deferred(); + const { clock, ping, watchdog } = createHarness({ + joinRefresh: () => refresh.promise, + }); + watchdog.activate(activation()); + + const tick = clock.runNext(); + await flushMicrotasks(); + expect(ping).not.toHaveBeenCalled(); + + refresh.resolve(); + await tick; + + expect(ping).toHaveBeenCalledTimes(1); + expect(clock.pendingCount).toBe(1); + }); + + it('serializes pings and schedules the next tick only after completion', async () => { + const firstPing = deferred(); + const ping = jest + .fn, [Readonly]>() + .mockReturnValueOnce(firstPing.promise) + .mockResolvedValue(undefined); + const { clock, watchdog } = createHarness({ ping }); + watchdog.activate(activation()); + + const firstTick = clock.runNext(); + await flushMicrotasks(); + + expect(ping).toHaveBeenCalledTimes(1); + expect(clock.pendingCount).toBe(0); + await expect(clock.runNext()).resolves.toBe(false); + + firstPing.resolve(); + await firstTick; + expect(clock.pendingCount).toBe(1); + + await clock.runNext(); + expect(ping).toHaveBeenCalledTimes(2); + }); + + it('does not ping or reschedule after the last lease stops during refresh', async () => { + const refresh = deferred(); + const { clock, ping, watchdog } = createHarness({ + joinRefresh: () => refresh.promise, + }); + const lease = activation(); + watchdog.activate(lease); + + const tick = clock.runNext(); + await flushMicrotasks(); + watchdog.deactivate(lease); + refresh.resolve(); + await tick; + + expect(ping).not.toHaveBeenCalled(); + expect(clock.pendingCount).toBe(0); + }); + + it('does not ping when cleanup happens during the final wire-ownership check', async () => { + const finalOwnership = deferred(); + let ownershipChecks = 0; + const { clock, ping, watchdog } = createHarness({ + isWireActive: () => { + ownershipChecks += 1; + return ownershipChecks === 1 ? true : finalOwnership.promise; + }, + }); + const lease = activation(); + watchdog.activate(lease); + + const tick = clock.runNext(); + await flushMicrotasks(); + expect(ownershipChecks).toBe(2); + + watchdog.deactivate(lease); + finalOwnership.resolve(true); + await tick; + + expect(ping).not.toHaveBeenCalled(); + expect(clock.pendingCount).toBe(0); + }); + + it('cleans every principal for one session and can clean all sessions', async () => { + const { clock, ping, watchdog } = createHarness(); + watchdog.activate(activation()); + watchdog.activate( + activation({ + leaseRef: 'lease-b', + principal: 'principal-b', + }) + ); + watchdog.activate( + activation({ + leaseRef: 'lease-c', + principal: 'principal-c', + sessionKey: 'session-b', + }) + ); + + expect(clock.pendingCount).toBe(3); + + watchdog.cleanupSession('session-a'); + expect(clock.pendingCount).toBe(1); + + await clock.runNext(); + expect(ping).toHaveBeenCalledWith({ + principal: 'principal-c', + sessionKey: 'session-b', + }); + + watchdog.cleanupAll(); + expect(clock.pendingCount).toBe(0); + }); + + it('continues scheduling after a ping failure without overlapping work', async () => { + const { clock, ping, watchdog } = createHarness({ + ping: jest + .fn, [Readonly]>() + .mockRejectedValueOnce(new Error('classified-watchdog-failure')) + .mockResolvedValue(undefined), + }); + watchdog.activate(activation()); + + await clock.runNext(); + + expect(ping).toHaveBeenCalledTimes(1); + expect(clock.pendingCount).toBe(1); + await clock.runNext(); + expect(ping).toHaveBeenCalledTimes(2); + }); +}); diff --git a/apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.ts b/apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.ts new file mode 100644 index 000000000..4b4bf07ea --- /dev/null +++ b/apps/electron-backend/src/app/services/stalker-session/stalker-watchdog.ts @@ -0,0 +1,285 @@ +export const STALKER_WATCHDOG_INTERVALS = { + defaultMs: 25_000, + jitterRatio: 0.1, + maximumMs: 5 * 60_000, + minimumMs: 10_000, +} as const; + +const PROFILE_INTERVAL_FIELDS = ['timeslot', 'watchdog_timeout'] as const; +const DECIMAL_SECONDS = /^(?:\d+(?:\.\d*)?|\.\d+)$/; +const DEFAULT_WATCHDOG_JITTER = () => Math.random() * 2 - 1; + +export type StalkerWatchdogClockCallback = () => void | Promise; + +export interface StalkerWatchdogClock { + schedule( + callback: StalkerWatchdogClockCallback, + delayMs: number + ): () => void; +} + +export interface StalkerWatchdogTarget { + readonly principal: string; + readonly sessionKey: string; +} + +export interface StalkerWatchdogActivation extends StalkerWatchdogTarget { + readonly leaseRef: string; + readonly profile?: Readonly>; +} + +export type StalkerWatchdogDeactivation = Pick< + StalkerWatchdogActivation, + 'leaseRef' | 'principal' | 'sessionKey' +>; + +export interface StalkerWatchdogCallbacks { + readonly isWireActive: ( + target: Readonly + ) => boolean | Promise; + /** + * Returns the current refresh promise when one exists. The callback must + * never start a refresh: a watchdog tick only joins manager-owned work. + */ + readonly joinRefresh: ( + target: Readonly + ) => Promise | undefined; + readonly ping: (target: Readonly) => Promise; +} + +export interface StalkerWatchdogDependencies extends StalkerWatchdogCallbacks { + readonly clock?: StalkerWatchdogClock; + /** + * Returns a deterministic jitter unit in [-1, 1]. Values outside the + * range are clamped; non-finite values mean no jitter. + */ + readonly jitter?: () => number; +} + +interface StalkerWatchdogEntry { + readonly baseIntervalMs: number; + cancelScheduled?: () => void; + inFlight?: Promise; + readonly key: string; + readonly leaseRefs: Set; + readonly target: Readonly; +} + +const SYSTEM_WATCHDOG_CLOCK: StalkerWatchdogClock = { + schedule(callback, delayMs) { + const timer = setTimeout(() => { + void callback(); + }, delayMs); + return () => clearTimeout(timer); + }, +}; + +export function resolveStalkerWatchdogIntervalMs( + profile?: Readonly> | null, + jitterUnit = 0 +): number { + const baseIntervalMs = resolveBaseIntervalMs(profile); + const normalizedJitter = normalizeJitter(jitterUnit); + const jitteredIntervalMs = + baseIntervalMs * + (1 + STALKER_WATCHDOG_INTERVALS.jitterRatio * normalizedJitter); + return Math.round(clampIntervalMs(jitteredIntervalMs)); +} + +export class StalkerWatchdog { + readonly #callbacks: StalkerWatchdogCallbacks; + readonly #clock: StalkerWatchdogClock; + readonly #entries = new Map(); + readonly #jitter: () => number; + + constructor(dependencies: StalkerWatchdogDependencies) { + this.#callbacks = dependencies; + this.#clock = dependencies.clock ?? SYSTEM_WATCHDOG_CLOCK; + this.#jitter = dependencies.jitter ?? DEFAULT_WATCHDOG_JITTER; + } + + activate(activation: StalkerWatchdogActivation): void { + const key = createEntryKey(activation); + const existing = this.#entries.get(key); + if (existing) { + existing.leaseRefs.add(activation.leaseRef); + return; + } + + const entry: StalkerWatchdogEntry = { + baseIntervalMs: resolveStalkerWatchdogIntervalMs( + activation.profile + ), + key, + leaseRefs: new Set([activation.leaseRef]), + target: Object.freeze({ + principal: activation.principal, + sessionKey: activation.sessionKey, + }), + }; + this.#entries.set(key, entry); + this.#schedule(entry); + } + + deactivate(deactivation: StalkerWatchdogDeactivation): void { + const entry = this.#entries.get(createEntryKey(deactivation)); + if (!entry) { + return; + } + entry.leaseRefs.delete(deactivation.leaseRef); + if (entry.leaseRefs.size === 0) { + this.#stop(entry); + } + } + + cleanupSession(sessionKey: string): void { + for (const entry of this.#entries.values()) { + if (entry.target.sessionKey === sessionKey) { + this.#stop(entry); + } + } + } + + cleanupAll(): void { + for (const entry of this.#entries.values()) { + this.#stop(entry); + } + } + + #schedule(entry: StalkerWatchdogEntry): void { + if (!this.#isCurrent(entry) || entry.cancelScheduled) { + return; + } + const delayMs = applyJitter(entry.baseIntervalMs, this.#readJitter()); + entry.cancelScheduled = this.#clock.schedule(() => { + entry.cancelScheduled = undefined; + return this.#runSerializedTick(entry); + }, delayMs); + } + + async #runSerializedTick(entry: StalkerWatchdogEntry): Promise { + if (entry.inFlight) { + await entry.inFlight; + return; + } + + const operation = this.#runTick(entry); + entry.inFlight = operation; + try { + await operation; + } finally { + if (entry.inFlight === operation) { + entry.inFlight = undefined; + } + this.#schedule(entry); + } + } + + async #runTick(entry: StalkerWatchdogEntry): Promise { + try { + if (!(await this.#mayPing(entry))) { + return; + } + + const refresh = this.#callbacks.joinRefresh(entry.target); + if (refresh) { + await refresh; + } + + if (!(await this.#mayPing(entry))) { + return; + } + + await this.#callbacks.ping(entry.target); + } catch { + // The manager-owned callbacks classify/report failures. A failed + // tick remains non-fatal so a later interval can recover. + } + } + + async #mayPing(entry: StalkerWatchdogEntry): Promise { + if (!this.#isCurrent(entry)) { + return false; + } + const wireActive = await this.#callbacks.isWireActive(entry.target); + return wireActive && this.#isCurrent(entry); + } + + #stop(entry: StalkerWatchdogEntry): void { + if (this.#entries.get(entry.key) !== entry) { + return; + } + this.#entries.delete(entry.key); + entry.cancelScheduled?.(); + entry.cancelScheduled = undefined; + entry.leaseRefs.clear(); + } + + #isCurrent(entry: StalkerWatchdogEntry): boolean { + return ( + this.#entries.get(entry.key) === entry && entry.leaseRefs.size > 0 + ); + } + + #readJitter(): number { + try { + return this.#jitter(); + } catch { + return 0; + } + } +} + +function resolveBaseIntervalMs( + profile?: Readonly> | null +): number { + for (const field of PROFILE_INTERVAL_FIELDS) { + const seconds = parsePositiveSeconds(profile?.[field]); + if (seconds !== undefined) { + return clampIntervalMs(seconds * 1_000); + } + } + return STALKER_WATCHDOG_INTERVALS.defaultMs; +} + +function parsePositiveSeconds(value: unknown): number | undefined { + if (typeof value === 'number') { + return Number.isFinite(value) && value > 0 ? value : undefined; + } + if (typeof value !== 'string') { + return undefined; + } + const normalized = value.trim(); + if (!DECIMAL_SECONDS.test(normalized)) { + return undefined; + } + const seconds = Number(normalized); + return Number.isFinite(seconds) && seconds > 0 ? seconds : undefined; +} + +function normalizeJitter(value: number): number { + if (!Number.isFinite(value)) { + return 0; + } + return Math.max(-1, Math.min(1, value)); +} + +function applyJitter(baseIntervalMs: number, jitterUnit: number): number { + const jitteredIntervalMs = + baseIntervalMs * + (1 + + STALKER_WATCHDOG_INTERVALS.jitterRatio * + normalizeJitter(jitterUnit)); + return Math.round(clampIntervalMs(jitteredIntervalMs)); +} + +function clampIntervalMs(value: number): number { + return Math.max( + STALKER_WATCHDOG_INTERVALS.minimumMs, + Math.min(STALKER_WATCHDOG_INTERVALS.maximumMs, value) + ); +} + +function createEntryKey(target: StalkerWatchdogTarget): string { + return JSON.stringify([target.sessionKey, target.principal]); +}