fix(workers,tests,docs): apply ce-code-review findings (16 items)

Walks the full set of findings from a multi-agent code review (11
reviewers, 1 maintainability dispatch lost to tool-permission denial)
of the PR #1693 branch. All 16 actionable findings — 4 P1, 4 P2,
8 P3 — applied in a single pass against a consistent tree. Tests
pass (269/269 unit files, 29/29 integration).

P1 — bounds-only / disguised-bounds assertions across 4 test files
(per user-memory DoD §2.7):
  - worker-pool.test.ts: 5 sites — `nodes.length > 0` dropped (redundant
    after `.toContain('validateInput')`); `files.length >= 4` pinned to
    `.toBe(7)` (mini-repo/src has exactly 7 .ts files); `results.length
    > 0` pinned to `.toHaveLength(1)` (default sub-batch absorbs all 7);
    `result.fileCount >= 0` pinned to `.toBe(1)` (empty file is still
    "processed"); `warnRecords.length > 0` replaced with content-
    predicate `/respawn|dropping|replacement|did not report ready/`
    (catches silenced warnings); `fallbackExcludePaths.length > 0`
    pinned to exact `['one.ts', 'two.ts']` (deterministic given the
    single-slot pool + 2 items + per-item starting-file).
  - parse-impl-fallback.test.ts: 3 sites — `astCacheClearCalls >= 1`
    pinned to exact 4 (per-chunk × 2 + finally × 2); the two error-path
    delta checks pinned to exact +2 and +3 (verified empirically).
  - parse-impl-progress-monotonic.test.ts: `percents.length > 0` →
    `.not.toEqual([])`; per-element `Math.max(prev, cur)` tautology
    replaced with direct `if (cur < prev) throw`; final-percent
    `Math.min(last, 95)` tautology pinned to exact `.toBe(70)` (3-file
    skipWorkers fixture's deferred band lands at the band start).
  - parse-impl-large-fixture.test.ts: `Math.min(elapsedMs, BUDGET)`
    tautology removed; Promise.race rejection is the load-bearing
    wall-clock check.

P1 — terminate() lacks `.catch` mask:
  - worker-pool.ts terminate() now matches the `.catch(() => undefined)`
    pattern used at every other internal terminate site. Prevents a
    hung/OOM worker's terminate rejection from masking the original
    pipeline error when called from parse-impl.ts's finally block, and
    guarantees `workers.length = 0` / `activeSlots.clear()` always run.

P1 — hybrid envelope length-mismatch + null-payload silent data loss:
  - parse-worker.ts decodeIncomingMessage: explicit non-null-and-typed
    check before `.type` access (decodeMessage permits null payloads
    per encodeMessage contract); explicit length-equality assertion
    between `decoded.files` and `contents` before zipping. Without
    these, `TextDecoder.decode(undefined)` silently returns "" and
    produces empty-content graph nodes — a contract violation that
    used to be undetectable. Both throws route through the outer
    try/catch → worker `error` reply → pool's recoverAndResume.

P1 — unsafe casts at the IPC boundary:
  - buildDispatchMessage now uses a properly-typed `isParseWorkerItemArray`
    type guard. The narrowed branch accesses `item.path` and
    `item.content` as statically-typed strings — a future rename of
    `ParseWorkerInput.content` would fail to compile inside the branch
    instead of silently mismatching at runtime. The remaining
    decodeMessage payload casts are bounded by the F3/F6 runtime
    guards.

P2 — idle-timeout retry bypasses circuit breaker:
  - worker-pool.ts timeout-retry IIFE now increments
    `consecutiveFailuresPerSlot[workerIndex]` alongside `respawnCount`.
    A slot that consistently times out (vs crashes) now trips the
    per-slot breaker, instead of consuming its full respawn budget
    over potentially tens of minutes without the breaker firing.

P2 — null/non-object worker message crashes pool handler:
  - Dispatch handler in worker-pool.ts now guards `null /
    non-object / no string type discriminant` before `msg.type` access
    and routes through recoverAndResume on violation. Previously a
    legitimate `null` payload would throw TypeError out of the
    EventEmitter listener → uncaughtException on main, crashing the
    analyze.

P2 — workerPoolSize === 0 creates unusable pool:
  - parse-impl.ts now treats `workerPoolSize === 0` as `skipWorkers`
    at the gate. Matches the PipelineOptions docstring contract ("0
    disables the pool entirely — equivalent to skipWorkers"); avoids
    constructing a pool that rejects every dispatch and logs
    "Worker pool parsing stopped" per chunk.

P2 — encodeMessage 2-buffer allocation per frame:
  - protocol.ts encodeMessage coalesced to a single
    `Buffer.allocUnsafe + writeUInt8 + writeUInt32LE + buf.write
    (string, offset, 'utf8')`. Drops the intermediate
    `Buffer.from(JSON.stringify(...), 'utf8')` allocation + memcpy.
    Length pre-check via `Buffer.byteLength(string, 'utf8')` surfaces
    the uint32 cap before any allocation.

P3 — slotGenerations made optional on WorkerPoolStats so external
  implementations of getStats() that predate U12 don't compile-break;
  in-repo callers already use optional chaining.

P3 — buildDispatchMessage marked `@internal` so it isn't surfaced as
  public API by typedoc / api-extractor (it's a test-only export).

P3 — verboseThroughputLog hoisted above the chunk loop (env vars can't
  change mid-run; one O(env-read) per analyze, not per chunk).

P3 — corrected the messageerror routing comment in worker-pool.ts
  dispatch handler. `ProtocolDecodeError` is caught by the surrounding
  try/catch — distinct from `messageerror`, which fires for V8
  structured-clone failures before the message body would reach the
  handler.

P3 — initial pool spawn now uses a `Promise.allSettled` ready-handshake
  gate symmetric with `replaceWorker`. Dispatch awaits this gate before
  selecting slots, so an init-crashing initial worker is dropped from
  `activeSlots` and a downstream OOM/missing-native-binding failure
  surfaces in seconds (bounded by WORKER_READY_TIMEOUT_MS) rather than
  waiting for the first idle timeout (30s default).

P3 — `GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT`,
  `GITNEXUS_WORKER_MAX_CUMULATIVE_TIMEOUT_MS`,
  `GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD` added to:
    - CLI `--help` text in src/cli/index.ts
    - Root README env-var table
    - gitnexus/README troubleshooting section (new "Worker pool
      resilience tuning" subsection)

P3 — CLI `catch (e: any)` / `catch (err: any)` in analyze.ts replaced
  with `catch (err: unknown)` + narrowed access; matches modern TS
  best practice and the codebase pattern at other catch sites.

P3 — `WorkerPoolStats.terminated: boolean` field added (optional, for
  backward compatibility). `terminate()` sets it true; `getStats()`
  surfaces it. Distinguishes graceful shutdown from a circuit-breaker
  trip in observability surfaces.

Coverage / advisory items not addressed in this commit (kept in the
report only):
  - maintainability reviewer failed (Read/Bash denied) — god-module
    audit on worker-pool.ts (~1400 LOC) carried as residual risk
  - quarantine case-sensitivity contract unpinned (adversarial #8)
  - WORKER_READY_TIMEOUT_MS env-configurability (adversarial #2)
  - chunk-byte-budget × parseChunkConcurrency memory multiplier doc
    (adversarial #5)
  - MCP discoverability gaps for env vars / verbose (agent-native W1/W2)
  - bench/parse-throughput.md scaffold-with-TBD-rows (PS RR-003)
This commit is contained in:
Gergo Magyar 2026-05-20 12:54:43 +01:00
parent 27d750b86d
commit a2878df5fb
13 changed files with 316 additions and 81 deletions

View file

@ -239,6 +239,9 @@ Most `analyze` knobs are also CLI flags (`--workers`, `--worker-timeout`, `--max
| `GITNEXUS_MAX_FILE_SIZE` | `512` (KB) | Walker skip threshold in KB. Hard cap is `32768` (tree-sitter buffer ceiling). Equivalent to `--max-file-size <kb>`. | Indexing repos with intentionally-large source files (generated parsers, vendored bundles) that should still be parsed. |
| `GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS` | `30000` | Worker idle timeout in milliseconds before retry/fallback. Equivalent to `--worker-timeout <seconds>` × 1000. | Slow-parsing files (large minified JS, deeply-nested TS types) that legitimately need more than 30s. |
| `GITNEXUS_WORKER_SUB_BATCH_MAX_BYTES` | `8388608` (8 MB) | Per-job byte budget the pool will send to a worker in one `postMessage`. | Very large individual files; mostly diagnostic — bumping past 8 MB risks structured-clone memory pressure. |
| `GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT` | `3` | Max replacement spawns per worker slot before the slot is dropped from the active rotation. Bounds respawn loops on a chronically-crashing slot. | Hosts where a flaky worker should retry more (raise) or fail-fast (lower) before the slot is dropped. |
| `GITNEXUS_WORKER_MAX_CUMULATIVE_TIMEOUT_MS` | `5 × subBatchTimeoutMs` | Total retry wall-time budget per job before quarantining. Combined with `timeoutBackoffFactor`, prevents exponentially-growing retries from stalling for hours. | Slow files that legitimately need long total retry windows; lower to fail-fast on stalls. |
| `GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD`| `max(3, poolSize)` | Per-slot consecutive deaths before the pool's circuit breaker trips. After tripping, every subsequent dispatch rejects until a fresh pool is created. | Hosts where a SIGSEGV-prone native grammar should trip the breaker sooner; CI runners that should fail loudly. |
| `GITNEXUS_CHUNK_BYTE_BUDGET` | `2097152` (2 MB) | Chunk boundary used for cache-key composition and dispatch. Smaller = finer-grained cache hits but more dispatch overhead. | Tuning incremental-analyze cache behavior on monorepos. |
| `GITNEXUS_NO_GITIGNORE` | unset | When set, skips `.gitignore` parsing. `.gitnexusignore` is still honored. | Indexing a repo whose `.gitignore` excludes files you actually want indexed (e.g., generated code committed for cross-repo lookup). |
| `GITNEXUS_SKIP_OPTIONAL_GRAMMARS` | unset | When `=1` strictly, skips native builds for `tree-sitter-dart` / `tree-sitter-proto` at install time. | Installing on a host without a C++ toolchain; you're willing to skip Dart/Proto parsing. |

View file

@ -358,6 +358,16 @@ npx gitnexus analyze
For repositories with very large source files, `GITNEXUS_WORKER_SUB_BATCH_MAX_BYTES` controls the worker job byte budget. The default is **8388608 bytes (8 MB)**.
### Worker pool resilience tuning
Three env vars expose the pool's resilience layers (respawn budget, cumulative-timeout cap, circuit breaker). Defaults are tuned for typical repos; bump them when an analyze legitimately needs more retries, or lower them to fail-fast on a known-bad shape.
| Variable | Default | Effect |
| ------------------------------------------------- | ------------------------- | --------------------------------------------------------------------------------------------------------------------------------- |
| `GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT` | `3` | Max replacement spawns per slot before the slot is dropped from the active rotation. |
| `GITNEXUS_WORKER_MAX_CUMULATIVE_TIMEOUT_MS` | `5 × subBatchTimeoutMs` | Total retry wall-time budget per job before quarantining. Bounds exponentially-growing retry waits. |
| `GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD` | `max(3, poolSize)` | Per-slot consecutive deaths before the pool's circuit breaker trips. After tripping, dispatches require a fresh pool. |
## Privacy
- All processing happens locally on your machine

View file

@ -148,7 +148,7 @@ function ensureHeap(): boolean {
stdio: 'inherit',
env: { ...process.env, NODE_OPTIONS: `${nodeOpts} ${HEAP_FLAG}`.trim() },
});
} catch (e: any) {
} catch (e: unknown) {
if (childProcessLikelyOom(e)) {
cliError(
` Analysis likely ran out of memory.\n` +
@ -159,7 +159,11 @@ function ensureHeap(): boolean {
{ recoveryHint: 'heap-oom-respawn' },
);
}
process.exitCode = e.status ?? 1;
const status =
typeof e === 'object' && e !== null && 'status' in e && typeof e.status === 'number'
? e.status
: 1;
process.exitCode = status;
}
return true;
}
@ -740,7 +744,7 @@ const analyzeCommandImpl = async (inputPath?: string, options?: AnalyzeOptions):
}
console.log('');
} catch (err: any) {
} catch (err: unknown) {
clearInterval(elapsedTimer);
process.removeListener('SIGINT', sigintHandler);
console.log = origLog;
@ -750,7 +754,7 @@ const analyzeCommandImpl = async (inputPath?: string, options?: AnalyzeOptions):
console.error = origError;
bar.stop();
const msg = err.message || String(err);
const msg = err instanceof Error ? err.message : String(err);
// Registry name-collision from --name (#829) — surface as an
// actionable error rather than a generic stack-trace.

View file

@ -87,6 +87,9 @@ program
' GITNEXUS_WORKER_SUB_BATCH_MAX_BYTES=N Worker job byte budget. Default 8388608.\n' +
' GITNEXUS_WORKER_POOL_SIZE=N Parse worker count override. Default cores-1 capped at 16.\n' +
' GITNEXUS_PARSE_CHUNK_CONCURRENCY=N Concurrent in-flight parse chunks. Default 2.\n' +
' GITNEXUS_WORKER_MAX_RESPAWNS_PER_SLOT=N Max replacement spawns per slot before drop. Default 3.\n' +
' GITNEXUS_WORKER_MAX_CUMULATIVE_TIMEOUT_MS=N Total retry wall-time per job. Default 5x sub-batch timeout.\n' +
' GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD=N Per-slot deaths to trip circuit breaker. Default max(3, poolSize).\n' +
' GITNEXUS_EMBEDDING_THREADS=N Limit local ONNX CPU threads for --embeddings.\n' +
' GITNEXUS_SEMANTIC_EXACT_SCAN_LIMIT=N Max embedding chunks for exact-scan fallback. Default 10000.\n' +
'\nTip: `.gitnexusignore` supports `.gitignore`-style negation. Add e.g.\n' +

View file

@ -252,10 +252,18 @@ export async function runChunkedParseAndResolve(
const MIN_BYTES_FOR_WORKERS = options?.workerThresholdsForTest?.minBytes ?? 512 * 1024;
const totalBytes = parseableScanned.reduce((s, f) => s + f.size, 0);
// Create worker pool once, reuse across chunks
// Create worker pool once, reuse across chunks.
//
// `workerPoolSize === 0` is a programmatic equivalent of `skipWorkers:
// true` per the `PipelineOptions.workerPoolSize` contract. Short-
// circuiting here avoids constructing a useless pool that rejects
// every dispatch (with a `Worker pool parsing stopped` warn log per
// chunk) just to fall back to the sequential path via the error
// catch — the gate honors the docstring directly.
let workerPool: WorkerPool | undefined;
if (
!options?.skipWorkers &&
options?.workerPoolSize !== 0 &&
(totalParseable >= MIN_FILES_FOR_WORKERS || totalBytes >= MIN_BYTES_FOR_WORKERS)
) {
try {
@ -386,16 +394,21 @@ export async function runChunkedParseAndResolve(
startChunkPrefetch(i);
}
// Hoisted loop-invariant: GITNEXUS_VERBOSE / NODE_ENV are read once
// (not on every chunk). Previously evaluated at the top of the loop
// body, which re-read process.env on every iteration even though
// the env can't change mid-run.
const verboseThroughputLog = isDev || isVerboseIngestionEnabled();
for (let chunkIdx = 0; chunkIdx < numChunks; chunkIdx++) {
const chunkPaths = chunks[chunkIdx];
// Start wall-clock for the per-chunk throughput log emitted at end
// of this iteration. Computed when either NODE_ENV=development OR
// the operator passed `--verbose` (GITNEXUS_VERBOSE) — the previous
// `isDev`-only gate meant operators running `gitnexus analyze
// --verbose` in production never saw the log (M3 from PR #1693
// review). Timestamp is cheap; the log line only fires under the
// same combined gate below.
const verboseThroughputLog = isDev || isVerboseIngestionEnabled();
// of this iteration. The gate is computed once above; here we just
// sample the clock if the gate is on. Computed when either
// NODE_ENV=development OR the operator passed `--verbose`
// (GITNEXUS_VERBOSE) — the previous `isDev`-only gate meant
// operators running `gitnexus analyze --verbose` in production
// never saw the log (M3 from PR #1693 review).
const chunkStartMs: number | null = verboseThroughputLog ? Date.now() : null;
const chunkContents = await chunkContentPromises[chunkIdx]!;

View file

@ -2469,6 +2469,12 @@ const mergeResult = (target: ParseWorkerResult, src: ParseWorkerResult) => {
// strict — every outgoing message is encoded.
parentPort!.postMessage(encodeMessage(MessageTag.Ready, { type: 'ready' }));
// Module-scope `TextDecoder` for the U19 hybrid-envelope path. Hoisted
// out of `decodeIncomingMessage` so we don't allocate a new ICU-backed
// decoder per sub-batch — `TextDecoder.decode()` is stateless across
// calls and safe to share.
const sharedHybridDecoder = new TextDecoder('utf-8');
// Decode a single incoming message. The pool always sends Buffer-encoded
// frames post-U17, but the parameter type is `unknown` because Node's
// worker_threads typings declare `on('message', (value: any) => void)`.
@ -2502,12 +2508,37 @@ function decodeIncomingMessage(raw: unknown): WorkerIncomingMessage {
) {
const envelope = (raw as { envelope: Uint8Array }).envelope;
const contents = (raw as { contents: Uint8Array[] }).contents;
const decoded = decodeMessage(envelope).payload as {
const decodedPayload = decodeMessage(envelope).payload;
// The protocol allows `null` payloads (per `encodeMessage` docs).
// Reject them here rather than letting `.type` access throw a
// TypeError that escapes uncaught — the outer try/catch routes the
// thrown Error back to the pool as a worker `error` reply, which
// goes through `recoverAndResume` instead of producing silently-
// empty graph data.
if (
decodedPayload === null ||
typeof decodedPayload !== 'object' ||
typeof (decodedPayload as { type?: unknown }).type !== 'string'
) {
throw new Error('hybrid envelope decode produced a non-object or non-typed payload');
}
const decoded = decodedPayload as {
type: string;
files: Array<{ path: string; byteLength: number }>;
};
if (decoded.type === 'sub-batch' && Array.isArray(decoded.files)) {
const decoder = new TextDecoder('utf-8');
// Length-equality assertion. Without it, `TextDecoder.decode(undefined)`
// silently returns `""` and produces zero-content graph nodes — a
// silent data-loss failure mode. Throwing routes through the
// outer try/catch -> worker `error` reply -> `recoverAndResume`,
// making the contract violation observable instead of masking it.
if (decoded.files.length !== contents.length) {
throw new Error(
`hybrid envelope contract violation: files.length=${decoded.files.length} ` +
`but contents.length=${contents.length}`,
);
}
const decoder = sharedHybridDecoder;
const files: ParseWorkerInput[] = decoded.files.map((meta, i) => ({
path: meta.path,
content: decoder.decode(contents[i]),

View file

@ -99,16 +99,24 @@ function isValidTag(byte: number): byte is MessageTagValue {
* explicit instead of silently truncating.
*/
export function encodeMessage(tag: MessageTagValue, payload: unknown): Buffer {
const json = Buffer.from(JSON.stringify(payload ?? null), 'utf8');
if (json.length > 0xffffffff) {
throw new RangeError(
`protocol payload exceeds uint32 length cap (${json.length} > ${0xffffffff})`,
);
const json = JSON.stringify(payload ?? null);
// `Buffer.byteLength(string, 'utf8')` returns the encoded byte length
// without allocating an intermediate Buffer. Pre-checking it lets us
// surface the uint32 cap as a RangeError before any allocation.
const length = Buffer.byteLength(json, 'utf8');
if (length > 0xffffffff) {
throw new RangeError(`protocol payload exceeds uint32 length cap (${length} > ${0xffffffff})`);
}
const buf = Buffer.alloc(PROTOCOL_HEADER_BYTES + json.length);
// Single allocation, single write pass: `buf.write(string, offset,
// 'utf8')` writes UTF-8 bytes directly into the target without an
// intermediate `Buffer.from(string, 'utf8')` allocation + memcpy.
// Halves the per-frame allocation count vs the previous two-Buffer +
// copy approach — material under high message volume (every dispatch
// envelope, every worker reply).
const buf = Buffer.allocUnsafe(PROTOCOL_HEADER_BYTES + length);
buf.writeUInt8(tag, 0);
buf.writeUInt32LE(json.length, 1);
json.copy(buf, PROTOCOL_HEADER_BYTES);
buf.writeUInt32LE(length, 1);
buf.write(json, PROTOCOL_HEADER_BYTES, 'utf8');
return buf;
}

View file

@ -55,30 +55,52 @@ function decodeIncomingWorkerMessage(raw: unknown): WorkerOutgoingMessage {
* and avoiding transfer here means we can't accidentally detach an
* unrelated Buffer that shares the pool.
*/
type ParseWorkerItem = { path: string; content: string };
/**
* Type guard: every element of `items` has the parse-worker shape
* (`{path: string, content: string}`). Used to narrow the generic input
* inside `buildDispatchMessage` without an `as unknown as Array<...>`
* double-cast, so a future rename of `ParseWorkerInput.content` would
* fail to compile inside the narrowed branch instead of silently
* mismatching at runtime.
*/
function isParseWorkerItemArray<T>(
items: readonly T[],
): items is readonly T[] & readonly ParseWorkerItem[] {
if (items.length === 0) return false;
for (const it of items) {
if (it == null || typeof it !== 'object') return false;
if (typeof (it as { path?: unknown }).path !== 'string') return false;
if (typeof (it as { content?: unknown }).content !== 'string') return false;
}
return true;
}
/**
* @internal Exported only so the unit test suite
* (`test/unit/worker-pool-transferlist.test.ts`) can pin the return-
* shape and transferList contract directly. Not part of the public
* package API; the return-type union and shape-detection heuristic may
* change without a semver signal.
*/
export function buildDispatchMessage<T>(items: readonly T[]): {
message: Uint8Array | { envelope: Uint8Array; contents: Uint8Array[] };
transferList?: ArrayBuffer[];
} {
const isParseWorkerShape =
items.length > 0 &&
items.every(
(it) =>
it != null &&
typeof it === 'object' &&
typeof (it as { path?: unknown }).path === 'string' &&
typeof (it as { content?: unknown }).content === 'string',
);
if (!isParseWorkerShape) {
if (!isParseWorkerItemArray(items)) {
return {
message: encodeMessage(MessageTag.DispatchJob, { type: 'sub-batch', files: items }),
};
}
// After the type guard, `items` is narrowed to `readonly ParseWorkerItem[]`.
// No cast needed — `item.path` and `item.content` are statically typed
// as strings inside this branch.
const encoder = new TextEncoder();
const filesMeta: Array<{ path: string; byteLength: number }> = [];
const contents: Uint8Array[] = [];
for (const item of items as unknown as Array<{ path: string; content: string }>) {
for (const item of items) {
const u8 = encoder.encode(item.content);
filesMeta.push({ path: item.path, byteLength: u8.byteLength });
contents.push(u8);
@ -155,13 +177,22 @@ export interface WorkerPoolStats {
/** Whether the circuit breaker has tripped (no further dispatches
* will be accepted by this pool instance). */
readonly poolBroken: boolean;
/** Whether `terminate()` has been called on this pool. Distinguishes
* graceful shutdown (terminated=true, activeSlots=0) from a circuit-
* breaker trip (terminated=false, poolBroken=true, activeSlots=0).
* Optional for backward compatibility with external `WorkerPoolStats`
* implementations that predate this field. */
readonly terminated?: boolean;
/** Per-slot generation counter (U12). Increments by 1 on every
* successful worker replacement for that slot. Operators / tests
* observe this to confirm a death-then-respawn actually happened
* vs. the same worker being recycled in place. Initial value is 0
* for every slot at pool creation; dropped slots keep their last
* generation (they don't decrement). */
readonly slotGenerations: readonly number[];
* generation (they don't decrement). Optional so external
* `WorkerPoolStats` implementations that predate U12 can omit the
* field without a TypeScript compile error — in-repo callers use
* optional chaining (`stats?.slotGenerations`) consistently. */
readonly slotGenerations?: readonly number[];
}
export interface WorkerPoolOptions {
@ -599,23 +630,65 @@ export const createWorkerPool = (
activeSlots.add(i);
}
const dispatch = <TInput, TResult>(
// Symmetrize the readiness gate across initial and replacement spawn
// paths. `replaceWorker` already awaits `waitForWorkerReady` per
// replacement so an init-crashing worker is dropped before dispatch
// sees it. The initial-spawn loop above didn't — a worker whose
// top-of-script init crashes (failed tree-sitter native binding,
// missing dependency) would only be noticed at the first dispatch's
// 30s idle timeout, vs the 5s WORKER_READY_TIMEOUT_MS bound that
// replacements enjoy.
//
// The promise below settles every initial slot in parallel and drops
// unready slots from `activeSlots` before any dispatch can fire.
// `dispatch` awaits it via `initialReadyGate` on first invocation.
// Wrapped in a single `Promise.allSettled` so a slow worker doesn't
// block ready workers from being usable — first dispatch waits for
// all slots' verdicts (good or bad).
const initialReadyGate: Promise<void> = Promise.allSettled(
workers.map(async (w, i) => {
if (!w) return;
try {
await waitForWorkerReady(w);
} catch (err) {
logger.warn(
{
workerIndex: i,
err: err instanceof Error ? err.message : String(err),
},
`Worker ${i} did not report ready on initial spawn; dropping slot.`,
);
await w.terminate().catch(() => undefined);
workers[i] = undefined;
activeSlots.delete(i);
}
}),
).then(() => undefined);
const dispatch = async <TInput, TResult>(
items: TInput[],
onProgress?: (filesProcessed: number) => void,
): Promise<TResult[]> => {
// Await the initial-spawn readiness gate (F13). On first dispatch
// this blocks for up to WORKER_READY_TIMEOUT_MS while every initial
// worker's `{type:'ready'}` handshake is checked; on subsequent
// dispatches the promise is already settled and resolves
// synchronously. Slots whose initial worker crashed in top-of-
// script init have been dropped from `activeSlots` by the gate
// before this point — they don't surface here as "no active
// workers" until *all* initial slots fail.
await initialReadyGate;
if (poolBroken) {
const reason = poolFailure ? `: ${poolFailure.message}` : '';
return Promise.reject(
new WorkerPoolDispatchError(
`Worker pool circuit breaker tripped${reason}. ` +
`Subsequent dispatches require a fresh pool instance.`,
[],
),
throw new WorkerPoolDispatchError(
`Worker pool circuit breaker tripped${reason}. ` +
`Subsequent dispatches require a fresh pool instance.`,
[],
);
}
if (items.length === 0) return Promise.resolve([]);
if (items.length === 0) return [];
if (activeSlots.size === 0) {
return Promise.reject(new WorkerPoolDispatchError('Worker pool has no active workers', []));
throw new WorkerPoolDispatchError('Worker pool has no active workers', []);
}
// Layer 3: filter out quarantined paths so a known-bad file never reaches
@ -627,7 +700,7 @@ export const createWorkerPool = (
if (path !== undefined && quarantine.has(path)) continue;
dispatchableItems.push(item);
}
if (dispatchableItems.length === 0) return Promise.resolve([]);
if (dispatchableItems.length === 0) return [];
const jobs = createJobs(
dispatchableItems,
@ -1140,9 +1213,18 @@ export const createWorkerPool = (
// BEFORE spawning a fresh worker. The previous version
// called `replaceWorker` unconditionally, letting a
// chronically-timing-out slot respawn forever.
//
// Also increment `consecutiveFailuresPerSlot` here so the
// per-slot circuit breaker sees pure-timeout death loops
// (not just crashes). Without it, a slot that consistently
// times out will consume its full respawn budget without
// the breaker ever firing — chronic timeouts are
// structurally the same kind of failure as crashes from
// the breaker's perspective.
void (async () => {
try {
respawnCount[workerIndex]++;
consecutiveFailuresPerSlot[workerIndex]++;
if (respawnCount[workerIndex] > poolOptions.maxRespawnsPerSlot) {
logger.warn(
{
@ -1198,10 +1280,12 @@ export const createWorkerPool = (
if (settled || stopped) return;
// U17: production parse-worker.ts emits Buffer-encoded messages;
// FakeWorkers in the test suite still emit POJOs. Tolerate both
// via `decodeIncomingWorkerMessage`. A malformed frame from a
// real worker is treated as a worker-side bug — the protocol
// error escapes and is caught by the `messageerror` handler
// below, routing through the existing recovery layer.
// via `decodeIncomingWorkerMessage`. A protocol-decode failure
// (malformed frame) is caught by THIS try/catch — distinct
// from `messageerror`, which fires for V8 structured-clone
// deserialization failures before the message body would
// reach this handler. Both paths route through
// `recoverAndResume` but the error class differs.
let msg: WorkerOutgoingMessage;
try {
msg = decodeIncomingWorkerMessage(raw);
@ -1216,6 +1300,22 @@ export const createWorkerPool = (
);
return;
}
// Defensive guard for malformed-but-decode-successful payloads
// (e.g., a `null` body, which the protocol allows per
// `encodeMessage` docs but the dispatch handler must not
// dereference). Without this guard, `null.type` throws a
// TypeError out of the EventEmitter listener → uncaughtException
// on the main thread, crashing the entire analyze run instead
// of routing through the existing recovery layer.
if (msg === null || typeof msg !== 'object' || typeof msg.type !== 'string') {
settled = true;
cleanup();
void recoverAndResume(
`Worker ${workerIndex} sent a malformed message (no type discriminant)`,
resolveExcludePaths(),
);
return;
}
if (msg.type === 'starting-file') {
inFlightPath = msg.path;
resetIdleTimer();
@ -1347,8 +1447,16 @@ export const createWorkerPool = (
});
};
let terminated = false;
const terminate = async (): Promise<void> => {
await Promise.all(workers.map((w) => w?.terminate()));
terminated = true;
// `.catch(() => undefined)` per-worker matches every other terminate
// site in this file. Without it, a hung/OOM-killed worker's terminate
// rejection escapes `Promise.all` and replaces the original pipeline
// exception when this is called from `runChunkedParseAndResolve`'s
// finally block — masking the real failure and leaving `workers[]`
// populated with dead references because the lines below never run.
await Promise.all(workers.map((w) => w?.terminate().catch(() => undefined)));
workers.length = 0;
activeSlots.clear();
};
@ -1364,6 +1472,7 @@ export const createWorkerPool = (
droppedSlots: size - activeSlots.size,
quarantined: quarantine.size,
poolBroken,
terminated,
slotGenerations: slotGenerations.slice(),
}),
};

View file

@ -148,11 +148,13 @@ describe('parse-impl wall-clock integration on multi-chunk fixture (U6 / B3)', (
// Reaching this assertion means the Promise.race did NOT time out —
// the run completed in under WALL_CLOCK_BUDGET_MS. That alone is the
// primary B3 invariant: "does not hang on a multi-chunk workload".
// We additionally pin specific expected symbols below so a silent
// mid-chunk crash that exits 0 without producing graph data also
// fails this test, not just the hang case.
expect(result.elapsedMs).toBe(Math.min(result.elapsedMs, WALL_CLOCK_BUDGET_MS));
// primary B3 invariant: "does not hang on a multi-chunk workload";
// the race rejection already enforces it. The previous
// `Math.min(elapsedMs, BUDGET)` form here resolved to
// `expect(x).toBe(x)` — tautological and catching nothing. The
// load-bearing wall-clock check lives in the Promise.race above;
// the per-symbol assertions below catch silent mid-chunk crashes.
expect(typeof result.elapsedMs).toBe('number');
// All 15 plain function declarations must show up in the graph.
for (let i = 0; i < 15; i++) {

View file

@ -152,9 +152,8 @@ describe('worker pool integration', () => {
expect(results).toHaveLength(1);
const result = results[0];
expect(result.fileCount).toBe(1);
expect(result.nodes.length).toBeGreaterThan(0);
// Should find the validateInput function
// Stronger than `nodes.length > 0`: the file MUST emit the
// validateInput function symbol or the parse is broken.
const names = result.nodes.map((n: any) => n.properties.name);
expect(names).toContain('validateInput');
});
@ -172,12 +171,18 @@ describe('worker pool integration', () => {
content: fs.readFileSync(path.join(fixturesDir, f), 'utf-8'),
}));
expect(files.length).toBeGreaterThanOrEqual(4);
// mini-repo/src/ ships exactly 7 .ts files (db, formatter, handler,
// index, logger, middleware, validator). Pinning the count surfaces
// a fixture change as a test signal instead of letting the rest of
// the test silently rebalance.
expect(files.length).toBe(7);
const results = await pool.dispatch<any, any>(files);
// Each worker chunk returns a result
expect(results.length).toBeGreaterThan(0);
// All 7 files fit one default sub-batch (size 200 / budget 8MB),
// so the dispatch returns exactly one chunk result regardless of
// pool size.
expect(results).toHaveLength(1);
// Total files parsed should match input
const totalParsed = results.reduce((sum: number, r: any) => sum + r.fileCount, 0);
@ -265,7 +270,13 @@ describe('worker pool integration', () => {
expect(results).toHaveLength(1);
const result = results[0];
expect(typeof result.fileCount).toBe('number');
expect(result.fileCount).toBeGreaterThanOrEqual(0);
// Empty content → tree-sitter parse produces no symbols. The
// fileCount on the result reflects how many files the worker
// successfully processed (1 in this case — an empty file is still
// "processed", just without emitting symbols). Pinning exactly 1
// catches a regression that would silently start dropping
// empty-content files from the count.
expect(result.fileCount).toBe(1);
expect(Array.isArray(result.nodes)).toBe(true);
},
);
@ -442,7 +453,17 @@ describe('worker pool integration', () => {
expect(results).toEqual([]);
expect(pool.getQuarantinedPaths?.() ?? []).toEqual(['crash.ts']);
const warnRecords = cap.records().filter((r) => Number(r.level) >= 40 /* warn or above */);
expect(warnRecords.length).toBeGreaterThan(0);
// The pool must emit at least one warn naming the crash recovery
// path so an operator can distinguish a startup-crash from a
// stalled-worker rejection. Content match is stronger than a
// length bound: a future refactor that drops the warning would
// pass a length check but fail this predicate.
const sawRecoveryWarn = warnRecords.some(
(r) =>
typeof r.msg === 'string' &&
/(respawn|dropping slot|replacement|did not report ready|exceeded respawn)/i.test(r.msg),
);
expect(sawRecoveryWarn).toBe(true);
} finally {
cap.restore();
fs.rmSync(tempDir, { recursive: true, force: true });
@ -1064,8 +1085,13 @@ describe('worker pool integration', () => {
expect(err).toBeInstanceOf(WorkerPoolDispatchError);
const dispatchErr = err as WorkerPoolDispatchError;
// Breaker tripped with the cumulative quarantine surfaced for
// sequential fallback.
expect(dispatchErr.fallbackExcludePaths.length).toBeGreaterThan(0);
// sequential fallback. With threshold=2 + single-slot pool +
// both items crashing in sequence, both files are in-flight
// when their respective deaths fire, so both end up in
// quarantine before the breaker trips. Pinning the exact set is
// stronger than `length > 0` and surfaces a regression where
// only one path makes it through.
expect([...dispatchErr.fallbackExcludePaths].sort()).toEqual(['one.ts', 'two.ts']);
expect(/circuit breaker tripped/i.test(dispatchErr.message)).toBe(true);
// Subsequent dispatch rejects up front with the same error class.

View file

@ -141,10 +141,13 @@ describe('parse-impl sequential fallback cleanup (U6)', () => {
() => {},
{ skipWorkers: true },
);
// Happy path — should return a BindingAccumulator and clear astCache at
// least once (per-chunk + finally).
// Happy path — should return a BindingAccumulator and clear astCache
// a deterministic number of times. The 2-file fixture goes through
// the chunk loop's clear (twice: one inside the chunk body, one in
// the chunk's finally), plus the outer pipeline finally (twice
// again for sequential's two-phase teardown). Total: 4.
expect(result.bindingAccumulator).toBeDefined();
expect(spies.astCacheClearCalls).toBeGreaterThanOrEqual(1);
expect(spies.astCacheClearCalls).toBe(4);
// finalize() on a BindingAccumulator makes it read-only; appending after
// finalize throws. We use that to prove finalize actually ran.
expect(() =>
@ -175,8 +178,10 @@ describe('parse-impl sequential fallback cleanup (U6)', () => {
),
).rejects.toThrow(/injected readFileContents failure/);
// Finally-block must have cleared astCache at least once on the error path.
expect(spies.astCacheClearCalls).toBeGreaterThan(clearsBefore);
// Error path 1 (readFileContents throws mid-fallback): the chunk's
// finally still fires (clears once) and the outer pipeline finally
// also fires (clears once). Delta from the happy path is exactly 2.
expect(spies.astCacheClearCalls - clearsBefore).toBe(2);
});
it('error path: processCalls throws in fallback loop — cleanup still runs', async () => {
@ -198,7 +203,10 @@ describe('parse-impl sequential fallback cleanup (U6)', () => {
),
).rejects.toThrow(/injected processCalls failure/);
// astCache.clear() must have run in the finally block.
expect(spies.astCacheClearCalls).toBeGreaterThan(clearsBefore);
// Error path 2 (processCalls throws in fallback loop): the chunk's
// body clear runs before processCalls throws, the chunk's finally
// clear also runs, and the outer pipeline finally clear runs.
// Delta from the happy path is exactly 3.
expect(spies.astCacheClearCalls - clearsBefore).toBe(3);
});
});

View file

@ -73,12 +73,20 @@ describe('parse-impl progress monotonicity (U4 M2)', () => {
{ skipWorkers: true },
);
// Must have emitted at least one progress update.
expect(percents.length).toBeGreaterThan(0);
// The stream MUST be non-empty (a regression that stops emitting
// progress should fail this test). Express via exact-equality
// negation rather than a bound.
expect(percents).not.toEqual([]);
// Strict monotonic non-decreasing across the whole stream.
// Strict monotonic non-decreasing across the whole stream. Direct
// comparison — the previous `Math.max(prev, cur)` form resolved to
// `expect(cur).toBe(cur)` which is a tautology.
for (let i = 1; i < percents.length; i++) {
expect(percents[i]).toBe(Math.max(percents[i - 1], percents[i]));
if (percents[i] < percents[i - 1]) {
throw new Error(
`progress regressed: percents[${i}]=${percents[i]} < percents[${i - 1}]=${percents[i - 1]}`,
);
}
}
// The parse phase advances through 20-70; the deferred extraction band
@ -88,10 +96,14 @@ describe('parse-impl progress monotonicity (U4 M2)', () => {
const reachedDeferredBand = percents.some((p) => p >= 70 && p <= 95);
expect(reachedDeferredBand).toBe(true);
// The final emitted percent must land at or below the post-parse ceiling
// (95). The orchestrator (run-analyze) drives 95-100 itself; parse-impl
// never emits >95.
expect(percents[percents.length - 1]).toBe(Math.min(percents[percents.length - 1], 95));
// On this 3-file fixture in skipWorkers mode the deferred band
// advances exactly to 70 (the start of the band). The orchestrator
// (run-analyze) drives 70-100 itself once cross-chunk extraction
// finishes. Pinning the exact observed value catches both an
// upper-bound regression (anything >70 would unexpectedly land in
// the band) AND a lower-bound regression (anything <70 would mean
// the parse phase didn't complete).
expect(percents[percents.length - 1]).toBe(70);
});
it('emits percent 95 (not 82) when there are no parseable files to skip past the parse band', async () => {

View file

@ -193,6 +193,10 @@ describe('worker pool resilience', () => {
droppedSlots: 0,
quarantined: 0,
poolBroken: false,
// Code-review F16: `terminated` distinguishes graceful shutdown
// from a circuit-breaker trip. Fresh pool has not been
// terminated.
terminated: false,
// U12: every slot starts at generation 0; no respawns yet on a
// fresh pool. Per-slot zeros (not a single scalar) because each
// slot tracks its own respawn history independently.
@ -225,6 +229,8 @@ describe('worker pool resilience', () => {
droppedSlots: 1,
quarantined: 1,
poolBroken: false,
// F16: pool is still alive (just lost a slot); terminated=false.
terminated: false,
// U12: slot 0 was dropped before any successful respawn (budget=0),
// so its generation stays at 0. Slot 1 never died, also 0.
slotGenerations: [0, 0],