fix(workers): exclude in-flight file on worker error/exit, not just singleton timeout

WorkerPoolDispatchError previously surfaced the stalled path only for the
singleton-timeout final-fail branch. Worker `error` and `exit` events (and
the msg-channel `error` reply) fell back to plain `Error`, so the sequential
fallback re-attempted every file in the active job — re-hanging on the same
pathological file when the worker crashed mid-parse.

Lift the in-flight-file inference into `inFlightExcludePath(job, lastProgress)`
and wire it into the three remaining in-pool failure sites. `lastProgress` is
already in `runWorker` scope, so `items[lastProgress]` (the next file the
worker was about to acknowledge) is the best single guess at the culprit;
earlier files are still re-tried sequentially. Returns `[]` when no path is
determinable (`lastProgress >= items.length`, or path missing/non-string) so
sequential retries the whole job.

Replacement-worker startup failures stay plain `Error` (no job context); the
result-before-flush protocol bug stays plain `Error` (code fault, not file).

Tests cover the three new exclusion paths plus a negative test confirming
non-WorkerPoolDispatchError throws fall through to full sequential retry.
This commit is contained in:
Gergo Magyar 2026-05-19 08:08:53 +01:00
parent 614f3803d3
commit 724bd2c415
2 changed files with 154 additions and 4 deletions

View file

@ -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<TInput>(job: WorkerJob<TInput>, lastProgress: number): string[] {
if (lastProgress >= job.items.length) return [];
const path = itemPath(job.items[lastProgress]);
return path ? [path] : [];
}
function createJobs<TInput>(
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,
),
);
}

View file

@ -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);
});
});