diff --git a/src/cli.ts b/src/cli.ts index 23a45c5..0241338 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 981cfa9..e321702 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 { 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'; @@ -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 ParseWorkerError) { + throw new CliError(err.message, ExitCode.ERROR); + } if (err instanceof UnsafeRootError) throw new CliError(err.message, ExitCode.ERROR); throw err; } diff --git a/src/engine/pool-guard.ts b/src/engine/pool-guard.ts new file mode 100644 index 0000000..a3a8baa --- /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 0000000..de61b87 --- /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 0000000..c276d90 --- /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 9837770..6759563 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. @@ -37,6 +38,8 @@ export interface ParseOptions { grammarsDir?: string; /** Heap budget (MiB) checked as parse results accumulate; 0/unset skips. */ memoryBudgetMb?: number; + /** Worker module path. Set by tests; production resolves the bundled worker. */ + workerFile?: string; } const DEFAULT_INLINE_THRESHOLD = 24; @@ -51,7 +54,7 @@ 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 = resolveWorkerFile(options.workerFile); const useInline = options.inline === true || jobs <= 1 || @@ -106,51 +109,89 @@ async function parsePooled( filename: workerFile, maxThreads: jobs, minThreads: 1, + // Bound a stuck worker.terminate() so destroy cannot hang forever. + terminateTimeout: 1_000, ...(workerHeapMb ? { resourceLimits: { maxOldGenerationSizeMb: workerHeapMb } } : {}), }); - 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) { - 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.`, + // 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); + return parsePoolGuard.using(pool as unknown as ManagedPool, async () => { + 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) { + // 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); } - throw err; - } finally { - await pool.destroy(); + }); +} + +/** 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 f23006a..555088c 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; }; }