mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-11 03:38:07 +00:00
fix(lbug): serialize singleton connection to stop --pdg analyze double-free
LadybugDB is single-writer and its Connection is NOT safe for concurrent
query execution. The WAL-checkpoint driver (5s setInterval) issued
`conn.query('CHECKPOINT')` on the same module-singleton `conn` the analyze
pipeline used for COPY. With --pdg the extra BasicBlock / REACHING_DEF / CDG /
POST_DOMINATE / TAINTED / CALL_SUMMARY / TAINT_PATH table COPYs outlast the
5s tick, so a checkpoint executed concurrently with an in-flight COPY on one
connection -> two libuv workers mutate shared native state -> heap corruption
("double free or corruption (out)" / SIGABRT, detected at the final
"Saving metadata..." free).
Fix: add conn-lock.ts (`withConnLock`, a promise-chain mutex) and run every
singleton-`conn` helper's full query + result-drain inside it: queryAndDrain
(when targetConn === conn), executePrepared, executeWithReusedStatement,
flushWAL, tryFlushWAL, getLbugStats, deleteAllInterprocTaintPaths,
deleteAllCallSummaries. Add an `if (inflight) return` reentrancy guard to the
driver tick so overdue ticks don't stack checkpoints. streamQuery is
intentionally NOT wrapped (read path, re-entrant per-row callback).
Reproduced the crash with concurrent queries on one raw Connection (serial =
stable); verified the fix drives the same overlap through the locked adapter
without crashing.
Tests: conn-lock serialization (no overlap / FIFO / throw-releases) and
driver reentrancy guard.
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
221069785b
commit
36ef9c27c7
5 changed files with 361 additions and 112 deletions
48
gitnexus/src/core/lbug/conn-lock.ts
Normal file
48
gitnexus/src/core/lbug/conn-lock.ts
Normal file
|
|
@ -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<void> = Promise.resolve();
|
||||
|
||||
export const withConnLock = async <T>(fn: () => Promise<T>): Promise<T> => {
|
||||
const prior = tail;
|
||||
let release!: () => void;
|
||||
tail = new Promise<void>((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();
|
||||
};
|
||||
|
|
@ -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<void> => {
|
||||
const queryResult = await targetConn.query(cypher);
|
||||
await drainQueryResult(queryResult);
|
||||
const run = async (): Promise<void> => {
|
||||
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<string, any>,
|
||||
): Promise<any[]> => {
|
||||
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<Record<string, any>>,
|
||||
): Promise<void> => {
|
||||
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<void> => {
|
||||
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<void> => {
|
|||
* whether to retry.
|
||||
*/
|
||||
export const tryFlushWAL = async (): Promise<boolean> => {
|
||||
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 };
|
||||
});
|
||||
};
|
||||
|
||||
// ============================================================================
|
||||
|
|
|
|||
|
|
@ -164,6 +164,14 @@ export const startWalCheckpointDriver = (
|
|||
|
||||
const tick = async (): Promise<void> => {
|
||||
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) => {
|
||||
|
|
|
|||
73
gitnexus/test/unit/conn-lock.test.ts
Normal file
73
gitnexus/test/unit/conn-lock.test.ts
Normal file
|
|
@ -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<void> => {
|
||||
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<void> => {
|
||||
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);
|
||||
});
|
||||
});
|
||||
66
gitnexus/test/unit/wal-checkpoint-driver-reentrancy.test.ts
Normal file
66
gitnexus/test/unit/wal-checkpoint-driver-reentrancy.test.ts
Normal file
|
|
@ -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<boolean>((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();
|
||||
});
|
||||
});
|
||||
Loading…
Add table
Reference in a new issue