refactor(workers): extract quarantine into its own module (U13 partial)

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<string> = 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.
This commit is contained in:
Gergo Magyar 2026-05-20 10:24:47 +01:00
parent ac90a133d0
commit 832d97678f
3 changed files with 174 additions and 13 deletions

View file

@ -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<string>`; 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<string>();
return {
add: (path) => {
paths.add(path);
},
has: (path) => paths.has(path),
snapshot: () => Array.from(paths),
get size() {
return paths.size;
},
};
}

View file

@ -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<number> = new Set();
const quarantined: Set<string> = 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<TInput> | 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(),
}),

View file

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