GitNexus/gitnexus/test/unit/worker-pool-startup-stderr.test.ts
Gergő Magyar c4b69402e1
feat(workers): self-healing worker pool + deferred-resolution observability (#1741) (#1947)
* fix(workers): fail fast instead of silently degrading on worker-pool startup failure (#1741)

When an explicitly-sized worker pool (--workers <N>) fails to start because
every worker crashes during top-of-script init, the parse phase used to log a
swallowed `logger.warn` and silently fall back to the ~10x slower sequential
parser. In #1741 (rc99) that turned a worker-startup regression into a
123-minute "stuck" parse with no explanation.

This change:
- Surfaces the real crash: the pool now spawns workers with `{ stderr: true }`,
  tees + captures each worker's stderr, and attaches the tail to its
  readiness-failure messages (propagated via
  WorkerPoolInitializationError.readinessFailures). "did not report ready"
  now carries the underlying native-binding/import error.
- Gates the fallback: when --workers was explicit and fallback was not opted
  into, a total startup failure throws an actionable error instead of
  degrading. Auto-sized pools still fall back, but loudly (logger.error +
  progress warning). New --allow-sequential-fallback flag (+ i18n) opts back in.
- Adds env-gated worker bootstrap-stage logging (GITNEXUS_WORKER_BOOTSTRAP /
  --verbose): imports+grammars loaded -> ready sent -> first task received, so
  a slow/crashing startup is diagnosable.

Tests: all-workers-failed gating (fatal vs loud degrade), stderr surfacing,
and the updated lazy-cache fallback contract (opt-in flag + fail-fast).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* feat(ingestion): always-on slow-file watchdog for deferred call resolution (#1741)

The original #1741 symptom is a run that appears stuck at "Resolving calls
(all chunks)... (9000/18066 files)" — the progress bar freezes inside a single
file's call resolution and nothing reaches the log. Rich per-file deferred
diagnostics already exist, but only behind --verbose / GITNEXUS_PROFILE_DEFERRED,
so a plain `analyze` run gives the user a frozen bar and silence.

Add an always-on (not verbose-gated) per-file watchdog in
processCallsFromExtracted: when a single file's call resolution exceeds
alwaysOnSlowFileWarnMs() (default 15s, override GITNEXUS_SLOW_FILE_WARN_MS,
0 disables) it emits a throttled logger.warn naming the culprit file and the
files-resolved-so-far — turning the silent stall into one actionable line.
Throttled (>=30s between warnings) so a genuinely slow repo can't storm the log.
The watchdog is observation-only; resolution behavior is unchanged.

Note: deliberately did NOT add a heritage child x parent product cap — the name
lookups are O(1) (type-registry Map.get) and the product is bounded, so the
heritage build is not the bottleneck; a cap would risk dropping real edges for
no measured gain.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* test(ingestion): worker-vs-sequential parity guard for binding/edge collapse (#1741)

rc99 produced almost no bindings/edges (13 bindings vs rc91's 106,305) because
a worker-path failure left extracted results unmerged while the run still
reported success. Rather than an arbitrary "implausibly low" runtime threshold
(which false-positives on legitimately low-binding repos/languages), pin the
invariant directly: for the same repo, worker mode and sequential mode must
produce the same graph.

The test runs the ts-simple cross-file fixture through worker mode
(workerPoolSize + lowered threshold) and sequential mode (skipWorkers), and
asserts: usedWorkerPool is true/false respectively (guards the test itself
against a silent fallback masking divergence), identical CALLS/IMPORTS/DEFINES/
HAS_METHOD edge sets and Class/Function/Method defs, and non-zero CALLS/IMPORTS
(the rc99 collapse signature).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix(workers): arm fail-fast for env-sized pools + fix watchdog /0 denominator (#1741)

Addresses two review findings on the #1741 worker-startup PR:

- Fail-fast gate missed the env channel. `explicitWorkers` keyed only off
  the `--workers` flag, so a pool sized via `GITNEXUS_WORKER_POOL_SIZE`
  (with no `--workers`) silently degraded to sequential on a total
  worker-startup crash — reproducing the original #1741 symptom for
  env-channel operators. The gate now arms on a non-zero size from either
  channel, via a single-source `envWorkerPoolSize()` helper exported from
  worker-pool.ts (also rewired through resolveAutoPoolSize). The fatal
  message now names the channel actually used instead of "--workers undefined".

- Always-on slow-file watchdog printed "Resolved N/0 files". `resolvedTotal`
  was pre-counted only on the profile path, but the watchdog reads it on
  every run, so a plain `analyze` showed a bogus /0 denominator on exactly
  the unprofiled hang the watchdog exists to explain. Pre-count now runs
  whenever its result is read (profile path OR watchdog active).

Tests: strengthened the watchdog test to assert "1/1" (not "/0"); added
env-channel fail-fast/degrade cases and made the gating suite hermetic
against an ambient GITNEXUS_WORKER_POOL_SIZE.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* feat(workers): self-healing worker pool replaces the fail-fast flag (#1741)

Replaces the interim --allow-sequential-fallback flag with automatic,
bounded self-healing in the worker pool — industry-standard supervision
(OTP restart-intensity, systemd StartLimit, circuit-breaker, AWS jittered
backoff) translated to the Node worker_threads pool.

worker-pool.ts — bounded startup self-heal (the missing layer):
- A worker that crashes during top-of-script init is now RETRIED with
  capped, full-jitter backoff (BASE 250ms, CAP 2s) up to a small per-slot
  budget, so a transient blip heals itself with no operator action. The
  prior code dropped an unready initial slot on its first crash.
- A DETERMINISTIC crash-loop (>=2 fresh workers crash with the same
  normalized signature before any reaches ready — the #1741 missing
  native-binding case) is detected and short-circuited, so the pool gives
  up in ~1s instead of burning every slot's budget. Correctness rests on
  the STRUCTURAL signal (zero workers ever ready + budget exhausted), so a
  missed signature only costs a few seconds, never a misfire; even a
  stderr-less crash groups via its normalized "exited with code N" message.
- Backoff sleeps are cancellable (unref'd timer + abort on terminate), so
  terminate() can't be wedged for the backoff duration.
- WorkerPoolInitializationError now carries a crashClass for an accurate,
  flag-free message. The runtime respawn/breaker path is unchanged.

parse-impl.ts — collapse to automatic fail-fast:
- handleWorkerStartupFailure always logs the real cause then THROWS with
  the captured crash + `--workers 0` as the explicit sequential escape.
  No more degrade branch; no dependence on how the pool was sized. This is
  reached only after the bounded self-heal is exhausted, so it can't
  resurrect the #1741 silent 123-minute sequential grind. Construction
  failure (broken install) also fails fast instead of degrading silently.

Removed --allow-sequential-fallback end to end (CLI, run-analyze, pipeline,
i18n). --workers 0 remains the explicit "parse sequentially" path; one flag
removed, none added. Grounded in a research+critique pass; the critique's
hazards (N-parallel race, empty-stderr timing, non-cancellable sleep,
runtime-breaker regression) are addressed or scoped out by design.

Tests: startup self-heal (transient recovers; deterministic fails fast
without burning the budget); gating test rewritten to the fail-fast-always
contract; obsolete degrade test removed.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix(workers): ref + cancel startup backoff so transient retries aren't dropped (#1741 U1)

abortableSleep unref'd its backoff timer, so a transient startup retry could be
silently dropped if that timer was the last ref'd handle on the event loop —
the process could exit mid-recovery. Keep the timer ref'd (a pending retry is
necessary work) and register a cancel fn in a pool-scoped set; terminate() now
clears pending backoffs so it can't be wedged for the backoff cap. A normally
fired timer self-deregisters (clear-on-settle), so no timer lingers after a
slot's retry loop exits. Exposes pendingStartupTimers in getStats.

Tests: terminate-during-backoff cancels + spawns nothing after (R2); the
recovery test now asserts no startup timer lingers after settle (R1).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix(workers): route GITNEXUS_WORKER_POOL_SIZE=0 to sequential, not a phantom fail-fast (#1741 U2)

env=0 (no --workers) built a size-0 pool that threw a fabricated "retry budget
exhausted / native binding" crash. The shouldUseWorkers gate now routes env=0
to the sequential path before pool construction — but only when no explicit
--workers <N> was given, so an explicit positive size wins over an ambient
env=0. The route emits one log line so the undocumented (possibly accidental)
env=0 case is observable instead of a silent degrade.

envWorkerPoolSize is un-exported (module-internal sizing reader); a new
workerPoolDisabledByEnv() predicate serves the gate. Empty/whitespace env is
now treated as unset (auto formula), not 0 — an empty assignment is an accident,
not a request for zero workers. Reattached the detached resolveAutoPoolSize
JSDoc and corrected the stale docstring.

Tests: env=0 → sequential (no spawn); explicit --workers wins over env=0;
workerPoolDisabledByEnv unit (0=true, positive/empty/invalid=false); getStats
shape updated for pendingStartupTimers.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix(workers): make deterministic crash-loop detection conservative (#1741 U3)

The old tally counted crash EVENTS in a shared signature->count map, so a
simultaneous transient crash storm (e.g. spawn EAGAIN under fork pressure) or
a single slot crashing identically twice falsely tripped "deterministic" and
hard-aborted work that would have self-healed. Replace it: a crash counts
toward deterministic only after its signature REPRODUCES across a respawn on
the same slot, and the short-circuit fires once >=2 distinct slots reproduced
(or 1 for a size-1 pool). Every slot now gets >=1 self-heal attempt before any
short-circuit; the structural budget floor still bounds the worst case.

crashSignature now also collapses Windows backslash paths and bare (no-0x) hex
runs so the fast-path fires on those platforms; exported for unit testing.

Tests: simultaneous storm self-heals (the discriminator vs an attempt-0 rule);
distinct-per-attempt crashes classify transient-exhausted; single-slot
reproduction classifies deterministic; crashSignature normalization unit.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix(workers): class-aware startup failure hint + reattach detached JSDoc (#1741 U4)

The "often a missing/broken native binding" hint was appended to every
failure class, including a pool *construction* failure where no worker ever
ran (a missing build / bad worker path). Make the hint class-aware: keep it
for the readiness/init classes, use a construction-specific hint otherwise,
and surface the construction error (e.g. "Worker script not found: …")
verbatim. Reattach the waitForWorkerReady JSDoc that the stderr-capture block
had detached from its function. (The abortableSleep docstring was already
corrected in U1.)

Tests: construction message surfaces the real error + drops the native-binding
guess; deterministic/transient messages keep the hint (regression guard).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-05-31 13:52:04 +01:00

256 lines
9.9 KiB
TypeScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

/**
* Worker startup failure surfaces the real crash via captured stderr (#1741).
*
* Before this, when every worker crashed during top-of-script init the pool
* rejected dispatch with a generic "did not report ready within 5000ms" and
* the actual cause (e.g. a broken native binding) was lost to the worker's
* inherited stderr. The pool now spawns workers with `{ stderr: true }`,
* tees + captures each worker's stderr, and attaches the tail to its
* readiness-failure messages — which propagate on
* `WorkerPoolInitializationError.readinessFailures`.
*
* This test injects a fake worker that prints a crash to stderr and exits
* without ever reporting `ready`, then asserts the captured stderr reaches
* the dispatch error.
*/
import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest';
import { EventEmitter } from 'node:events';
import path from 'node:path';
import os from 'node:os';
import fs from 'node:fs';
import { pathToFileURL } from 'node:url';
import {
createWorkerPool,
WorkerPoolInitializationError,
} from '../../src/core/ingestion/workers/worker-pool.js';
const CRASH_STDERR =
"Error: Cannot find module 'tree-sitter-c-sharp/bindings/node'\n at parse-worker.ts:10\n";
/**
* Worker double that crashes during startup: emits a crash to its `stderr`
* stream, never sends `{type:'ready'}`, then exits non-zero. Mirrors a
* native-binding load failure in `parse-worker.ts`.
*/
class CrashingWorker extends EventEmitter {
readonly stderr = new EventEmitter();
constructor(crashText: string = CRASH_STDERR) {
super();
queueMicrotask(() => {
// stderr first so it's captured before the exit builds the message.
this.stderr.emit('data', Buffer.from(crashText));
this.emit('exit', 1);
});
}
postMessage(): void {}
async terminate(): Promise<number> {
return 1;
}
}
/** Worker double that starts cleanly: reports `{type:'ready'}` and never dies. */
class ReadyWorker extends EventEmitter {
readonly stderr = new EventEmitter();
constructor() {
super();
queueMicrotask(() => this.emit('message', { type: 'ready' }));
}
postMessage(): void {}
async terminate(): Promise<number> {
return 0;
}
}
let tempDir: string;
let workerUrl: URL;
let stderrSpy: ReturnType<typeof vi.spyOn>;
beforeEach(() => {
tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-startup-stderr-'));
const workerPath = path.join(tempDir, 'fake-worker.js');
fs.writeFileSync(workerPath, '// fake');
workerUrl = pathToFileURL(workerPath) as URL;
// The tee writes captured worker stderr to process.stderr; silence it so
// the (intentional) crash text doesn't pollute test output.
stderrSpy = vi.spyOn(process.stderr, 'write').mockReturnValue(true);
});
afterEach(() => {
stderrSpy.mockRestore();
try {
fs.rmSync(tempDir, { recursive: true, force: true });
} catch {
/* best-effort */
}
});
describe('worker pool — startup stderr surfacing (#1741)', () => {
it('attaches captured worker stderr to the WorkerPoolInitializationError', async () => {
const pool = createWorkerPool(workerUrl, 2, {
workerFactory: () => new CrashingWorker() as unknown as Worker,
});
const dispatch = pool.dispatch([{ path: 'src/a.ts', content: 'x' }]);
await expect(dispatch).rejects.toBeInstanceOf(WorkerPoolInitializationError);
const err = await dispatch.catch((e: unknown) => e as WorkerPoolInitializationError);
expect(err.readinessFailures.length).toBeGreaterThan(0);
// The real crash reason, recovered from the worker's stderr, is present.
const joined = err.readinessFailures.join('\n');
expect(joined).toContain('Worker stderr:');
expect(joined).toContain('tree-sitter-c-sharp');
// And the captured stderr was teed to process.stderr (visibility preserved).
expect(stderrSpy).toHaveBeenCalled();
await pool.terminate().catch(() => undefined);
});
});
describe('worker pool — startup self-healing (#1741)', () => {
it('self-heals a transient startup crash: respawns the slot and recovers', async () => {
let calls = 0;
const pool = createWorkerPool(workerUrl, 1, {
// First spawn crashes during init; the bounded retry respawns the slot
// and the second spawn comes up clean — recovery with no operator action.
workerFactory: () => {
calls++;
return (calls === 1 ? new CrashingWorker() : new ReadyWorker()) as unknown as Worker;
},
});
// Empty dispatch forces the initial-ready gate to settle without needing
// the full sub-batch protocol.
await pool.dispatch([]);
expect(calls).toBeGreaterThanOrEqual(2); // crashed once, respawned
expect(pool.getStats().activeSlots).toBe(1); // slot recovered and is live
expect(pool.getStats().poolBroken).toBe(false);
// R1 (clear-on-settle): the ref'd backoff timer self-cleared when it fired,
// so no startup timer lingers after the slot's retry loop exited.
expect(pool.getStats().pendingStartupTimers).toBe(0);
await pool.terminate().catch(() => undefined);
});
it('fails fast on a deterministic crash-loop without burning every slot budget', async () => {
let calls = 0;
const pool = createWorkerPool(workerUrl, 3, {
// Every worker crashes identically — a deterministic fault (the #1741
// missing-binding case). Retrying cannot help.
workerFactory: () => {
calls++;
return new CrashingWorker() as unknown as Worker;
},
});
const err = await pool
.dispatch([{ path: 'a.ts', content: 'x' }])
.catch((e: unknown) => e as WorkerPoolInitializationError);
expect(err).toBeInstanceOf(WorkerPoolInitializationError);
expect(err.crashClass).toBe('deterministic-startup');
// No short-circuit would mean 3 slots × (1 + STARTUP_RESTART_BUDGET=2) = 9
// spawns; the reproduced-across-respawn signal trips once ≥2 slots crash
// identically a second time, well before the full budget.
expect(calls).toBeLessThan(9);
await pool.terminate().catch(() => undefined);
});
it('a simultaneous transient crash storm self-heals — not misclassified as deterministic (R4)', async () => {
// The discriminator: all 3 slots crash IDENTICALLY on their first spawn, but
// each respawn comes up ready. A rule that tallied attempt-0 crashes by
// distinct slot would trip "deterministic" here and hard-abort; the
// reproduced-across-respawn rule lets every slot self-heal.
let calls = 0;
const pool = createWorkerPool(workerUrl, 3, {
workerFactory: () => {
calls++;
// Calls 1-3 are the initial spawns (all crash identically); 4+ are respawns.
return (calls <= 3 ? new CrashingWorker() : new ReadyWorker()) as unknown as Worker;
},
});
await pool.dispatch([]); // settle the initial-ready gate
expect(pool.getStats().activeSlots).toBe(3); // every slot recovered
expect(pool.getStats().poolBroken).toBe(false);
await pool.terminate().catch(() => undefined);
});
it('classifies distinct-per-attempt crashes as transient-exhausted, not deterministic (R5)', async () => {
// Each spawn crashes with a DISTINCT signature, so no signature reproduces
// across a respawn — the deterministic short-circuit never fires and the
// pool exhausts its per-slot budget.
let calls = 0;
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => {
calls++;
return new CrashingWorker(
`Error: distinct-failure-${'abcdefg'[calls] ?? 'z'}`,
) as unknown as Worker;
},
});
const err = await pool
.dispatch([{ path: 'a.ts', content: 'x' }])
.catch((e: unknown) => e as WorkerPoolInitializationError);
expect(err).toBeInstanceOf(WorkerPoolInitializationError);
expect(err.crashClass).toBe('transient-exhausted');
// 1 initial + STARTUP_RESTART_BUDGET (2) retries = 3 spawns, all distinct.
expect(calls).toBe(3);
await pool.terminate().catch(() => undefined);
});
it('classifies a single slot whose crash reproduces across a respawn as deterministic (R4/R5)', async () => {
// Size-1 pool, identical crash on attempt 0 and its respawn => the signature
// reproduced => deterministic. The slot is not run to full budget once the
// same crash survives a respawn.
let calls = 0;
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => {
calls++;
return new CrashingWorker() as unknown as Worker; // identical every spawn
},
});
const err = await pool
.dispatch([{ path: 'a.ts', content: 'x' }])
.catch((e: unknown) => e as WorkerPoolInitializationError);
expect(err.crashClass).toBe('deterministic-startup');
expect(calls).toBe(2); // attempt 0 + one respawn that reproduced
await pool.terminate().catch(() => undefined);
});
it('terminate() during startup cancels any pending backoff and spawns nothing after (R2)', async () => {
let calls = 0;
const pool = createWorkerPool(workerUrl, 1, {
// Crashes on every spawn, so the slot is in its bounded retry/backoff loop.
workerFactory: () => {
calls++;
return new CrashingWorker() as unknown as Worker;
},
});
// Let the first crash register and the slot enter its (ref'd) backoff.
await new Promise((resolve) => setTimeout(resolve, 0));
const callsBeforeTerminate = calls;
expect(callsBeforeTerminate).toBeGreaterThanOrEqual(1);
// terminate() must cancel the pending backoff (not wait it out) and stop the loop.
await pool.terminate();
expect(pool.getStats().terminated).toBe(true);
// No ref'd backoff timer left pinning the loop after terminate.
expect(pool.getStats().pendingStartupTimers).toBe(0);
expect(pool.getStats().activeSlots).toBe(0);
// No worker is spawned after terminate — the woken loop sees `terminated` and gives up.
await new Promise((resolve) => setTimeout(resolve, 10));
expect(calls).toBe(callsBeforeTerminate);
});
});