diff --git a/gitnexus/src/core/ingestion/workers/worker-pool.ts b/gitnexus/src/core/ingestion/workers/worker-pool.ts index 577bf8760..5823a5676 100644 --- a/gitnexus/src/core/ingestion/workers/worker-pool.ts +++ b/gitnexus/src/core/ingestion/workers/worker-pool.ts @@ -66,6 +66,13 @@ export interface WorkerPoolStats { /** Whether the circuit breaker has tripped (no further dispatches * will be accepted by this pool instance). */ readonly poolBroken: boolean; + /** Per-slot generation counter (U12). Increments by 1 on every + * successful worker replacement for that slot. Operators / tests + * observe this to confirm a death-then-respawn actually happened + * vs. the same worker being recycled in place. Initial value is 0 + * for every slot at pool creation; dropped slots keep their last + * generation (they don't decrement). */ + readonly slotGenerations: readonly number[]; } export interface WorkerPoolOptions { @@ -467,6 +474,16 @@ export const createWorkerPool = ( // failure streak instead of being masked by another slot's successes. // Reset to 0 on that slot's next successful job. const consecutiveFailuresPerSlot: number[] = new Array(size).fill(0); + // Per-slot generation counter (U12). Incremented on every successful + // worker replacement (see replaceWorker below). Handlers in the + // dispatch loop capture the slot's generation at attach time and + // short-circuit when they fire on a stale generation. Defensive layer + // on top of the existing `settled` flag + listener removal — protects + // against any future refactor that loosens cleanup() ordering or + // re-attaches handlers without resetting the per-job state. Exposed + // via getStats so operators (and tests) can verify a slot was + // actually replaced and not just the same worker recycled. + const slotGenerations: number[] = new Array(size).fill(0); let poolBroken = false; let poolFailure: Error | undefined; @@ -570,6 +587,12 @@ export const createWorkerPool = ( return false; } workers[workerIndex] = replacement; + // U12: bump the slot generation atomically with the worker swap so + // any late event from the OLD worker that somehow slipped past + // cleanup() carries a stale generation and short-circuits in the + // handler guard below. Increment AFTER `workers[workerIndex]` is + // updated so observers (getStats) see the new pair consistently. + slotGenerations[workerIndex]++; return true; }; @@ -1055,7 +1078,16 @@ export const createWorkerPool = ( }, job.timeoutMs); }; + // U12: capture the slot's generation at handler-attach time so any + // late event from a previous worker on this slot (which would carry + // an older generation) short-circuits below. Defensive — cleanup() + // already removes listeners synchronously when a death is observed, + // so under the current control flow no listener should fire on a + // stale generation. The guard catches future-refactor mistakes. + const slotGen = slotGenerations[workerIndex]; + const handler = (msg: WorkerOutgoingMessage) => { + if (slotGenerations[workerIndex] !== slotGen) return; if (settled || stopped) return; if (msg.type === 'starting-file') { inFlightPath = msg.path; @@ -1119,6 +1151,7 @@ export const createWorkerPool = ( }; const errorHandler = (err: Error) => { + if (slotGenerations[workerIndex] !== slotGen) return; if (!settled) { settled = true; cleanup(); @@ -1130,6 +1163,7 @@ export const createWorkerPool = ( }; const exitHandler = (code: number) => { + if (slotGenerations[workerIndex] !== slotGen) return; if (!settled) { settled = true; cleanup(); @@ -1154,6 +1188,7 @@ export const createWorkerPool = ( // and let the per-slot respawn budget and circuit breaker decide // whether to keep this slot in rotation. const messageErrorHandler = (err: Error) => { + if (slotGenerations[workerIndex] !== slotGen) return; if (!settled) { settled = true; cleanup(); @@ -1197,6 +1232,7 @@ export const createWorkerPool = ( droppedSlots: size - activeSlots.size, quarantined: quarantined.size, poolBroken, + slotGenerations: slotGenerations.slice(), }), }; }; diff --git a/gitnexus/test/unit/worker-pool-resilience.test.ts b/gitnexus/test/unit/worker-pool-resilience.test.ts index 0bb3b5492..b64137559 100644 --- a/gitnexus/test/unit/worker-pool-resilience.test.ts +++ b/gitnexus/test/unit/worker-pool-resilience.test.ts @@ -144,6 +144,10 @@ describe('worker pool resilience', () => { droppedSlots: 0, quarantined: 0, poolBroken: false, + // U12: every slot starts at generation 0; no respawns yet on a + // fresh pool. Per-slot zeros (not a single scalar) because each + // slot tracks its own respawn history independently. + slotGenerations: [0, 0, 0], }); void pool.terminate(); }); @@ -172,6 +176,9 @@ describe('worker pool resilience', () => { droppedSlots: 1, quarantined: 1, poolBroken: false, + // U12: slot 0 was dropped before any successful respawn (budget=0), + // so its generation stays at 0. Slot 1 never died, also 0. + slotGenerations: [0, 0], }); await pool.terminate(); }); diff --git a/gitnexus/test/unit/worker-pool-slot-generation.test.ts b/gitnexus/test/unit/worker-pool-slot-generation.test.ts new file mode 100644 index 000000000..b53991f23 --- /dev/null +++ b/gitnexus/test/unit/worker-pool-slot-generation.test.ts @@ -0,0 +1,203 @@ +/** + * U12 — Per-slot generation counter. + * + * worker-pool.ts now tracks a monotonic generation counter per slot, + * incremented on every successful worker replacement. The dispatch + * loop's handlers capture the slot's generation at attach time and + * short-circuit when they fire on a stale generation — defensive + * insurance against any future refactor that loosens cleanup() + * ordering or re-attaches handlers across the swap. + * + * In the current implementation, cleanup() synchronously removes + * listeners on a Worker instance the moment a death is observed, so + * no listener can naturally fire on a stale generation. The test + * surface is therefore the observable counter via `getStats()`: + * + * - Fresh pool: every slot starts at generation 0 + * - After a death + successful respawn: that slot's generation is 1 + * - After a death where the respawn budget is exhausted: that slot's + * generation stays at its last successful-respawn value (the + * drop-slot path does NOT bump generation, because no new worker + * came online for the slot) + */ +import { describe, it, expect, beforeEach, afterEach } from 'vitest'; +import { EventEmitter } from 'node:events'; +import path from 'node:path'; +import { pathToFileURL } from 'node:url'; +import fs from 'node:fs'; +import os from 'node:os'; +import { createWorkerPool } from '../../src/core/ingestion/workers/worker-pool.js'; + +type FakeAction = + | { kind: 'crash-after-starting'; startingPath: string; code: number } + | { kind: 'parse-ok'; files: { path: string }[] }; + +let nextActions: FakeAction[] = []; + +class FakeWorker extends EventEmitter { + constructor() { + super(); + queueMicrotask(() => { + this.emit('online'); + this.emit('message', { type: 'ready' }); + }); + } + postMessage(msg: unknown): void { + if (typeof msg !== 'object' || msg === null) return; + const m = msg as { type?: string }; + if (m.type !== 'sub-batch') return; + const action = nextActions.shift(); + if (!action) return; + queueMicrotask(() => this.run(action)); + } + private async run(action: FakeAction): Promise { + if (action.kind === 'crash-after-starting') { + this.emit('message', { type: 'starting-file', path: action.startingPath }); + this.emit('exit', action.code); + return; + } + if (action.kind === 'parse-ok') { + for (const f of action.files) { + this.emit('message', { type: 'starting-file', path: f.path }); + } + this.emit('message', { type: 'progress', filesProcessed: action.files.length }); + this.emit('message', { type: 'sub-batch-done' }); + await Promise.resolve(); + this.emit('message', { + type: 'result', + data: { fileCount: action.files.length, paths: action.files.map((f) => f.path) }, + }); + } + } + async terminate(): Promise { + this.emit('exit', 0); + return 0; + } +} + +let tempDir: string; +let workerUrl: URL; + +beforeEach(() => { + nextActions = []; + tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-slot-generation-')); + const workerPath = path.join(tempDir, 'fake-worker.js'); + fs.writeFileSync(workerPath, '// fake'); + workerUrl = pathToFileURL(workerPath) as URL; +}); + +afterEach(() => { + try { + fs.rmSync(tempDir, { recursive: true, force: true }); + } catch { + // best-effort cleanup + } +}); + +describe('worker pool slot-generation counter (U12)', () => { + it('starts every slot at generation 0 on a fresh pool', () => { + const pool = createWorkerPool(workerUrl, 4, { + workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker, + }); + try { + const stats = pool.getStats?.(); + expect(stats?.slotGenerations).toEqual([0, 0, 0, 0]); + } finally { + void pool.terminate(); + } + }); + + it('increments the slot generation exactly once on a successful respawn', async () => { + const pool = createWorkerPool(workerUrl, 1, { + workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker, + // Generous budgets so the replacement actually comes online (the + // happy path for the counter increment). + maxRespawnsPerSlot: 5, + consecutiveFailureThreshold: 10, + }); + + // Script: first dispatch crashes the worker; pool's replaceWorker + // creates a new FakeWorker (generation should bump to 1); the new + // worker handles the requeued remainder via parse-ok. + nextActions.push({ kind: 'crash-after-starting', startingPath: 'src/bad.ts', code: 134 }); + nextActions.push({ kind: 'parse-ok', files: [{ path: 'src/ok.ts' }] }); + + try { + await pool.dispatch<{ path: string; content: string }, unknown>([ + { path: 'src/bad.ts', content: '' }, + { path: 'src/ok.ts', content: '' }, + ]); + const stats = pool.getStats?.(); + expect(stats?.slotGenerations).toEqual([1]); + } finally { + await pool.terminate(); + } + }); + + it('leaves the slot generation unchanged when the respawn budget is exhausted (slot dropped, not replaced)', async () => { + const pool = createWorkerPool(workerUrl, 1, { + workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker, + // maxRespawnsPerSlot:0 means the first crash drops the slot without + // creating a replacement worker. No worker comes online for slot 0, + // so the generation MUST NOT bump (it's incremented in replaceWorker + // only AFTER a successful waitForWorkerReady). + maxRespawnsPerSlot: 0, + consecutiveFailureThreshold: 10, + }); + + nextActions.push({ kind: 'crash-after-starting', startingPath: 'src/bad.ts', code: 134 }); + + try { + // Dispatch rejects when all 1 slot is dropped — that's the + // breaker-tripped exhaustion path. The rejection is the EXPECTED + // outcome for this scenario; the load-bearing assertion is the + // post-rejection stats snapshot showing the generation did NOT + // bump (no successful respawn happened on the dropped slot). + await expect( + pool.dispatch<{ path: string; content: string }, unknown>([ + { path: 'src/bad.ts', content: '' }, + ]), + ).rejects.toBeDefined(); + const stats = pool.getStats?.(); + // Slot 0 was dropped before any successful respawn; generation + // stays at 0. droppedSlots == size confirms the slot is gone. + expect(stats?.slotGenerations).toEqual([0]); + expect(stats?.droppedSlots).toBe(1); + } finally { + await pool.terminate(); + } + }); + + it('increments each slot independently — one slot crashing does not affect another slot generation', async () => { + const pool = createWorkerPool(workerUrl, 2, { + workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker, + maxRespawnsPerSlot: 5, + consecutiveFailureThreshold: 10, + }); + + // Two files dispatched. The pool round-robins them across slots — + // exact assignment is implementation-detail, but on a 2-slot pool + // with 2 items the first item goes to one slot and the second to + // the other. We script BOTH possible orderings via a crash on the + // first action and parse-ok on the second; whichever slot got the + // bad file gets respawned (generation 1), the other stays at 0. + nextActions.push({ kind: 'crash-after-starting', startingPath: 'src/bad.ts', code: 134 }); + nextActions.push({ kind: 'parse-ok', files: [{ path: 'src/ok.ts' }] }); + nextActions.push({ kind: 'parse-ok', files: [{ path: 'src/bad.ts' }] }); + + try { + await pool.dispatch<{ path: string; content: string }, unknown>([ + { path: 'src/bad.ts', content: '' }, + { path: 'src/ok.ts', content: '' }, + ]); + const stats = pool.getStats?.(); + const gens = stats?.slotGenerations ?? []; + // Exactly one slot bumped to 1; the other stayed at 0. Sort to + // make the assertion order-independent across the round-robin + // assignment. + expect([...gens].sort()).toEqual([0, 1]); + } finally { + await pool.terminate(); + } + }); +});