From 3ffd6ad396c0e10fb12bf8176d56d45ad7c9ed5a Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Tue, 19 May 2026 09:28:02 +0100 Subject: [PATCH] feat(workers): resilient + scalable worker pool MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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` 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. --- .../src/core/ingestion/parsing-processor.ts | 24 +- .../core/ingestion/workers/parse-worker.ts | 7 + .../src/core/ingestion/workers/worker-pool.ts | 563 +++++++++++++++--- .../test/unit/worker-pool-resilience.test.ts | 344 +++++++++++ 4 files changed, 859 insertions(+), 79 deletions(-) create mode 100644 gitnexus/test/unit/worker-pool-resilience.test.ts diff --git a/gitnexus/src/core/ingestion/parsing-processor.ts b/gitnexus/src/core/ingestion/parsing-processor.ts index 73a72aa22..a276a0ab4 100644 --- a/gitnexus/src/core/ingestion/parsing-processor.ts +++ b/gitnexus/src/core/ingestion/parsing-processor.ts @@ -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; diff --git a/gitnexus/src/core/ingestion/workers/parse-worker.ts b/gitnexus/src/core/ingestion/workers/parse-worker.ts index da681b070..2166ca950 100644 --- a/gitnexus/src/core/ingestion/workers/parse-worker.ts +++ b/gitnexus/src/core/ingestion/workers/parse-worker.ts @@ -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