diff --git a/gitnexus/src/core/ingestion/workers/worker-pool.ts b/gitnexus/src/core/ingestion/workers/worker-pool.ts index 0e5c3e75b..2ef895f2e 100644 --- a/gitnexus/src/core/ingestion/workers/worker-pool.ts +++ b/gitnexus/src/core/ingestion/workers/worker-pool.ts @@ -142,6 +142,26 @@ function itemPath(item: unknown): string | undefined { return typeof path === 'string' ? path : undefined; } +/** + * Best-guess path of the file in flight when a worker dies mid-job. + * + * `lastProgress` is the number of files the worker has acknowledged via + * `progress` messages, so `items[lastProgress]` is the next file it was + * about to process — the most likely culprit when the worker crashes + * (OOM, native addon SIGSEGV) or reports an error. + * + * Excluding only this single path keeps the blast radius small: earlier + * files in the job get re-tried by the sequential fallback, and any + * pathological file gets the same skip treatment as the singleton- + * timeout path. Returns `[]` when no path is determinable so sequential + * retries the whole job. + */ +function inFlightExcludePath(job: WorkerJob, lastProgress: number): string[] { + if (lastProgress >= job.items.length) return []; + const path = itemPath(job.items[lastProgress]); + return path ? [path] : []; +} + function createJobs( items: TInput[], maxItems: number, @@ -435,7 +455,12 @@ export const createWorkerPool = ( } else if (msg.type === 'error') { settled = true; cleanup(); - void fail(new Error(`Worker ${workerIndex} error: ${msg.error}`)); + void fail( + new WorkerPoolDispatchError( + `Worker ${workerIndex} error: ${msg.error}`, + inFlightExcludePath(job, lastProgress), + ), + ); } else if (msg.type === 'result') { if (!waitingForFlush) { settled = true; @@ -456,7 +481,12 @@ export const createWorkerPool = ( if (!settled) { settled = true; cleanup(); - void fail(err); + void fail( + new WorkerPoolDispatchError( + `Worker ${workerIndex} error: ${err.message}`, + inFlightExcludePath(job, lastProgress), + ), + ); } }; @@ -464,9 +494,13 @@ export const createWorkerPool = ( if (!settled) { settled = true; cleanup(); + const excludes = inFlightExcludePath(job, lastProgress); + const inFlightSuffix = excludes.length > 0 ? ` (in-flight: ${excludes[0]})` : ''; void fail( - new Error( - `Worker ${workerIndex} exited with code ${code}. Likely OOM or native addon failure.`, + new WorkerPoolDispatchError( + `Worker ${workerIndex} exited with code ${code}. ` + + `Likely OOM or native addon failure${inFlightSuffix}.`, + excludes, ), ); } diff --git a/gitnexus/test/unit/parsing-worker-fallback.test.ts b/gitnexus/test/unit/parsing-worker-fallback.test.ts index ab434eb75..8525e263b 100644 --- a/gitnexus/test/unit/parsing-worker-fallback.test.ts +++ b/gitnexus/test/unit/parsing-worker-fallback.test.ts @@ -78,4 +78,120 @@ describe('processParsing worker fallback', () => { graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'stuck'), ).toBe(false); }); + + it('skips worker-error in-flight file during sequential fallback', async () => { + const graph = createKnowledgeGraph(); + const progressDetails: string[] = []; + const workerPool: WorkerPool = { + size: 1, + dispatch: vi.fn(async () => { + throw new WorkerPoolDispatchError('Worker 0 error: native crash', ['src/crashed.ts']); + }), + terminate: vi.fn(async () => undefined), + }; + + const result = await processParsing( + graph, + [ + { path: 'src/crashed.ts', content: 'export function crashed() { return 0; }\n' }, + { path: 'src/a.ts', content: 'export function a() { return 1; }\n' }, + ], + createSymbolTable(), + createASTCache(), + createASTCache(), + (_current, _total, detail) => { + progressDetails.push(detail); + }, + workerPool, + ); + + expect(result).toBeNull(); + expect(progressDetails).toContain('Skipping 1 worker-timeout file(s) in sequential fallback'); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'a'), + ).toBe(true); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'crashed'), + ).toBe(false); + }); + + it('skips worker-exit in-flight file during sequential fallback', async () => { + const graph = createKnowledgeGraph(); + const progressDetails: string[] = []; + const workerPool: WorkerPool = { + size: 1, + dispatch: vi.fn(async () => { + throw new WorkerPoolDispatchError( + 'Worker 0 exited with code 134. Likely OOM or native addon failure (in-flight: src/oom.ts).', + ['src/oom.ts'], + ); + }), + terminate: vi.fn(async () => undefined), + }; + + const result = await processParsing( + graph, + [ + { path: 'src/oom.ts', content: 'export function oom() { return 0; }\n' }, + { path: 'src/a.ts', content: 'export function a() { return 1; }\n' }, + ], + createSymbolTable(), + createASTCache(), + createASTCache(), + (_current, _total, detail) => { + progressDetails.push(detail); + }, + workerPool, + ); + + expect(result).toBeNull(); + expect(progressDetails).toContain('Skipping 1 worker-timeout file(s) in sequential fallback'); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'a'), + ).toBe(true); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'oom'), + ).toBe(false); + }); + + it('runs full sequential fallback when the worker pool throws a non-WorkerPoolDispatchError', async () => { + const graph = createKnowledgeGraph(); + const progressDetails: string[] = []; + const workerPool: WorkerPool = { + size: 1, + dispatch: vi.fn(async () => { + throw new Error('replacement worker failed'); + }), + terminate: vi.fn(async () => undefined), + }; + + const result = await processParsing( + graph, + [ + { path: 'src/keep.ts', content: 'export function keep() { return 0; }\n' }, + { path: 'src/a.ts', content: 'export function a() { return 1; }\n' }, + ], + createSymbolTable(), + createASTCache(), + createASTCache(), + (_current, _total, detail) => { + progressDetails.push(detail); + }, + workerPool, + ); + + expect(result).toBeNull(); + expect(progressDetails).toContain( + 'Sequential fallback after worker issue: replacement worker failed', + ); + expect( + progressDetails.some((d) => d.startsWith('Skipping ') && d.includes('worker-timeout file')), + ).toBe(false); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'a'), + ).toBe(true); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'keep'), + ).toBe(true); + }); });