feat(workers): resilient + scalable worker pool

Restructures `createWorkerPool` so a single bad file no longer kills the
pool for the rest of an analyze run. Five interlocking layers:

1. **Auto-respawn on error/exit** — worker death triggers `replaceWorker`
   on the same slot, bounded by `maxRespawnsPerSlot` (default 3). The slot
   is dropped from rotation when the budget is exhausted; other slots
   keep running.

2. **Circuit breaker** — replaces the permanent `poolBroken=true` with a
   consecutive-failure counter. The pool only trips after
   `consecutiveFailureThreshold` deaths (default `max(3, poolSize)`) with
   no successful job in between. A successful job resets the counter so
   transient bursts of bad files don't escalate.

3. **Session-scoped file quarantine** — paths identified as the in-flight
   file at the moment of a worker death are added to a `Set<string>` on
   the pool. `dispatch()` filters quarantined items up front (they never
   reach a worker again this pool lifetime). Exposed via the new
   `WorkerPool.getQuarantinedPaths()` so callers can log/route them.
   `processParsing` surfaces the per-chunk quarantine summary alongside
   the existing fallback-exclusion log.

4. **Authoritative in-flight tracking** — `parse-worker.ts` emits
   `{type:'starting-file', path}` before each file. The pool tracks this
   per slot and uses it for crash attribution, falling back to the
   `items[lastProgress]` heuristic only when no starting-file has been
   observed (very-early crash, older worker build). Closes the
   reorder/race concerns raised by reviewers C1 and R3 in the earlier
   review run.

5. **Per-job cumulative timeout budget** — each `WorkerJob` tracks the
   total wall time spent across attempts/splits/retries. When the budget
   is exhausted (default 5x `subBatchIdleTimeoutMs`), the pool surfaces
   the in-flight path instead of letting exponential backoff balloon
   into multi-hour stalls.

Cross-layer wiring: a new `wakeIdleSlots` helper kicks any non-busy live
slot when items are requeued (after a death or split-retry), so a dropped
slot doesn't strand work in the queue. `recoverAndResume` consolidates
the per-job teardown shared by the three in-pool death sites (`error`,
`exit`, msg-channel `error`).

New env knobs: `GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT`,
`GITNEXUS_WORKER_MAX_CUMULATIVE_TIMEOUT_MS`,
`GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD`.
New `WorkerPoolOptions.workerFactory` injection point for unit tests.

Tests: 12 new unit tests using a FakeWorker mock cover quarantine
seeding, slot-respawn, slot-drop after budget, breaker trip + reset,
and quarantine filtering. Plus option-resolution tests for the three
new env vars. All 19 worker-pool/-fallback/-options tests pass; full
unit suite 6040 passed / 30 skipped / 0 failed.
This commit is contained in:
Gergo Magyar 2026-05-19 09:28:02 +01:00
parent 25af32c121
commit 3ffd6ad396
4 changed files with 859 additions and 79 deletions

View file

@ -875,7 +875,7 @@ export const processParsing = async (
);
}
try {
return await processParsingWithWorkers(
const data = await processParsingWithWorkers(
graph,
files,
symbolTable,
@ -884,6 +884,28 @@ export const processParsing = async (
reportProgress,
outRawResults,
);
// Session-scoped quarantine (worker-pool resilience Layer 3): surface
// any files this pool has decided are unsafe for workers so the
// operator can see what was skipped. The pool already filtered them
// out of dispatch; we only need to log + progress-report. Quarantine
// is session-scoped per pool instance — a fresh `createWorkerPool`
// call clears it.
const quarantineSet = new Set(workerPool.getQuarantinedPaths?.() ?? []);
if (quarantineSet.size > 0) {
const quarantinedInChunk = files.filter((file) => quarantineSet.has(file.path));
if (quarantinedInChunk.length > 0) {
logger.warn(
{ quarantinedFiles: quarantinedInChunk.map((file) => file.path) },
`Worker quarantine: ${quarantinedInChunk.length} file(s) skipped in this chunk (cumulative pool quarantine: ${quarantineSet.size}).`,
);
reportProgress?.(
lastProgress,
files.length,
`${quarantinedInChunk.length} worker-quarantined file(s) skipped`,
);
}
}
return data;
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
let fallbackFiles = files;

View file

@ -1401,6 +1401,13 @@ const processFileGroup = (
// Skip files larger than the max tree-sitter buffer (32 MB)
if (getTreeSitterContentByteLength(file.content) > TREE_SITTER_MAX_BUFFER) continue;
// Authoritative in-flight signal for the pool: lets `WorkerPool` exclude
// exactly this file if the worker dies during parse/extract, instead of
// guessing from `items[lastProgress]` (which the language-grouped order
// here would defeat). The pool gracefully ignores this when running an
// older worker build that doesn't emit it.
if (parentPort) parentPort.postMessage({ type: 'starting-file', path: file.path });
// Vue SFC preprocessing: extract <script> block content
let parseContent = file.content;
let lineOffset = 0;

View file

@ -8,6 +8,12 @@ export interface WorkerPool {
/**
* Dispatch items across workers. Items are split into bounded jobs, each job
* is committed independently, and stalled jobs are split/retried locally.
*
* Files in {@link WorkerPool.getQuarantinedPaths} are filtered out before
* dispatch — they have already caused a worker death this pool lifetime and
* are not safe to re-attempt in workers. The caller is responsible for
* routing them (e.g. to sequential fallback); inspect the quarantine
* snapshot before and after each dispatch.
*/
dispatch<TInput, TResult>(
items: TInput[],
@ -17,8 +23,16 @@ export interface WorkerPool {
/** Terminate all workers. Must be called when done. */
terminate(): Promise<void>;
/** Number of workers in the pool */
/** Number of worker slots originally requested for the pool. */
readonly size: number;
/**
* Snapshot of paths quarantined by this pool instance. Populated when a
* worker dies with an authoritative in-flight file (Layer 4 starting-file
* message) or a singleton-timeout exclusion. Cleared only by pool teardown
* — quarantine is session-scoped per `createWorkerPool` invocation.
*/
getQuarantinedPaths(): readonly string[];
}
export interface WorkerPoolOptions {
@ -27,6 +41,34 @@ export interface WorkerPoolOptions {
subBatchIdleTimeoutMs?: number;
maxTimeoutRetries?: number;
timeoutBackoffFactor?: number;
/**
* Max replacement spawns per worker slot before the slot is dropped from
* the active rotation. Bounds respawn loops on a slot that consistently
* crashes the worker (likely a system-level fault rather than a single
* bad input). Default 3.
*/
maxRespawnsPerSlot?: number;
/**
* Hard ceiling on total wall time the pool will spend retrying / splitting
* any single job. Combined with `timeoutBackoffFactor`, this prevents
* exponentially-growing retry waits from accumulating into multi-hour
* stalls before the pool finally surfaces the bad file to sequential
* fallback. Default 5x `subBatchIdleTimeoutMs`.
*/
maxCumulativeTimeoutMs?: number;
/**
* Number of consecutive worker deaths (no successful job in between) that
* trip the pool circuit breaker. Once tripped, the pool rejects every
* subsequent `dispatch` with `WorkerPoolDispatchError` until a new pool is
* created. Default `Math.max(3, poolSize)`.
*/
consecutiveFailureThreshold?: number;
/**
* Test-only injection point for the Worker constructor. When provided,
* the pool uses this factory instead of `new Worker(workerUrl)`. Production
* code should leave this unset.
*/
workerFactory?: (workerUrl: URL) => Worker;
}
export class WorkerPoolDispatchError extends Error {
@ -45,7 +87,15 @@ type WorkerOutgoingMessage =
| { type: 'warning'; message: string }
| { type: 'sub-batch-done' }
| { type: 'error'; error: string }
| { type: 'result'; data: unknown };
| { type: 'result'; data: unknown }
/**
* Authoritative in-flight signal: worker is about to process this file.
* Pool records it per slot so worker death can be attributed exactly,
* instead of guessing from `items[lastProgress]` (which language-grouped
* worker processing defeats). Optional — older worker builds may not
* emit it; pool falls back to the heuristic when absent.
*/
| { type: 'starting-file'; path: string };
interface WorkerJob<TInput> {
startIndex: number;
@ -54,6 +104,13 @@ interface WorkerJob<TInput> {
attempt: number;
splitDepth: number;
timeoutMs: number;
/**
* Running total of timeoutMs across all attempts/splits/respawn-retries
* for this conceptual unit of work. Tracked separately from `timeoutMs`
* so we can bound the *total* wait the pool incurs on a single job, not
* just the current attempt. See {@link WorkerPoolOptions.maxCumulativeTimeoutMs}.
*/
cumulativeTimeoutMs: number;
}
interface WorkerJobResult<TResult> {
@ -71,6 +128,9 @@ const SUB_BATCH_MAX_BYTES = 8 * 1024 * 1024;
const DEFAULT_SUB_BATCH_IDLE_TIMEOUT_MS = 30_000;
const DEFAULT_TIMEOUT_RETRIES = 1;
const DEFAULT_TIMEOUT_BACKOFF_FACTOR = 2;
const DEFAULT_MAX_RESPAWNS_PER_SLOT = 3;
const DEFAULT_MAX_CUMULATIVE_TIMEOUT_FACTOR = 5;
const DEFAULT_CONSECUTIVE_FAILURE_THRESHOLD_FLOOR = 3;
function positiveInteger(value: unknown): number | undefined {
const parsed = typeof value === 'string' ? Number(value) : value;
@ -86,22 +146,47 @@ function nonNegativeInteger(value: unknown): number | undefined {
: undefined;
}
interface ResolvedWorkerPoolOptions {
subBatchSize: number;
subBatchMaxBytes: number;
subBatchIdleTimeoutMs: number;
maxTimeoutRetries: number;
timeoutBackoffFactor: number;
maxRespawnsPerSlot: number;
maxCumulativeTimeoutMs: number;
consecutiveFailureThreshold: number;
}
export function resolveWorkerPoolOptions(
options: WorkerPoolOptions = {},
): Required<WorkerPoolOptions> {
poolSize?: number,
): ResolvedWorkerPoolOptions {
const subBatchIdleTimeoutMs =
positiveInteger(options.subBatchIdleTimeoutMs) ??
positiveInteger(process.env.GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS) ??
DEFAULT_SUB_BATCH_IDLE_TIMEOUT_MS;
return {
subBatchSize: positiveInteger(options.subBatchSize) ?? SUB_BATCH_SIZE,
subBatchMaxBytes:
positiveInteger(options.subBatchMaxBytes) ??
positiveInteger(process.env.GITNEXUS_WORKER_SUB_BATCH_MAX_BYTES) ??
SUB_BATCH_MAX_BYTES,
subBatchIdleTimeoutMs:
positiveInteger(options.subBatchIdleTimeoutMs) ??
positiveInteger(process.env.GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS) ??
DEFAULT_SUB_BATCH_IDLE_TIMEOUT_MS,
subBatchIdleTimeoutMs,
maxTimeoutRetries: nonNegativeInteger(options.maxTimeoutRetries) ?? DEFAULT_TIMEOUT_RETRIES,
timeoutBackoffFactor:
positiveInteger(options.timeoutBackoffFactor) ?? DEFAULT_TIMEOUT_BACKOFF_FACTOR,
maxRespawnsPerSlot:
nonNegativeInteger(options.maxRespawnsPerSlot) ??
nonNegativeInteger(process.env.GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT) ??
DEFAULT_MAX_RESPAWNS_PER_SLOT,
maxCumulativeTimeoutMs:
positiveInteger(options.maxCumulativeTimeoutMs) ??
positiveInteger(process.env.GITNEXUS_WORKER_MAX_CUMULATIVE_TIMEOUT_MS) ??
subBatchIdleTimeoutMs * DEFAULT_MAX_CUMULATIVE_TIMEOUT_FACTOR,
consecutiveFailureThreshold:
positiveInteger(options.consecutiveFailureThreshold) ??
positiveInteger(process.env.GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD) ??
Math.max(DEFAULT_CONSECUTIVE_FAILURE_THRESHOLD_FLOOR, poolSize ?? 0),
};
}
@ -143,18 +228,18 @@ function itemPath(item: unknown): string | undefined {
}
/**
* Best-guess path of the file in flight when a worker dies mid-job.
* Best-guess path of the file in flight when a worker dies mid-job — used as
* the fallback when the authoritative `starting-file` message hasn't been
* observed yet (very early job-startup crash, or older worker build that
* doesn't emit the signal).
*
* `lastProgress` is the number of files the worker has acknowledged via
* `progress` messages, so `items[lastProgress]` is the next file it was
* about to process — the most likely culprit when the worker crashes
* (OOM, native addon SIGSEGV) or reports an error.
*
* Excluding only this single path keeps the blast radius small: earlier
* files in the job get re-tried by the sequential fallback, and any
* pathological file gets the same skip treatment as the singleton-
* timeout path. Returns `[]` when no path is determinable so sequential
* retries the whole job.
* Returns `[]` when no path is determinable so the caller retries the whole
* job.
*/
function inFlightExcludePath<TInput>(job: WorkerJob<TInput>, lastProgress: number): string[] {
if (lastProgress >= job.items.length) return [];
@ -182,6 +267,7 @@ function createJobs<TInput>(
attempt: 0,
splitDepth: 0,
timeoutMs,
cumulativeTimeoutMs: timeoutMs,
});
startIndex += batch.length;
batch = [];
@ -202,6 +288,26 @@ function createJobs<TInput>(
/**
* Create a pool of worker threads.
*
* Resilience model (PR #1693 / 1694):
* - Layer 1 (auto-respawn): a worker `error`/`exit` triggers a replacement on
* the same slot, bounded by {@link WorkerPoolOptions.maxRespawnsPerSlot}.
* The slot is dropped from the rotation when its budget is exhausted.
* - Layer 2 (circuit breaker): `consecutiveFailureThreshold` consecutive
* worker deaths (no successful job between) — OR all slots exhausting their
* respawn budget — trip the breaker. Every subsequent dispatch rejects
* with `WorkerPoolDispatchError` and the caller must build a new pool.
* - Layer 3 (quarantine): a path identified as the in-flight file at the
* time of a worker death is added to `quarantined` and filtered out of
* future dispatches. Snapshot via {@link WorkerPool.getQuarantinedPaths}.
* - Layer 4 (authoritative in-flight): the worker emits a `starting-file`
* message before each parse attempt; the pool prefers this for crash
* attribution and falls back to {@link inFlightExcludePath} only when no
* signal has been observed yet.
* - Layer 5 (cumulative timeout budget): each job tracks the total wall
* time spent across all attempts/splits/retries. When the budget is
* exhausted, the pool surfaces the in-flight path via `WorkerPoolDispatchError`
* instead of letting timeouts compound indefinitely.
*/
export const createWorkerPool = (
workerUrl: URL,
@ -216,13 +322,19 @@ export const createWorkerPool = (
}
const size = poolSize ?? Math.min(8, Math.max(1, os.cpus().length - 1));
const poolOptions = resolveWorkerPoolOptions(options);
const workers: Worker[] = [];
const poolOptions = resolveWorkerPoolOptions(options, size);
const spawnWorker = options?.workerFactory ?? ((url: URL) => new Worker(url));
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();
let consecutiveFailures = 0;
let poolBroken = false;
let poolFailure: Error | undefined;
for (let i = 0; i < size; i++) {
workers.push(new Worker(workerUrl));
workers[i] = spawnWorker(workerUrl);
activeSlots.add(i);
}
const dispatch = <TInput, TResult>(
@ -232,14 +344,31 @@ export const createWorkerPool = (
if (poolBroken) {
const reason = poolFailure ? `: ${poolFailure.message}` : '';
return Promise.reject(
new Error(`Worker pool is unavailable after a previous failure${reason}`),
new WorkerPoolDispatchError(
`Worker pool circuit breaker tripped${reason}. ` +
`Subsequent dispatches require a fresh pool instance.`,
[],
),
);
}
if (items.length === 0) return Promise.resolve([]);
if (workers.length === 0) return Promise.reject(new Error('Worker pool has no active workers'));
if (activeSlots.size === 0) {
return Promise.reject(new WorkerPoolDispatchError('Worker pool has no active workers', []));
}
// Layer 3: filter out quarantined paths so a known-bad file never reaches
// a worker again this pool lifetime. The caller queries
// `getQuarantinedPaths` after dispatch to route filtered items.
const dispatchableItems: TInput[] = [];
for (const item of items) {
const path = itemPath(item);
if (path !== undefined && quarantined.has(path)) continue;
dispatchableItems.push(item);
}
if (dispatchableItems.length === 0) return Promise.resolve([]);
const jobs = createJobs(
items,
dispatchableItems,
poolOptions.subBatchSize,
poolOptions.subBatchMaxBytes,
poolOptions.subBatchIdleTimeoutMs,
@ -247,47 +376,72 @@ export const createWorkerPool = (
return new Promise<TResult[]>((resolve, reject) => {
const results: WorkerJobResult<TResult>[] = [];
const inFlightProgress = new Array(workers.length).fill(0);
const inFlightProgress = new Array(size).fill(0);
// Tracks which slots are currently mid-job so the "wake idle slots"
// pass after a requeue doesn't double-dispatch to a busy slot.
const busySlots: Set<number> = new Set();
let completedFiles = 0;
let activeWorkers = 0;
let stopped = false;
let maxReported = 0;
const wakeIdleSlots = () => {
if (stopped || jobs.length === 0) return;
for (const slot of activeSlots) {
if (busySlots.has(slot)) continue;
if (jobs.length === 0) break;
runWorker(slot);
}
};
const reportProgress = () => {
if (!onProgress) return;
const inFlight = inFlightProgress.reduce((sum, value) => sum + value, 0);
const next = Math.min(items.length, Math.max(maxReported, completedFiles + inFlight));
const next = Math.min(
dispatchableItems.length,
Math.max(maxReported, completedFiles + inFlight),
);
if (next === maxReported) return;
maxReported = next;
onProgress(next);
};
const replaceWorker = async (workerIndex: number) => {
const worker = workers[workerIndex];
await worker?.terminate().catch(() => undefined);
if (stopped) return;
const replacement = new Worker(workerUrl);
const replaceWorker = async (workerIndex: number): Promise<boolean> => {
const existing = workers[workerIndex];
await existing?.terminate().catch(() => undefined);
workers[workerIndex] = undefined;
if (stopped) return false;
const replacement = spawnWorker(workerUrl);
try {
await waitForWorkerOnline(replacement);
} catch (err) {
await replacement.terminate().catch(() => undefined);
throw new Error(
`Replacement worker ${workerIndex} failed to start: ${err instanceof Error ? err.message : String(err)}`,
logger.warn(
{ workerIndex, error: err instanceof Error ? err.message : String(err) },
`Worker ${workerIndex} replacement failed to come online; dropping slot.`,
);
return false;
}
if (stopped) {
await replacement.terminate().catch(() => undefined);
return;
return false;
}
workers[workerIndex] = replacement;
return true;
};
const fail = async (err: Error) => {
// Terminal failure path: trip the pool circuit breaker and reject the
// outer dispatch promise with the cumulative exclude paths. This is the
// ONLY place that sets `poolBroken = true` — recoverable single-worker
// failures stay local to `handleWorkerDeath`.
const tripBreaker = async (err: WorkerPoolDispatchError) => {
poolBroken = true;
poolFailure = err;
if (stopped) return;
stopped = true;
await Promise.all(workers.map((worker) => worker.terminate().catch(() => undefined)));
await Promise.all(workers.map((worker) => worker?.terminate().catch(() => undefined)));
for (let i = 0; i < workers.length; i++) workers[i] = undefined;
activeSlots.clear();
reject(err);
};
@ -296,17 +450,151 @@ export const createWorkerPool = (
if (jobs.length === 0 && activeWorkers === 0) {
stopped = true;
results.sort((a, b) => a.startIndex - b.startIndex);
if (onProgress && maxReported < items.length) onProgress(items.length);
if (onProgress && maxReported < dispatchableItems.length)
onProgress(dispatchableItems.length);
resolve(results.map((result) => result.data));
}
};
// Re-queue the non-quarantined remainder of a dead worker's job so a
// healthy worker can finish the work. Earlier items in the dead job
// were never flushed back to the main thread, so they must be
// re-processed. The new job carries the existing job's startIndex so
// result ordering is preserved.
const requeueRemainder = (job: WorkerJob<TInput>, excluded: readonly string[]) => {
if (excluded.length === 0) {
jobs.unshift(job);
return;
}
const excludeSet = new Set(excluded);
const filtered = job.items.filter((item) => {
const p = itemPath(item);
return p === undefined || !excludeSet.has(p);
});
if (filtered.length === 0) return;
const requeueTimeoutMs = job.timeoutMs;
jobs.unshift({
startIndex: job.startIndex,
items: filtered,
estimatedBytes: filtered.reduce((sum, item) => sum + estimateItemBytes(item), 0),
attempt: job.attempt,
splitDepth: job.splitDepth,
timeoutMs: requeueTimeoutMs,
cumulativeTimeoutMs: job.cumulativeTimeoutMs + requeueTimeoutMs,
});
};
// Recoverable worker death — quarantine the in-flight path, attempt
// to respawn the slot, re-queue the rest of the job, and continue.
// Trips the circuit breaker only when consecutiveFailures crosses the
// threshold OR all slots have exhausted their respawn budget.
const handleWorkerDeath = async (
workerIndex: number,
reason: string,
excludePaths: readonly string[],
) => {
if (stopped) return;
consecutiveFailures++;
for (const p of excludePaths) {
if (p) quarantined.add(p);
}
if (consecutiveFailures >= poolOptions.consecutiveFailureThreshold) {
void tripBreaker(
new WorkerPoolDispatchError(
`${reason}. Pool circuit breaker tripped after ${consecutiveFailures} ` +
`consecutive failures (threshold: ${poolOptions.consecutiveFailureThreshold}).`,
Array.from(quarantined),
),
);
return;
}
respawnCount[workerIndex]++;
if (respawnCount[workerIndex] > poolOptions.maxRespawnsPerSlot) {
logger.warn(
{
workerIndex,
respawnCount: respawnCount[workerIndex],
maxRespawns: poolOptions.maxRespawnsPerSlot,
reason,
},
`Worker ${workerIndex} exceeded respawn budget; dropping slot.`,
);
const dead = workers[workerIndex];
await dead?.terminate().catch(() => undefined);
workers[workerIndex] = undefined;
activeSlots.delete(workerIndex);
if (activeSlots.size === 0) {
void tripBreaker(
new WorkerPoolDispatchError(
`${reason}. All ${size} worker slot(s) exhausted their respawn budget.`,
Array.from(quarantined),
),
);
return;
}
return;
}
logger.warn(
{
workerIndex,
respawnCount: respawnCount[workerIndex],
reason,
excludePaths,
},
`Worker ${workerIndex} died; respawning slot (attempt ${respawnCount[workerIndex]}/${poolOptions.maxRespawnsPerSlot}).`,
);
const respawned = await replaceWorker(workerIndex);
if (!respawned) {
activeSlots.delete(workerIndex);
if (activeSlots.size === 0) {
void tripBreaker(
new WorkerPoolDispatchError(
`${reason}. Replacement worker startup failed and no slots remain.`,
Array.from(quarantined),
),
);
}
return;
}
};
const requeueAfterTimeout = (
workerIndex: number,
job: WorkerJob<TInput>,
lastProgress: number,
inFlightPath: string | undefined,
): boolean => {
const nextTimeout = Math.ceil(job.timeoutMs * poolOptions.timeoutBackoffFactor);
const nextCumulative = job.cumulativeTimeoutMs + nextTimeout;
// Layer 5: respect the per-job cumulative timeout budget. Once
// exhausted, surface the in-flight file via WorkerPoolDispatchError
// instead of letting exponential backoff stall further.
if (nextCumulative > poolOptions.maxCumulativeTimeoutMs) {
const exhausted =
inFlightPath !== undefined
? [inFlightPath]
: itemPath(job.items[0])
? [itemPath(job.items[0]) as string]
: [];
logger.warn(
{
workerIndex,
cumulativeMs: job.cumulativeTimeoutMs,
nextCumulativeMs: nextCumulative,
maxCumulativeMs: poolOptions.maxCumulativeTimeoutMs,
exhausted,
},
`Worker ${workerIndex} parse job exhausted cumulative timeout budget. Surfacing in-flight file(s).`,
);
void handleWorkerDeath(
workerIndex,
`Worker ${workerIndex} parse job exhausted cumulative timeout budget ` +
`(${(nextCumulative / 1000).toFixed(0)}s > ${(poolOptions.maxCumulativeTimeoutMs / 1000).toFixed(0)}s cap)`,
exhausted,
);
return false;
}
if (job.items.length > 1) {
const midpoint = Math.ceil(job.items.length / 2);
@ -319,6 +607,7 @@ export const createWorkerPool = (
attempt: job.attempt,
splitDepth: job.splitDepth + 1,
timeoutMs: nextTimeout,
cumulativeTimeoutMs: nextCumulative,
};
const second: WorkerJob<TInput> = {
startIndex: job.startIndex + midpoint,
@ -327,6 +616,7 @@ export const createWorkerPool = (
attempt: job.attempt,
splitDepth: job.splitDepth + 1,
timeoutMs: nextTimeout,
cumulativeTimeoutMs: nextCumulative,
};
logger.warn(
{
@ -362,39 +652,91 @@ export const createWorkerPool = (
...job,
attempt: nextAttempt,
timeoutMs: nextTimeout,
cumulativeTimeoutMs: nextCumulative,
});
return true;
}
const stalledPath = itemPath(job.items[0]);
void fail(
new WorkerPoolDispatchError(
`Worker ${workerIndex} parse job idle timeout after ${job.timeoutMs / 1000}s ` +
`(single item${stalledPath ? `: ${stalledPath}` : ''}, ` +
`${job.estimatedBytes} bytes, last progress: ${lastProgress}). ` +
`Analyze will retry through sequential fallback. Increase with ` +
`--worker-timeout or GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS.`,
stalledPath ? [stalledPath] : [],
),
const stalledPath = inFlightPath ?? itemPath(job.items[0]);
const excludes = stalledPath ? [stalledPath] : [];
logger.warn(
{
workerIndex,
timeoutSec: job.timeoutMs / 1000,
stalledPath,
cumulativeMs: job.cumulativeTimeoutMs,
},
`Worker ${workerIndex} parse job idle timeout exhausted retries; quarantining file and respawning slot.`,
);
void handleWorkerDeath(
workerIndex,
`Worker ${workerIndex} parse job idle timeout after ${job.timeoutMs / 1000}s ` +
`(single item${stalledPath ? `: ${stalledPath}` : ''}, ` +
`${job.estimatedBytes} bytes, last progress: ${lastProgress})`,
excludes,
);
return false;
};
const runWorker = (workerIndex: number) => {
if (stopped) return;
if (!activeSlots.has(workerIndex)) return;
const job = jobs.shift();
if (!job) {
maybeDone();
return;
}
// Drop quarantined items that may have been re-queued before a death
// added them to quarantine — keeps the worker from ever seeing a
// known-bad file.
if (quarantined.size > 0) {
const dispatchable = job.items.filter((item) => {
const p = itemPath(item);
return p === undefined || !quarantined.has(p);
});
if (dispatchable.length === 0) {
// Whole job was quarantined; drop and try next.
runWorker(workerIndex);
return;
}
if (dispatchable.length !== job.items.length) {
job.items = dispatchable;
job.estimatedBytes = dispatchable.reduce(
(sum, item) => sum + estimateItemBytes(item),
0,
);
}
}
activeWorkers++;
busySlots.add(workerIndex);
inFlightProgress[workerIndex] = 0;
const worker = workers[workerIndex];
if (!worker) {
// Slot was dropped between scheduling and execution; requeue the
// job for another slot and bail.
activeWorkers--;
busySlots.delete(workerIndex);
jobs.unshift(job);
wakeIdleSlots();
maybeDone();
return;
}
let settled = false;
let waitingForFlush = false;
let idleTimer: ReturnType<typeof setTimeout> | null = null;
let lastProgress = 0;
// Authoritative in-flight file from the worker's `starting-file`
// message. Cleared on `progress` so a between-files crash falls
// back to the `items[lastProgress]` heuristic, which then points
// at the next file (the one about to start) — the right guess.
let inFlightPath: string | undefined;
const resolveExcludePaths = (): readonly string[] => {
if (inFlightPath !== undefined) return [inFlightPath];
return inFlightExcludePath(job, lastProgress);
};
const cleanup = () => {
if (idleTimer) clearTimeout(idleTimer);
@ -405,44 +747,100 @@ export const createWorkerPool = (
const finishJob = () => {
activeWorkers--;
busySlots.delete(workerIndex);
inFlightProgress[workerIndex] = 0;
runWorker(workerIndex);
maybeDone();
};
// Recover-and-resume flow shared by all in-pool worker death sites
// (`error`, `exit`, msg-channel error). Bridges the per-job teardown
// into the pool-level handleWorkerDeath recovery + breaker logic.
const recoverAndResume = async (reason: string, excludePaths: readonly string[]) => {
activeWorkers--;
busySlots.delete(workerIndex);
inFlightProgress[workerIndex] = 0;
requeueRemainder(job, excludePaths);
await handleWorkerDeath(workerIndex, reason, excludePaths);
if (stopped) return;
// Slot may have been dropped or respawned. Kick the current slot
// if still active, then wake any other idle live slots so the
// requeued remainder can be picked up immediately (without this,
// dropped-slot scenarios can deadlock when no other slot is
// currently busy and the next finishJob never fires).
if (activeSlots.has(workerIndex)) {
runWorker(workerIndex);
}
wakeIdleSlots();
maybeDone();
};
const resetIdleTimer = () => {
if (idleTimer) clearTimeout(idleTimer);
idleTimer = setTimeout(async () => {
idleTimer = setTimeout(() => {
if (!settled) {
settled = true;
cleanup();
inFlightProgress[workerIndex] = 0;
const shouldContinue = requeueAfterTimeout(workerIndex, job, lastProgress);
const stalledPath = inFlightPath;
const shouldContinue = requeueAfterTimeout(
workerIndex,
job,
lastProgress,
stalledPath,
);
if (!shouldContinue) {
activeWorkers--;
// handleWorkerDeath path was taken by requeueAfterTimeout;
// recover the slot (respawn if budget allows) and continue.
void (async () => {
activeWorkers--;
busySlots.delete(workerIndex);
if (stopped) return;
if (activeSlots.has(workerIndex)) {
const respawned = await replaceWorker(workerIndex);
if (!respawned) activeSlots.delete(workerIndex);
}
if (stopped) return;
if (activeSlots.has(workerIndex)) {
runWorker(workerIndex);
}
wakeIdleSlots();
maybeDone();
})();
return;
}
try {
await replaceWorker(workerIndex);
} catch (err) {
void fail(err instanceof Error ? err : new Error(String(err)));
return;
} finally {
activeWorkers--;
}
reportProgress();
runWorker(workerIndex);
maybeDone();
// Timeout-retry path: spawn a fresh worker on this slot to
// pick up the next attempt.
void (async () => {
try {
const respawned = await replaceWorker(workerIndex);
if (!respawned) {
activeSlots.delete(workerIndex);
}
} finally {
activeWorkers--;
busySlots.delete(workerIndex);
}
if (stopped) return;
reportProgress();
if (activeSlots.has(workerIndex)) runWorker(workerIndex);
wakeIdleSlots();
maybeDone();
})();
}
}, job.timeoutMs);
};
const handler = (msg: WorkerOutgoingMessage) => {
if (settled || stopped) return;
if (msg.type === 'progress') {
if (msg.type === 'starting-file') {
inFlightPath = msg.path;
resetIdleTimer();
} else if (msg.type === 'progress') {
const bounded = Math.min(job.items.length, Math.max(0, msg.filesProcessed));
inFlightProgress[workerIndex] = bounded;
lastProgress = bounded;
inFlightPath = undefined;
resetIdleTimer();
reportProgress();
} else if (msg.type === 'warning') {
@ -455,23 +853,30 @@ export const createWorkerPool = (
} else if (msg.type === 'error') {
settled = true;
cleanup();
void fail(
new WorkerPoolDispatchError(
`Worker ${workerIndex} error: ${msg.error}`,
inFlightExcludePath(job, lastProgress),
),
void recoverAndResume(
`Worker ${workerIndex} error: ${msg.error}`,
resolveExcludePaths(),
);
} else if (msg.type === 'result') {
if (!waitingForFlush) {
settled = true;
cleanup();
void fail(new Error(`Worker ${workerIndex} protocol error: result before flush`));
void tripBreaker(
new WorkerPoolDispatchError(
`Worker ${workerIndex} protocol error: result before flush`,
Array.from(quarantined),
),
);
return;
}
settled = true;
cleanup();
results.push({ startIndex: job.startIndex, data: msg.data as TResult });
completedFiles += job.items.length;
// Layer 2: a successful job resets the consecutive-failure
// counter so transient bursts of bad files don't trip the
// breaker prematurely.
consecutiveFailures = 0;
reportProgress();
finishJob();
}
@ -481,11 +886,9 @@ export const createWorkerPool = (
if (!settled) {
settled = true;
cleanup();
void fail(
new WorkerPoolDispatchError(
`Worker ${workerIndex} error: ${err.message}`,
inFlightExcludePath(job, lastProgress),
),
void recoverAndResume(
`Worker ${workerIndex} error: ${err.message}`,
resolveExcludePaths(),
);
}
};
@ -494,14 +897,12 @@ export const createWorkerPool = (
if (!settled) {
settled = true;
cleanup();
const excludes = inFlightExcludePath(job, lastProgress);
const excludes = resolveExcludePaths();
const inFlightSuffix = excludes.length > 0 ? ` (in-flight: ${excludes[0]})` : '';
void fail(
new WorkerPoolDispatchError(
`Worker ${workerIndex} exited with code ${code}. ` +
`Likely OOM or native addon failure${inFlightSuffix}.`,
excludes,
),
void recoverAndResume(
`Worker ${workerIndex} exited with code ${code}. ` +
`Likely OOM or native addon failure${inFlightSuffix}.`,
excludes,
);
}
};
@ -517,14 +918,20 @@ export const createWorkerPool = (
worker.postMessage({ type: 'sub-batch', files: job.items });
};
for (let i = 0; i < workers.length; i++) runWorker(i);
for (const slotIndex of activeSlots) runWorker(slotIndex);
});
};
const terminate = async (): Promise<void> => {
await Promise.all(workers.map((w) => w.terminate()));
await Promise.all(workers.map((w) => w?.terminate()));
workers.length = 0;
activeSlots.clear();
};
return { dispatch, terminate, size };
return {
dispatch,
terminate,
size,
getQuarantinedPaths: () => Array.from(quarantined),
};
};

View file

@ -0,0 +1,344 @@
import { describe, expect, it, vi, beforeEach } 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,
WorkerPoolDispatchError,
resolveWorkerPoolOptions,
} from '../../src/core/ingestion/workers/worker-pool.js';
/**
* Minimal `node:worker_threads` Worker double for unit-testing the pool's
* resilience layers (auto-respawn, circuit breaker, quarantine, retry
* budget). Tests script behaviour via `nextActions`: each action runs on
* the next dispatched sub-batch postMessage. `'crash'` and `'exit'` mimic
* real worker failures; `'parse-ok'` mimics a healthy completion.
*/
type FakeWorkerAction =
| { kind: 'parse-ok'; files: { path: string }[]; result?: unknown }
| { kind: 'crash-exit'; code: number; afterStartingFiles?: number }
| { kind: 'crash-error'; message: string; afterStartingFiles?: number };
const nextActions: FakeWorkerAction[] = [];
let workerInstances: FakeWorker[] = [];
class FakeWorker extends EventEmitter {
readonly seenMessages: unknown[] = [];
constructor() {
super();
workerInstances.push(this);
// Real Worker fires 'online' asynchronously after the runtime is ready;
// replicate so `waitForWorkerOnline` resolves.
queueMicrotask(() => this.emit('online'));
}
postMessage(msg: unknown): void {
this.seenMessages.push(msg);
if (typeof msg !== 'object' || msg === null) return;
const m = msg as { type?: string; files?: { path: string }[] };
if (m.type !== 'sub-batch') return;
const action = nextActions.shift();
if (!action) {
// No script set; behave as a hung worker (no reply) — the idle timer
// will eventually fire. Tests should always script enough actions.
return;
}
queueMicrotask(() => this.runAction(action, m.files ?? []));
}
private async runAction(action: FakeWorkerAction, files: { path: string }[]): Promise<void> {
if (action.kind === 'parse-ok') {
for (const file of action.files) {
this.emit('message', { type: 'starting-file', path: file.path });
}
this.emit('message', { type: 'progress', filesProcessed: action.files.length });
this.emit('message', { type: 'sub-batch-done' });
// sub-batch-done triggers the pool to post {type:'flush'} which we
// ignore in postMessage above (only 'sub-batch' triggers actions).
// For the result, wait one microtask so the flush is observed.
await Promise.resolve();
this.emit('message', {
type: 'result',
data: action.result ?? { fileCount: action.files.length },
});
return;
}
if (action.kind === 'crash-exit') {
const upTo = Math.min(action.afterStartingFiles ?? 0, files.length);
for (let i = 0; i < upTo; i++) {
this.emit('message', { type: 'starting-file', path: files[i].path });
}
this.emit('exit', action.code);
return;
}
if (action.kind === 'crash-error') {
const upTo = Math.min(action.afterStartingFiles ?? 0, files.length);
for (let i = 0; i < upTo; i++) {
this.emit('message', { type: 'starting-file', path: files[i].path });
}
this.emit('error', new Error(action.message));
return;
}
}
async terminate(): Promise<number> {
this.emit('exit', 0);
return 0;
}
removeListener(event: string | symbol, listener: (...args: unknown[]) => void): this {
return super.removeListener(event, listener);
}
}
// Create a real on-disk worker script so createWorkerPool's existsSync gate
// passes. The script is never actually executed because we inject
// FakeWorker via workerFactory; it just has to exist as a file path.
let tempDir: string;
let workerUrl: URL;
beforeEach(() => {
nextActions.length = 0;
workerInstances = [];
tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-pool-resilience-'));
const workerPath = path.join(tempDir, 'fake-worker.js');
fs.writeFileSync(workerPath, '// fake');
workerUrl = pathToFileURL(workerPath) as URL;
});
describe('worker pool resilience', () => {
it('seeds an empty quarantine on a fresh pool', () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker,
});
expect(pool.getQuarantinedPaths()).toEqual([]);
void pool.terminate();
});
it('quarantines the in-flight file on worker exit and respawns the slot', async () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker,
consecutiveFailureThreshold: 5,
maxRespawnsPerSlot: 3,
});
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({
kind: 'parse-ok',
files: [{ path: 'src/good.ts' }],
result: { fileCount: 1 },
});
const results = await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/bad.ts', content: '' },
{ path: 'src/good.ts', content: '' },
]);
expect(results).toEqual([{ fileCount: 1 }]);
expect(pool.getQuarantinedPaths()).toEqual(['src/bad.ts']);
// First FakeWorker died; second is the respawn. Total = 2.
expect(workerInstances.length).toBe(2);
await pool.terminate();
});
it('drops a slot after maxRespawnsPerSlot exceeded and continues on other slots', async () => {
const pool = createWorkerPool(workerUrl, 2, {
workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker,
consecutiveFailureThreshold: 10,
maxRespawnsPerSlot: 1,
});
// Slot 0 dies twice, exceeding budget=1; slot 1 succeeds with the
// requeued remainder.
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({
kind: 'parse-ok',
files: [{ path: 'src/c.ts' }, { path: 'src/d.ts' }],
result: { fileCount: 2 },
});
const results = await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/a.ts', content: '' },
{ path: 'src/b.ts', content: '' },
{ path: 'src/c.ts', content: '' },
{ path: 'src/d.ts', content: '' },
]);
expect(results).toEqual([{ fileCount: 2 }]);
// Two bad files quarantined; both pre-crash 'starting-file' targets.
expect(pool.getQuarantinedPaths().sort()).toEqual(['src/a.ts', 'src/b.ts']);
await pool.terminate();
});
it('trips the circuit breaker after consecutiveFailureThreshold deaths', async () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker,
consecutiveFailureThreshold: 2,
maxRespawnsPerSlot: 5,
});
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
await expect(
pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/x.ts', content: '' },
{ path: 'src/y.ts', content: '' },
]),
).rejects.toBeInstanceOf(WorkerPoolDispatchError);
expect(pool.getQuarantinedPaths().sort()).toEqual(['src/x.ts', 'src/y.ts']);
// Subsequent dispatch on a tripped pool rejects without running anything.
await expect(
pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/z.ts', content: '' },
]),
).rejects.toBeInstanceOf(WorkerPoolDispatchError);
await pool.terminate();
});
it('resets consecutive-failure counter on a successful job', async () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker,
consecutiveFailureThreshold: 2,
maxRespawnsPerSlot: 5,
});
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({
kind: 'parse-ok',
files: [{ path: 'src/recovered.ts' }],
result: { fileCount: 1 },
});
const r1 = await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/bad.ts', content: '' },
{ path: 'src/recovered.ts', content: '' },
]);
expect(r1).toEqual([{ fileCount: 1 }]);
// Second dispatch: another death. Counter was reset by the prior success,
// so this single failure should not trip the breaker (threshold=2).
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({
kind: 'parse-ok',
files: [{ path: 'src/ok.ts' }],
result: { fileCount: 1 },
});
const r2 = await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/bad2.ts', content: '' },
{ path: 'src/ok.ts', content: '' },
]);
expect(r2).toEqual([{ fileCount: 1 }]);
expect(pool.getQuarantinedPaths().sort()).toEqual(['src/bad.ts', 'src/bad2.ts']);
await pool.terminate();
});
it('filters already-quarantined paths from new dispatches', async () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker,
consecutiveFailureThreshold: 5,
maxRespawnsPerSlot: 3,
});
// First dispatch: quarantine 'src/poison.ts'
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({
kind: 'parse-ok',
files: [{ path: 'src/a.ts' }],
result: { fileCount: 1 },
});
await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/poison.ts', content: '' },
{ path: 'src/a.ts', content: '' },
]);
expect(pool.getQuarantinedPaths()).toEqual(['src/poison.ts']);
// Second dispatch including the quarantined file: pool filters before
// workers see it. The action should never be popped because the only
// dispatchable item is src/b.ts.
nextActions.push({
kind: 'parse-ok',
files: [{ path: 'src/b.ts' }],
result: { fileCount: 1 },
});
const results = await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/poison.ts', content: '' },
{ path: 'src/b.ts', content: '' },
]);
expect(results).toEqual([{ fileCount: 1 }]);
// The most recent sub-batch the pool dispatched is the dispatch-2
// payload. With poison already in the quarantine when dispatch 2 ran,
// the pool must have filtered it out before reaching a worker.
const allSubBatches = workerInstances
.flatMap((w) => w.seenMessages)
.filter(
(m): m is { type: string; files: { path: string }[] } =>
typeof m === 'object' && m !== null && (m as { type?: string }).type === 'sub-batch',
);
const lastSubBatch = allSubBatches[allSubBatches.length - 1];
expect(lastSubBatch.files.map((f) => f.path)).toEqual(['src/b.ts']);
await pool.terminate();
});
it('returns an empty result without dispatching when every item is quarantined', async () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new FakeWorker() as unknown as import('node:worker_threads').Worker,
consecutiveFailureThreshold: 5,
maxRespawnsPerSlot: 3,
});
nextActions.push({ kind: 'crash-exit', code: 134, afterStartingFiles: 1 });
nextActions.push({
kind: 'parse-ok',
files: [{ path: 'src/a.ts' }],
result: { fileCount: 1 },
});
await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/poison.ts', content: '' },
{ path: 'src/a.ts', content: '' },
]);
const baselineWorkers = workerInstances.length;
const results = await pool.dispatch<{ path: string; content: string }, unknown>([
{ path: 'src/poison.ts', content: '' },
]);
expect(results).toEqual([]);
expect(workerInstances.length).toBe(baselineWorkers);
await pool.terminate();
});
});
describe('worker pool option resolution', () => {
it('resolves maxRespawnsPerSlot from explicit options', () => {
const opts = resolveWorkerPoolOptions({ maxRespawnsPerSlot: 7 }, 4);
expect(opts.maxRespawnsPerSlot).toBe(7);
});
it('defaults consecutiveFailureThreshold to max(3, poolSize)', () => {
expect(resolveWorkerPoolOptions({}, 1).consecutiveFailureThreshold).toBe(3);
expect(resolveWorkerPoolOptions({}, 8).consecutiveFailureThreshold).toBe(8);
});
it('defaults maxCumulativeTimeoutMs to 5x subBatchIdleTimeoutMs', () => {
const opts = resolveWorkerPoolOptions({ subBatchIdleTimeoutMs: 1000 }, 1);
expect(opts.maxCumulativeTimeoutMs).toBe(5000);
});
it('reads GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT env override', () => {
vi.stubEnv('GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT', '2');
try {
expect(resolveWorkerPoolOptions({}, 1).maxRespawnsPerSlot).toBe(2);
} finally {
vi.unstubAllEnvs();
}
});
it('reads GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD env override', () => {
vi.stubEnv('GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD', '12');
try {
expect(resolveWorkerPoolOptions({}, 1).consecutiveFailureThreshold).toBe(12);
} finally {
vi.unstubAllEnvs();
}
});
});