fix(server): coordinate worker SIGTERM cancellation with completion (#2264 P3)

The worker SIGTERM handler unconditionally sent {type:'error','Analysis cancelled'}
and didn't coordinate with the message handler that sends complete, so a cancel near
the finish line could report a cancelled job complete, or a late SIGTERM could flip
an already-complete job to failed.

Add a single terminal-outcome claim (createTerminalClaim) shared by the message
handler and the SIGTERM handler: whoever claims it first reports its terminal
message; the other skips its terminal send. Single-threaded JS makes the
check-and-set atomic. The cleanup + process.exit still run regardless.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JBJomjoTdBV2eveDVq4JMm
This commit is contained in:
Gergo Magyar 2026-06-21 13:10:42 +00:00
parent 1565e14f69
commit 7e1e5828a2
3 changed files with 142 additions and 32 deletions

View file

@ -20,18 +20,27 @@ export interface WorkerAnalysisDeps {
runFullAnalysis: typeof import('../core/run-analyze.js').runFullAnalysis;
assertAnalysisFinalized: typeof import('../storage/repo-manager.js').assertAnalysisFinalized;
send: (msg: WorkerMessage) => void;
/**
* Claim the single terminal-outcome slot. Returns `true` for the first caller
* (which may then send its `complete`/`error`) and `false` for every caller
* after — so a SIGTERM cancellation and a near-simultaneous completion can't
* both report a terminal outcome (#2264 P3). See {@link createTerminalClaim}.
*/
claimTerminal: () => boolean;
}
/**
* Run the analysis and report the outcome to the parent over IPC. Always reports
* exactly one terminal message (`complete` or `error`) and never throws — the
* caller schedules `process.exit` after this resolves.
* Run the analysis and report the outcome to the parent over IPC. Reports at most
* one terminal message (`complete` or `error`) — and none if a cancellation
* already claimed the terminal slot — and never throws; the caller schedules
* `process.exit` after this resolves.
*/
export async function runWorkerAnalysis(
repoPath: string,
options: AnalyzeOptions,
deps: WorkerAnalysisDeps,
): Promise<void> {
let terminal: WorkerMessage;
try {
const result = await deps.runFullAnalysis(
repoPath,
@ -55,10 +64,31 @@ export async function runWorkerAnalysis(
// Send a JSON-safe projection, NOT the raw result: the IPC channel is
// default-JSON serialization and `result.pipelineResult` carries the live
// KnowledgeGraph. See analyze-worker-ipc.ts.
deps.send({ type: 'complete', result: projectAnalyzeResultForIpc(result) });
terminal = { type: 'complete', result: projectAnalyzeResultForIpc(result) };
} catch (err: unknown) {
// Report the failure to the parent over IPC (the parent surfaces the message).
const message = err instanceof Error ? err.message : 'Analysis failed';
deps.send({ type: 'error', message });
terminal = { type: 'error', message };
}
// P3 (#2264): only report if a SIGTERM cancellation hasn't already claimed the
// terminal slot — otherwise a cancel near the finish line would report the
// analysis as `complete` over the top of the cancellation.
if (deps.claimTerminal()) deps.send(terminal);
}
/**
* Create the single-use terminal-outcome claim shared by the worker's message
* handler and its SIGTERM handler. The first call returns `true`; every later
* call returns `false`. This is the coordination point that prevents a cancel and
* a completion from both reporting a terminal status (#2264 P3). Single-threaded
* JS guarantees the check-and-set is atomic (no preemption mid-call).
*/
export function createTerminalClaim(): () => boolean {
let claimed = false;
return () => {
if (claimed) return false;
claimed = true;
return true;
};
}

View file

@ -13,7 +13,7 @@
import { runFullAnalysis, type AnalyzeOptions } from '../core/run-analyze.js';
import { type AnalyzeResultIpc } from './analyze-worker-ipc.js';
import { runWorkerAnalysis } from './analyze-worker-core.js';
import { runWorkerAnalysis, createTerminalClaim } from './analyze-worker-core.js';
import { assertAnalysisFinalized } from '../storage/repo-manager.js';
import { closeLbug } from '../core/lbug/lbug-adapter.js';
@ -53,6 +53,12 @@ function send(msg: WorkerMessage) {
process.send?.(msg);
}
// Single terminal-outcome slot shared by the message handler and the SIGTERM
// handler: whoever claims it first reports its complete/error; the other skips its
// terminal send, so a cancel near the finish line can't also report success and a
// late SIGTERM can't flip an already-reported job (#2264 P3).
const claimTerminal = createTerminalClaim();
// Catch uncaught exceptions and unhandled rejections — report them to the parent
// over IPC (the same channel the analysis path uses), then exit. The report runs
// in `try` and the exit in `finally` so a throw from send() on a closed channel
@ -86,7 +92,12 @@ process.on('unhandledRejection', (reason: unknown) => {
// `finally` so it always fires.
const SIGTERM_CLEANUP_TIMEOUT_MS = 2000;
process.on('SIGTERM', () => {
send({ type: 'error', message: 'Analysis cancelled (worker received SIGTERM)' });
// Only report the cancellation if the analysis hasn't already reported a
// terminal outcome (#2264 P3) — otherwise this would flip an already-complete
// job to failed. The cleanup + exit below run regardless.
if (claimTerminal()) {
send({ type: 'error', message: 'Analysis cancelled (worker received SIGTERM)' });
}
void Promise.race([
closeLbug({ skipNativeClose: true }).catch((err: unknown) => {
const message =
@ -112,6 +123,7 @@ process.on('message', async (msg: StartMessage) => {
runFullAnalysis,
assertAnalysisFinalized,
send,
claimTerminal,
});
} finally {
// LadybugDB's native module prevents clean exit — force it (same reason the

View file

@ -1,30 +1,38 @@
/**
* Unit tests for the analyze-worker core seam (#2264 P2). The worker must NOT
* report `complete` for a half-finalized repo (meta.json written but the global
* registry entry missing) — it must surface that as an error, mirroring the CLI's
* assertAnalysisFinalized guard. Driven via the side-effect-free
* `runWorkerAnalysis` seam with injected fakes, so no fork()/process.on side
* effects of the entry module are touched.
* Unit tests for the analyze-worker core seam (#2264).
*
* P2: the worker must NOT report `complete` for a half-finalized repo (meta.json
* written but the global registry entry missing) — it must surface that as an
* error, mirroring the CLI's assertAnalysisFinalized guard.
*
* P3: a SIGTERM cancellation and a near-simultaneous completion must not both
* report a terminal outcome — the `claimTerminal` slot coordinates them.
*
* Driven via the side-effect-free `runWorkerAnalysis` seam with injected fakes, so
* no fork()/process.on side effects of the entry module are touched.
*/
import { describe, it, expect, vi } from 'vitest';
import {
runWorkerAnalysis,
createTerminalClaim,
type WorkerAnalysisDeps,
} from '../../src/server/analyze-worker-core.js';
import type { AnalyzeResult } from '../../src/core/run-analyze.js';
import type { WorkerMessage } from '../../src/server/analyze-worker.js';
describe('runWorkerAnalysis — worker finalize guard (#2264 P2)', () => {
const baseResult: AnalyzeResult = {
repoName: 'repo',
repoPath: '/repo',
stats: {},
alreadyUpToDate: false,
ftsRepairedOnly: false,
};
const baseResult: AnalyzeResult = {
repoName: 'repo',
repoPath: '/repo',
stats: {},
alreadyUpToDate: false,
ftsRepairedOnly: false,
};
const okRun: WorkerAnalysisDeps['runFullAnalysis'] = vi.fn(async () => baseResult);
const okRun: WorkerAnalysisDeps['runFullAnalysis'] = vi.fn(async () => baseResult);
const okFinalize: WorkerAnalysisDeps['assertAnalysisFinalized'] = vi.fn(async () => undefined);
const alwaysClaim: WorkerAnalysisDeps['claimTerminal'] = () => true;
describe('runWorkerAnalysis — finalize guard (#2264 P2)', () => {
it('reports error (not complete) when finalization fails for an unregistered repo', async () => {
const send = vi.fn<(msg: WorkerMessage) => void>();
const assertAnalysisFinalized: WorkerAnalysisDeps['assertAnalysisFinalized'] = vi.fn(
@ -33,7 +41,16 @@ describe('runWorkerAnalysis — worker finalize guard (#2264 P2)', () => {
},
);
await runWorkerAnalysis('/repo', {}, { runFullAnalysis: okRun, assertAnalysisFinalized, send });
await runWorkerAnalysis(
'/repo',
{},
{
runFullAnalysis: okRun,
assertAnalysisFinalized,
send,
claimTerminal: alwaysClaim,
},
);
expect(send).toHaveBeenCalledWith({
type: 'error',
@ -44,11 +61,17 @@ describe('runWorkerAnalysis — worker finalize guard (#2264 P2)', () => {
it('reports complete exactly once when finalization succeeds', async () => {
const send = vi.fn<(msg: WorkerMessage) => void>();
const assertAnalysisFinalized: WorkerAnalysisDeps['assertAnalysisFinalized'] = vi.fn(
async () => undefined,
);
await runWorkerAnalysis('/repo', {}, { runFullAnalysis: okRun, assertAnalysisFinalized, send });
await runWorkerAnalysis(
'/repo',
{},
{
runFullAnalysis: okRun,
assertAnalysisFinalized: okFinalize,
send,
claimTerminal: alwaysClaim,
},
);
const completes = send.mock.calls.filter((c) => c[0].type === 'complete');
expect(completes).toHaveLength(1);
@ -59,17 +82,62 @@ describe('runWorkerAnalysis — worker finalize guard (#2264 P2)', () => {
const failingRun: WorkerAnalysisDeps['runFullAnalysis'] = vi.fn(async () => {
throw new Error('boom');
});
const assertAnalysisFinalized: WorkerAnalysisDeps['assertAnalysisFinalized'] = vi.fn(
async () => undefined,
);
// Fresh local mock (not the shared okFinalize) so the "never called" assertion
// reflects only this test.
const finalize = vi.fn<WorkerAnalysisDeps['assertAnalysisFinalized']>(async () => undefined);
await runWorkerAnalysis(
'/repo',
{},
{ runFullAnalysis: failingRun, assertAnalysisFinalized, send },
{
runFullAnalysis: failingRun,
assertAnalysisFinalized: finalize,
send,
claimTerminal: alwaysClaim,
},
);
expect(send).toHaveBeenCalledWith({ type: 'error', message: 'boom' });
expect(assertAnalysisFinalized).not.toHaveBeenCalled();
expect(finalize).not.toHaveBeenCalled();
});
});
describe('runWorkerAnalysis — terminal-claim coordination (#2264 P3)', () => {
it('sends NO terminal message when the slot is already claimed (cancellation won)', async () => {
const send = vi.fn<(msg: WorkerMessage) => void>();
const alreadyClaimed: WorkerAnalysisDeps['claimTerminal'] = () => false;
await runWorkerAnalysis(
'/repo',
{},
{
runFullAnalysis: okRun,
assertAnalysisFinalized: okFinalize,
send,
claimTerminal: alreadyClaimed,
},
);
const terminals = send.mock.calls.filter(
(c) => c[0].type === 'complete' || c[0].type === 'error',
);
expect(terminals).toHaveLength(0);
});
});
describe('createTerminalClaim (#2264 P3)', () => {
it('returns true for the first claim and false for every claim after', () => {
const claim = createTerminalClaim();
expect(claim()).toBe(true);
expect(claim()).toBe(false);
expect(claim()).toBe(false);
});
it('gives independent claims separate slots', () => {
const a = createTerminalClaim();
const b = createTerminalClaim();
expect(a()).toBe(true);
expect(b()).toBe(true);
expect(a()).toBe(false);
});
});