diff --git a/gitnexus/src/core/ingestion/workers/quarantine.ts b/gitnexus/src/core/ingestion/workers/quarantine.ts new file mode 100644 index 000000000..21a738809 --- /dev/null +++ b/gitnexus/src/core/ingestion/workers/quarantine.ts @@ -0,0 +1,59 @@ +/** + * Quarantine layer (Layer 3 of the worker-pool resilience model). + * + * Tracks paths that caused a worker death this pool lifetime and must + * not be re-dispatched to a worker. Session-scoped — created once per + * `createWorkerPool` invocation and discarded with the pool. + * + * This module is the first piece of the U13 layer-extraction work. The + * doc-review's A10 finding flagged the full 5-module split as + * abstraction-without-multi-consumer-demand, so the rest of the + * extraction is deferred until a real second consumer emerges (e.g., a + * non-parse worker pool that reuses the same resilience layers). + * Extracting the smallest self-contained layer first validates the + * factory + interface pattern with minimal risk: behavior is unchanged, + * the worker-pool.ts public API is unchanged, and existing tests act as + * the regression net. + */ + +/** + * Operations a {@link createQuarantine} instance exposes to the worker + * pool. Intentionally tiny — anything more would invite the abstraction + * overhead doc-review A10 cautioned against. Snapshot returns a fresh + * `string[]` (not a `Set` or iterator) so callers can pass it directly + * to `WorkerPoolDispatchError` without an `Array.from` dance and so + * mutations to the returned array can't accidentally leak back into the + * internal set. + */ +export interface Quarantine { + /** Mark `path` as known-bad for the remainder of this pool's life. */ + add(path: string): void; + /** Whether `path` has been quarantined. */ + has(path: string): boolean; + /** Defensive copy of every quarantined path. */ + snapshot(): string[]; + /** How many distinct paths are currently quarantined. */ + readonly size: number; +} + +/** + * Construct a fresh quarantine. Each `createWorkerPool` invocation gets + * its own instance — quarantines never outlive the pool that created + * them. The implementation is a thin wrapper around `Set`; the + * named interface exists to make the resilience layer addressable as a + * unit (named module, dedicated tests) instead of an inline Set field + * tangled into 1100+ LOC of pool plumbing. + */ +export function createQuarantine(): Quarantine { + const paths = new Set(); + return { + add: (path) => { + paths.add(path); + }, + has: (path) => paths.has(path), + snapshot: () => Array.from(paths), + get size() { + return paths.size; + }, + }; +} diff --git a/gitnexus/src/core/ingestion/workers/worker-pool.ts b/gitnexus/src/core/ingestion/workers/worker-pool.ts index 5823a5676..c87f65449 100644 --- a/gitnexus/src/core/ingestion/workers/worker-pool.ts +++ b/gitnexus/src/core/ingestion/workers/worker-pool.ts @@ -4,6 +4,7 @@ import fs from 'node:fs'; import { fileURLToPath } from 'node:url'; import { logger } from '../../logger.js'; +import { createQuarantine } from './quarantine.js'; export interface WorkerPool { /** * Dispatch items across workers. Items are split into bounded jobs, each job @@ -468,7 +469,12 @@ export const createWorkerPool = ( const workers: (Worker | undefined)[] = new Array(size); const respawnCount: number[] = new Array(size).fill(0); const activeSlots: Set = new Set(); - const quarantined: Set = new Set(); + // Layer 3 (quarantine): tracked via the dedicated `quarantine.ts` + // module so the resilience layer is addressable as a unit (named + // interface, isolated tests) rather than an inline Set tangled into + // 1100+ LOC of pool plumbing. Public worker-pool API is unchanged — + // `getQuarantinedPaths()` still returns the same defensive copy. + const quarantine = createQuarantine(); // Per-slot consecutive-failure counter (F6): replaces the prior pool-wide // scalar so a chronically-failing slot trips the breaker on its own // failure streak instead of being masked by another slot's successes. @@ -517,7 +523,7 @@ export const createWorkerPool = ( const dispatchableItems: TInput[] = []; for (const item of items) { const path = itemPath(item); - if (path !== undefined && quarantined.has(path)) continue; + if (path !== undefined && quarantine.has(path)) continue; dispatchableItems.push(item); } if (dispatchableItems.length === 0) return Promise.resolve([]); @@ -658,7 +664,7 @@ export const createWorkerPool = ( } const firstPath = itemPath(job.items[0]); if (firstPath !== undefined) { - quarantined.add(firstPath); + quarantine.add(firstPath); logger.warn( { startIndex: job.startIndex, firstPath, deaths }, `Conceptual job ${job.startIndex} died ${deaths} times unattributably; ` + @@ -707,7 +713,7 @@ export const createWorkerPool = ( if (stopped) return; consecutiveFailuresPerSlot[workerIndex]++; for (const p of excludePaths) { - if (p) quarantined.add(p); + if (p) quarantine.add(p); } if (consecutiveFailuresPerSlot[workerIndex] >= poolOptions.consecutiveFailureThreshold) { tripBreaker( @@ -715,7 +721,7 @@ export const createWorkerPool = ( `${reason}. Pool circuit breaker tripped: slot ${workerIndex} hit ` + `${consecutiveFailuresPerSlot[workerIndex]} consecutive failures ` + `(threshold: ${poolOptions.consecutiveFailureThreshold}).`, - Array.from(quarantined), + quarantine.snapshot(), ), ); return; @@ -739,7 +745,7 @@ export const createWorkerPool = ( tripBreaker( new WorkerPoolDispatchError( `${reason}. All ${size} worker slot(s) exhausted their respawn budget.`, - Array.from(quarantined), + quarantine.snapshot(), ), ); return; @@ -762,7 +768,7 @@ export const createWorkerPool = ( tripBreaker( new WorkerPoolDispatchError( `${reason}. Replacement worker startup failed and no slots remain.`, - Array.from(quarantined), + quarantine.snapshot(), ), ); } @@ -909,10 +915,10 @@ export const createWorkerPool = ( // queued jobs are fully quarantined back-to-back). let job: WorkerJob | undefined; while ((job = jobs.shift()) !== undefined) { - if (quarantined.size === 0) break; + if (quarantine.size === 0) break; const dispatchable = job.items.filter((item) => { const p = itemPath(item); - return p === undefined || !quarantined.has(p); + return p === undefined || !quarantine.has(p); }); if (dispatchable.length === 0) continue; if (dispatchable.length !== job.items.length) { @@ -1064,7 +1070,7 @@ export const createWorkerPool = ( tripBreaker( new WorkerPoolDispatchError( `Worker pool exhausted all slots during idle-timeout retry.`, - Array.from(quarantined), + quarantine.snapshot(), ), ); return; @@ -1120,7 +1126,7 @@ export const createWorkerPool = ( tripBreaker( new WorkerPoolDispatchError( `Worker ${workerIndex} protocol error: result before flush`, - Array.from(quarantined), + quarantine.snapshot(), ), ); return; @@ -1225,12 +1231,12 @@ export const createWorkerPool = ( dispatch, terminate, size, - getQuarantinedPaths: () => Array.from(quarantined), + getQuarantinedPaths: () => quarantine.snapshot(), getStats: () => ({ size, activeSlots: activeSlots.size, droppedSlots: size - activeSlots.size, - quarantined: quarantined.size, + quarantined: quarantine.size, poolBroken, slotGenerations: slotGenerations.slice(), }), diff --git a/gitnexus/test/unit/workers/quarantine.test.ts b/gitnexus/test/unit/workers/quarantine.test.ts new file mode 100644 index 000000000..9e8c61825 --- /dev/null +++ b/gitnexus/test/unit/workers/quarantine.test.ts @@ -0,0 +1,96 @@ +/** + * U13 (partial) — Isolated tests for the extracted quarantine layer. + * + * Worker-pool resilience integration tests (`worker-pool-resilience`, + * `worker-pool-windows-quarantine`, `worker-pool.test.ts`) already + * exercise the quarantine through the full pool. This file pins the + * module's interface CONTRACT directly so a future change to the + * quarantine surface — extra methods, signature drift, snapshot + * shape — surfaces a focused failure here instead of cascading into + * the larger integration suite. + * + * Notably: the `snapshot()` return type is `string[]`, not `Set` or + * iterator. Tests pin that callers can both mutate the returned array + * (it's a defensive copy) AND pass it directly to consumers expecting + * `string[]` (the `WorkerPoolDispatchError` fallback-exclude-paths + * shape). + */ +import { describe, it, expect } from 'vitest'; +import { createQuarantine } from '../../../src/core/ingestion/workers/quarantine.js'; + +describe('quarantine module (U13 partial)', () => { + it('starts empty', () => { + const q = createQuarantine(); + expect(q.size).toBe(0); + expect(q.snapshot()).toEqual([]); + expect(q.has('any/path.ts')).toBe(false); + }); + + it('records exact-string paths via add() and reports them via has() + size', () => { + const q = createQuarantine(); + q.add('src/bad.ts'); + expect(q.size).toBe(1); + expect(q.has('src/bad.ts')).toBe(true); + expect(q.has('src/other.ts')).toBe(false); + }); + + it('deduplicates repeated add() calls — size grows by exactly one distinct path', () => { + const q = createQuarantine(); + q.add('src/bad.ts'); + q.add('src/bad.ts'); + q.add('src/bad.ts'); + expect(q.size).toBe(1); + }); + + it('preserves separator-style for round-trip (no normalization — matches U9 / M5 contract)', () => { + // Pins the contract worker-pool-windows-quarantine.test.ts asserts + // at the pool level: the quarantine layer treats paths as opaque + // strings. `src\\bad.ts` and `src/bad.ts` are distinct entries. + const q = createQuarantine(); + q.add('src\\bad.ts'); + expect(q.has('src\\bad.ts')).toBe(true); + expect(q.has('src/bad.ts')).toBe(false); + expect(q.size).toBe(1); + }); + + it('snapshot() returns a defensive copy — mutation does not leak back into the quarantine', () => { + const q = createQuarantine(); + q.add('src/a.ts'); + q.add('src/b.ts'); + const snap = q.snapshot(); + expect(snap.sort()).toEqual(['src/a.ts', 'src/b.ts']); + + // Mutate the returned array; the internal state must not change. + snap.length = 0; + snap.push('src/never-added.ts'); + expect(q.size).toBe(2); + expect(q.has('src/a.ts')).toBe(true); + expect(q.has('src/never-added.ts')).toBe(false); + }); + + it('snapshots are independent — successive calls return fresh arrays', () => { + const q = createQuarantine(); + q.add('src/a.ts'); + const first = q.snapshot(); + const second = q.snapshot(); + expect(first).toEqual(second); + expect(first).not.toBe(second); + }); + + it('reflects subsequent add() calls in later snapshots', () => { + const q = createQuarantine(); + q.add('src/a.ts'); + expect(q.snapshot()).toEqual(['src/a.ts']); + q.add('src/b.ts'); + expect(q.snapshot().sort()).toEqual(['src/a.ts', 'src/b.ts']); + }); + + it('size is a getter, not a stale property — reflects state at access time', () => { + const q = createQuarantine(); + expect(q.size).toBe(0); + q.add('src/a.ts'); + expect(q.size).toBe(1); + q.add('src/b.ts'); + expect(q.size).toBe(2); + }); +});