diff --git a/gitnexus/src/cli/analyze.ts b/gitnexus/src/cli/analyze.ts index 429482a34..d399c5082 100644 --- a/gitnexus/src/cli/analyze.ts +++ b/gitnexus/src/cli/analyze.ts @@ -13,7 +13,8 @@ import os from 'os'; import { spawn } from 'child_process'; import v8 from 'v8'; import cliProgress from 'cli-progress'; -import { closeLbugBeforeExit, isLbugReady } from '../core/lbug/lbug-adapter.js'; +import { isLbugReady } from '../core/lbug/lbug-adapter.js'; +import { boundedCheckpointBeforeExit } from '../core/lbug/shutdown-helpers.js'; import { isLbugCheckpointIoError, isWalCorruptionError, @@ -1179,22 +1180,17 @@ const analyzeCommandImpl = async ( aborted = true; bar.stop(); console.log('\n Interrupted — cleaning up...'); - // process.exit(130) follows, so skip the native close (LadybugDB destructor - // can double-free after --pdg writes, #2264); closeLbugBeforeExit's CHECKPOINT - // still flushes the WAL. But that CHECKPOINT queues behind the connection lock - // held by an in-flight COPY, so a single Ctrl-C during a long --pdg COPY would - // otherwise appear hung until the COPY releases. Bound it with a short timeout - // so the interrupt stays responsive; the WAL replays on the next analyze - // (#2264 review P3). A second Ctrl-C (`if (aborted) process.exit(1)` above) - // remains the immediate escape hatch. - const SIGINT_CLEANUP_TIMEOUT_MS = 2000; - void Promise.race([ - closeLbugBeforeExit().catch(() => {}), - new Promise((resolve) => setTimeout(resolve, SIGINT_CLEANUP_TIMEOUT_MS)), - ]).finally(async () => { - const { flushLoggerSync } = await import('../core/logger.js'); - flushLoggerSync(); - process.exit(130); + // Bounded CHECKPOINT-then-exit (#2264 review P3): skip the native close (the + // LadybugDB destructor can double-free after --pdg writes), but don't hang + // behind a long --pdg COPY holding the connection lock — bound it so a single + // Ctrl-C stays responsive; the WAL replays on the next analyze. A second + // Ctrl-C (`if (aborted) process.exit(1)` above) remains the escape hatch. + void boundedCheckpointBeforeExit({ + exitCode: 130, + beforeExit: async () => { + const { flushLoggerSync } = await import('../core/logger.js'); + flushLoggerSync(); + }, }); }; process.on('SIGINT', sigintHandler); diff --git a/gitnexus/src/core/lbug/shutdown-helpers.ts b/gitnexus/src/core/lbug/shutdown-helpers.ts new file mode 100644 index 000000000..542709d72 --- /dev/null +++ b/gitnexus/src/core/lbug/shutdown-helpers.ts @@ -0,0 +1,53 @@ +/** + * Shared bounded "checkpoint, then exit" cleanup for interrupt/cancel signals + * (#2264). The CLI SIGINT handler and the forked worker's SIGTERM handler both + * need to: CHECKPOINT the WAL for durability (skipping the native close — see + * closeLbugBeforeExit), but NOT hang behind an in-flight COPY that holds the + * connection lock, and then exit. Bounding the CHECKPOINT with a short timeout + * keeps a single Ctrl-C / cancel responsive; the WAL replays on the next analyze. + */ +import { closeLbugBeforeExit } from './lbug-adapter.js'; + +/** Default cap so a CHECKPOINT queued behind a long COPY can't wedge the signal. */ +export const DEFAULT_EXIT_CLEANUP_TIMEOUT_MS = 2000; + +export interface BoundedCheckpointExitOptions { + /** Exit code to terminate with (130 for SIGINT, 0 for a worker SIGTERM). */ + exitCode: number; + /** Cap on the CHECKPOINT; defaults to {@link DEFAULT_EXIT_CLEANUP_TIMEOUT_MS}. */ + timeoutMs?: number; + /** Report a CHECKPOINT failure (e.g. over IPC) rather than swallowing it. */ + onFlushError?: (err: unknown) => void; + /** Run just before exit (e.g. flush the logger synchronously). */ + beforeExit?: () => void | Promise; + /** @internal test seam — defaults to {@link closeLbugBeforeExit}. */ + checkpoint?: () => Promise; + /** @internal test seam — defaults to `process.exit`. */ + exit?: (code: number) => void; +} + +/** + * Best-effort CHECKPOINT bounded by a timeout, then exit. Never rejects — the + * exit always fires (in `finally`) even if the CHECKPOINT throws. Fire-and-forget + * from a signal handler (`void boundedCheckpointBeforeExit({...})`). + */ +export async function boundedCheckpointBeforeExit( + opts: BoundedCheckpointExitOptions, +): Promise { + const timeoutMs = opts.timeoutMs ?? DEFAULT_EXIT_CLEANUP_TIMEOUT_MS; + const checkpoint = opts.checkpoint ?? closeLbugBeforeExit; + const exit = opts.exit ?? ((code: number) => process.exit(code)); + let timer: ReturnType | undefined; + try { + await Promise.race([ + checkpoint().catch((err: unknown) => opts.onFlushError?.(err)), + new Promise((resolve) => { + timer = setTimeout(resolve, timeoutMs); + }), + ]); + } finally { + if (timer !== undefined) clearTimeout(timer); + await opts.beforeExit?.(); + exit(opts.exitCode); + } +} diff --git a/gitnexus/src/server/analyze-worker.ts b/gitnexus/src/server/analyze-worker.ts index 33faa986d..f8d583e71 100644 --- a/gitnexus/src/server/analyze-worker.ts +++ b/gitnexus/src/server/analyze-worker.ts @@ -15,7 +15,7 @@ import { runFullAnalysis, type AnalyzeOptions } from '../core/run-analyze.js'; import { type AnalyzeResultIpc } from './analyze-worker-ipc.js'; import { runWorkerAnalysis, createTerminalClaim } from './analyze-worker-core.js'; import { assertAnalysisFinalized } from '../storage/repo-manager.js'; -import { closeLbugBeforeExit } from '../core/lbug/lbug-adapter.js'; +import { boundedCheckpointBeforeExit } from '../core/lbug/shutdown-helpers.js'; interface StartMessage { type: 'start'; @@ -82,15 +82,11 @@ process.on('unhandledRejection', (reason: unknown) => { }); // Handle cancellation / timeout shutdown (analyze-job.ts `cancelJob` sends -// SIGTERM). Mirror the CLI SIGINT path (#2264): a best-effort CHECKPOINT that -// SKIPS the native close. A real conn.close()/db.close() here can double-free in -// LadybugDB's ClientContext destructor after --pdg writes, AND it would block -// behind the in-flight COPY's connection lock — so a single cancel could abort -// the worker or hang until the COPY releases. Bound the CHECKPOINT with a short -// timeout; process.exit reclaims the native handles regardless. A CHECKPOINT -// failure is reported to the parent over IPC, not swallowed; the exit lives in -// `finally` so it always fires. -const SIGTERM_CLEANUP_TIMEOUT_MS = 2000; +// SIGTERM). Bounded CHECKPOINT-then-exit shared with the CLI SIGINT path (#2264): +// skip the native close (the LadybugDB destructor can double-free after --pdg +// writes), but don't block behind the in-flight COPY's connection lock — so a +// single cancel can't abort or hang the worker. A CHECKPOINT failure is reported +// to the parent over IPC, not swallowed; the exit always fires. process.on('SIGTERM', () => { // Only report the cancellation if the analysis hasn't already reported a // terminal outcome (#2264 P3) — otherwise this would flip an already-complete @@ -98,14 +94,14 @@ process.on('SIGTERM', () => { if (claimTerminal()) { send({ type: 'error', message: 'Analysis cancelled (worker received SIGTERM)' }); } - void Promise.race([ - closeLbugBeforeExit().catch((err: unknown) => { + void boundedCheckpointBeforeExit({ + exitCode: 0, + onFlushError: (err: unknown) => { const message = err instanceof Error ? err.message : 'Worker checkpoint failed during SIGTERM'; send({ type: 'error', message }); - }), - new Promise((resolve) => setTimeout(resolve, SIGTERM_CLEANUP_TIMEOUT_MS)), - ]).finally(() => process.exit(0)); + }, + }); }); // Listen for start command from parent — guarded against re-entry diff --git a/gitnexus/test/unit/shutdown-helpers.test.ts b/gitnexus/test/unit/shutdown-helpers.test.ts new file mode 100644 index 000000000..e4606383c --- /dev/null +++ b/gitnexus/test/unit/shutdown-helpers.test.ts @@ -0,0 +1,62 @@ +/** + * Unit tests for the shared bounded checkpoint-then-exit cleanup (#2264) used by + * the CLI SIGINT handler and the worker SIGTERM handler. The checkpoint + exit are + * injected so the helper is testable without touching the real LadybugDB close or + * the real process.exit. + */ +import { describe, it, expect, vi } from 'vitest'; +import { boundedCheckpointBeforeExit } from '../../src/core/lbug/shutdown-helpers.js'; + +describe('boundedCheckpointBeforeExit (#2264)', () => { + it('checkpoints, runs beforeExit, then exits with the given code (in order)', async () => { + const order: string[] = []; + const exit = vi.fn<(code: number) => void>((c) => { + order.push(`exit:${c}`); + }); + + await boundedCheckpointBeforeExit({ + exitCode: 130, + checkpoint: vi.fn(async () => { + order.push('checkpoint'); + }), + beforeExit: () => { + order.push('beforeExit'); + }, + exit, + }); + + expect(exit).toHaveBeenCalledWith(130); + expect(order).toEqual(['checkpoint', 'beforeExit', 'exit:130']); + }); + + it('reports a checkpoint failure via onFlushError and still exits', async () => { + const onFlushError = vi.fn<(err: unknown) => void>(); + const exit = vi.fn<(code: number) => void>(); + const err = new Error('checkpoint boom'); + + await boundedCheckpointBeforeExit({ + exitCode: 0, + checkpoint: vi.fn(async () => { + throw err; + }), + onFlushError, + exit, + }); + + expect(onFlushError).toHaveBeenCalledWith(err); + expect(exit).toHaveBeenCalledWith(0); + }); + + it('exits via the timeout when the checkpoint hangs', async () => { + const exit = vi.fn<(code: number) => void>(); + + await boundedCheckpointBeforeExit({ + exitCode: 0, + timeoutMs: 0, + checkpoint: () => new Promise(() => {}), // never resolves + exit, + }); + + expect(exit).toHaveBeenCalledWith(0); + }); +});