mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-03 02:21:44 +00:00
* 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>
256 lines
9.9 KiB
TypeScript
256 lines
9.9 KiB
TypeScript
/**
|
||
* 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);
|
||
});
|
||
});
|