From 7e1e5828a2292ee85bbe82ec37198bb4be2bb8ff Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Sun, 21 Jun 2026 13:10:42 +0000 Subject: [PATCH] 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) Claude-Session: https://claude.ai/code/session_01JBJomjoTdBV2eveDVq4JMm --- gitnexus/src/server/analyze-worker-core.ts | 40 +++++- gitnexus/src/server/analyze-worker.ts | 16 ++- .../test/unit/analyze-worker-core.test.ts | 118 ++++++++++++++---- 3 files changed, 142 insertions(+), 32 deletions(-) diff --git a/gitnexus/src/server/analyze-worker-core.ts b/gitnexus/src/server/analyze-worker-core.ts index 09b8b0c11..acb31e6d1 100644 --- a/gitnexus/src/server/analyze-worker-core.ts +++ b/gitnexus/src/server/analyze-worker-core.ts @@ -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 { + 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; + }; } diff --git a/gitnexus/src/server/analyze-worker.ts b/gitnexus/src/server/analyze-worker.ts index da44b3409..90575f22e 100644 --- a/gitnexus/src/server/analyze-worker.ts +++ b/gitnexus/src/server/analyze-worker.ts @@ -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 diff --git a/gitnexus/test/unit/analyze-worker-core.test.ts b/gitnexus/test/unit/analyze-worker-core.test.ts index d2961f650..a5a3f20d9 100644 --- a/gitnexus/test/unit/analyze-worker-core.test.ts +++ b/gitnexus/test/unit/analyze-worker-core.test.ts @@ -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(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); }); });