From a2878df5fb5ec620df7618b8a7a921cc37684e7d Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Wed, 20 May 2026 12:54:43 +0100 Subject: [PATCH] fix(workers,tests,docs): apply ce-code-review findings (16 items) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- README.md | 3 + gitnexus/README.md | 10 ++ gitnexus/src/cli/analyze.ts | 12 +- gitnexus/src/cli/index.ts | 3 + .../ingestion/pipeline-phases/parse-impl.ts | 29 ++- .../core/ingestion/workers/parse-worker.ts | 35 +++- .../src/core/ingestion/workers/protocol.ts | 24 ++- .../src/core/ingestion/workers/worker-pool.ts | 167 +++++++++++++++--- .../parse-impl-large-fixture.test.ts | 12 +- gitnexus/test/integration/worker-pool.test.ts | 46 +++-- .../test/unit/parse-impl-fallback.test.ts | 22 ++- .../parse-impl-progress-monotonic.test.ts | 28 ++- .../test/unit/worker-pool-resilience.test.ts | 6 + 13 files changed, 316 insertions(+), 81 deletions(-) diff --git a/README.md b/README.md index 7f78b9a08..f1ce3ba01 100644 --- a/README.md +++ b/README.md @@ -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 `. | 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 ` × 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. | diff --git a/gitnexus/README.md b/gitnexus/README.md index 55cdc12f7..1c6f13f3c 100644 --- a/gitnexus/README.md +++ b/gitnexus/README.md @@ -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 diff --git a/gitnexus/src/cli/analyze.ts b/gitnexus/src/cli/analyze.ts index 0e001b845..dd1d71d8d 100644 --- a/gitnexus/src/cli/analyze.ts +++ b/gitnexus/src/cli/analyze.ts @@ -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. diff --git a/gitnexus/src/cli/index.ts b/gitnexus/src/cli/index.ts index dbd8b0592..7bf2278c8 100644 --- a/gitnexus/src/cli/index.ts +++ b/gitnexus/src/cli/index.ts @@ -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' + diff --git a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts index 3678539e4..40853f45b 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts @@ -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]!; diff --git a/gitnexus/src/core/ingestion/workers/parse-worker.ts b/gitnexus/src/core/ingestion/workers/parse-worker.ts index cd1307693..00de843e9 100644 --- a/gitnexus/src/core/ingestion/workers/parse-worker.ts +++ b/gitnexus/src/core/ingestion/workers/parse-worker.ts @@ -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]), diff --git a/gitnexus/src/core/ingestion/workers/protocol.ts b/gitnexus/src/core/ingestion/workers/protocol.ts index e8a603531..8ff27a779 100644 --- a/gitnexus/src/core/ingestion/workers/protocol.ts +++ b/gitnexus/src/core/ingestion/workers/protocol.ts @@ -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; } diff --git a/gitnexus/src/core/ingestion/workers/worker-pool.ts b/gitnexus/src/core/ingestion/workers/worker-pool.ts index 38a7e7651..4b05cd9b3 100644 --- a/gitnexus/src/core/ingestion/workers/worker-pool.ts +++ b/gitnexus/src/core/ingestion/workers/worker-pool.ts @@ -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( + 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(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 = ( + // 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 = 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 ( items: TInput[], onProgress?: (filesProcessed: number) => void, ): Promise => { + // 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 => { - 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(), }), }; diff --git a/gitnexus/test/integration/parse-impl-large-fixture.test.ts b/gitnexus/test/integration/parse-impl-large-fixture.test.ts index c51cdb400..29ebae7bf 100644 --- a/gitnexus/test/integration/parse-impl-large-fixture.test.ts +++ b/gitnexus/test/integration/parse-impl-large-fixture.test.ts @@ -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++) { diff --git a/gitnexus/test/integration/worker-pool.test.ts b/gitnexus/test/integration/worker-pool.test.ts index 04b4372fd..cbcca250f 100644 --- a/gitnexus/test/integration/worker-pool.test.ts +++ b/gitnexus/test/integration/worker-pool.test.ts @@ -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(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. diff --git a/gitnexus/test/unit/parse-impl-fallback.test.ts b/gitnexus/test/unit/parse-impl-fallback.test.ts index 869016552..089d9aef7 100644 --- a/gitnexus/test/unit/parse-impl-fallback.test.ts +++ b/gitnexus/test/unit/parse-impl-fallback.test.ts @@ -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); }); }); diff --git a/gitnexus/test/unit/parse-impl-progress-monotonic.test.ts b/gitnexus/test/unit/parse-impl-progress-monotonic.test.ts index 6418c743f..ce21157f7 100644 --- a/gitnexus/test/unit/parse-impl-progress-monotonic.test.ts +++ b/gitnexus/test/unit/parse-impl-progress-monotonic.test.ts @@ -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 () => { diff --git a/gitnexus/test/unit/worker-pool-resilience.test.ts b/gitnexus/test/unit/worker-pool-resilience.test.ts index e2dea546e..6bc14387b 100644 --- a/gitnexus/test/unit/worker-pool-resilience.test.ts +++ b/gitnexus/test/unit/worker-pool-resilience.test.ts @@ -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],