mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-11 03:38:07 +00:00
refactor(lbug): share the bounded checkpoint-then-exit cleanup (SIGINT/SIGTERM) (#2264)
The CLI SIGINT handler (analyze.ts) and the worker SIGTERM handler (analyze-worker.ts) had near-identical Promise.race([closeLbugBeforeExit, timeout]).finally(exit) blocks with separately-hardcoded 2s timeouts. Extract boundedCheckpointBeforeExit into a shared shutdown-helpers module — parameterized by exit code, an optional flush-error reporter (worker reports over IPC), and an optional beforeExit hook (CLI flushes the logger). checkpoint + exit are injectable test seams, so it's unit-tested without the real LadybugDB close or process.exit. 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:
parent
2e5353837a
commit
aa7003da19
4 changed files with 139 additions and 32 deletions
|
|
@ -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<void>((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);
|
||||
|
|
|
|||
53
gitnexus/src/core/lbug/shutdown-helpers.ts
Normal file
53
gitnexus/src/core/lbug/shutdown-helpers.ts
Normal file
|
|
@ -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<void>;
|
||||
/** @internal test seam — defaults to {@link closeLbugBeforeExit}. */
|
||||
checkpoint?: () => Promise<void>;
|
||||
/** @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<void> {
|
||||
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<typeof setTimeout> | undefined;
|
||||
try {
|
||||
await Promise.race([
|
||||
checkpoint().catch((err: unknown) => opts.onFlushError?.(err)),
|
||||
new Promise<void>((resolve) => {
|
||||
timer = setTimeout(resolve, timeoutMs);
|
||||
}),
|
||||
]);
|
||||
} finally {
|
||||
if (timer !== undefined) clearTimeout(timer);
|
||||
await opts.beforeExit?.();
|
||||
exit(opts.exitCode);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<void>((resolve) => setTimeout(resolve, SIGTERM_CLEANUP_TIMEOUT_MS)),
|
||||
]).finally(() => process.exit(0));
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
// Listen for start command from parent — guarded against re-entry
|
||||
|
|
|
|||
62
gitnexus/test/unit/shutdown-helpers.test.ts
Normal file
62
gitnexus/test/unit/shutdown-helpers.test.ts
Normal file
|
|
@ -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<void>(() => {}), // never resolves
|
||||
exit,
|
||||
});
|
||||
|
||||
expect(exit).toHaveBeenCalledWith(0);
|
||||
});
|
||||
});
|
||||
Loading…
Add table
Reference in a new issue