diff --git a/gitnexus/src/core/lbug/conn-lock.ts b/gitnexus/src/core/lbug/conn-lock.ts new file mode 100644 index 000000000..bb10d8113 --- /dev/null +++ b/gitnexus/src/core/lbug/conn-lock.ts @@ -0,0 +1,48 @@ +/** + * Serialize every operation on the shared singleton LadybugDB connection. + * + * LadybugDB is single-writer and its `Connection` is NOT safe for concurrent + * query execution: dispatching two queries on one connection at the same time + * lets two libuv workers mutate shared native engine state at once, corrupting + * the heap. This surfaced as `double free or corruption (out)` / SIGSEGV at the + * end of `analyze --pdg`, where the periodic WAL-checkpoint driver + * (`wal-checkpoint-driver.ts`) fired `CHECKPOINT` on the same connection a + * long-running PDG-table COPY was still using. `--pdg` makes those COPYs outlast + * the driver's 5 s tick, so the overlap (rare without `--pdg`) becomes reliable. + * + * Every singleton-`conn` helper in `lbug-adapter.ts` runs its full query + + * result-drain inside this lock, so the checkpoint driver, the bulk COPY, the + * embedding writeback, and the PDG edge deletes are mutually exclusive — the + * property that makes a strictly-serial workload stable. + * + * Implementation: a promise chain. Each caller installs a fresh unresolved tail, + * awaits the previous holder's tail, runs, then releases its own in `finally` + * (so a thrown op never wedges the connection). FIFO and re-entrancy-unsafe by + * design — a wrapped helper MUST NOT call another wrapped helper. + */ +let tail: Promise = Promise.resolve(); + +export const withConnLock = async (fn: () => Promise): Promise => { + const prior = tail; + let release!: () => void; + tail = new Promise((resolve) => { + release = resolve; + }); + await prior; + try { + return await fn(); + } finally { + release(); + } +}; + +/** + * Test-only: reset the lock chain to a fresh resolved tail. Production code has + * no reason to call this — a leaked-but-resolved tail is harmless — but unit + * tests want a clean chain per case. + * + * @internal + */ +export const __resetConnLockForTests = (): void => { + tail = Promise.resolve(); +}; diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index eb04ec5a4..2ba9aede9 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -6,6 +6,7 @@ import { finished } from 'stream/promises'; import path from 'path'; import lbug from '@ladybugdb/core'; import { closeQueryResults } from './query-result-utils.js'; +import { withConnLock } from './conn-lock.js'; import { KnowledgeGraph } from '../graph/types.js'; import { NODE_TABLES, @@ -178,6 +179,18 @@ export const splitRelCsvByLabelPair = async ( let db: lbug.Database | null = null; let conn: lbug.Connection | null = null; + +// Serialize every operation on the shared singleton `conn`. LadybugDB's +// Connection is single-writer and is NOT safe for concurrent query execution; +// the periodic WAL-checkpoint driver overlapping a long `--pdg` COPY on this +// connection corrupted native state (`double free or corruption`). Each +// singleton-`conn` helper below runs its full query + drain inside withConnLock. +// Invariant: a wrapped helper MUST NOT call another wrapped helper (re-entry +// self-deadlocks); all current holders are leaf-level. `streamQuery` is +// deliberately NOT wrapped — its per-row callback can re-enter the adapter and +// it only runs on the read path where the checkpoint driver is inactive. +// See conn-lock.ts for the full rationale. + let currentDbPath: string | null = null; let currentDbReadOnly = false; let ftsLoaded = false; @@ -472,8 +485,15 @@ const readQueryRows = async ( }; const queryAndDrain = async (targetConn: lbug.Connection, cypher: string): Promise => { - const queryResult = await targetConn.query(cypher); - await drainQueryResult(queryResult); + const run = async (): Promise => { + const queryResult = await targetConn.query(cypher); + await drainQueryResult(queryResult); + }; + // Serialize only when this runs on the shared singleton connection (the bulk + // node/relationship COPY captures `writeConn = conn`). Per-file / temp + // connections are distinct objects with no shared native state, so they must + // not block on — or be blocked by — the singleton's lock. + return targetConn === conn ? withConnLock(run) : run(); }; const READ_ONLY_SHADOW_REPLAY_PROBE = 'MATCH (n) RETURN n LIMIT 1'; @@ -1535,23 +1555,27 @@ export const executePrepared = async ( cypher: string, params: Record, ): Promise => { - if (!conn) { + const c = conn; + if (!c) { throw new Error('LadybugDB not initialized. Call initLbug first.'); } - const stmt = await conn.prepare(cypher); - if (!stmt.isSuccess()) { - const errMsg = await stmt.getErrorMessage(); - throw new Error(`Prepare failed: ${errMsg}`); - } - const queryResult = await conn.execute(stmt, params); - return await readQueryRows(queryResult); + return withConnLock(async () => { + const stmt = await c.prepare(cypher); + if (!stmt.isSuccess()) { + const errMsg = await stmt.getErrorMessage(); + throw new Error(`Prepare failed: ${errMsg}`); + } + const queryResult = await c.execute(stmt, params); + return await readQueryRows(queryResult); + }); }; export const executeWithReusedStatement = async ( cypher: string, paramsList: Array>, ): Promise => { - if (!conn) { + const c = conn; + if (!c) { throw new Error('LadybugDB not initialized. Call initLbug first.'); } if (paramsList.length === 0) return; @@ -1559,39 +1583,50 @@ export const executeWithReusedStatement = async ( const SUB_BATCH_SIZE = 4; for (let i = 0; i < paramsList.length; i += SUB_BATCH_SIZE) { const subBatch = paramsList.slice(i, i + SUB_BATCH_SIZE); - const stmt = await conn.prepare(cypher); - if (!stmt.isSuccess()) { - const errMsg = await stmt.getErrorMessage(); - throw new Error(`Prepare failed: ${errMsg}`); - } - try { - for (const params of subBatch) { - await drainQueryResult(await conn.execute(stmt, params)); + // One critical section per sub-batch: the prepare + its executes run with + // exclusive access to the connection (so the WAL checkpoint driver cannot + // interleave a CHECKPOINT mid-batch), while the lock is released between + // sub-batches to let the driver checkpoint during a long writeback. + await withConnLock(async () => { + const stmt = await c.prepare(cypher); + if (!stmt.isSuccess()) { + const errMsg = await stmt.getErrorMessage(); + throw new Error(`Prepare failed: ${errMsg}`); } - } catch (e) { - const msg = e instanceof Error ? e.message : String(e); - const queryPreview = cypher.replace(/\s+/g, ' ').slice(0, 120); - throw new Error( - `Batch execution failed for rows ${i + 1}-${i + subBatch.length}: ${msg} (${queryPreview})`, - ); - } - // Note: LadybugDB PreparedStatement doesn't require explicit close() + try { + for (const params of subBatch) { + await drainQueryResult(await c.execute(stmt, params)); + } + } catch (e) { + const msg = e instanceof Error ? e.message : String(e); + const queryPreview = cypher.replace(/\s+/g, ' ').slice(0, 120); + throw new Error( + `Batch execution failed for rows ${i + 1}-${i + subBatch.length}: ${msg} (${queryPreview})`, + ); + } + // Note: LadybugDB PreparedStatement doesn't require explicit close() + }); } }; export const getLbugStats = async (): Promise<{ nodes: number; edges: number }> => { - if (!conn) return { nodes: 0, edges: 0 }; + const c = conn; + if (!c) return { nodes: 0, edges: 0 }; + // Called during analyze finalize while the WAL-checkpoint driver is still + // running; each count read takes the connection lock so it cannot execute + // concurrently with a driver CHECKPOINT. Per-query locking lets the driver + // checkpoint between table counts rather than waiting for the whole sweep. let totalNodes = 0; for (const tableName of NODE_TABLES) { try { - const queryResult = await conn.query( - `MATCH (n:${escapeTableName(tableName)}) RETURN count(n) AS cnt`, - ); - const nodeRows = await readQueryRows(queryResult); - if (nodeRows.length > 0) { - totalNodes += Number(nodeRows[0]?.cnt ?? nodeRows[0]?.[0] ?? 0); - } + totalNodes += await withConnLock(async () => { + const queryResult = await c.query( + `MATCH (n:${escapeTableName(tableName)}) RETURN count(n) AS cnt`, + ); + const nodeRows = await readQueryRows(queryResult); + return nodeRows.length > 0 ? Number(nodeRows[0]?.cnt ?? nodeRows[0]?.[0] ?? 0) : 0; + }); } catch { // ignore } @@ -1599,13 +1634,13 @@ export const getLbugStats = async (): Promise<{ nodes: number; edges: number }> let totalEdges = 0; try { - const queryResult = await conn.query( - `MATCH ()-[r:${REL_TABLE_NAME}]->() RETURN count(r) AS cnt`, - ); - const edgeRows = await readQueryRows(queryResult); - if (edgeRows.length > 0) { - totalEdges = Number(edgeRows[0]?.cnt ?? edgeRows[0]?.[0] ?? 0); - } + totalEdges = await withConnLock(async () => { + const queryResult = await c.query( + `MATCH ()-[r:${REL_TABLE_NAME}]->() RETURN count(r) AS cnt`, + ); + const edgeRows = await readQueryRows(queryResult); + return edgeRows.length > 0 ? Number(edgeRows[0]?.cnt ?? edgeRows[0]?.[0] ?? 0) : 0; + }); } catch { // ignore } @@ -1768,10 +1803,13 @@ export const fetchExistingEmbeddingHashes = async ( * @see safeClose — CHECKPOINT + connection/database close */ export const flushWAL = async (): Promise => { - if (!conn) return; + const c = conn; + if (!c) return; try { - const checkpointResult = await conn.query('CHECKPOINT'); - await drainQueryResult(checkpointResult); + await withConnLock(async () => { + const checkpointResult = await c.query('CHECKPOINT'); + await drainQueryResult(checkpointResult); + }); } catch (err) { logger.debug( `GitNexus: LadybugDB CHECKPOINT skipped/failed during WAL flush: ${summarizeError(err)}`, @@ -1794,9 +1832,15 @@ export const flushWAL = async (): Promise => { * whether to retry. */ export const tryFlushWAL = async (): Promise => { - if (!conn) return false; - const checkpointResult = await conn.query('CHECKPOINT'); - await drainQueryResult(checkpointResult); + const c = conn; + if (!c) return false; + // Runs on the periodic WAL-checkpoint driver. The lock makes this CHECKPOINT + // wait for any in-flight COPY / writeback on the singleton connection instead + // of executing concurrently with it (the `analyze --pdg` heap-corruption bug). + await withConnLock(async () => { + const checkpointResult = await c.query('CHECKPOINT'); + await drainQueryResult(checkpointResult); + }); return true; }; @@ -2036,44 +2080,49 @@ export const deleteAllCommunitiesAndProcesses = async (): Promise<{ * plain DELETE on the typed CodeRelation rows — endpoints are untouched. */ export const deleteAllInterprocTaintPaths = async (): Promise<{ edgesDeleted: number }> => { - if (!conn) { + const c = conn; + if (!c) { throw new Error('LadybugDB not initialized. Call initLbug first.'); } - let edgesDeleted = 0; - let countResult: lbug.QueryResult | lbug.QueryResult[] | undefined; - try { - countResult = await conn.query( - `MATCH ()-[r:CodeRelation]->() WHERE r.type = 'TAINT_PATH' RETURN count(r) AS cnt`, - ); - const result = Array.isArray(countResult) ? countResult[0] : countResult; - const rows = await result.getAll(); - const count = Number(rows[0]?.cnt ?? rows[0]?.[0] ?? 0); - if (count > 0) { - await conn.query(`MATCH ()-[r:CodeRelation]->() WHERE r.type = 'TAINT_PATH' DELETE r`); - edgesDeleted = count; - } - } catch (err) { - // A missing table on a freshly-initialized DB is the benign, expected case - // (the count query above is what throws) — stay silent. Any OTHER failure - // (lock, disk, native error) would leave stale TAINT_PATH rows that the - // subsequent re-extract then DUPLICATES (CodeRelation has no PK), so it - // must ABORT the writeback (#2084 review P2-5): re-throw so the caller's - // crash-recovery dirty flag forces a clean full rebuild on the next run, - // rather than silently writing duplicate cross-function findings. - const msg = err instanceof Error ? err.message : String(err); - if (/no table|not exist|not found|does not exist|Table .* does not exist/i.test(msg)) { + // count + DELETE run as one critical section on the singleton connection so a + // concurrent WAL-checkpoint cannot corrupt native state mid-delete (#pdg). + return withConnLock(async () => { + let edgesDeleted = 0; + let countResult: lbug.QueryResult | lbug.QueryResult[] | undefined; + try { + countResult = await c.query( + `MATCH ()-[r:CodeRelation]->() WHERE r.type = 'TAINT_PATH' RETURN count(r) AS cnt`, + ); + const result = Array.isArray(countResult) ? countResult[0] : countResult; + const rows = await result.getAll(); + const count = Number(rows[0]?.cnt ?? rows[0]?.[0] ?? 0); + if (count > 0) { + await c.query(`MATCH ()-[r:CodeRelation]->() WHERE r.type = 'TAINT_PATH' DELETE r`); + edgesDeleted = count; + } + } catch (err) { + // A missing table on a freshly-initialized DB is the benign, expected case + // (the count query above is what throws) — stay silent. Any OTHER failure + // (lock, disk, native error) would leave stale TAINT_PATH rows that the + // subsequent re-extract then DUPLICATES (CodeRelation has no PK), so it + // must ABORT the writeback (#2084 review P2-5): re-throw so the caller's + // crash-recovery dirty flag forces a clean full rebuild on the next run, + // rather than silently writing duplicate cross-function findings. + const msg = err instanceof Error ? err.message : String(err); + if (/no table|not exist|not found|does not exist|Table .* does not exist/i.test(msg)) { + if (countResult) await closeQueryResults(countResult); + return { edgesDeleted }; + } if (countResult) await closeQueryResults(countResult); - return { edgesDeleted }; + throw new Error( + `[taint-interproc] failed to clear existing TAINT_PATH edges before incremental ` + + `re-write (${msg}) — aborting to avoid duplicate cross-function findings; ` + + `the next run will full-rebuild`, + ); } if (countResult) await closeQueryResults(countResult); - throw new Error( - `[taint-interproc] failed to clear existing TAINT_PATH edges before incremental ` + - `re-write (${msg}) — aborting to avoid duplicate cross-function findings; ` + - `the next run will full-rebuild`, - ); - } - if (countResult) await closeQueryResults(countResult); - return { edgesDeleted }; + return { edgesDeleted }; + }); }; /** @@ -2088,42 +2137,47 @@ export const deleteAllInterprocTaintPaths = async (): Promise<{ edgesDeleted: nu * an unchanged function's summary from being lost. */ export const deleteAllCallSummaries = async (): Promise<{ edgesDeleted: number }> => { - if (!conn) { + const c = conn; + if (!c) { throw new Error('LadybugDB not initialized. Call initLbug first.'); } - let edgesDeleted = 0; - let countResult: lbug.QueryResult | lbug.QueryResult[] | undefined; - try { - countResult = await conn.query( - `MATCH ()-[r:CodeRelation]->() WHERE r.type = 'CALL_SUMMARY' RETURN count(r) AS cnt`, - ); - const result = Array.isArray(countResult) ? countResult[0] : countResult; - const rows = await result.getAll(); - const count = Number(rows[0]?.cnt ?? rows[0]?.[0] ?? 0); - if (count > 0) { - await conn.query(`MATCH ()-[r:CodeRelation]->() WHERE r.type = 'CALL_SUMMARY' DELETE r`); - edgesDeleted = count; - } - } catch (err) { - // A missing table on a freshly-initialized DB is the benign, expected case - // (the count query is what throws) — stay silent. Any OTHER failure would - // leave stale rows that the re-extract then DUPLICATES (CodeRelation has no - // PK), so it must ABORT the writeback: re-throw so the caller's crash- - // recovery dirty flag forces a clean full rebuild on the next run. - const msg = err instanceof Error ? err.message : String(err); - if (/no table|not exist|not found|does not exist|Table .* does not exist/i.test(msg)) { + // count + DELETE run as one critical section on the singleton connection so a + // concurrent WAL-checkpoint cannot corrupt native state mid-delete (#pdg). + return withConnLock(async () => { + let edgesDeleted = 0; + let countResult: lbug.QueryResult | lbug.QueryResult[] | undefined; + try { + countResult = await c.query( + `MATCH ()-[r:CodeRelation]->() WHERE r.type = 'CALL_SUMMARY' RETURN count(r) AS cnt`, + ); + const result = Array.isArray(countResult) ? countResult[0] : countResult; + const rows = await result.getAll(); + const count = Number(rows[0]?.cnt ?? rows[0]?.[0] ?? 0); + if (count > 0) { + await c.query(`MATCH ()-[r:CodeRelation]->() WHERE r.type = 'CALL_SUMMARY' DELETE r`); + edgesDeleted = count; + } + } catch (err) { + // A missing table on a freshly-initialized DB is the benign, expected case + // (the count query is what throws) — stay silent. Any OTHER failure would + // leave stale rows that the re-extract then DUPLICATES (CodeRelation has no + // PK), so it must ABORT the writeback: re-throw so the caller's crash- + // recovery dirty flag forces a clean full rebuild on the next run. + const msg = err instanceof Error ? err.message : String(err); + if (/no table|not exist|not found|does not exist|Table .* does not exist/i.test(msg)) { + if (countResult) await closeQueryResults(countResult); + return { edgesDeleted }; + } if (countResult) await closeQueryResults(countResult); - return { edgesDeleted }; + throw new Error( + `[call-summary] failed to clear existing CALL_SUMMARY edges before incremental ` + + `re-write (${msg}) — aborting to avoid duplicate summaries; ` + + `the next run will full-rebuild`, + ); } if (countResult) await closeQueryResults(countResult); - throw new Error( - `[call-summary] failed to clear existing CALL_SUMMARY edges before incremental ` + - `re-write (${msg}) — aborting to avoid duplicate summaries; ` + - `the next run will full-rebuild`, - ); - } - if (countResult) await closeQueryResults(countResult); - return { edgesDeleted }; + return { edgesDeleted }; + }); }; // ============================================================================ diff --git a/gitnexus/src/core/lbug/wal-checkpoint-driver.ts b/gitnexus/src/core/lbug/wal-checkpoint-driver.ts index 57dc6b2d6..08e48e6f2 100644 --- a/gitnexus/src/core/lbug/wal-checkpoint-driver.ts +++ b/gitnexus/src/core/lbug/wal-checkpoint-driver.ts @@ -164,6 +164,14 @@ export const startWalCheckpointDriver = ( const tick = async (): Promise => { if (stopped) return; + // Reentrancy guard: setInterval keeps firing on its fixed cadence even when + // the previous checkpoint has not settled (a CHECKPOINT can outlast the + // period during a large `--pdg` writeback). Without this, each overdue tick + // would queue another CHECKPOINT — they now serialize on the connection lock + // (lbug-adapter `withConnLock`), but letting them pile up is still pointless + // work and widens the window for a backlog at stop(). Skip while one is in + // flight; the next tick covers any WAL accumulated in the meantime. + if (inflight) return; inflight = runCheckpointWithRetry() .then(() => undefined) .catch((err) => { diff --git a/gitnexus/test/unit/conn-lock.test.ts b/gitnexus/test/unit/conn-lock.test.ts new file mode 100644 index 000000000..cbce6ac0c --- /dev/null +++ b/gitnexus/test/unit/conn-lock.test.ts @@ -0,0 +1,73 @@ +/** + * Unit tests for the LadybugDB connection serialization lock (conn-lock.ts). + * + * This lock is the fix for the `analyze --pdg` native crash: the WAL-checkpoint + * driver's periodic CHECKPOINT was executing on the shared singleton connection + * concurrently with a long-running COPY, and LadybugDB's single-writer + * Connection corrupts native heap state under concurrent query execution + * (`double free or corruption (out)` / SIGSEGV). These tests assert the + * lock's one-at-a-time guarantee deterministically, with no native engine — + * the property that makes the otherwise-crashing overlap safe. + */ +import { afterEach, describe, expect, it } from 'vitest'; +import { withConnLock, __resetConnLockForTests } from '../../src/core/lbug/conn-lock.js'; + +afterEach(() => { + __resetConnLockForTests(); +}); + +describe('withConnLock — connection serialization', () => { + it('runs critical sections one at a time in FIFO order with no interleave', async () => { + const events: string[] = []; + const section = (id: string, yields: number) => async (): Promise => { + events.push(`${id}:enter`); + // Yield to the microtask queue repeatedly. Without serialization a later + // section's `enter` would slip in between these yields. + for (let i = 0; i < yields; i++) await Promise.resolve(); + events.push(`${id}:exit`); + }; + + // B and C are launched while A (which yields the most) is mid-flight. + await Promise.all([ + withConnLock(section('A', 5)), + withConnLock(section('B', 0)), + withConnLock(section('C', 0)), + ]); + + expect(events).toEqual(['A:enter', 'A:exit', 'B:enter', 'B:exit', 'C:enter', 'C:exit']); + }); + + it('never lets two critical sections overlap under heavy concurrency', async () => { + let active = 0; + const observedMax: number[] = []; + const op = () => async (): Promise => { + active++; + observedMax.push(active); + await Promise.resolve(); + await Promise.resolve(); + active--; + }; + + await Promise.all(Array.from({ length: 25 }, () => withConnLock(op()))); + + // The concurrency count observed at the top of every critical section was + // always exactly 1 — i.e. no two ran at once. This is precisely what stops + // the checkpoint driver from racing a COPY on the native connection. + expect(Math.max(...observedMax)).toBe(1); + }); + + it('releases the lock when a critical section throws (no permanent wedge)', async () => { + await expect( + withConnLock(async () => { + throw new Error('boom'); + }), + ).rejects.toThrow('boom'); + + // A failed op must not strand the lock — the next caller still acquires it. + await expect(withConnLock(async () => 'recovered')).resolves.toBe('recovered'); + }); + + it('returns the wrapped operation result', async () => { + await expect(withConnLock(async () => 42)).resolves.toBe(42); + }); +}); diff --git a/gitnexus/test/unit/wal-checkpoint-driver-reentrancy.test.ts b/gitnexus/test/unit/wal-checkpoint-driver-reentrancy.test.ts new file mode 100644 index 000000000..289a5a23b --- /dev/null +++ b/gitnexus/test/unit/wal-checkpoint-driver-reentrancy.test.ts @@ -0,0 +1,66 @@ +/** + * Reentrancy-guard test for the manual WAL checkpoint driver. + * + * `setInterval` fires on a fixed cadence regardless of whether the previous + * checkpoint has settled. During a large `--pdg` writeback a CHECKPOINT can + * outlast the period; without a guard each overdue tick would launch ANOTHER + * concurrent CHECKPOINT on the singleton connection. The guard (`if (inflight) + * return`) ensures at most one checkpoint is ever in flight. + * + * `tryFlushWAL` is mocked so we can hold a checkpoint "in flight" and drive the + * interval with fake timers — no native engine involved. + */ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +vi.mock('../../src/core/lbug/lbug-adapter.js', () => ({ + tryFlushWAL: vi.fn(), +})); + +import { startWalCheckpointDriver } from '../../src/core/lbug/wal-checkpoint-driver.js'; +import { tryFlushWAL } from '../../src/core/lbug/lbug-adapter.js'; + +const mockedTryFlush = vi.mocked(tryFlushWAL); + +describe('startWalCheckpointDriver — reentrancy guard', () => { + let originalEnv: string | undefined; + + beforeEach(() => { + originalEnv = process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT; + delete process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT; // default = enabled + vi.useFakeTimers(); + mockedTryFlush.mockReset(); + }); + + afterEach(() => { + vi.useRealTimers(); + if (originalEnv === undefined) delete process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT; + else process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT = originalEnv; + }); + + it('starts only one checkpoint while a prior one is still in flight, then resumes', async () => { + let resolveFirst!: () => void; + mockedTryFlush + // First checkpoint is held open until we resolve it. + .mockImplementationOnce( + () => + new Promise((resolve) => { + resolveFirst = () => resolve(true); + }), + ) + // Any later checkpoint completes immediately. + .mockImplementation(() => Promise.resolve(true)); + + const driver = startWalCheckpointDriver({ periodMs: 10 }); + + // ~5 ticks fire while the first checkpoint is still pending. + await vi.advanceTimersByTimeAsync(55); + expect(mockedTryFlush).toHaveBeenCalledTimes(1); + + // Let the first settle; subsequent ticks may now fire a new checkpoint. + resolveFirst(); + await vi.advanceTimersByTimeAsync(25); + expect(mockedTryFlush.mock.calls.length).toBeGreaterThanOrEqual(2); + + await driver.stop(); + }); +});