import { logger } from "./logger.js"; import type { createPinboard } from "./server.js"; const DEFAULT_EJECT_AFTER_FAILURES = 5; const DEFAULT_BASE_EJECTION_MS = 30_000; const DEFAULT_MAX_EJECTION_WINDOWS = 10; const DEFAULT_MAX_EJECTION_RATIO = 0.5; export namespace createOutlierDetector { export type Deps = { /** Live worker count the ejection cap is computed against. */ workerCount(): Promise; now?: (() => number) | undefined; }; export type Instance = { /** SessionIds currently withheld from the placement candidate set. */ ejected(): string[]; /** Serve/wake outcome for a placed worker: ok resets its streak, a * failure counts toward ejection. */ report(sessionId: string, ok: boolean): Promise; }; } type Entry = { failures: number; ejections: number; until: number | null }; /** * Envoy-style outlier ejection: consecutive activation failures eject a * worker from the placement candidate set for a doubling cooldown, capped so * a fleet-wide failure never empties the set. */ export const createOutlierDetector = ( deps: createOutlierDetector.Deps, options: createPinboard.OutlierOptions, ): createOutlierDetector.Instance => { const now = deps.now ?? (() => Date.now()); const ejectAfterFailures = options.ejectAfterFailures ?? DEFAULT_EJECT_AFTER_FAILURES; const baseEjectionMs = options.baseEjectionMs ?? DEFAULT_BASE_EJECTION_MS; const maxEjectionMs = options.maxEjectionMs ?? baseEjectionMs * DEFAULT_MAX_EJECTION_WINDOWS; const maxEjectionRatio = options.maxEjectionRatio ?? DEFAULT_MAX_EJECTION_RATIO; if (!Number.isInteger(ejectAfterFailures) || ejectAfterFailures < 1) { throw new Error( `ejectAfterFailures must be a positive integer: ${String(options.ejectAfterFailures)}`, ); } if (!Number.isFinite(baseEjectionMs) || baseEjectionMs <= 0) { throw new Error( `baseEjectionMs must be a positive finite number: ${String(options.baseEjectionMs)}`, ); } if (!Number.isFinite(maxEjectionMs) || maxEjectionMs < baseEjectionMs) { throw new Error( `maxEjectionMs must be >= baseEjectionMs: ${String(options.maxEjectionMs)}`, ); } if (!(maxEjectionRatio > 0 && maxEjectionRatio <= 1)) { throw new Error( `maxEjectionRatio must be in (0, 1]: ${String(options.maxEjectionRatio)}`, ); } const entries = new Map(); const active = (entry: Entry): boolean => entry.until !== null && entry.until > now(); const eject = async (sessionId: string, entry: Entry): Promise => { let ejectedCount = 0; for (const other of entries.values()) { if (active(other)) ejectedCount += 1; } const total = await deps.workerCount(); const allowed = Math.min(Math.floor(total * maxEjectionRatio), total - 1); if (ejectedCount + 1 > allowed) return; entry.ejections += 1; entry.failures = 0; const windowMs = Math.min( baseEjectionMs * 2 ** (entry.ejections - 1), maxEjectionMs, ); entry.until = now() + windowMs; logger.warn( { operation: "outlier", sessionId, ejections: entry.ejections, windowMs }, "worker ejected from placement after consecutive activation failures", ); }; return { ejected: () => { const out: string[] = []; for (const [sessionId, entry] of entries) { if (active(entry)) out.push(sessionId); } return out; }, report: async (sessionId, ok) => { let entry = entries.get(sessionId); if (ok) { if (entry !== undefined && !active(entry)) entries.delete(sessionId); return; } if (entry === undefined) { entry = { failures: 0, ejections: 0, until: null }; entries.set(sessionId, entry); } if (entry.until !== null) { // Cooldown running: in-flight stragglers don't extend it; an expired // cooldown means this failure is the half-open probe — re-eject. if (!active(entry)) await eject(sessionId, entry); return; } entry.failures += 1; if (entry.failures >= ejectAfterFailures) await eject(sessionId, entry); }, }; };