From 601acb6dd5610251670955a942b394b842e89988 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 9 Oct 2026 22:12:29 +0000 Subject: [PATCH 1/2] fix(engine): stop parse workers on failure and interrupt A crashed parse worker made pool shutdown throw an uncaught stack, and a parent that exited without kill left child workers reparented. Stop workers on failure, SIGINT, SIGTERM, and normal exit, and report a one-line remedy. Signed-off-by: Cursor Agent Co-authored-by: vibgrate-team --- ARCHITECTURE.md | 4 +- DOCS.md | 8 + src/commands/build.ts | 5 +- src/core-open/run-core-scan.ts | 8 +- src/engine/pool.ts | 292 ++++++++++++++++++++++++++-- test/fixtures/parse-pool-harness.ts | 43 ++++ test/fixtures/parse-pool-worker.mjs | 32 +++ test/parse-pool-shutdown.test.ts | 250 ++++++++++++++++++++++++ 8 files changed, 626 insertions(+), 16 deletions(-) create mode 100644 test/fixtures/parse-pool-harness.ts create mode 100644 test/fixtures/parse-pool-worker.mjs create mode 100644 test/parse-pool-shutdown.test.ts diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 475150fb..34e3516c 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -108,7 +108,9 @@ meaningful. (`engine/discover.ts` sorts). - **Parallel but ordered** — the worker pool (`engine/pool.ts` + `engine/parse-worker.ts`) parses files concurrently for speed, then results are - re-ordered deterministically before they enter the graph. + re-ordered deterministically before they enter the graph. The pool stops its + workers when parsing finishes, when parsing fails, and on SIGINT/SIGTERM, so + a build or scan does not leave parse workers behind. If you touch anything that ends up in `graph.json` or a report, add a test that asserts stable output. diff --git a/DOCS.md b/DOCS.md index b6b753ce..2af0f4ef 100644 --- a/DOCS.md +++ b/DOCS.md @@ -5101,6 +5101,14 @@ and `0` always means "disabled". | `VG_JOBS` | CPU cores − 1 | Default parse worker count when `--jobs` isn't passed. Fewer workers = lower peak memory (each worker loads its own grammar set). | | `VG_WORKER_HEAP_MB` | platform default | Per-worker old-generation heap cap, so one runaway parse can't take the whole machine. | +If a parse worker crashes, runs out of memory, or the command is interrupted +(Ctrl+C or SIGTERM), `vg` stops the workers before it exits. The failure is one +line on stderr — what failed, and whether to re-run with `--jobs 1`, exclude +files, or raise `VG_WORKER_HEAP_MB` — not a stack trace. A finished build stops +its workers the same way. Check with a process list (`ps`) filtered for this +command's pid: after `vg build` or `vg scan` returns, no parse worker should +still be running. Graph output is unchanged. + Skips are deterministic functions of the input (file size, file count) — never of observed memory — so identical input still produces an identical `graph.json`. To give the build more room instead of limiting it, raise the diff --git a/src/commands/build.ts b/src/commands/build.ts index 981cfa95..53921911 100644 --- a/src/commands/build.ts +++ b/src/commands/build.ts @@ -22,6 +22,7 @@ import { renderReport } from '../engine/report.js'; import { renderHtml } from '../engine/html.js'; import { UsageError, mergeExcludes } from '../engine/discover.js'; import { ResourceLimitError } from '../engine/limits.js'; +import { ParseWorkerFailure } from '../engine/pool.js'; import { UnsafeRootError } from '../core-open/utils/root-safety.js'; import { CliError, ExitCode, usageError } from '../util/exit.js'; import { resolveSelfJsEntry } from '../util/cli-invocation.js'; @@ -150,7 +151,9 @@ export async function runBuild( } catch (err) { bar?.done(); if (err instanceof UsageError) throw usageError(err.message); - if (err instanceof ResourceLimitError) throw new CliError(err.message, ExitCode.ERROR); + if (err instanceof ResourceLimitError || err instanceof ParseWorkerFailure) { + throw new CliError(err.message, ExitCode.ERROR); + } if (err instanceof UnsafeRootError) throw new CliError(err.message, ExitCode.ERROR); throw err; } diff --git a/src/core-open/run-core-scan.ts b/src/core-open/run-core-scan.ts index a6d6ebc5..9e2b1023 100644 --- a/src/core-open/run-core-scan.ts +++ b/src/core-open/run-core-scan.ts @@ -764,8 +764,12 @@ export async function runCoreScan( { projects: allProjects, solutions, extended }, ); progress.completeStep('map', detail || 'done'); - } catch { - progress.completeStep('map', 'skipped (map build failed)'); + } catch (err) { + // One line, no stack: a worker crash or resource limit should say what + // failed and what to do, then the scan continues without the map. + const line = (err instanceof Error ? err.message : String(err)).split('\n')[0]?.trim() || 'map build failed'; + const detail = line.length > 240 ? `${line.slice(0, 237)}...` : line; + progress.completeStep('map', `skipped (${detail})`); } } diff --git a/src/engine/pool.ts b/src/engine/pool.ts index 9837770c..ad5e8b8b 100644 --- a/src/engine/pool.ts +++ b/src/engine/pool.ts @@ -22,6 +22,10 @@ import type { ParseTask } from './parse-worker.js'; * single-threaded inline path is always available (and used for small repos or * when the worker module isn't resolvable, e.g. under ts-only test runners), * producing byte-identical output to the pooled path. + * + * Workers are stopped on success, on failure, and on SIGINT/SIGTERM. A crashed + * worker must not dump a stack or leave a Node process behind after `vg build` + * or `vg scan` returns. */ export interface ParseOptions { @@ -37,6 +41,261 @@ export interface ParseOptions { grammarsDir?: string; /** Heap budget (MiB) checked as parse results accumulate; 0/unset skips. */ memoryBudgetMb?: number; + /** + * Parse-worker module. Tests pass a fixture; production resolves the bundled + * `parse-worker.js` next to this file. + */ + workerFile?: string; + /** + * `child_process` runs each worker as its own OS process so a crash can be + * killed by pid. The default `worker_threads` pool is what `vg build` uses. + */ + workerRuntime?: 'worker_threads' | 'child_process'; +} + +/** A parse worker died or was stopped. The message is the whole user-facing error. */ +export class ParseWorkerFailure extends Error { + readonly isParseWorkerFailure = true; + constructor(message: string) { + super(message); + this.name = 'ParseWorkerFailure'; + } +} + +/** How long `destroy()` may take before workers are killed outright. */ +const POOL_SHUTDOWN_MS = 2_000; + +interface KillableWorker { + terminate?: () => Promise; + unref?: () => void; + /** Set when the worker is a `child_process` (tinypool `ProcessWorker`). */ + process?: { kill: (signal?: NodeJS.Signals | number) => boolean }; +} + +interface ManagedPool { + threads: KillableWorker[]; + destroy: () => Promise; + cancelPendingTasks: () => void; + on: (event: 'error', listener: (err: unknown) => void) => void; +} + +let activePool: ManagedPool | null = null; +let detachSignals: (() => void) | null = null; +let exitHookInstalled = false; + +/** Workers still tracked by the live parse pool. Zero after shutdown. */ +export function activeParseWorkerCount(): number { + if (!activePool) return 0; + try { + return activePool.threads.length; + } catch { + return 0; + } +} + +function installExitHook(): void { + if (exitHookInstalled) return; + exitHookInstalled = true; + // `exit` is synchronous. Child workers are not reaped with the parent, so + // kill them here — including when a signal handler calls `process.exit`. + process.on('exit', () => { + if (activePool) killWorkers(activePool); + }); +} + +function killWorkers(pool: ManagedPool): void { + let workers: KillableWorker[] = []; + try { + workers = pool.threads; + } catch { + return; + } + for (const worker of workers) killWorker(worker); +} + +function killWorker(worker: KillableWorker): void { + const child = worker.process; + if (child && typeof child.kill === 'function') { + try { + child.kill('SIGKILL'); + } catch { + // Already reaped. + } + return; + } + try { + worker.unref?.(); + } catch { + // Already gone. + } + try { + void worker.terminate?.(); + } catch { + // Already gone. + } +} + +function onStopSignal(signal: NodeJS.Signals): void { + const pool = activePool; + if (!pool) return; + killWorkers(pool); + const code = signal === 'SIGINT' ? 130 : 143; + const word = signal === 'SIGINT' ? 'interrupted' : 'terminated'; + const message = `vg: ${word} while parsing. Parse workers were stopped.\n`; + let exited = false; + const finish = (): void => { + if (exited) return; + exited = true; + process.exit(code); + }; + try { + process.stderr.write(message, finish); + } catch { + finish(); + return; + } + // If stderr never drains, still exit so workers cannot outlive the command. + setTimeout(finish, 50); +} + +function armPool(pool: ManagedPool): void { + activePool = pool; + installExitHook(); + if (detachSignals) return; + const onInt = (): void => onStopSignal('SIGINT'); + const onTerm = (): void => onStopSignal('SIGTERM'); + process.on('SIGINT', onInt); + process.on('SIGTERM', onTerm); + detachSignals = () => { + process.removeListener('SIGINT', onInt); + process.removeListener('SIGTERM', onTerm); + }; +} + +function disarmPool(): void { + // Drop the signal listeners before clearing the pool so a signal delivered + // in this window still finds the workers and exits, instead of being swallowed. + if (detachSignals) { + detachSignals(); + detachSignals = null; + } + activePool = null; +} + +function isTinypoolShutdownBug(err: unknown): boolean { + return err instanceof TypeError && err.message.includes('removeListener'); +} + +async function closePool(pool: ManagedPool): Promise { + // tinypool's `destroy()` waits on `events.once(worker, 'exit')`. If the + // worker emits `error` first, that helper throws an uncaught + // `removeListener` TypeError and the process dumps a stack. Swallow only + // that bug; the `run()` rejection is the error the caller sees. + const onUncaught = (err: Error): void => { + if (isTinypoolShutdownBug(err)) { + killWorkers(pool); + return; + } + process.removeListener('uncaughtException', onUncaught); + process.stderr.write(`${err.stack ?? err.message}\n`); + process.exit(1); + }; + process.on('uncaughtException', onUncaught); + try { + try { + pool.cancelPendingTasks(); + } catch { + // Queue already drained. + } + let timer: ReturnType | undefined; + const timedOut = new Promise((resolve) => { + timer = setTimeout(() => { + killWorkers(pool); + resolve(); + }, POOL_SHUTDOWN_MS); + }); + try { + await Promise.race([ + pool.destroy().then( + () => undefined, + () => { + killWorkers(pool); + }, + ), + timedOut, + ]); + } finally { + if (timer) clearTimeout(timer); + } + killWorkers(pool); + } finally { + process.removeListener('uncaughtException', onUncaught); + } +} + +/** + * Flags that make a worker a second copy of the tool (tsx, a preload, an + * inspector) rather than a parse process. Heap and V8 flags are kept so a + * raised `--max-old-space-size` still applies inside the pool. + */ +function workerExecArgv(): string[] { + const dropValue = new Set(['-r', '--require', '--import', '--loader', '--experimental-loader']); + const argv = process.execArgv; + const out: string[] = []; + for (let i = 0; i < argv.length; i++) { + const arg = argv[i] ?? ''; + const inspect = + arg === '--inspect' || + arg === '--inspect-brk' || + arg.startsWith('--inspect=') || + arg.startsWith('--inspect-brk') || + arg.startsWith('--inspect-port'); + if (inspect) continue; + if (dropValue.has(arg)) { + i += 1; + continue; + } + if ( + arg.startsWith('--require=') || + arg.startsWith('--import=') || + arg.startsWith('--loader=') || + arg.startsWith('--experimental-loader=') || + arg.includes('tsx') + ) { + continue; + } + out.push(arg); + } + return out; +} + +function workerEnvironment(): Record { + const env: Record = {}; + for (const key of Object.keys(process.env)) { + const value = process.env[key]; + if (typeof value === 'string') env[key] = value; + } + return env; +} + +function actionableParseError(err: unknown, workerHeapMb: number | undefined): Error { + if (err instanceof ResourceLimitError || err instanceof ParseWorkerFailure) return err; + if (isWorkerOom(err)) { + return new ResourceLimitError( + `graph build stopped: a parse worker exceeded its ${workerHeapMb ?? '?'} MiB heap cap ` + + `(VG_WORKER_HEAP_MB). Raise the cap, exclude the offending files (--exclude), or ` + + `run single-threaded with --jobs 1.`, + ); + } + const raw = err instanceof Error ? err.message : String(err); + const line = raw.split('\n')[0]?.trim() || 'unknown error'; + const brief = line.length > 200 ? `${line.slice(0, 197)}...` : line; + const cap = workerHeapMb ? ` (currently ${workerHeapMb} MiB)` : ''; + return new ParseWorkerFailure( + `graph build stopped: a parse worker failed (${brief}). ` + + 'Re-run with --jobs 1 to parse in this process, exclude the offending files with --exclude, ' + + `or raise the per-worker heap with VG_WORKER_HEAP_MB${cap}.`, + ); } const DEFAULT_INLINE_THRESHOLD = 24; @@ -51,12 +310,15 @@ export async function parseFiles( // workers caps peak memory too (each worker holds its own grammar set). const jobs = Math.max(1, options.jobs ?? envJobs() ?? (Math.min(cores - 1, files.length) || 1)); - const workerFile = resolveWorkerFile(); + const workerFile = options.workerFile ?? resolveWorkerFile(); const useInline = options.inline === true || jobs <= 1 || - files.length < threshold || - workerFile === null; + workerFile === null || + // An explicit worker module (tests) always uses the pool. Production + // stays inline below the threshold, where spinning workers up costs more + // than it saves. + (options.workerFile === undefined && files.length < threshold); if (useInline) { // Inline runs in this process — apply the override directly. @@ -106,8 +368,17 @@ async function parsePooled( filename: workerFile, maxThreads: jobs, minThreads: 1, + runtime: options.workerRuntime ?? 'worker_threads', + terminateTimeout: 1_000, + execArgv: workerExecArgv(), + env: workerEnvironment(), ...(workerHeapMb ? { resourceLimits: { maxOldGenerationSizeMb: workerHeapMb } } : {}), }); + // Without this listener, a worker crash makes `destroy()` throw an uncaught + // TypeError (`emitter.removeListener is not a function`) and the process + // prints a stack instead of the message from `run()`. + pool.on('error', () => undefined); + armPool(pool); try { // More, smaller buckets than threads → finer live progress + better load // balancing. Round-robin keeps shards balanced; the final sort makes the @@ -132,16 +403,13 @@ async function parsePooled( ); return results.flat(); } catch (err) { - if (isWorkerOom(err)) { - throw new ResourceLimitError( - `graph build stopped: a parse worker exceeded its ${workerHeapMb ?? '?'} MiB heap cap ` + - `(VG_WORKER_HEAP_MB). Raise the cap, exclude the offending files (--exclude), or ` + - `run single-threaded with --jobs 1.`, - ); - } - throw err; + throw actionableParseError(err, workerHeapMb); } finally { - await pool.destroy(); + try { + await closePool(pool); + } finally { + disarmPool(); + } } } diff --git a/test/fixtures/parse-pool-harness.ts b/test/fixtures/parse-pool-harness.ts new file mode 100644 index 00000000..01d55157 --- /dev/null +++ b/test/fixtures/parse-pool-harness.ts @@ -0,0 +1,43 @@ +/** + * Scripted check for issue #301. Runs the same `parseFiles` pool `vg build` + * and `vg scan` use, then exits. The parent test asserts the process is gone + * and that no parse worker is left in the process list. + * + * tsx test/fixtures/parse-pool-harness.ts [worker_threads|child_process] + */ +import { fileURLToPath } from 'node:url'; +import type { DiscoveredFile } from '../../src/engine/discover.js'; +import { parseFiles } from '../../src/engine/pool.js'; + +const mode = process.argv[2] || 'crash'; +const runtime = process.argv[3] === 'child_process' ? 'child_process' : 'worker_threads'; +process.env.VG_POOL_WORKER_MODE = mode === 'crash' || mode === 'hang' ? mode : ''; + +const workerFile = fileURLToPath(new URL('./parse-pool-worker.mjs', import.meta.url)); +const lang = { id: 'ts' } as DiscoveredFile['lang']; +const files: DiscoveredFile[] = ['b.ts', 'a.ts'].map((rel) => ({ rel, abs: rel, lang })); + +let announced = false; +try { + await parseFiles(files, { + jobs: 2, + workerFile, + workerRuntime: runtime, + memoryBudgetMb: mode === 'budget' ? 1 : 0, + onProgress: () => { + if (!announced && mode === 'hang') { + announced = true; + process.stdout.write('READY\n'); + } + }, + }); + if (mode === 'hang' || mode === 'crash' || mode === 'budget') { + process.stderr.write(`expected ${mode} to fail\n`); + process.exit(2); + } + process.stdout.write('OK\n'); +} catch (err) { + const message = err instanceof Error ? err.message : String(err); + process.stderr.write(`${message}\n`); + process.exit(1); +} diff --git a/test/fixtures/parse-pool-worker.mjs b/test/fixtures/parse-pool-worker.mjs new file mode 100644 index 00000000..d3c1a135 --- /dev/null +++ b/test/fixtures/parse-pool-worker.mjs @@ -0,0 +1,32 @@ +/** + * Fixture worker for parse-pool shutdown tests. Not a grammar parser. + * `VG_POOL_WORKER_MODE` selects crash / hang; otherwise it echoes tasks in + * payload order so the parent can assert a stable sort. + */ +export default async function run(payload) { + const mode = process.env.VG_POOL_WORKER_MODE || ''; + if (mode === 'crash') { + // Uncaught, off the task promise. This is the failure that used to make + // pool shutdown throw `emitter.removeListener is not a function`. + setTimeout(() => { + throw new Error('parse worker crashed'); + }, 20); + await new Promise(() => {}); + } + if (mode === 'hang') { + await new Promise(() => {}); + } + const tasks = (payload && payload.tasks) || []; + return tasks.map((task) => ({ + rel: task.rel, + lang: task.lang, + hash: 'h', + bytes: 1, + defs: [], + calls: [], + imports: [], + heritage: [], + typeRefs: [], + guards: [], + })); +} diff --git a/test/parse-pool-shutdown.test.ts b/test/parse-pool-shutdown.test.ts new file mode 100644 index 00000000..500b7e81 --- /dev/null +++ b/test/parse-pool-shutdown.test.ts @@ -0,0 +1,250 @@ +import { createRequire } from 'node:module'; +import { spawn, spawnSync, type ChildProcess } from 'node:child_process'; +import * as path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { describe, expect, it } from 'vitest'; +import type { DiscoveredFile } from '../src/engine/discover.js'; +import { parseFiles, activeParseWorkerCount } from '../src/engine/pool.js'; + +const require = createRequire(import.meta.url); +const tsx = require.resolve('tsx/cli'); +const here = path.dirname(fileURLToPath(import.meta.url)); +const harness = path.join(here, 'fixtures/parse-pool-harness.ts'); +const workerFile = path.join(here, 'fixtures/parse-pool-worker.mjs'); + +interface ProcRow { + pid: number; + ppid: number; + stat: string; + args: string; +} + +function listProcesses(): ProcRow[] { + const out = spawnSync('ps', ['-eo', 'pid,ppid,stat,args'], { encoding: 'utf8' }); + const rows: ProcRow[] = []; + for (const line of (out.stdout || '').split('\n').slice(1)) { + const match = line.trim().match(/^(\d+)\s+(\d+)\s+(\S+)\s+(.*)$/); + if (!match) continue; + rows.push({ pid: Number(match[1]), ppid: Number(match[2]), stat: match[3], args: match[4] }); + } + return rows; +} + +function isWorkerProc(row: ProcRow): boolean { + if (row.stat.includes('Z')) return false; + return row.args.includes('entry/process.js') || row.args.includes('parse-pool-worker.mjs'); +} + +function workerPids(): Set { + return new Set(listProcesses().filter(isWorkerProc).map((row) => row.pid)); +} + +function descendants(pid: number): number[] { + const rows = listProcesses(); + const out: number[] = []; + const queue = [pid]; + const seen = new Set([pid]); + while (queue.length) { + const cur = queue.pop(); + if (cur === undefined) break; + for (const row of rows) { + if (row.ppid === cur && !seen.has(row.pid) && !row.stat.includes('Z')) { + seen.add(row.pid); + out.push(row.pid); + queue.push(row.pid); + } + } + } + return out; +} + +function isAlive(pid: number): boolean { + const row = listProcesses().find((item) => item.pid === pid); + if (!row || row.stat.includes('Z')) return false; + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +function delay(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +interface Running { + child: ChildProcess; + output: { stdout: string; stderr: string }; +} + +function startHarness(mode: string, runtime: 'worker_threads' | 'child_process'): Running { + const child = spawn(process.execPath, [tsx, harness, mode, runtime], { + stdio: ['ignore', 'pipe', 'pipe'], + env: { + ...process.env, + NO_COLOR: '1', + FORCE_COLOR: '0', + VIBGRATE_NO_KERNEL: '1', + }, + }); + const output = { stdout: '', stderr: '' }; + child.stdout?.on('data', (chunk: Buffer) => { + output.stdout += chunk.toString('utf8'); + }); + child.stderr?.on('data', (chunk: Buffer) => { + output.stderr += chunk.toString('utf8'); + }); + return { child, output }; +} + +function waitExit( + running: Running, + timeoutMs = 15_000, +): Promise<{ code: number | null; signal: NodeJS.Signals | null; stdout: string; stderr: string }> { + const { child, output } = running; + return new Promise((resolve, reject) => { + const finish = (code: number | null, signal: NodeJS.Signals | null): void => { + resolve({ code, signal, stdout: output.stdout, stderr: output.stderr }); + }; + if (child.exitCode !== null || child.signalCode !== null) { + setImmediate(() => finish(child.exitCode, child.signalCode)); + return; + } + const timer = setTimeout(() => { + child.kill('SIGKILL'); + reject(new Error(`harness timed out\nstdout:\n${output.stdout}\nstderr:\n${output.stderr}`)); + }, timeoutMs); + child.once('error', (err) => { + clearTimeout(timer); + reject(err); + }); + child.once('exit', (code, signal) => { + clearTimeout(timer); + // The signal path writes stderr and then exits; give the pipe a moment. + setTimeout(() => finish(code, signal), 100); + }); + }); +} + +async function waitForStdout(running: Running, needle: string): Promise { + const started = Date.now(); + while (!running.output.stdout.includes(needle)) { + if (Date.now() - started > 10_000) { + throw new Error(`timed out waiting for ${needle}\nstdout:\n${running.output.stdout}\nstderr:\n${running.output.stderr}`); + } + const { child } = running; + if (child.exitCode !== null || child.signalCode !== null) { + throw new Error( + `harness exited before ${needle} (code ${child.exitCode} signal ${child.signalCode})\n${running.output.stderr}`, + ); + } + await delay(30); + } +} + +const files: DiscoveredFile[] = ['b.ts', 'a.ts'].map((rel) => ({ + rel, + abs: rel, + lang: { id: 'ts' } as DiscoveredFile['lang'], +})); + +describe('parse worker shutdown', () => { + it('keeps pooled parse output ordered and stable', async () => { + delete process.env.VG_POOL_WORKER_MODE; + const opts = { jobs: 2, workerFile, memoryBudgetMb: 0 as const }; + const once = await parseFiles(files, opts); + const twice = await parseFiles(files, opts); + expect(once.map((row) => row.rel)).toEqual(['a.ts', 'b.ts']); + expect(JSON.stringify(once)).toBe(JSON.stringify(twice)); + expect(activeParseWorkerCount()).toBe(0); + }); + + it('a memory-budget failure during parse exits the pool', async () => { + delete process.env.VG_POOL_WORKER_MODE; + const run = parseFiles(files, { jobs: 2, workerFile, memoryBudgetMb: 1 }); + const timeout = new Promise((_resolve, reject) => { + setTimeout(() => reject(new Error('parse hung after a budget failure')), 8_000); + }); + await expect(Promise.race([run, timeout])).rejects.toThrow(/VG_MEMORY_BUDGET_MB/); + expect(activeParseWorkerCount()).toBe(0); + }); + + it.each(['worker_threads', 'child_process'] as const)( + 'a finished %s pool exits and leaves no workers', + async (runtime) => { + const before = workerPids(); + const running = startHarness('ok', runtime); + const result = await waitExit(running); + expect(result.code).toBe(0); + expect(result.stdout).toContain('OK'); + expect(result.stderr).not.toContain('removeListener'); + await delay(200); + const leaked = [...workerPids()].filter((pid) => !before.has(pid) && isAlive(pid)); + expect(leaked).toEqual([]); + }, + ); + + it.each(['worker_threads', 'child_process'] as const)( + 'a crashing %s worker exits with an actionable error and leaves no workers', + async (runtime) => { + const before = workerPids(); + const running = startHarness('crash', runtime); + const result = await waitExit(running); + expect(result.code).toBe(1); + expect(result.stderr).toContain('graph build stopped: a parse worker failed'); + expect(result.stderr).toContain('--jobs 1'); + expect(result.stderr).toContain('--exclude'); + expect(result.stderr).toContain('VG_WORKER_HEAP_MB'); + expect(result.stderr).not.toContain('removeListener'); + // worker_threads share the parent's stderr. A shutdown bug prints a + // `removeListener` stack from node:internal; the actionable line must + // be the whole failure. child_process workers print their own crash + // before exiting — that text is the worker, and the parent must still + // exit and reap it. + if (runtime === 'worker_threads') { + expect(result.stderr).not.toContain('node:internal'); + expect(result.stderr).not.toMatch(/\n\s+at /); + } + await delay(300); + const leaked = [...workerPids()].filter((pid) => !before.has(pid) && isAlive(pid)); + expect(leaked).toEqual([]); + expect(isAlive(running.child.pid ?? -1)).toBe(false); + }, + ); + + it.each([ + ['worker_threads', 'SIGTERM', 143, 'terminated while parsing'], + ['child_process', 'SIGTERM', 143, 'terminated while parsing'], + ['child_process', 'SIGINT', 130, 'interrupted while parsing'], + ] as const)( + 'a hung %s pool stops workers on %s', + async (runtime, signal, code, message) => { + const before = workerPids(); + const running = startHarness('hang', runtime); + await waitForStdout(running, 'READY'); + let kids: number[] = []; + if (runtime === 'child_process') { + const started = Date.now(); + while (kids.length === 0 && Date.now() - started < 5_000) { + kids = descendants(running.child.pid ?? -1); + if (kids.length === 0) await delay(50); + } + expect(kids.length).toBeGreaterThan(0); + } + running.child.kill(signal); + const result = await waitExit(running); + expect(result.signal).toBeNull(); + expect(result.code).toBe(code); + expect(result.stderr).toContain(message); + expect(result.stderr).not.toContain('removeListener'); + expect(result.stderr).not.toMatch(/\n\s+at /); + await delay(300); + const still = kids.filter((pid) => isAlive(pid)); + expect(still).toEqual([]); + const leaked = [...workerPids()].filter((pid) => !before.has(pid) && isAlive(pid)); + expect(leaked).toEqual([]); + expect(isAlive(running.child.pid ?? -1)).toBe(false); + }, + ); +}); From 1b80bbf2fb6f498590a91699dc891c308fe6fd0a Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sat, 10 Oct 2026 13:54:02 +0000 Subject: [PATCH 2/2] fix(engine): stop parse workers when build or scan ends Ctrl-C, SIGTERM, a worker failure, and a normal finish all shut the parse pool down, so Node workers are not left running afterwards. A worker that fails still prints what happened and how to re-run (--jobs 1 or --exclude). Signed-off-by: Cursor Agent Co-authored-by: vibgrate-team --- ARCHITECTURE.md | 4 +- CHANGELOG.md | 5 + DOCS.md | 8 - src/cli.ts | 9 +- src/commands/build.ts | 4 +- src/core-open/run-core-scan.ts | 8 +- src/engine/pool-guard.ts | 247 ++++++++++++++++++ src/engine/pool-teardown.fixture.ts | 31 +++ src/engine/pool-teardown.test.ts | 381 ++++++++++++++++++++++++++++ src/engine/pool.ts | 375 ++++++--------------------- src/reporting/commands/scan.ts | 106 ++++---- test/fixtures/parse-pool-harness.ts | 43 ---- test/fixtures/parse-pool-worker.mjs | 32 --- test/parse-pool-shutdown.test.ts | 250 ------------------ 14 files changed, 808 insertions(+), 695 deletions(-) create mode 100644 src/engine/pool-guard.ts create mode 100644 src/engine/pool-teardown.fixture.ts create mode 100644 src/engine/pool-teardown.test.ts delete mode 100644 test/fixtures/parse-pool-harness.ts delete mode 100644 test/fixtures/parse-pool-worker.mjs delete mode 100644 test/parse-pool-shutdown.test.ts diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 34e3516c..475150fb 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -108,9 +108,7 @@ meaningful. (`engine/discover.ts` sorts). - **Parallel but ordered** — the worker pool (`engine/pool.ts` + `engine/parse-worker.ts`) parses files concurrently for speed, then results are - re-ordered deterministically before they enter the graph. The pool stops its - workers when parsing finishes, when parsing fails, and on SIGINT/SIGTERM, so - a build or scan does not leave parse workers behind. + re-ordered deterministically before they enter the graph. If you touch anything that ends up in `graph.json` or a report, add a test that asserts stable output. diff --git a/CHANGELOG.md b/CHANGELOG.md index 887e092f..a36c361b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,6 +42,11 @@ backward compatible. ### Fixed +- **Interrupted or failed `vg build` and `vg scan` stop their parse workers.** + Ctrl-C, a worker failure, and a normal finish all shut the parse pool down, + so Node workers are not left running afterwards. A worker that fails still + prints what happened and how to re-run (`--jobs 1` or `--exclude`). + - **SBOM component order no longer follows scan or filesystem order.** `vg sbom export` and CycloneDX/SPDX graph export sort a component by its Package URL when one is written, otherwise by package name, then by version. diff --git a/DOCS.md b/DOCS.md index 2af0f4ef..b6b753ce 100644 --- a/DOCS.md +++ b/DOCS.md @@ -5101,14 +5101,6 @@ and `0` always means "disabled". | `VG_JOBS` | CPU cores − 1 | Default parse worker count when `--jobs` isn't passed. Fewer workers = lower peak memory (each worker loads its own grammar set). | | `VG_WORKER_HEAP_MB` | platform default | Per-worker old-generation heap cap, so one runaway parse can't take the whole machine. | -If a parse worker crashes, runs out of memory, or the command is interrupted -(Ctrl+C or SIGTERM), `vg` stops the workers before it exits. The failure is one -line on stderr — what failed, and whether to re-run with `--jobs 1`, exclude -files, or raise `VG_WORKER_HEAP_MB` — not a stack trace. A finished build stops -its workers the same way. Check with a process list (`ps`) filtered for this -command's pid: after `vg build` or `vg scan` returns, no parse worker should -still be running. Graph output is unchanged. - Skips are deterministic functions of the input (file size, file count) — never of observed memory — so identical input still produces an identical `graph.json`. To give the build more room instead of limiting it, raise the diff --git a/src/cli.ts b/src/cli.ts index 23a45c54..0241338b 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -46,6 +46,7 @@ import { ConfigFileError } from './core-open/config.js'; import { LockfileParseError } from './core-open/utils/lockfile-parse.js'; import { UnsafeRootError } from './core-open/utils/root-safety.js'; import { CliError, ExitCode, usageError } from './util/exit.js'; +import { ParseWorkerError } from './engine/pool.js'; import { c, info, disableColor, exitAfterFlush } from './util/output.js'; // Drift-reporting commands (merged from the Vibgrate CLI). These run on the @@ -438,7 +439,13 @@ function handleError(err: unknown): never { /* stdout closed */ } }; - if (err instanceof CliError || err instanceof LockfileParseError || err instanceof ConfigFileError || err instanceof UnsafeRootError) { + if ( + err instanceof CliError || + err instanceof ParseWorkerError || + err instanceof LockfileParseError || + err instanceof ConfigFileError || + err instanceof UnsafeRootError + ) { const code = err instanceof CliError ? err.code : ExitCode.ERROR; emitHostError(err.message); info(c.red(`error: ${err.message}`)); diff --git a/src/commands/build.ts b/src/commands/build.ts index 53921911..e321702b 100644 --- a/src/commands/build.ts +++ b/src/commands/build.ts @@ -22,7 +22,7 @@ import { renderReport } from '../engine/report.js'; import { renderHtml } from '../engine/html.js'; import { UsageError, mergeExcludes } from '../engine/discover.js'; import { ResourceLimitError } from '../engine/limits.js'; -import { ParseWorkerFailure } from '../engine/pool.js'; +import { ParseWorkerError } from '../engine/pool.js'; import { UnsafeRootError } from '../core-open/utils/root-safety.js'; import { CliError, ExitCode, usageError } from '../util/exit.js'; import { resolveSelfJsEntry } from '../util/cli-invocation.js'; @@ -151,7 +151,7 @@ export async function runBuild( } catch (err) { bar?.done(); if (err instanceof UsageError) throw usageError(err.message); - if (err instanceof ResourceLimitError || err instanceof ParseWorkerFailure) { + if (err instanceof ResourceLimitError || err instanceof ParseWorkerError) { throw new CliError(err.message, ExitCode.ERROR); } if (err instanceof UnsafeRootError) throw new CliError(err.message, ExitCode.ERROR); diff --git a/src/core-open/run-core-scan.ts b/src/core-open/run-core-scan.ts index 9e2b1023..a6d6ebc5 100644 --- a/src/core-open/run-core-scan.ts +++ b/src/core-open/run-core-scan.ts @@ -764,12 +764,8 @@ export async function runCoreScan( { projects: allProjects, solutions, extended }, ); progress.completeStep('map', detail || 'done'); - } catch (err) { - // One line, no stack: a worker crash or resource limit should say what - // failed and what to do, then the scan continues without the map. - const line = (err instanceof Error ? err.message : String(err)).split('\n')[0]?.trim() || 'map build failed'; - const detail = line.length > 240 ? `${line.slice(0, 237)}...` : line; - progress.completeStep('map', `skipped (${detail})`); + } catch { + progress.completeStep('map', 'skipped (map build failed)'); } } diff --git a/src/engine/pool-guard.ts b/src/engine/pool-guard.ts new file mode 100644 index 00000000..a3a8baa2 --- /dev/null +++ b/src/engine/pool-guard.ts @@ -0,0 +1,247 @@ +/** + * Lifetime of the parse-worker pool. + * + * tinypool keeps its threads alive until `destroy()` settles. A worker stuck + * in native code, a SIGINT/SIGTERM, or `process.exit` from the progress UI + * (which runs before this function's `finally`) used to skip that call, so a + * failed or interrupted `vg build` / `vg scan` left Node workers running and + * could swallow the error the user needed. This guard always stops the pool, + * then lets the original failure through. + */ + +/** A thread the pool can still be holding after `destroy()` gives up. */ +export interface PoolThread { + terminate?: () => Promise; + unref?: () => void; +} + +/** The slice of tinypool this guard needs. Tests pass a fake. */ +export interface ManagedPool { + destroy(): Promise; + cancelPendingTasks?(): void; + threads?: readonly PoolThread[]; + on?(event: 'error', listener: (err: unknown) => void): void; +} + +export interface PoolGuardHost { + prependListener(signal: NodeJS.Signals, listener: () => void): void; + removeListener(signal: NodeJS.Signals, listener: () => void): void; + /** Fires when the loop would otherwise drain with workers still tracked. */ + onBeforeExit?(listener: () => void): () => void; + /** + * Install a wrapper around the host's hard exit. The returned function + * restores the previous exit. The wrapper must not itself exit. + */ + installExit(handler: (code?: number) => void): () => void; + /** Hard-exit. Called only after workers have been asked to stop. */ + exit(code: number): void; + writeError(message: string): void; +} + +export const PARSE_INTERRUPTED_MESSAGE = + 'interrupted. Parse workers were stopped. Re-run the command to continue.'; + +export const PARSE_TERMINATED_MESSAGE = + 'terminated. Parse workers were stopped. Re-run the command to continue.'; + +const DEFAULT_BUDGET_MS = 2_000; + +function delay(ms: number, ref: boolean): Promise { + return new Promise((resolve) => { + const timer = setTimeout(resolve, ms); + if (!ref) timer.unref?.(); + }); +} + +async function stopThread(thread: PoolThread, budgetMs: number): Promise { + try { + thread.unref?.(); + } catch { + /* already gone */ + } + if (!thread.terminate) return; + await Promise.race([ + thread.terminate().then( + () => undefined, + () => undefined, + ), + delay(budgetMs, false), + ]); + try { + thread.unref?.(); + } catch { + /* already gone */ + } +} + +async function stopPool(pool: ManagedPool, budgetMs: number): Promise { + const threads = [...(pool.threads ?? [])]; + for (const thread of threads) { + try { + thread.unref?.(); + } catch { + /* already gone */ + } + } + try { + pool.cancelPendingTasks?.(); + } catch { + /* already closing */ + } + await Promise.race([ + pool.destroy().then( + () => undefined, + () => undefined, + ), + delay(budgetMs, false), + ]); + await Promise.all(threads.map((thread) => stopThread(thread, Math.min(budgetMs, 500)))); +} + +export class ParsePoolGuard { + private readonly pools = new Set(); + private readonly stopping = new Map>(); + /** Bodies currently inside `using`. beforeExit must not stop those. */ + private active = 0; + private armed = false; + private exitStarted = false; + private pendingCode = 0; + private restoreExit: (() => void) | null = null; + private removeBeforeExit: (() => void) | null = null; + + private readonly onSigint = (): void => { + this.requestExit(130); + }; + + private readonly onSigterm = (): void => { + this.requestExit(143); + }; + + private readonly onBeforeExit = (): void => { + // A live `using()` body is still parsing. Its message ports normally keep + // the loop alive; if they don't, exiting is what stops unref'd threads. + // Only a pool that nobody is waiting on needs an explicit stop here. + if (this.active > 0 || this.pools.size === 0 || this.exitStarted) return; + void this.releaseAll(); + }; + + constructor( + private readonly host: PoolGuardHost, + private readonly budgetMs = DEFAULT_BUDGET_MS, + ) {} + + /** True once a signal or a hard exit has taken over shutdown. */ + get isExiting(): boolean { + return this.exitStarted; + } + + track(pool: ManagedPool): void { + this.pools.add(pool); + this.arm(); + } + + /** Run `body` with the pool tracked, and always stop the pool afterwards. */ + async using(pool: ManagedPool, body: () => Promise): Promise { + this.track(pool); + this.active += 1; + try { + return await body(); + } finally { + this.active -= 1; + await this.release(pool); + } + } + + async release(pool: ManagedPool): Promise { + const existing = this.stopping.get(pool); + if (existing) return existing; + this.pools.delete(pool); + const job = stopPool(pool, this.budgetMs).finally(() => { + this.stopping.delete(pool); + if (this.pools.size === 0 && this.stopping.size === 0 && !this.exitStarted) this.disarm(); + }); + this.stopping.set(pool, job); + return job; + } + + private async releaseAll(): Promise { + const pools = [...this.pools]; + await Promise.all(pools.map((pool) => this.release(pool))); + } + + private arm(): void { + if (this.armed) return; + this.armed = true; + this.host.prependListener('SIGINT', this.onSigint); + this.host.prependListener('SIGTERM', this.onSigterm); + this.removeBeforeExit = this.host.onBeforeExit?.(this.onBeforeExit) ?? null; + this.restoreExit = this.host.installExit((code) => { + this.requestExit(typeof code === 'number' ? code : 0); + }); + } + + private disarm(): void { + if (!this.armed) return; + this.armed = false; + this.host.removeListener('SIGINT', this.onSigint); + this.host.removeListener('SIGTERM', this.onSigterm); + this.removeBeforeExit?.(); + this.removeBeforeExit = null; + this.restoreExit?.(); + this.restoreExit = null; + } + + private requestExit(code: number): void { + if (!this.exitStarted) this.pendingCode = code; + if (this.exitStarted) return; + this.exitStarted = true; + if (code === 130) this.host.writeError(`error: ${PARSE_INTERRUPTED_MESSAGE}\n`); + else if (code === 143) this.host.writeError(`error: ${PARSE_TERMINATED_MESSAGE}\n`); + const exitBudget = this.budgetMs + 500; + void Promise.race([this.releaseAll(), delay(exitBudget, true)]).finally(() => { + this.disarm(); + this.host.exit(this.pendingCode); + }); + } +} + +export function nodePoolGuardHost(): PoolGuardHost { + return { + prependListener(signal, listener) { + process.prependListener(signal, listener); + }, + removeListener(signal, listener) { + process.removeListener(signal, listener); + }, + onBeforeExit(listener) { + process.on('beforeExit', listener); + return () => { + process.removeListener('beforeExit', listener); + }; + }, + installExit(handler) { + const previous = process.exit; + const wrapped = ((code?: number): never => { + handler(code); + return undefined as never; + }) as typeof process.exit; + process.exit = wrapped; + return () => { + if (process.exit === wrapped) process.exit = previous; + }; + }, + exit(code) { + process.exit(code); + }, + writeError(message) { + try { + process.stderr.write(message); + } catch { + /* stderr already closed */ + } + }, + }; +} + +/** Process-wide guard used by the parse pool. Armed only while a pool is live. */ +export const parsePoolGuard = new ParsePoolGuard(nodePoolGuardHost()); diff --git a/src/engine/pool-teardown.fixture.ts b/src/engine/pool-teardown.fixture.ts new file mode 100644 index 00000000..de61b87e --- /dev/null +++ b/src/engine/pool-teardown.fixture.ts @@ -0,0 +1,31 @@ +/** + * Child entry for the parse-pool teardown check. Spawned by + * pool-teardown.test.ts — not a test, and not a CLI command. + * + * argv: + */ +import type { DiscoveredFile } from './discover.js'; +import type { LanguageDef } from './languages.js'; +import { parseFiles } from './pool.js'; + +const mode = process.argv[2]; +const workerFile = process.argv[3]; +if (!mode || !workerFile) { + process.stderr.write('usage: pool-teardown.fixture.ts \n'); + process.exit(2); +} +process.env.VG_POOL_TEST_MODE = mode; + +const lang = { id: 'typescript' } as LanguageDef; +const files: DiscoveredFile[] = [ + { rel: 'a.ts', abs: '/tmp/vg-pool-a.ts', lang }, + { rel: 'b.ts', abs: '/tmp/vg-pool-b.ts', lang }, +]; + +try { + await parseFiles(files, { jobs: 2, inlineThreshold: 0, workerFile }); +} catch (err) { + const message = err instanceof Error ? err.message : String(err); + process.stderr.write(`${message}\n`); + process.exit(1); +} diff --git a/src/engine/pool-teardown.test.ts b/src/engine/pool-teardown.test.ts new file mode 100644 index 00000000..c276d90f --- /dev/null +++ b/src/engine/pool-teardown.test.ts @@ -0,0 +1,381 @@ +import { spawn, type ChildProcess } from 'node:child_process'; +import * as fs from 'node:fs'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { describe, it, expect, afterEach } from 'vitest'; +import { ResourceLimitError } from './limits.js'; +import { ParseWorkerError, toParseWorkerError } from './pool.js'; +import { + PARSE_INTERRUPTED_MESSAGE, + PARSE_TERMINATED_MESSAGE, + ParsePoolGuard, + type ManagedPool, + type PoolGuardHost, + type PoolThread, +} from './pool-guard.js'; + +const pkgRoot = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '../..'); +const fixture = path.join(pkgRoot, 'src/engine/pool-teardown.fixture.ts'); + +function fakeThread(): PoolThread & { terminated: number; unrefed: number } { + const thread = { + terminated: 0, + unrefed: 0, + terminate(): Promise { + thread.terminated += 1; + return Promise.resolve(0); + }, + unref(): void { + thread.unrefed += 1; + }, + }; + return thread; +} + +function fakeHost(): PoolGuardHost & { + emit(signal: NodeJS.Signals): void; + callExit(code?: number): void; + fireBeforeExit(): void; + exits: number[]; + errors: string[]; + listenerCount(signal: NodeJS.Signals): number; +} { + const listeners = new Map void>>(); + const exits: number[] = []; + const errors: string[] = []; + let onExit: ((code?: number) => void) | null = null; + let beforeExit: (() => void) | null = null; + return { + exits, + errors, + prependListener(signal, listener) { + const list = listeners.get(signal) ?? []; + list.unshift(listener); + listeners.set(signal, list); + }, + removeListener(signal, listener) { + listeners.set( + signal, + (listeners.get(signal) ?? []).filter((item) => item !== listener), + ); + }, + onBeforeExit(listener) { + beforeExit = listener; + return () => { + if (beforeExit === listener) beforeExit = null; + }; + }, + installExit(handler) { + onExit = handler; + return () => { + onExit = null; + }; + }, + exit(code) { + exits.push(code); + }, + writeError(message) { + errors.push(message); + }, + emit(signal) { + for (const listener of [...(listeners.get(signal) ?? [])]) listener(); + }, + callExit(code) { + onExit?.(code); + }, + fireBeforeExit() { + beforeExit?.(); + }, + listenerCount(signal) { + return listeners.get(signal)?.length ?? 0; + }, + }; +} + +function livePool(thread: PoolThread, destroy: () => Promise): ManagedPool & { cancelled: number } { + const pool = { + cancelled: 0, + threads: [thread], + cancelPendingTasks() { + pool.cancelled += 1; + }, + destroy, + }; + return pool; +} + +describe('parse worker errors stay actionable', () => { + it('turns an out-of-memory worker into a resource error that says how to re-run', () => { + const err = toParseWorkerError( + Object.assign(new Error('JavaScript heap out of memory'), { code: 'ERR_WORKER_OUT_OF_MEMORY' }), + 256, + ); + expect(err).toBeInstanceOf(ResourceLimitError); + expect(err.message).toContain('256'); + expect(err.message).toContain('VG_WORKER_HEAP_MB'); + expect(err.message).toContain('--jobs 1'); + expect(err.message).toContain('--exclude'); + }); + + it('keeps an existing resource error and wraps other failures without file paths', () => { + const existing = new ResourceLimitError('already actionable'); + expect(toParseWorkerError(existing)).toBe(existing); + + const wrapped = toParseWorkerError(new Error('boom-parse at /home/dev/repo/src/app.ts\nstack')); + expect(wrapped).toBeInstanceOf(ParseWorkerError); + expect(wrapped.message).toContain('boom-parse'); + expect(wrapped.message).toContain('--jobs 1'); + expect(wrapped.message).toContain('--exclude'); + expect(wrapped.message).not.toContain('/home/dev'); + expect(wrapped.message).not.toContain('\n'); + }); +}); + +describe('parse pool guard', () => { + it('stops the pool after a normal run and drops its signal handlers', async () => { + const host = fakeHost(); + const guard = new ParsePoolGuard(host, 30); + const thread = fakeThread(); + let destroyed = 0; + const pool = livePool(thread, async () => { + destroyed += 1; + }); + await expect(guard.using(pool, async () => 'ok')).resolves.toBe('ok'); + expect(destroyed).toBe(1); + expect(thread.terminated).toBeGreaterThan(0); + expect(host.exits).toEqual([]); + expect(host.listenerCount('SIGINT')).toBe(0); + expect(host.listenerCount('SIGTERM')).toBe(0); + }); + + it('stops the pool when the run fails and still throws the original error', async () => { + const host = fakeHost(); + const guard = new ParsePoolGuard(host, 30); + const thread = fakeThread(); + let destroyed = 0; + const pool = livePool(thread, async () => { + destroyed += 1; + }); + const err = new Error('keep-me'); + await expect( + guard.using(pool, async () => { + throw err; + }), + ).rejects.toBe(err); + expect(destroyed).toBe(1); + expect(thread.terminated).toBeGreaterThan(0); + expect(host.exits).toEqual([]); + expect(host.errors).toEqual([]); + }); + + it('stops threads when destroy never settles', async () => { + const host = fakeHost(); + const guard = new ParsePoolGuard(host, 40); + const thread = fakeThread(); + const pool = livePool(thread, () => new Promise(() => undefined)); + await expect(guard.using(pool, async () => 7)).resolves.toBe(7); + expect(pool.cancelled).toBe(1); + expect(thread.terminated).toBeGreaterThan(0); + expect(thread.unrefed).toBeGreaterThan(0); + expect(host.listenerCount('SIGINT')).toBe(0); + }); + + it('tears the pool down on SIGINT and exits 130 with an actionable line', async () => { + const host = fakeHost(); + const guard = new ParsePoolGuard(host, 40); + const thread = fakeThread(); + const pool = livePool(thread, () => new Promise(() => undefined)); + void guard.using(pool, () => new Promise(() => undefined)); + await new Promise((resolve) => setTimeout(resolve, 0)); + host.emit('SIGINT'); + host.emit('SIGINT'); + await waitFor(() => host.exits.length > 0); + expect(host.exits).toEqual([130]); + expect(host.errors).toEqual([`error: ${PARSE_INTERRUPTED_MESSAGE}\n`]); + expect(pool.cancelled).toBe(1); + expect(thread.terminated).toBeGreaterThan(0); + expect(host.listenerCount('SIGINT')).toBe(0); + }); + + it('tears the pool down on SIGTERM and exits 143', async () => { + const host = fakeHost(); + const guard = new ParsePoolGuard(host, 40); + const thread = fakeThread(); + const pool = livePool(thread, async () => undefined); + void guard.using(pool, () => new Promise(() => undefined)); + await new Promise((resolve) => setTimeout(resolve, 0)); + host.emit('SIGTERM'); + await waitFor(() => host.exits.length > 0); + expect(host.exits).toEqual([143]); + expect(host.errors[0]).toContain(PARSE_TERMINATED_MESSAGE); + }); + + it('stops the pool when the process exits and does not add an interrupt line', async () => { + const host = fakeHost(); + const guard = new ParsePoolGuard(host, 40); + const thread = fakeThread(); + const pool = livePool(thread, async () => undefined); + void guard.using(pool, () => new Promise(() => undefined)); + await new Promise((resolve) => setTimeout(resolve, 0)); + host.callExit(1); + await waitFor(() => host.exits.length > 0); + expect(host.exits).toEqual([1]); + expect(host.errors).toEqual([]); + expect(thread.terminated).toBeGreaterThan(0); + expect(guard.isExiting).toBe(true); + }); + + it('stops a pool that is still tracked when the event loop would drain', async () => { + const host = fakeHost(); + const guard = new ParsePoolGuard(host, 40); + const thread = fakeThread(); + let destroyed = 0; + const pool = livePool(thread, async () => { + destroyed += 1; + }); + guard.track(pool); + host.fireBeforeExit(); + await waitFor(() => destroyed === 1 && thread.terminated > 0); + expect(thread.terminated).toBeGreaterThan(0); + expect(host.listenerCount('SIGINT')).toBe(0); + }); +}); + +describe('parse pool process teardown', () => { + const children: ChildProcess[] = []; + const dirs: string[] = []; + + afterEach(() => { + for (const child of children) { + if (child.pid && child.exitCode === null && !child.killed) { + try { + child.kill('SIGKILL'); + } catch { + /* already gone */ + } + } + } + children.length = 0; + while (dirs.length) fs.rmSync(dirs.pop()!, { recursive: true, force: true }); + }); + + it('stops real workers on success, failure, and SIGINT', async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'vg-pool-')); + dirs.push(dir); + const worker = path.join(dir, 'worker.mjs'); + fs.writeFileSync( + worker, + [ + "import * as fs from 'node:fs';", + 'export default async function run() {', + ' const mode = process.env.VG_POOL_TEST_MODE;', + " if (mode === 'hang') {", + ' const ready = process.env.VG_POOL_READY;', + " if (ready) fs.writeFileSync(ready, 'ready');", + ' const stamp = process.env.VG_POOL_HEARTBEAT;', + ' if (stamp) {', + ' const beat = () => { try { fs.writeFileSync(stamp, String(Date.now())); } catch { /* parent left */ } };', + ' beat();', + ' setInterval(beat, 40);', + ' }', + ' await new Promise(() => undefined);', + ' }', + " if (mode === 'fail') throw new Error('boom-parse');", + ' return [];', + '}', + '', + ].join('\n'), + ); + + const ok = await runMode(worker, 'ok'); + expect(ok.code).toBe(0); + expect(ok.stderr).not.toContain('error:'); + expect(alive(ok.pid)).toBe(false); + + const failed = await runMode(worker, 'fail'); + expect(failed.code).toBe(1); + expect(failed.stderr).toContain('boom-parse'); + expect(failed.stderr).toContain('--jobs 1'); + expect(failed.stderr).toContain('--exclude'); + expect(failed.stderr).not.toContain(dir); + expect(alive(failed.pid)).toBe(false); + + const heartbeat = path.join(dir, 'beat'); + const ready = path.join(dir, 'ready'); + const interrupted = await runMode(worker, 'hang', heartbeat, ready); + expect(interrupted.code).toBe(130); + expect(interrupted.stderr).toContain(PARSE_INTERRUPTED_MESSAGE); + expect(interrupted.stderr).not.toContain('Terminating worker thread'); + expect(alive(interrupted.pid)).toBe(false); + const stamped = fs.readFileSync(heartbeat, 'utf8'); + await new Promise((resolve) => setTimeout(resolve, 200)); + expect(fs.readFileSync(heartbeat, 'utf8')).toBe(stamped); + }, 40_000); + + async function runMode( + worker: string, + mode: 'ok' | 'fail' | 'hang', + heartbeat?: string, + readyFile?: string, + ): Promise<{ code: number; stderr: string; pid: number }> { + const child = spawn(process.execPath, ['--import', 'tsx', fixture, mode, worker], { + cwd: pkgRoot, + env: { + ...process.env, + NO_COLOR: '1', + VG_POOL_TEST_MODE: mode, + ...(heartbeat ? { VG_POOL_HEARTBEAT: heartbeat } : {}), + ...(readyFile ? { VG_POOL_READY: readyFile } : {}), + }, + stdio: ['ignore', 'pipe', 'pipe'], + }); + children.push(child); + const pid = child.pid; + if (!pid) throw new Error('parse pool fixture failed to start'); + let stderr = ''; + child.stderr?.on('data', (chunk: Buffer) => { + stderr += chunk.toString('utf8'); + }); + child.stdout?.on('data', () => undefined); + const exited = waitExit(child, 20_000); + if (mode === 'hang') { + await waitFor(() => (readyFile ? fs.existsSync(readyFile) : false), 20_000); + child.kill('SIGINT'); + } + const code = await exited; + return { code, stderr, pid }; + } +}); + +function alive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +function waitExit(child: ChildProcess, ms: number): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(() => { + child.kill('SIGKILL'); + reject(new Error('parse pool fixture did not exit')); + }, ms); + child.once('exit', (code, signal) => { + clearTimeout(timer); + if (signal === 'SIGINT') resolve(130); + else if (signal) reject(new Error(`fixture died on ${signal}`)); + else resolve(code ?? 1); + }); + }); +} + +async function waitFor(pred: () => boolean, ms = 2_000): Promise { + const start = Date.now(); + while (!pred()) { + if (Date.now() - start > ms) throw new Error('timed out waiting for parse pool teardown'); + await new Promise((resolve) => setTimeout(resolve, 15)); + } +} diff --git a/src/engine/pool.ts b/src/engine/pool.ts index ad5e8b8b..67595637 100644 --- a/src/engine/pool.ts +++ b/src/engine/pool.ts @@ -9,6 +9,7 @@ import type { DiscoveredFile } from './discover.js'; import type { FileParse } from './types.js'; import { stampWarning, WARNING_CODES } from '../core-open/warnings.js'; import type { ParseTask } from './parse-worker.js'; +import { parsePoolGuard, type ManagedPool } from './pool-guard.js'; /** * Parse a set of discovered files into FileParse tables. @@ -22,10 +23,6 @@ import type { ParseTask } from './parse-worker.js'; * single-threaded inline path is always available (and used for small repos or * when the worker module isn't resolvable, e.g. under ts-only test runners), * producing byte-identical output to the pooled path. - * - * Workers are stopped on success, on failure, and on SIGINT/SIGTERM. A crashed - * worker must not dump a stack or leave a Node process behind after `vg build` - * or `vg scan` returns. */ export interface ParseOptions { @@ -41,261 +38,8 @@ export interface ParseOptions { grammarsDir?: string; /** Heap budget (MiB) checked as parse results accumulate; 0/unset skips. */ memoryBudgetMb?: number; - /** - * Parse-worker module. Tests pass a fixture; production resolves the bundled - * `parse-worker.js` next to this file. - */ + /** Worker module path. Set by tests; production resolves the bundled worker. */ workerFile?: string; - /** - * `child_process` runs each worker as its own OS process so a crash can be - * killed by pid. The default `worker_threads` pool is what `vg build` uses. - */ - workerRuntime?: 'worker_threads' | 'child_process'; -} - -/** A parse worker died or was stopped. The message is the whole user-facing error. */ -export class ParseWorkerFailure extends Error { - readonly isParseWorkerFailure = true; - constructor(message: string) { - super(message); - this.name = 'ParseWorkerFailure'; - } -} - -/** How long `destroy()` may take before workers are killed outright. */ -const POOL_SHUTDOWN_MS = 2_000; - -interface KillableWorker { - terminate?: () => Promise; - unref?: () => void; - /** Set when the worker is a `child_process` (tinypool `ProcessWorker`). */ - process?: { kill: (signal?: NodeJS.Signals | number) => boolean }; -} - -interface ManagedPool { - threads: KillableWorker[]; - destroy: () => Promise; - cancelPendingTasks: () => void; - on: (event: 'error', listener: (err: unknown) => void) => void; -} - -let activePool: ManagedPool | null = null; -let detachSignals: (() => void) | null = null; -let exitHookInstalled = false; - -/** Workers still tracked by the live parse pool. Zero after shutdown. */ -export function activeParseWorkerCount(): number { - if (!activePool) return 0; - try { - return activePool.threads.length; - } catch { - return 0; - } -} - -function installExitHook(): void { - if (exitHookInstalled) return; - exitHookInstalled = true; - // `exit` is synchronous. Child workers are not reaped with the parent, so - // kill them here — including when a signal handler calls `process.exit`. - process.on('exit', () => { - if (activePool) killWorkers(activePool); - }); -} - -function killWorkers(pool: ManagedPool): void { - let workers: KillableWorker[] = []; - try { - workers = pool.threads; - } catch { - return; - } - for (const worker of workers) killWorker(worker); -} - -function killWorker(worker: KillableWorker): void { - const child = worker.process; - if (child && typeof child.kill === 'function') { - try { - child.kill('SIGKILL'); - } catch { - // Already reaped. - } - return; - } - try { - worker.unref?.(); - } catch { - // Already gone. - } - try { - void worker.terminate?.(); - } catch { - // Already gone. - } -} - -function onStopSignal(signal: NodeJS.Signals): void { - const pool = activePool; - if (!pool) return; - killWorkers(pool); - const code = signal === 'SIGINT' ? 130 : 143; - const word = signal === 'SIGINT' ? 'interrupted' : 'terminated'; - const message = `vg: ${word} while parsing. Parse workers were stopped.\n`; - let exited = false; - const finish = (): void => { - if (exited) return; - exited = true; - process.exit(code); - }; - try { - process.stderr.write(message, finish); - } catch { - finish(); - return; - } - // If stderr never drains, still exit so workers cannot outlive the command. - setTimeout(finish, 50); -} - -function armPool(pool: ManagedPool): void { - activePool = pool; - installExitHook(); - if (detachSignals) return; - const onInt = (): void => onStopSignal('SIGINT'); - const onTerm = (): void => onStopSignal('SIGTERM'); - process.on('SIGINT', onInt); - process.on('SIGTERM', onTerm); - detachSignals = () => { - process.removeListener('SIGINT', onInt); - process.removeListener('SIGTERM', onTerm); - }; -} - -function disarmPool(): void { - // Drop the signal listeners before clearing the pool so a signal delivered - // in this window still finds the workers and exits, instead of being swallowed. - if (detachSignals) { - detachSignals(); - detachSignals = null; - } - activePool = null; -} - -function isTinypoolShutdownBug(err: unknown): boolean { - return err instanceof TypeError && err.message.includes('removeListener'); -} - -async function closePool(pool: ManagedPool): Promise { - // tinypool's `destroy()` waits on `events.once(worker, 'exit')`. If the - // worker emits `error` first, that helper throws an uncaught - // `removeListener` TypeError and the process dumps a stack. Swallow only - // that bug; the `run()` rejection is the error the caller sees. - const onUncaught = (err: Error): void => { - if (isTinypoolShutdownBug(err)) { - killWorkers(pool); - return; - } - process.removeListener('uncaughtException', onUncaught); - process.stderr.write(`${err.stack ?? err.message}\n`); - process.exit(1); - }; - process.on('uncaughtException', onUncaught); - try { - try { - pool.cancelPendingTasks(); - } catch { - // Queue already drained. - } - let timer: ReturnType | undefined; - const timedOut = new Promise((resolve) => { - timer = setTimeout(() => { - killWorkers(pool); - resolve(); - }, POOL_SHUTDOWN_MS); - }); - try { - await Promise.race([ - pool.destroy().then( - () => undefined, - () => { - killWorkers(pool); - }, - ), - timedOut, - ]); - } finally { - if (timer) clearTimeout(timer); - } - killWorkers(pool); - } finally { - process.removeListener('uncaughtException', onUncaught); - } -} - -/** - * Flags that make a worker a second copy of the tool (tsx, a preload, an - * inspector) rather than a parse process. Heap and V8 flags are kept so a - * raised `--max-old-space-size` still applies inside the pool. - */ -function workerExecArgv(): string[] { - const dropValue = new Set(['-r', '--require', '--import', '--loader', '--experimental-loader']); - const argv = process.execArgv; - const out: string[] = []; - for (let i = 0; i < argv.length; i++) { - const arg = argv[i] ?? ''; - const inspect = - arg === '--inspect' || - arg === '--inspect-brk' || - arg.startsWith('--inspect=') || - arg.startsWith('--inspect-brk') || - arg.startsWith('--inspect-port'); - if (inspect) continue; - if (dropValue.has(arg)) { - i += 1; - continue; - } - if ( - arg.startsWith('--require=') || - arg.startsWith('--import=') || - arg.startsWith('--loader=') || - arg.startsWith('--experimental-loader=') || - arg.includes('tsx') - ) { - continue; - } - out.push(arg); - } - return out; -} - -function workerEnvironment(): Record { - const env: Record = {}; - for (const key of Object.keys(process.env)) { - const value = process.env[key]; - if (typeof value === 'string') env[key] = value; - } - return env; -} - -function actionableParseError(err: unknown, workerHeapMb: number | undefined): Error { - if (err instanceof ResourceLimitError || err instanceof ParseWorkerFailure) return err; - if (isWorkerOom(err)) { - return new ResourceLimitError( - `graph build stopped: a parse worker exceeded its ${workerHeapMb ?? '?'} MiB heap cap ` + - `(VG_WORKER_HEAP_MB). Raise the cap, exclude the offending files (--exclude), or ` + - `run single-threaded with --jobs 1.`, - ); - } - const raw = err instanceof Error ? err.message : String(err); - const line = raw.split('\n')[0]?.trim() || 'unknown error'; - const brief = line.length > 200 ? `${line.slice(0, 197)}...` : line; - const cap = workerHeapMb ? ` (currently ${workerHeapMb} MiB)` : ''; - return new ParseWorkerFailure( - `graph build stopped: a parse worker failed (${brief}). ` + - 'Re-run with --jobs 1 to parse in this process, exclude the offending files with --exclude, ' + - `or raise the per-worker heap with VG_WORKER_HEAP_MB${cap}.`, - ); } const DEFAULT_INLINE_THRESHOLD = 24; @@ -310,15 +54,12 @@ export async function parseFiles( // workers caps peak memory too (each worker holds its own grammar set). const jobs = Math.max(1, options.jobs ?? envJobs() ?? (Math.min(cores - 1, files.length) || 1)); - const workerFile = options.workerFile ?? resolveWorkerFile(); + const workerFile = resolveWorkerFile(options.workerFile); const useInline = options.inline === true || jobs <= 1 || - workerFile === null || - // An explicit worker module (tests) always uses the pool. Production - // stays inline below the threshold, where spinning workers up costs more - // than it saves. - (options.workerFile === undefined && files.length < threshold); + files.length < threshold || + workerFile === null; if (useInline) { // Inline runs in this process — apply the override directly. @@ -368,57 +109,89 @@ async function parsePooled( filename: workerFile, maxThreads: jobs, minThreads: 1, - runtime: options.workerRuntime ?? 'worker_threads', + // Bound a stuck worker.terminate() so destroy cannot hang forever. terminateTimeout: 1_000, - execArgv: workerExecArgv(), - env: workerEnvironment(), ...(workerHeapMb ? { resourceLimits: { maxOldGenerationSizeMb: workerHeapMb } } : {}), }); - // Without this listener, a worker crash makes `destroy()` throw an uncaught - // TypeError (`emitter.removeListener is not a function`) and the process - // prints a stack instead of the message from `run()`. + // A worker that dies with no in-flight task emits 'error'. Without a + // listener that becomes an uncaught exception and skips pool teardown. pool.on('error', () => undefined); - armPool(pool); - try { - // More, smaller buckets than threads → finer live progress + better load - // balancing. Round-robin keeps shards balanced; the final sort makes the - // output independent of bucket count, so determinism is unaffected. - const total = files.length; - const buckets = chunk( - files.map((f) => ({ rel: f.rel, abs: f.abs, lang: f.lang.id })), - Math.min(total, jobs * 8), - ); - let done = 0; - onProgress?.(0, total); - const results = await Promise.all( - buckets.map((b) => - (pool.run({ tasks: b, grammarsDir }) as Promise).then((r) => { - done += b.length; - onProgress?.(done, total); - // Results accumulate in *this* process; guard its heap as they land. - checkMemoryBudget('parse', memoryBudgetMb); - return r; - }), - ), - ); - return results.flat(); - } catch (err) { - throw actionableParseError(err, workerHeapMb); - } finally { + return parsePoolGuard.using(pool as unknown as ManagedPool, async () => { try { - await closePool(pool); - } finally { - disarmPool(); + // More, smaller buckets than threads → finer live progress + better load + // balancing. Round-robin keeps shards balanced; the final sort makes the + // output independent of bucket count, so determinism is unaffected. + const total = files.length; + const buckets = chunk( + files.map((f) => ({ rel: f.rel, abs: f.abs, lang: f.lang.id })), + Math.min(total, jobs * 8), + ); + let done = 0; + onProgress?.(0, total); + const results = await Promise.all( + buckets.map((b) => + (pool.run({ tasks: b, grammarsDir }) as Promise).then((r) => { + done += b.length; + onProgress?.(done, total); + // Results accumulate in *this* process; guard its heap as they land. + checkMemoryBudget('parse', memoryBudgetMb); + return r; + }), + ), + ); + return results.flat(); + } catch (err) { + // A signal is already shutting the pool down and will exit. Parking + // keeps this rejection from being reported as a second, unrelated error + // or from letting the command finish successfully. + if (parsePoolGuard.isExiting) await new Promise(() => undefined); + throw toParseWorkerError(err, workerHeapMb); } + }); +} + +/** A parse worker failed. The message says what happened and how to re-run. */ +export class ParseWorkerError extends Error { + readonly isParseWorkerError = true; + constructor(message: string) { + super(message); + this.name = 'ParseWorkerError'; } } +/** Turn a worker failure into an error the CLI can print as-is. */ +export function toParseWorkerError(err: unknown, workerHeapMb?: number): Error { + if (err instanceof ResourceLimitError || err instanceof ParseWorkerError) return err; + if (isWorkerOom(err)) { + return new ResourceLimitError( + `graph build stopped: a parse worker exceeded its ${workerHeapMb ?? '?'} MiB heap cap ` + + `(VG_WORKER_HEAP_MB). Raise the cap, exclude the offending files (--exclude), or ` + + `run single-threaded with --jobs 1.`, + ); + } + return new ParseWorkerError( + `graph build stopped: a parse worker failed (${publicDetail(err)}). ` + + `Re-run with --jobs 1 to parse in this process, or skip files with --exclude.`, + ); +} + +/** One line, no absolute paths — caller-facing, not a stack trace. */ +function publicDetail(err: unknown): string { + const raw = err instanceof Error && err.message ? err.message : 'the worker stopped unexpectedly'; + let line = raw.replace(/\s+/g, ' ').trim(); + line = line.replace(/file:\/\/\S+/g, 'a file'); + line = line.replace(/(?:[A-Za-z]:\\(?:[^\\\s]+\\)*[^\\\s]+)|(?:\/(?:[\w.+@~-]+\/)+[\w.+@~-]+)/g, 'a file'); + if (!line) return 'the worker stopped unexpectedly'; + return line.length > 240 ? `${line.slice(0, 239)}…` : line; +} + function isWorkerOom(err: unknown): boolean { const e = err as { code?: string; message?: string } | null; return e?.code === 'ERR_WORKER_OUT_OF_MEMORY' || /out of memory/i.test(e?.message ?? ''); } -function resolveWorkerFile(): string | null { +function resolveWorkerFile(override?: string): string | null { + if (override) return fs.existsSync(override) ? override : null; // Only the compiled .js worker is runnable by a bare worker_thread. Under a // TS-only runner (vitest/tsx) the .js won't exist → fall back to inline. const here = path.dirname(fileURLToPath(import.meta.url)); diff --git a/src/reporting/commands/scan.ts b/src/reporting/commands/scan.ts index f23006a0..555088c5 100644 --- a/src/reporting/commands/scan.ts +++ b/src/reporting/commands/scan.ts @@ -703,57 +703,65 @@ export const scanCommand = new Command('scan') let iacRun = null as SecurityRunResult | null; if (wantGraph) { scanOpts.postScan = async (report, ctx) => { - const result = await buildGraph({ - root: rootDir, - exclude: opts.exclude, - onParseProgress: (done, total) => report(done, total, 'parsing'), - }); - builtGraph = result.graph; - const written = writeArtifacts(result.graph, { root: rootDir }); - if (written.architecturePolicyError) console.error(chalk.red(`\narchitecture policy: ${written.architecturePolicyError}`)); - // Freshness snapshot → lets `vg serve`/`vg ask` auto-refresh this map - // when the working tree drifts (see engine/freshness.ts). - writeSnapshot(rootDir, result.graph.provenance.corpusHash, result.fileStats, { - exclude: opts.exclude, - }); - // Refine before format/write — architecture is already on these objects. - // --no-graph / map failed / --max-privacy never reach here. - // AST roles first so @Entity / @Controller confirm (or override a - // filename suffix) before boundary violations are collected. - refineArchitectureWithAstRoles( - { projects: ctx.projects, solutions: ctx.solutions, extended: ctx.extended }, - result.fileRoles, - ); - refineArchitectureWithGraph( - { projects: ctx.projects, solutions: ctx.solutions, extended: ctx.extended }, - architectureGraphView(result.graph), - ); - // Infrastructure packs run here, after the map, so facts bind to the - // node ids the graph just assigned — and before runCoreScan assembles - // the artifact, so the section is on it for every formatter and the - // upload. `ctx.extended` is the very object the artifact carries. - const { counts } = result.graph.meta; - let detail = `${counts.nodes.toLocaleString()} nodes · ${counts.edges.toLocaleString()} edges`; - if (wantIac) { - try { - // Provision the module on first use the way `vg build` does — - // bounded, consent-respecting, never under --offline (plan §2.8). - iacRun = await runSecurityPacks({ - root: rootDir, - exclude: opts.exclude, - packs: IAC_PACKS, - provision: { offline: Boolean(opts.offline) }, - }); - } catch { - iacRun = { status: 'abstained' }; - } - if (iacRun.status === 'ok') { - ctx.extended.security = iacRun.section; - const n = iacRun.section.findings.length; - detail += ` · ${n} infrastructure finding${n === 1 ? '' : 's'}`; + try { + const result = await buildGraph({ + root: rootDir, + exclude: opts.exclude, + onParseProgress: (done, total) => report(done, total, 'parsing'), + }); + builtGraph = result.graph; + const written = writeArtifacts(result.graph, { root: rootDir }); + if (written.architecturePolicyError) console.error(chalk.red(`\narchitecture policy: ${written.architecturePolicyError}`)); + // Freshness snapshot → lets `vg serve`/`vg ask` auto-refresh this map + // when the working tree drifts (see engine/freshness.ts). + writeSnapshot(rootDir, result.graph.provenance.corpusHash, result.fileStats, { + exclude: opts.exclude, + }); + // Refine before format/write — architecture is already on these objects. + // --no-graph / map failed / --max-privacy never reach here. + // AST roles first so @Entity / @Controller confirm (or override a + // filename suffix) before boundary violations are collected. + refineArchitectureWithAstRoles( + { projects: ctx.projects, solutions: ctx.solutions, extended: ctx.extended }, + result.fileRoles, + ); + refineArchitectureWithGraph( + { projects: ctx.projects, solutions: ctx.solutions, extended: ctx.extended }, + architectureGraphView(result.graph), + ); + // Infrastructure packs run here, after the map, so facts bind to the + // node ids the graph just assigned — and before runCoreScan assembles + // the artifact, so the section is on it for every formatter and the + // upload. `ctx.extended` is the very object the artifact carries. + const { counts } = result.graph.meta; + let detail = `${counts.nodes.toLocaleString()} nodes · ${counts.edges.toLocaleString()} edges`; + if (wantIac) { + try { + // Provision the module on first use the way `vg build` does — + // bounded, consent-respecting, never under --offline (plan §2.8). + iacRun = await runSecurityPacks({ + root: rootDir, + exclude: opts.exclude, + packs: IAC_PACKS, + provision: { offline: Boolean(opts.offline) }, + }); + } catch { + iacRun = { status: 'abstained' }; + } + if (iacRun.status === 'ok') { + ctx.extended.security = iacRun.section; + const n = iacRun.section.findings.length; + detail += ` · ${n} infrastructure finding${n === 1 ? '' : 's'}`; + } } + return detail; + } catch (err) { + // Fail-soft at the step, but say why the map stopped. The pool has + // already torn its workers down; this is the message the user sees. + const message = err instanceof Error ? err.message : String(err); + console.error(chalk.red(`error: ${message}`)); + throw err; } - return detail; }; } diff --git a/test/fixtures/parse-pool-harness.ts b/test/fixtures/parse-pool-harness.ts deleted file mode 100644 index 01d55157..00000000 --- a/test/fixtures/parse-pool-harness.ts +++ /dev/null @@ -1,43 +0,0 @@ -/** - * Scripted check for issue #301. Runs the same `parseFiles` pool `vg build` - * and `vg scan` use, then exits. The parent test asserts the process is gone - * and that no parse worker is left in the process list. - * - * tsx test/fixtures/parse-pool-harness.ts [worker_threads|child_process] - */ -import { fileURLToPath } from 'node:url'; -import type { DiscoveredFile } from '../../src/engine/discover.js'; -import { parseFiles } from '../../src/engine/pool.js'; - -const mode = process.argv[2] || 'crash'; -const runtime = process.argv[3] === 'child_process' ? 'child_process' : 'worker_threads'; -process.env.VG_POOL_WORKER_MODE = mode === 'crash' || mode === 'hang' ? mode : ''; - -const workerFile = fileURLToPath(new URL('./parse-pool-worker.mjs', import.meta.url)); -const lang = { id: 'ts' } as DiscoveredFile['lang']; -const files: DiscoveredFile[] = ['b.ts', 'a.ts'].map((rel) => ({ rel, abs: rel, lang })); - -let announced = false; -try { - await parseFiles(files, { - jobs: 2, - workerFile, - workerRuntime: runtime, - memoryBudgetMb: mode === 'budget' ? 1 : 0, - onProgress: () => { - if (!announced && mode === 'hang') { - announced = true; - process.stdout.write('READY\n'); - } - }, - }); - if (mode === 'hang' || mode === 'crash' || mode === 'budget') { - process.stderr.write(`expected ${mode} to fail\n`); - process.exit(2); - } - process.stdout.write('OK\n'); -} catch (err) { - const message = err instanceof Error ? err.message : String(err); - process.stderr.write(`${message}\n`); - process.exit(1); -} diff --git a/test/fixtures/parse-pool-worker.mjs b/test/fixtures/parse-pool-worker.mjs deleted file mode 100644 index d3c1a135..00000000 --- a/test/fixtures/parse-pool-worker.mjs +++ /dev/null @@ -1,32 +0,0 @@ -/** - * Fixture worker for parse-pool shutdown tests. Not a grammar parser. - * `VG_POOL_WORKER_MODE` selects crash / hang; otherwise it echoes tasks in - * payload order so the parent can assert a stable sort. - */ -export default async function run(payload) { - const mode = process.env.VG_POOL_WORKER_MODE || ''; - if (mode === 'crash') { - // Uncaught, off the task promise. This is the failure that used to make - // pool shutdown throw `emitter.removeListener is not a function`. - setTimeout(() => { - throw new Error('parse worker crashed'); - }, 20); - await new Promise(() => {}); - } - if (mode === 'hang') { - await new Promise(() => {}); - } - const tasks = (payload && payload.tasks) || []; - return tasks.map((task) => ({ - rel: task.rel, - lang: task.lang, - hash: 'h', - bytes: 1, - defs: [], - calls: [], - imports: [], - heritage: [], - typeRefs: [], - guards: [], - })); -} diff --git a/test/parse-pool-shutdown.test.ts b/test/parse-pool-shutdown.test.ts deleted file mode 100644 index 500b7e81..00000000 --- a/test/parse-pool-shutdown.test.ts +++ /dev/null @@ -1,250 +0,0 @@ -import { createRequire } from 'node:module'; -import { spawn, spawnSync, type ChildProcess } from 'node:child_process'; -import * as path from 'node:path'; -import { fileURLToPath } from 'node:url'; -import { describe, expect, it } from 'vitest'; -import type { DiscoveredFile } from '../src/engine/discover.js'; -import { parseFiles, activeParseWorkerCount } from '../src/engine/pool.js'; - -const require = createRequire(import.meta.url); -const tsx = require.resolve('tsx/cli'); -const here = path.dirname(fileURLToPath(import.meta.url)); -const harness = path.join(here, 'fixtures/parse-pool-harness.ts'); -const workerFile = path.join(here, 'fixtures/parse-pool-worker.mjs'); - -interface ProcRow { - pid: number; - ppid: number; - stat: string; - args: string; -} - -function listProcesses(): ProcRow[] { - const out = spawnSync('ps', ['-eo', 'pid,ppid,stat,args'], { encoding: 'utf8' }); - const rows: ProcRow[] = []; - for (const line of (out.stdout || '').split('\n').slice(1)) { - const match = line.trim().match(/^(\d+)\s+(\d+)\s+(\S+)\s+(.*)$/); - if (!match) continue; - rows.push({ pid: Number(match[1]), ppid: Number(match[2]), stat: match[3], args: match[4] }); - } - return rows; -} - -function isWorkerProc(row: ProcRow): boolean { - if (row.stat.includes('Z')) return false; - return row.args.includes('entry/process.js') || row.args.includes('parse-pool-worker.mjs'); -} - -function workerPids(): Set { - return new Set(listProcesses().filter(isWorkerProc).map((row) => row.pid)); -} - -function descendants(pid: number): number[] { - const rows = listProcesses(); - const out: number[] = []; - const queue = [pid]; - const seen = new Set([pid]); - while (queue.length) { - const cur = queue.pop(); - if (cur === undefined) break; - for (const row of rows) { - if (row.ppid === cur && !seen.has(row.pid) && !row.stat.includes('Z')) { - seen.add(row.pid); - out.push(row.pid); - queue.push(row.pid); - } - } - } - return out; -} - -function isAlive(pid: number): boolean { - const row = listProcesses().find((item) => item.pid === pid); - if (!row || row.stat.includes('Z')) return false; - try { - process.kill(pid, 0); - return true; - } catch { - return false; - } -} - -function delay(ms: number): Promise { - return new Promise((resolve) => setTimeout(resolve, ms)); -} - -interface Running { - child: ChildProcess; - output: { stdout: string; stderr: string }; -} - -function startHarness(mode: string, runtime: 'worker_threads' | 'child_process'): Running { - const child = spawn(process.execPath, [tsx, harness, mode, runtime], { - stdio: ['ignore', 'pipe', 'pipe'], - env: { - ...process.env, - NO_COLOR: '1', - FORCE_COLOR: '0', - VIBGRATE_NO_KERNEL: '1', - }, - }); - const output = { stdout: '', stderr: '' }; - child.stdout?.on('data', (chunk: Buffer) => { - output.stdout += chunk.toString('utf8'); - }); - child.stderr?.on('data', (chunk: Buffer) => { - output.stderr += chunk.toString('utf8'); - }); - return { child, output }; -} - -function waitExit( - running: Running, - timeoutMs = 15_000, -): Promise<{ code: number | null; signal: NodeJS.Signals | null; stdout: string; stderr: string }> { - const { child, output } = running; - return new Promise((resolve, reject) => { - const finish = (code: number | null, signal: NodeJS.Signals | null): void => { - resolve({ code, signal, stdout: output.stdout, stderr: output.stderr }); - }; - if (child.exitCode !== null || child.signalCode !== null) { - setImmediate(() => finish(child.exitCode, child.signalCode)); - return; - } - const timer = setTimeout(() => { - child.kill('SIGKILL'); - reject(new Error(`harness timed out\nstdout:\n${output.stdout}\nstderr:\n${output.stderr}`)); - }, timeoutMs); - child.once('error', (err) => { - clearTimeout(timer); - reject(err); - }); - child.once('exit', (code, signal) => { - clearTimeout(timer); - // The signal path writes stderr and then exits; give the pipe a moment. - setTimeout(() => finish(code, signal), 100); - }); - }); -} - -async function waitForStdout(running: Running, needle: string): Promise { - const started = Date.now(); - while (!running.output.stdout.includes(needle)) { - if (Date.now() - started > 10_000) { - throw new Error(`timed out waiting for ${needle}\nstdout:\n${running.output.stdout}\nstderr:\n${running.output.stderr}`); - } - const { child } = running; - if (child.exitCode !== null || child.signalCode !== null) { - throw new Error( - `harness exited before ${needle} (code ${child.exitCode} signal ${child.signalCode})\n${running.output.stderr}`, - ); - } - await delay(30); - } -} - -const files: DiscoveredFile[] = ['b.ts', 'a.ts'].map((rel) => ({ - rel, - abs: rel, - lang: { id: 'ts' } as DiscoveredFile['lang'], -})); - -describe('parse worker shutdown', () => { - it('keeps pooled parse output ordered and stable', async () => { - delete process.env.VG_POOL_WORKER_MODE; - const opts = { jobs: 2, workerFile, memoryBudgetMb: 0 as const }; - const once = await parseFiles(files, opts); - const twice = await parseFiles(files, opts); - expect(once.map((row) => row.rel)).toEqual(['a.ts', 'b.ts']); - expect(JSON.stringify(once)).toBe(JSON.stringify(twice)); - expect(activeParseWorkerCount()).toBe(0); - }); - - it('a memory-budget failure during parse exits the pool', async () => { - delete process.env.VG_POOL_WORKER_MODE; - const run = parseFiles(files, { jobs: 2, workerFile, memoryBudgetMb: 1 }); - const timeout = new Promise((_resolve, reject) => { - setTimeout(() => reject(new Error('parse hung after a budget failure')), 8_000); - }); - await expect(Promise.race([run, timeout])).rejects.toThrow(/VG_MEMORY_BUDGET_MB/); - expect(activeParseWorkerCount()).toBe(0); - }); - - it.each(['worker_threads', 'child_process'] as const)( - 'a finished %s pool exits and leaves no workers', - async (runtime) => { - const before = workerPids(); - const running = startHarness('ok', runtime); - const result = await waitExit(running); - expect(result.code).toBe(0); - expect(result.stdout).toContain('OK'); - expect(result.stderr).not.toContain('removeListener'); - await delay(200); - const leaked = [...workerPids()].filter((pid) => !before.has(pid) && isAlive(pid)); - expect(leaked).toEqual([]); - }, - ); - - it.each(['worker_threads', 'child_process'] as const)( - 'a crashing %s worker exits with an actionable error and leaves no workers', - async (runtime) => { - const before = workerPids(); - const running = startHarness('crash', runtime); - const result = await waitExit(running); - expect(result.code).toBe(1); - expect(result.stderr).toContain('graph build stopped: a parse worker failed'); - expect(result.stderr).toContain('--jobs 1'); - expect(result.stderr).toContain('--exclude'); - expect(result.stderr).toContain('VG_WORKER_HEAP_MB'); - expect(result.stderr).not.toContain('removeListener'); - // worker_threads share the parent's stderr. A shutdown bug prints a - // `removeListener` stack from node:internal; the actionable line must - // be the whole failure. child_process workers print their own crash - // before exiting — that text is the worker, and the parent must still - // exit and reap it. - if (runtime === 'worker_threads') { - expect(result.stderr).not.toContain('node:internal'); - expect(result.stderr).not.toMatch(/\n\s+at /); - } - await delay(300); - const leaked = [...workerPids()].filter((pid) => !before.has(pid) && isAlive(pid)); - expect(leaked).toEqual([]); - expect(isAlive(running.child.pid ?? -1)).toBe(false); - }, - ); - - it.each([ - ['worker_threads', 'SIGTERM', 143, 'terminated while parsing'], - ['child_process', 'SIGTERM', 143, 'terminated while parsing'], - ['child_process', 'SIGINT', 130, 'interrupted while parsing'], - ] as const)( - 'a hung %s pool stops workers on %s', - async (runtime, signal, code, message) => { - const before = workerPids(); - const running = startHarness('hang', runtime); - await waitForStdout(running, 'READY'); - let kids: number[] = []; - if (runtime === 'child_process') { - const started = Date.now(); - while (kids.length === 0 && Date.now() - started < 5_000) { - kids = descendants(running.child.pid ?? -1); - if (kids.length === 0) await delay(50); - } - expect(kids.length).toBeGreaterThan(0); - } - running.child.kill(signal); - const result = await waitExit(running); - expect(result.signal).toBeNull(); - expect(result.code).toBe(code); - expect(result.stderr).toContain(message); - expect(result.stderr).not.toContain('removeListener'); - expect(result.stderr).not.toMatch(/\n\s+at /); - await delay(300); - const still = kids.filter((pid) => isAlive(pid)); - expect(still).toEqual([]); - const leaked = [...workerPids()].filter((pid) => !before.has(pid) && isAlive(pid)); - expect(leaked).toEqual([]); - expect(isAlive(running.child.pid ?? -1)).toBe(false); - }, - ); -});