feat(workers): per-slot generation counter for late-event protection (U12)

Adds a monotonic per-slot generation counter to createWorkerPool's
state. Each successful worker replacement (replaceWorker) bumps the
slot's counter exactly once — atomically with the workers[slotIndex]
swap, so observers (getStats) see the new (worker, generation) pair
consistently. Handler closures in the dispatch loop capture the
slot's generation at attach time and short-circuit when they fire
on a stale generation.

In the current implementation, cleanup() synchronously removes
listeners on a Worker instance the moment a death is observed, so
no listener naturally fires on a stale generation — the guard is a
defensive layer protecting against any future refactor that loosens
cleanup() ordering or re-attaches handlers across the swap. The
load-bearing observable is the slotGenerations[] array exposed via
WorkerPoolStats so operators (and tests) can confirm a slot was
actually replaced and not just the same worker recycled.

Implementation:
  - const slotGenerations: number[] = new Array(size).fill(0) in
    createWorkerPool's per-pool state, alongside respawnCount and
    consecutiveFailuresPerSlot.
  - replaceWorker: slotGenerations[workerIndex]++ AFTER the
    workers[workerIndex] = replacement swap (only on the success
    branch — drop-slot paths leave the counter unchanged).
  - runWorker dispatch loop: const slotGen = slotGenerations[workerIndex]
    captured before handler attachment; every handler (handler /
    errorHandler / exitHandler / messageErrorHandler) starts with
    `if (slotGenerations[workerIndex] !== slotGen) return`.
  - WorkerPoolStats gains `readonly slotGenerations: readonly number[]`.
  - getStats() returns slotGenerations.slice() so callers can't mutate
    pool state by writing to the returned array.

Two existing toEqual snapshots in worker-pool-resilience.test.ts
extended with the new slotGenerations field (both expect all-zeros —
neither test scenario triggers a respawn).

New test file (worker-pool-slot-generation.test.ts, 4 tests):
  1. Fresh pool: every slot at generation 0.
  2. Successful crash + respawn: generation bumps to 1 exactly once.
  3. Crash that drops the slot (maxRespawnsPerSlot:0): generation
     stays at 0 because no successful respawn happened. The dispatch
     rejection on breaker trip is the expected outcome here; the
     load-bearing assertion is the post-rejection stats.
  4. Multi-slot independence: one slot crashing bumps only that
     slot's generation, not the other. Order-independent via sort()
     because the round-robin assignment isn't pinned by contract.

All assertions exact .toEqual / .toBe per DoD §2.7.
This commit is contained in:
Gergo Magyar 2026-05-20 09:18:21 +01:00
parent f7120f8547
commit b702fc643d
3 changed files with 246 additions and 0 deletions

View file

@ -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(),
}),
};
};

View file

@ -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();
});

View file

@ -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<void> {
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<number> {
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();
}
});
});