From 832d97678fae996f6693668f3a0a484b1e244b89 Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Wed, 20 May 2026 10:24:47 +0100 Subject: [PATCH] refactor(workers): extract quarantine into its own module (U13 partial) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Honest partial U13: extract the quarantine resilience layer (Layer 3 of the 5-layer model) into a dedicated module with a small explicit interface. The full 5-module split that the original plan named was flagged by doc-review A10 as abstraction-without-multi-consumer-demand ("Each has exactly one consumer: worker-pool.ts. None of these layers is imported elsewhere in the codebase pre-extraction, and the plan doesn't identify any future consumer.") This commit ships the smallest self-contained layer as a named module to validate the factory + interface pattern with minimal risk. The remaining four layers (respawn-budget, cumulative-timeout, circuit-breaker, slot-attribution) stay inline until a real second consumer emerges (e.g., a non-parse worker pool that reuses the same resilience layers). Module shape (`workers/quarantine.ts`, ~30 LOC): interface Quarantine { add(path: string): void; has(path: string): boolean; snapshot(): string[]; // defensive copy readonly size: number; // getter, reflects state at access time } function createQuarantine(): Quarantine Replaces in `worker-pool.ts`: - `const quarantined: Set = new Set()` -> `createQuarantine()` - `quarantined.has(p)` -> `quarantine.has(p)` (2 sites) - `quarantined.add(p)` -> `quarantine.add(p)` (2 sites) - `quarantined.size` -> `quarantine.size` (2 sites) - `Array.from(quarantined)` -> `quarantine.snapshot()` (6 sites) Public worker-pool.ts API is unchanged — `getQuarantinedPaths()` still returns the same defensive `string[]` copy. The behavioral contract is preserved: paths are quarantined as opaque strings (the U9 / M5 non-normalization contract still holds — see the new dedicated test). Tests: - 8 isolated unit tests for the quarantine module — pins the interface contract (empty start, add/has/size, dedup on repeated add, no separator normalization, snapshot defensive copy + freshness, size-getter live behavior). - All 86 existing worker-pool tests pass unchanged — they exercise the quarantine through the pool and act as the regression net for behavior preservation. Why not the full 5-module extraction in this commit: doc-review A10's concern is real — a single-consumer abstraction adds module-boundary overhead (5 sets of imports, 5 dedicated test files, 5 interfaces to keep in sync with worker-pool) without any structural benefit until a second consumer materializes. Extracting one validates the pattern; the remaining four can be moved on demand. --- .../src/core/ingestion/workers/quarantine.ts | 59 ++++++++++++ .../src/core/ingestion/workers/worker-pool.ts | 32 ++++--- gitnexus/test/unit/workers/quarantine.test.ts | 96 +++++++++++++++++++ 3 files changed, 174 insertions(+), 13 deletions(-) create mode 100644 gitnexus/src/core/ingestion/workers/quarantine.ts create mode 100644 gitnexus/test/unit/workers/quarantine.test.ts 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); + }); +});