mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-08-28 05:25:25 +00:00
* test(lbug): lock PARALLEL=false as a tested correctness invariant (#2203)
The parallel CSV reader (Kuzu-derived, default PARALLEL=true) cannot parse
quoted fields with embedded newlines (kuzudb/kuzu#5778); our content/text
columns hold source code, so PARALLEL=false is mandatory for correctness.
Add a live-DB multiline-quoted round-trip that fails if it is ever flipped,
plus a static guard on the generated COPY queries. Export COPY_CSV_OPTS /
getCopyQuery for the static assertion; document the invariant at the source.
* feat(lbug): expose node/rel phase boundary via onNodePhaseComplete hook (#2203)
streamAllCSVsToDisk now fires an optional onNodePhaseComplete(nodeFiles)
callback right after node CSVs are flushed and before the relationship pass
writes any rel_*.csv — the boundary the COPY-overlap leg needs. The node-file
manifest construction is hoisted above the rel pass and reused in the return,
so output is byte-for-byte identical when no callback is supplied (verified by
the emit bench fingerprint and the splitRelCsvByLabelPair differential oracle).
The callback is not awaited, so the rel pass runs concurrently with the
caller's node COPY.
* perf(lbug): overlap node COPY with relationship emit (#2203)
The deferred parallelism leg of #2203. LadybugDB is single-writer and its
parallel CSV reader is unsafe for our multiline content (kuzudb/kuzu#5778), so
the only safe parallelism is pipeline-overlap: node COPY (uses conn, never the
rel files) runs concurrently with the relationship emit pass (writes rel_*.csv,
never conn). Node COPY now starts at streamAllCSVsToDisk's onNodePhaseComplete
boundary while the rel pass keeps writing; the relationship COPY still waits for
node COPY (FK precondition), so DB load order and content are unchanged.
- Extract copyNodeCSVs; start it in the hook (overlap) or after emit (serial).
- GITNEXUS_SERIAL_LBUG_LOAD=1 forces the legacy strictly-sequential path
(operator escape hatch + differential-test oracle).
- Settle the in-flight node-COPY promise on emit failure (no unhandled
rejection); rethrow node-COPY errors at the FK barrier.
- Preserve the PDG manifest merge + collision guards (node merge at the hook,
rel merge before rel COPY) and all retry/fallback/cleanup behavior.
- PROF_LBUG_LOAD gains mode=overlap|serial; copy-nodes becomes the residual
node-COPY time after emit (trends to 0 as overlap hides it).
* test(lbug): differential gate — overlap load === serial load (#2203)
Loads one fixture (multiple node tables, multiple edge pairs, multiline File
content + BasicBlock text) into two fresh DBs — once via the default node-COPY
‖ rel-emit overlap, once via GITNEXUS_SERIAL_LBUG_LOAD=1 — and asserts the two
databases are content-equivalent: identical per-table node counts, per-type
edge counts, byte-for-byte multiline content/text, and identical
insertedRels/skippedRels/warnings. This is the issue's byte-identical-content
acceptance gate for the parallelism leg.
* fix(review): apply autofix feedback
- csv-generator: onNodePhaseComplete doc-contract now matches reality (a sync
throw is allowed and is how loadGraphToLbug surfaces the manifest collision
guard) — drops the inaccurate 'must not throw synchronously' line.
- lbug-adapter: copyNodeCSVs totalSteps is the node-table count (drop the +1
rel-step holdover; the rel COPY has its own progress line).
- lbug-adapter: on emit+node-COPY double-failure, log the swallowed node-COPY
error before rethrowing the emit error (diagnosability).
- lbug-load-prof test: assert mode=overlap on the default path.
* fix(test): use mkdtemp for secure temp dirs (CodeQL js/insecure-temporary-file)
CodeQL flagged lbug-load-overlap.test.ts writing a file into a predictable
os.tmpdir() path. Create the base temp dir with fs.mkdtemp (atomic, random
suffix) in both new live-DB tests, and switch to the gitnexus-lbug- prefix that
TEST_FIXTURE_PREFIXES recognizes so the Windows stale-sidecar sweep covers
these fixtures.
* fix(lbug): check PDG manifest rel-pair collision before node COPY (#2203)
Found by Codex in tri-review. The manifest rel-pair collision guard ran after
the FK barrier (after node COPY committed), so on that should-never-happen
error branch the overlap path left orphan node rows AND the
GITNEXUS_SERIAL_LBUG_LOAD escape hatch diverged from the legacy 'validate
manifest before any COPY' behavior. Move the rel merge + collision check ahead
of beginNodeCopy/the barrier: the serial path now detects a collision before
committing any node rows (legacy parity restored — the escape hatch is a
faithful oracle again), and the overlap path detects it as early as csvResult
is available. The node-collision guard already ran before node COPY (in the
hook).
* test(lbug): cover rel-emit failure with node COPY in flight (#2203)
Resolves a P1 review gap on PR #2226: the overlap's catch(emitErr) branch
(settle the in-flight node-COPY promise, then rethrow the emit error) was
untested. Fault-injects via a vi.mock of streamAllCSVsToDisk that fires
onNodePhaseComplete (starting a real node COPY on a live DB) then throws,
asserting loadGraphToLbug rejects with the emit error and no unhandled
rejection leaks. Also covers the both-fail case (node COPY error is logged,
emit error still wins). Listener removed in finally; macrotask queue flushed
before the assertion so it can't pass vacuously.
* test(lbug): cover node-COPY hard-failure rethrow at the FK barrier (#2203)
Resolves the second P1 review gap on PR #2226. Mocks emit to fire
onNodePhaseComplete with a nodeFiles entry pointing at a missing CSV (a
bind-time COPY error that IGNORE_ERRORS does not suppress) and otherwise
succeed, so copyNodeCSVs throws, the error is captured in nodeCopyError, and
loadGraphToLbug rethrows it at the FK barrier — asserted via rejects /COPY
failed for File/.
* test(lbug): cover PDG manifest rel-pair collision in overlap + serial (#2203)
Resolves the P2 gap behind the Codex tri-review finding: the manifest rel-pair
collision guard (moved ahead of node COPY in ad195582) had no test. A leaky
graph with a structural BasicBlock->BasicBlock edge (routed by id-prefix, no
BasicBlock nodes — isolating the rel-pair clash from the node-CSV one) plus a
PdgEmitSink manifest declaring the same pair makes loadGraphToLbug reject with
the rel-pair collision error, asserted on both the overlap (default) and serial
(GITNEXUS_SERIAL_LBUG_LOAD=1) paths.
* perf(lbug): yield the event loop periodically during relationship emit (#2203)
Resolves a P2 review finding on PR #2226: the relationship-emit loop ran long
synchronous stretches between write-stream drain awaits, which could starve the
overlapped node-COPY callbacks on fast I/O and erode the node-COPY-||-rel-emit
overlap. Yield via setImmediate every REL_YIELD_EVERY (5000) edges so the node
COPY and drains get scheduling time. Scheduling-only — emit bench fingerprint
unchanged (byte-identical), csv-pipeline determinism + overlap differential
green.
* refactor(lbug): extract shared copyCsvWithRetry helper (#2203)
Resolves a P2 maintainability finding on PR #2226: the COPY + IGNORE_ERRORS
retry block was duplicated in copyNodeCSVs and the inline relationship-COPY
loop. Extract copyCsvWithRetry(conn, query, onError); the callback receives the
RAW retry error so each site keeps its own message shape + slice length (node
throws, slices 200; relationship warns + records the failed pair, slices 80).
Behavior-preserving — guarded by the live-DB round-trips plus the new
node-COPY-failure and overlap error-path tests.
* docs(lbug): document loadGraphToLbug non-transactionality (#2203)
Resolves the advisory review finding on PR #2226: loadGraphToLbug runs
independent COPYs with no surrounding transaction, so a mid-load failure leaves
a partial DB and recovery is a --force re-analyze. Make that contract explicit
on the function so callers don't assume atomicity.
This commit is contained in:
parent
ff0124e067
commit
b895a20415
8 changed files with 922 additions and 87 deletions
|
|
@ -47,6 +47,15 @@ const orderedRelationships = (
|
|||
/** Flush buffered rows to disk every N rows */
|
||||
const FLUSH_EVERY = 500;
|
||||
|
||||
/**
|
||||
* Yield the event loop every N relationship rows during the emit pass (#2226 F4)
|
||||
* so a concurrent node COPY (the overlap in loadGraphToLbug) and write-stream
|
||||
* drain callbacks get scheduling time during long synchronous emit stretches.
|
||||
* Scheduling-only — never changes row content or order (byte-identical). Tuning
|
||||
* constant, not load-bearing.
|
||||
*/
|
||||
const REL_YIELD_EVERY = 5000;
|
||||
|
||||
// ============================================================================
|
||||
// CSV ESCAPE UTILITIES
|
||||
// ============================================================================
|
||||
|
|
@ -282,11 +291,26 @@ export interface StreamedCSVResult {
|
|||
* Stream all CSV data directly to disk files.
|
||||
* Iterates graph nodes exactly ONCE — routes each node to the right writer.
|
||||
* File contents are lazy-read from disk with a generous LRU cache.
|
||||
*
|
||||
* `onNodePhaseComplete` (optional, #2203 parallelism leg): fired exactly once,
|
||||
* right after every node CSV is fully flushed to disk and BEFORE the
|
||||
* relationship pass starts writing any `rel_*.csv`. It receives the finished
|
||||
* node-file manifest so the caller can begin `COPY`-ing nodes while this
|
||||
* function keeps generating relationship CSVs (the only single-writer-safe
|
||||
* overlap — node `COPY` ‖ relationship emit). It is intentionally NOT awaited:
|
||||
* the relationship pass proceeds concurrently with whatever the caller
|
||||
* schedules. A synchronous throw from the callback is allowed and propagates out
|
||||
* of this function (rejecting the returned promise) — it is raised before the
|
||||
* relationship pass begins, so no `rel_*.csv` is written; `loadGraphToLbug` uses
|
||||
* this to surface its PDG-manifest collision guard. The callback must NOT, however,
|
||||
* schedule un-awaited async work that can reject unobserved. Absent ⇒ today's
|
||||
* behavior, byte-for-byte.
|
||||
*/
|
||||
export const streamAllCSVsToDisk = async (
|
||||
graph: KnowledgeGraph,
|
||||
repoPath: string,
|
||||
csvDir: string,
|
||||
onNodePhaseComplete?: (nodeFiles: Map<NodeTableName, { csvPath: string; rows: number }>) => void,
|
||||
): Promise<StreamedCSVResult> => {
|
||||
// Deterministic (id-sorted) node/relationship row order when enabled;
|
||||
// default off = today's graph-insertion order (byte-identical).
|
||||
|
|
@ -615,29 +639,11 @@ export const streamAllCSVsToDisk = async (
|
|||
];
|
||||
await Promise.all(allWriters.map((w) => w.finish()));
|
||||
|
||||
// --- Stream relationships directly to per-FROM→TO-label-pair files ---
|
||||
// (#2203 U2) Route every edge to its pair file in this single pass. The old
|
||||
// monolithic relations.csv — and its line-by-line re-read + per-edge regex
|
||||
// re-split in loadGraphToLbug — are gone, so the ~1M-edge set is written and
|
||||
// read once instead of twice. The router applies the SAME label-derivation +
|
||||
// validTables filter as the legacy splitRelCsvByLabelPair, so the per-pair
|
||||
// files are byte-identical (asserted by the differential test).
|
||||
const relRouter = new RelPairRouter(csvDir, REL_CSV_HEADER, new Set<string>(NODE_TABLES));
|
||||
try {
|
||||
for (const rel of orderedRelationships(graph, sortOutput)) {
|
||||
const pending = relRouter.route(rel.sourceId, rel.targetId, buildRelRow(rel));
|
||||
if (pending) await pending;
|
||||
}
|
||||
await relRouter.close();
|
||||
} catch (err) {
|
||||
relRouter.destroy();
|
||||
// Rethrow the real stream error (EMFILE / disk-full) rather than the generic
|
||||
// AbortError a pending drain-await rejects with — mirrors the retained
|
||||
// splitRelCsvByLabelPair's `throw streamError ?? err`.
|
||||
throw relRouter.lastError ?? err;
|
||||
}
|
||||
|
||||
// Build result map — only include tables that have rows
|
||||
// Build the node-file manifest now (all writers are flushed; `.rows` is
|
||||
// final). Hoisted above the relationship pass so `onNodePhaseComplete` can
|
||||
// hand the caller a complete node manifest to start COPY-ing while we keep
|
||||
// generating relationship CSVs below (#2203 overlap). The same map is
|
||||
// returned, so the result is unchanged when no callback is supplied.
|
||||
const nodeFiles = new Map<NodeTableName, { csvPath: string; rows: number }>();
|
||||
const tableMap: [NodeTableName, BufferedCSVWriter][] = [
|
||||
['File', fileWriter],
|
||||
|
|
@ -666,6 +672,37 @@ export const streamAllCSVsToDisk = async (
|
|||
}
|
||||
}
|
||||
|
||||
// Node CSVs are on disk; relationship CSVs have not been touched yet. Hand
|
||||
// the manifest to the caller (not awaited — the rel pass runs concurrently).
|
||||
onNodePhaseComplete?.(nodeFiles);
|
||||
|
||||
// --- Stream relationships directly to per-FROM→TO-label-pair files ---
|
||||
// (#2203 U2) Route every edge to its pair file in this single pass. The old
|
||||
// monolithic relations.csv — and its line-by-line re-read + per-edge regex
|
||||
// re-split in loadGraphToLbug — are gone, so the ~1M-edge set is written and
|
||||
// read once instead of twice. The router applies the SAME label-derivation +
|
||||
// validTables filter as the legacy splitRelCsvByLabelPair, so the per-pair
|
||||
// files are byte-identical (asserted by the differential test).
|
||||
const relRouter = new RelPairRouter(csvDir, REL_CSV_HEADER, new Set<string>(NODE_TABLES));
|
||||
try {
|
||||
let emitted = 0;
|
||||
for (const rel of orderedRelationships(graph, sortOutput)) {
|
||||
const pending = relRouter.route(rel.sourceId, rel.targetId, buildRelRow(rel));
|
||||
if (pending) await pending;
|
||||
// Periodically hand the event loop back so the overlapped node COPY and
|
||||
// write-stream drains run instead of starving behind this synchronous
|
||||
// loop (#2226 F4). No effect on emitted bytes — pure scheduling.
|
||||
if (++emitted % REL_YIELD_EVERY === 0) await new Promise((r) => setImmediate(r));
|
||||
}
|
||||
await relRouter.close();
|
||||
} catch (err) {
|
||||
relRouter.destroy();
|
||||
// Rethrow the real stream error (EMFILE / disk-full) rather than the generic
|
||||
// AbortError a pending drain-await rejects with — mirrors the retained
|
||||
// splitRelCsvByLabelPair's `throw streamError ?? err`.
|
||||
throw relRouter.lastError ?? err;
|
||||
}
|
||||
|
||||
return {
|
||||
nodeFiles,
|
||||
relsByPair: relRouter.byPair,
|
||||
|
|
|
|||
|
|
@ -16,7 +16,7 @@ import {
|
|||
STALE_HASH_SENTINEL,
|
||||
NodeTableName,
|
||||
} from './schema.js';
|
||||
import { streamAllCSVsToDisk } from './csv-generator.js';
|
||||
import { streamAllCSVsToDisk, type StreamedCSVResult } from './csv-generator.js';
|
||||
import type { PdgEmitManifest } from './pdg-emit-sink.js';
|
||||
import { getNodeLabel as deriveNodeLabel, type WriteStreamFactory } from './rel-pair-routing.js';
|
||||
import type { CachedEmbedding } from '../embeddings/types.js';
|
||||
|
|
@ -876,6 +876,73 @@ const doInitLbug = async (dbPath: string, readOnly: boolean = false) => {
|
|||
|
||||
export type LbugProgressCallback = (message: string) => void;
|
||||
|
||||
/**
|
||||
* Run a COPY, retrying once with IGNORE_ERRORS=true (which skips row-level
|
||||
* errors) on first failure. On a second failure, hand the RAW retry error to
|
||||
* `onError` — each call site formats + slices its own message (#2226 F5: node
|
||||
* COPY slices to 200 chars and throws; relationship COPY slices to 80 and warns,
|
||||
* so the helper must not pre-format and lose that distinction). `onError` may
|
||||
* throw to propagate the failure.
|
||||
*/
|
||||
const copyCsvWithRetry = async (
|
||||
targetConn: lbug.Connection,
|
||||
copyQuery: string,
|
||||
onError: (retryErr: unknown) => void,
|
||||
): Promise<void> => {
|
||||
try {
|
||||
await queryAndDrain(targetConn, copyQuery);
|
||||
} catch {
|
||||
try {
|
||||
const retryQuery = copyQuery.replace(
|
||||
'auto_detect=false)',
|
||||
'auto_detect=false, IGNORE_ERRORS=true)',
|
||||
);
|
||||
await queryAndDrain(targetConn, retryQuery);
|
||||
} catch (retryErr) {
|
||||
onError(retryErr);
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
* Bulk-COPY every node CSV sequentially on the single writable connection
|
||||
* (LadybugDB allows one write txn at a time). Extracted from loadGraphToLbug so
|
||||
* it can run either at the node-phase boundary — overlapping the relationship
|
||||
* emit pass (#2203) — or after emit in the serial escape-hatch path. Each COPY
|
||||
* keeps the IGNORE_ERRORS=true retry; a hard failure throws (no node rows ⇒ the
|
||||
* relationship COPY would dangle on missing endpoints).
|
||||
*/
|
||||
const copyNodeCSVs = async (
|
||||
targetConn: lbug.Connection,
|
||||
nodeFileEntries: [NodeTableName, { csvPath: string; rows: number }][],
|
||||
log: (message: string) => void,
|
||||
totalSteps: number,
|
||||
): Promise<void> => {
|
||||
let stepsDone = 0;
|
||||
for (const [table, { csvPath, rows }] of nodeFileEntries) {
|
||||
stepsDone++;
|
||||
log(`Loading nodes ${stepsDone}/${totalSteps}: ${table} (${rows.toLocaleString()} rows)`);
|
||||
|
||||
const copyQuery = getCopyQuery(table, normalizeCopyPath(csvPath));
|
||||
await copyCsvWithRetry(targetConn, copyQuery, (retryErr) => {
|
||||
const retryMsg = retryErr instanceof Error ? retryErr.message : String(retryErr);
|
||||
throw new Error(`COPY failed for ${table}: ${retryMsg.slice(0, 200)}`);
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
* Persist a KnowledgeGraph: stream CSVs, then bulk-COPY nodes (overlapped with
|
||||
* relationship emit — see the body) and relationships.
|
||||
*
|
||||
* NOT TRANSACTIONAL (#2226). Each `COPY` commits independently and there is no
|
||||
* surrounding transaction, so a failure partway through — a node `COPY` that
|
||||
* throws at the FK barrier, a relationship `COPY` failure, or a `pdgEmitManifest`
|
||||
* collision raised after node rows have already committed in the overlap path —
|
||||
* leaves a partially-loaded DB. The caller surfaces the error; recovery is a
|
||||
* `--force` re-analyze (a full rebuild), not a partial retry. Callers must not
|
||||
* assume the DB is either fully loaded or untouched after a rejection.
|
||||
*/
|
||||
export const loadGraphToLbug = async (
|
||||
graph: KnowledgeGraph,
|
||||
repoPath: string,
|
||||
|
|
@ -904,36 +971,99 @@ export const loadGraphToLbug = async (
|
|||
// the gap that the DB-persistence path is un-timed today (the analyze
|
||||
// "emit" number is the scope-resolution emit bucket, not this COPY path).
|
||||
const PROF = process.env.PROF_LBUG_LOAD === '1';
|
||||
// Escape hatch / differential oracle (#2203): force the legacy strictly-serial
|
||||
// load order (emit everything, THEN COPY nodes, THEN COPY rels) instead of the
|
||||
// default node-COPY ‖ rel-emit overlap. Lets an operator revert the behavior at
|
||||
// runtime, and lets a test load the same graph both ways and assert identical
|
||||
// persisted content.
|
||||
const SERIAL = process.env.GITNEXUS_SERIAL_LBUG_LOAD === '1';
|
||||
const mark = (): bigint => (PROF ? process.hrtime.bigint() : 0n);
|
||||
const span = (a: bigint, b: bigint): string => (Number(b - a) / 1e6).toFixed(1);
|
||||
const tStart = mark();
|
||||
|
||||
const csvDir = resolveNativeSafeStorageDir(storagePath, 'csv');
|
||||
|
||||
log('Streaming CSVs to disk...');
|
||||
const csvResult = await streamAllCSVsToDisk(graph, repoPath, csvDir);
|
||||
// The single writable connection (LadybugDB is single-writer). Captured as a
|
||||
// const so the node-COPY closure has a non-null reference — TS cannot narrow
|
||||
// the reassignable module-level `conn` across the callback boundary.
|
||||
const writeConn = conn;
|
||||
const validTables = new Set<string>(NODE_TABLES as readonly string[]);
|
||||
|
||||
// Merge the streamed PDG-emit CSVs (#2202) into the COPY plan so the
|
||||
// BasicBlock node table + per-pair PDG edges (CFG / REACHING_DEF / CDG /
|
||||
// POST_DOMINATE / TAINTED / SANITIZES) load through the SAME node + per-pair
|
||||
// COPY loops as the structural CSVs. The graph held zero BasicBlocks when
|
||||
// streaming, so `streamAllCSVsToDisk` produced none of these — the manifest
|
||||
// is the sole source and there is no double-COPY. Absent ⇒ no-op.
|
||||
if (pdgEmitManifest) {
|
||||
// Merge the streamed PDG-emit node CSVs (#2202) into a node-file map. Collision
|
||||
// guard: a BasicBlock in the in-memory graph during a streamed run is an
|
||||
// invariant violation (streamAllCSVsToDisk would also emit basicblock.csv), so
|
||||
// fail loudly rather than drop rows (#2202 review #3). Runs at the node-phase
|
||||
// boundary so the manifest BasicBlock table COPYs with the structural CSVs.
|
||||
const mergeManifestNodeFiles = (
|
||||
nodeFilesMap: Map<NodeTableName, { csvPath: string; rows: number }>,
|
||||
): void => {
|
||||
if (!pdgEmitManifest) return;
|
||||
for (const [table, meta] of pdgEmitManifest.nodeFiles) {
|
||||
// A collision means a BasicBlock leaked into the in-memory graph during a
|
||||
// streamed run (streamAllCSVsToDisk then emitted a structural basicblock.csv).
|
||||
// That is a streaming-invariant violation — fail loudly rather than
|
||||
// silently overwrite one CSV with the other and drop its rows (#2202 review #3).
|
||||
if (csvResult.nodeFiles.has(table)) {
|
||||
if (nodeFilesMap.has(table)) {
|
||||
throw new Error(
|
||||
`Streaming PDG manifest collides with a structural node CSV for "${table}" — ` +
|
||||
`the in-memory graph should hold zero ${table} nodes when streaming. ` +
|
||||
`A ${table} node leaked into the graph during a streamed emit.`,
|
||||
);
|
||||
}
|
||||
csvResult.nodeFiles.set(table, meta);
|
||||
nodeFilesMap.set(table, meta);
|
||||
}
|
||||
};
|
||||
|
||||
// Node COPY is the only DB write that can overlap relationship CSV emit: the
|
||||
// rel pass writes new rel_*.csv files and never touches `conn`, while node COPY
|
||||
// uses `conn` and never touches the rel files. We start node COPY at the
|
||||
// node-phase boundary and let the rel pass run concurrently — the only
|
||||
// single-writer-safe parallelism (#2203). The rel COPY still waits for node
|
||||
// COPY (FK precondition), so the DB load order is unchanged.
|
||||
let nodeCopyPromise: Promise<void> | undefined;
|
||||
let nodeCopyError: unknown;
|
||||
const beginNodeCopy = (
|
||||
nodeFilesMap: Map<NodeTableName, { csvPath: string; rows: number }>,
|
||||
): void => {
|
||||
mergeManifestNodeFiles(nodeFilesMap);
|
||||
const entries = [...nodeFilesMap.entries()];
|
||||
// copyNodeCSVs logs node progress as step/total; it processes only node
|
||||
// tables (the rel COPY has its own "Loading edges" progress line), so the
|
||||
// denominator is the node-table count — not +1 reserving a rel step.
|
||||
// .catch captures the failure so an overlapped (mid-emit) rejection cannot
|
||||
// surface as an unhandled rejection; it is rethrown at the FK barrier below.
|
||||
nodeCopyPromise = copyNodeCSVs(writeConn, entries, log, entries.length).catch((e) => {
|
||||
nodeCopyError = e;
|
||||
});
|
||||
};
|
||||
|
||||
log('Streaming CSVs to disk...');
|
||||
let csvResult: StreamedCSVResult;
|
||||
try {
|
||||
csvResult = SERIAL
|
||||
? await streamAllCSVsToDisk(graph, repoPath, csvDir)
|
||||
: await streamAllCSVsToDisk(graph, repoPath, csvDir, beginNodeCopy);
|
||||
} catch (emitErr) {
|
||||
// Relationship emit failed. In overlap mode a node COPY may be in flight —
|
||||
// settle it (the .catch above means this never rejects) before rethrowing so
|
||||
// it cannot leak as an unhandled rejection.
|
||||
if (nodeCopyPromise) await nodeCopyPromise;
|
||||
// If node COPY ALSO failed, emitErr wins the throw — log the swallowed node
|
||||
// error so a half-loaded DB isn't misattributed to the emit failure alone.
|
||||
if (nodeCopyError) {
|
||||
logger.warn(
|
||||
{ err: nodeCopyError },
|
||||
'[lbug-load] node COPY also failed while relationship emit was failing',
|
||||
);
|
||||
}
|
||||
throw emitErr;
|
||||
}
|
||||
const tCsv = mark();
|
||||
|
||||
// Merge the streamed PDG-emit per-pair rel CSVs (#2202) into the COPY plan —
|
||||
// collision-guarded. Done BEFORE node COPY so the serial escape hatch detects a
|
||||
// manifest/structural pair collision before committing any node rows (legacy
|
||||
// parity with the pre-overlap path), and the overlap path detects it as early
|
||||
// as csvResult is available. When a manifest is present, streaming was on and
|
||||
// the in-memory graph held zero BasicBlocks, so a structural collision means a
|
||||
// streaming-invariant violation — fail loudly rather than load corrupt data.
|
||||
if (pdgEmitManifest) {
|
||||
for (const [pairKey, meta] of pdgEmitManifest.relsByPair) {
|
||||
if (csvResult.relsByPair.has(pairKey)) {
|
||||
throw new Error(
|
||||
|
|
@ -945,38 +1075,18 @@ export const loadGraphToLbug = async (
|
|||
csvResult.totalValidRels += meta.rows;
|
||||
}
|
||||
}
|
||||
const tCsv = mark();
|
||||
|
||||
const validTables = new Set<string>(NODE_TABLES as readonly string[]);
|
||||
// Serial path: all CSVs are on disk and node COPY has not started — start it
|
||||
// here so the barrier below blocks on it exactly as the legacy path did.
|
||||
if (SERIAL) beginNodeCopy(csvResult.nodeFiles);
|
||||
|
||||
// Bulk COPY all node CSVs (sequential — LadybugDB allows only one write txn at a time)
|
||||
const nodeFiles = [...csvResult.nodeFiles.entries()];
|
||||
const totalSteps = nodeFiles.length + 1; // +1 for relationships
|
||||
let stepsDone = 0;
|
||||
|
||||
for (const [table, { csvPath, rows }] of nodeFiles) {
|
||||
stepsDone++;
|
||||
log(`Loading nodes ${stepsDone}/${totalSteps}: ${table} (${rows.toLocaleString()} rows)`);
|
||||
|
||||
const normalizedPath = normalizeCopyPath(csvPath);
|
||||
const copyQuery = getCopyQuery(table, normalizedPath);
|
||||
|
||||
try {
|
||||
await queryAndDrain(conn, copyQuery);
|
||||
} catch (err) {
|
||||
try {
|
||||
const retryQuery = copyQuery.replace(
|
||||
'auto_detect=false)',
|
||||
'auto_detect=false, IGNORE_ERRORS=true)',
|
||||
);
|
||||
await queryAndDrain(conn, retryQuery);
|
||||
} catch (retryErr) {
|
||||
const retryMsg = retryErr instanceof Error ? retryErr.message : String(retryErr);
|
||||
throw new Error(`COPY failed for ${table}: ${retryMsg.slice(0, 200)}`);
|
||||
}
|
||||
}
|
||||
// FK barrier: node rows must exist before the relationship COPY resolves their
|
||||
// endpoints. In overlap mode most of node COPY was hidden behind rel emit, so
|
||||
// this await is the *residual* node-COPY time (≈0 when fully overlapped).
|
||||
if (nodeCopyPromise) await nodeCopyPromise;
|
||||
if (nodeCopyError) {
|
||||
throw nodeCopyError instanceof Error ? nodeCopyError : new Error(String(nodeCopyError));
|
||||
}
|
||||
|
||||
const tCopyNodes = mark();
|
||||
|
||||
// Bulk COPY relationships. They were already routed to per-FROM→TO-label-pair
|
||||
|
|
@ -999,28 +1109,19 @@ export const loadGraphToLbug = async (
|
|||
pairIdx++;
|
||||
const [fromLabel, toLabel] = pairKey.split('|');
|
||||
const normalizedPath = normalizeCopyPath(pairCsvPath);
|
||||
// PARALLEL=false is load-bearing here too — see COPY_CSV_OPTS (#2203 / kuzudb/kuzu#5778).
|
||||
const copyQuery = `COPY ${REL_TABLE_NAME} FROM "${normalizedPath}" (from="${fromLabel}", to="${toLabel}", HEADER=true, ESCAPE='"', DELIM=',', QUOTE='"', PARALLEL=false, auto_detect=false)`;
|
||||
|
||||
if (pairIdx % 5 === 0 || rows > 1000) {
|
||||
log(`Loading edges: ${pairIdx}/${relsByPair.size} types (${fromLabel} -> ${toLabel})`);
|
||||
}
|
||||
|
||||
try {
|
||||
await queryAndDrain(conn, copyQuery);
|
||||
} catch (err) {
|
||||
try {
|
||||
const retryQuery = copyQuery.replace(
|
||||
'auto_detect=false)',
|
||||
'auto_detect=false, IGNORE_ERRORS=true)',
|
||||
);
|
||||
await queryAndDrain(conn, retryQuery);
|
||||
} catch (retryErr) {
|
||||
const retryMsg = retryErr instanceof Error ? retryErr.message : String(retryErr);
|
||||
warnings.push(`${fromLabel}->${toLabel} (${rows} edges): ${retryMsg.slice(0, 80)}`);
|
||||
failedPairEdges += rows;
|
||||
failedPairCsvPaths.add(pairCsvPath);
|
||||
}
|
||||
}
|
||||
await copyCsvWithRetry(conn, copyQuery, (retryErr) => {
|
||||
const retryMsg = retryErr instanceof Error ? retryErr.message : String(retryErr);
|
||||
warnings.push(`${fromLabel}->${toLabel} (${rows} edges): ${retryMsg.slice(0, 80)}`);
|
||||
failedPairEdges += rows;
|
||||
failedPairCsvPaths.add(pairCsvPath);
|
||||
});
|
||||
// Only delete if not in failedPairCsvPaths (needed for fallback)
|
||||
if (!failedPairCsvPaths.has(pairCsvPath)) {
|
||||
try {
|
||||
|
|
@ -1077,8 +1178,13 @@ export const loadGraphToLbug = async (
|
|||
const tEnd = mark();
|
||||
let totalNodeRows = 0;
|
||||
for (const [, { rows }] of csvResult.nodeFiles) totalNodeRows += rows;
|
||||
// `mode` records which load path ran. In overlap mode `csv-emit` is the wall
|
||||
// to streamAllCSVsToDisk's return (node COPY overlapped part of it) and
|
||||
// `copy-nodes` is the RESIDUAL node-COPY await after emit returned — it
|
||||
// trends to 0 as the overlap hides node COPY behind relationship emit. In
|
||||
// serial mode the buckets carry their legacy, disjoint meaning.
|
||||
logger.warn(
|
||||
`[lbug-load prof] csv-emit=${span(tStart, tCsv)}ms ` +
|
||||
`[lbug-load prof] mode=${SERIAL ? 'serial' : 'overlap'} csv-emit=${span(tStart, tCsv)}ms ` +
|
||||
`copy-nodes=${span(tCsv, tCopyNodes)}ms copy-rels=${span(tCopyNodes, tCopyRels)}ms ` +
|
||||
`fallback=${span(tCopyRels, tFallback)}ms total=${span(tStart, tEnd)}ms ` +
|
||||
`(${totalNodeRows} nodes, ${insertedRels} rels)`,
|
||||
|
|
@ -1092,7 +1198,18 @@ export const loadGraphToLbug = async (
|
|||
// Source code content is full of backslashes which confuse the auto-detection.
|
||||
// We MUST explicitly set ESCAPE='"' to use RFC 4180 escaping, and disable auto_detect to prevent
|
||||
// LadybugDB from overriding our settings based on sample rows.
|
||||
const COPY_CSV_OPTS = `(HEADER=true, ESCAPE='"', DELIM=',', QUOTE='"', PARALLEL=false, auto_detect=false)`;
|
||||
//
|
||||
// PARALLEL=false IS LOAD-BEARING FOR CORRECTNESS — DO NOT FLIP IT (#2203).
|
||||
// LadybugDB's parallel CSV reader (Kuzu-derived; default PARALLEL=true) splits the
|
||||
// file into byte ranges parsed concurrently, and CANNOT determine line boundaries
|
||||
// when a quoted field contains an embedded newline — it errors with "Quoted newlines
|
||||
// are not supported in parallel CSV reader. Please specify PARALLEL=FALSE", or worse,
|
||||
// mis-parses silently (upstream kuzudb/kuzu#5778, still open). Our `content`/`text`
|
||||
// columns hold source code, so quoted multiline fields are guaranteed. PARALLEL=false
|
||||
// is therefore required, not conservative. The multiline-quoted round-trip in
|
||||
// test/integration/copy-parallel-invariant.test.ts fails loudly if this is ever flipped.
|
||||
// Exported so that test asserts the invariant statically as well.
|
||||
export const COPY_CSV_OPTS = `(HEADER=true, ESCAPE='"', DELIM=',', QUOTE='"', PARALLEL=false, auto_detect=false)`;
|
||||
|
||||
// Multi-language table names that were created with backticks in CODE_ELEMENT_BASE
|
||||
// and must always be referenced with backticks in queries
|
||||
|
|
@ -1170,7 +1287,7 @@ const TABLES_WITH_EXPORTED = new Set<string>([
|
|||
'CodeElement',
|
||||
]);
|
||||
|
||||
const getCopyQuery = (table: NodeTableName, filePath: string): string => {
|
||||
export const getCopyQuery = (table: NodeTableName, filePath: string): string => {
|
||||
const t = escapeTableName(table);
|
||||
if (table === 'File') {
|
||||
return `COPY ${t}(id, name, filePath, content) FROM "${filePath}" ${COPY_CSV_OPTS}`;
|
||||
|
|
|
|||
130
gitnexus/test/integration/copy-parallel-invariant.test.ts
Normal file
130
gitnexus/test/integration/copy-parallel-invariant.test.ts
Normal file
|
|
@ -0,0 +1,130 @@
|
|||
/**
|
||||
* Integration test: PARALLEL=false is a load-bearing correctness invariant for
|
||||
* the bulk-COPY persistence path (#2203 — the "parallelized emit" follow-up).
|
||||
*
|
||||
* Issue #2203 asked us to investigate parallelizing emit. The most obvious
|
||||
* lever — LadybugDB's intra-COPY parallel CSV reader (`PARALLEL=true`, the
|
||||
* Kuzu default) — is UNSAFE for our data: that reader splits the file into byte
|
||||
* ranges parsed concurrently and cannot find line boundaries when a quoted
|
||||
* field contains an embedded newline (upstream kuzudb/kuzu#5778, still open;
|
||||
* error text "Quoted newlines are not supported in parallel CSV reader. Please
|
||||
* specify PARALLEL=FALSE"). Our `content`/`text` columns hold source code, so
|
||||
* quoted multiline fields are guaranteed. PARALLEL=false is therefore MANDATORY,
|
||||
* not conservative — this test locks that in two ways:
|
||||
*
|
||||
* 1. Live-DB proof — a node whose `text` carries embedded newlines AND quotes
|
||||
* round-trips byte-exact through the real csv-emit → COPY → query path.
|
||||
* If anyone flips PARALLEL=true, the parallel reader mis-parses this row and
|
||||
* the assertion fails loudly. (An edge also round-trips, exercising the rel
|
||||
* COPY path, which uses the same PARALLEL=false option.)
|
||||
* 2. Static guard — the generated COPY query strings still carry PARALLEL=false,
|
||||
* giving a crisp failure independent of a live DB.
|
||||
*
|
||||
* Needs a real LadybugDB connection (initLbug), so it lives under integration.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, afterAll } from 'vitest';
|
||||
import fs from 'fs/promises';
|
||||
import path from 'path';
|
||||
import os from 'os';
|
||||
import { NODE_TABLES } from 'gitnexus-shared';
|
||||
import { buildTestGraph } from '../helpers/test-graph.js';
|
||||
|
||||
let tmpBase: string;
|
||||
let storagePath: string;
|
||||
let dbPath: string;
|
||||
|
||||
// A BasicBlock `text` with the exact hazard the parallel reader cannot handle:
|
||||
// embedded newlines INSIDE a field that also contains double-quotes. After
|
||||
// escapeCSVField this becomes a quoted multiline CSV field.
|
||||
const HAZARD_TEXT = 'const msg = "line one";\nconst other = "she said \\"hi\\"";\nreturn msg;';
|
||||
const BB1 = 'BasicBlock:src/hazard.ts:0';
|
||||
const BB2 = 'BasicBlock:src/hazard.ts:1';
|
||||
|
||||
beforeAll(async () => {
|
||||
// mkdtemp (not a predictable os.tmpdir join) + the `gitnexus-lbug-` prefix
|
||||
// that TEST_FIXTURE_PREFIXES recognizes for the stale-sidecar sweep.
|
||||
tmpBase = await fs.mkdtemp(path.join(os.tmpdir(), 'gitnexus-lbug-copy-parallel-'));
|
||||
storagePath = path.join(tmpBase, '.gitnexus');
|
||||
dbPath = path.join(storagePath, 'lbug');
|
||||
await fs.mkdir(dbPath, { recursive: true });
|
||||
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
await adapter.initLbug(dbPath);
|
||||
|
||||
const graph = buildTestGraph(
|
||||
[
|
||||
{
|
||||
id: BB1,
|
||||
label: 'BasicBlock',
|
||||
name: '',
|
||||
filePath: 'src/hazard.ts',
|
||||
startLine: 1,
|
||||
endLine: 3,
|
||||
extra: { text: HAZARD_TEXT },
|
||||
},
|
||||
{
|
||||
id: BB2,
|
||||
label: 'BasicBlock',
|
||||
name: '',
|
||||
filePath: 'src/hazard.ts',
|
||||
startLine: 4,
|
||||
endLine: 4,
|
||||
extra: { text: 'sink(msg);' },
|
||||
},
|
||||
],
|
||||
[{ sourceId: BB1, targetId: BB2, type: 'CFG', reason: 'cfg-edge' }],
|
||||
);
|
||||
|
||||
await adapter.loadGraphToLbug(graph, tmpBase, storagePath);
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
try {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
await adapter.closeLbug();
|
||||
} catch {
|
||||
/* may not have opened */
|
||||
}
|
||||
if (tmpBase) {
|
||||
for (let attempt = 0; attempt < 5; attempt++) {
|
||||
try {
|
||||
await fs.rm(tmpBase, { recursive: true, force: true });
|
||||
return;
|
||||
} catch {
|
||||
if (attempt < 4) await new Promise((r) => setTimeout(r, 200 * (attempt + 1)));
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
describe('PARALLEL=false correctness invariant (#2203 / kuzudb/kuzu#5778)', () => {
|
||||
it('a multiline-quoted content field round-trips byte-exact through COPY', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const rows = await adapter.executeQuery(
|
||||
`MATCH (n:BasicBlock {id: '${BB1}'}) RETURN n.text AS text`,
|
||||
);
|
||||
expect(rows).toHaveLength(1);
|
||||
// Byte-exact: the embedded newlines and the doubled quotes survived the
|
||||
// quoted-field round-trip. PARALLEL=true would have mis-split this row.
|
||||
expect(rows[0].text).toBe(HAZARD_TEXT);
|
||||
});
|
||||
|
||||
it('an edge round-trips through the relationship COPY path (same PARALLEL=false)', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const rows = await adapter.executeQuery(
|
||||
"MATCH (:BasicBlock)-[r:CodeRelation {type: 'CFG'}]->(:BasicBlock) RETURN count(r) AS c",
|
||||
);
|
||||
expect(Number(rows[0].c)).toBe(1);
|
||||
});
|
||||
|
||||
it('every generated node COPY query statically carries PARALLEL=false', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
expect(adapter.COPY_CSV_OPTS).toContain('PARALLEL=false');
|
||||
expect(adapter.COPY_CSV_OPTS).not.toContain('PARALLEL=true');
|
||||
for (const table of NODE_TABLES) {
|
||||
const q = adapter.getCopyQuery(table, '/tmp/x.csv');
|
||||
expect(q, `${table} COPY must keep PARALLEL=false`).toContain('PARALLEL=false');
|
||||
expect(q, `${table} COPY must not enable PARALLEL`).not.toContain('PARALLEL=true');
|
||||
}
|
||||
});
|
||||
});
|
||||
|
|
@ -6,6 +6,7 @@
|
|||
*/
|
||||
import { describe, it, expect, beforeAll, beforeEach, afterAll } from 'vitest';
|
||||
import fs from 'fs/promises';
|
||||
import { readdirSync } from 'node:fs';
|
||||
import { finished } from 'stream/promises';
|
||||
import path from 'path';
|
||||
import { createTempDir, type TestDBHandle } from '../helpers/test-db.js';
|
||||
|
|
@ -565,3 +566,69 @@ describe('streamAllCSVsToDisk — direct per-pair emit matches the split oracle'
|
|||
}
|
||||
});
|
||||
});
|
||||
|
||||
// The overlap leg (#2203) needs to start COPY-ing nodes while relationship CSVs
|
||||
// are still being written. streamAllCSVsToDisk exposes that boundary via an
|
||||
// onNodePhaseComplete callback. These tests pin the contract: it fires once,
|
||||
// after node CSVs exist and before any rel CSV does, and supplying it does not
|
||||
// change the emitted output.
|
||||
describe('onNodePhaseComplete hook (#2203 overlap boundary)', () => {
|
||||
const hookGraph = () =>
|
||||
buildTestGraph(
|
||||
[
|
||||
{ id: 'File:src/index.ts', label: 'File', name: 'index.ts', filePath: 'src/index.ts' },
|
||||
{
|
||||
id: 'Function:src/index.ts:main:1',
|
||||
label: 'Function',
|
||||
name: 'main',
|
||||
filePath: 'src/index.ts',
|
||||
startLine: 1,
|
||||
endLine: 3,
|
||||
},
|
||||
],
|
||||
[
|
||||
{
|
||||
sourceId: 'File:src/index.ts',
|
||||
targetId: 'Function:src/index.ts:main:1',
|
||||
type: 'DEFINES',
|
||||
},
|
||||
],
|
||||
);
|
||||
|
||||
it('fires exactly once, after node CSVs are flushed and before any rel CSV exists', async () => {
|
||||
const hookCsvDir = path.join(tmpHandle.dbPath, 'csv-hook-timing');
|
||||
let calls = 0;
|
||||
let nodeCsvsPresent = false;
|
||||
let relCsvsPresent = true;
|
||||
let handedKeys: string[] = [];
|
||||
|
||||
const result = await streamAllCSVsToDisk(hookGraph(), repoDir, hookCsvDir, (nodeFiles) => {
|
||||
calls++;
|
||||
handedKeys = [...nodeFiles.keys()].sort();
|
||||
const entries = readdirSync(hookCsvDir);
|
||||
nodeCsvsPresent = entries.includes('file.csv') && entries.includes('function.csv');
|
||||
relCsvsPresent = entries.some((f) => f.startsWith('rel_'));
|
||||
});
|
||||
|
||||
expect(calls).toBe(1);
|
||||
expect(nodeCsvsPresent).toBe(true);
|
||||
// No relationship CSV may exist yet — the rel pass starts after the hook.
|
||||
expect(relCsvsPresent).toBe(false);
|
||||
// The manifest handed to the callback is the one returned to the caller.
|
||||
expect(handedKeys).toEqual([...result.nodeFiles.keys()].sort());
|
||||
expect(result.totalValidRels).toBe(1);
|
||||
});
|
||||
|
||||
it('supplying the callback does not change the node manifest (no behavior change)', async () => {
|
||||
const withDir = path.join(tmpHandle.dbPath, 'csv-hook-with');
|
||||
const withoutDir = path.join(tmpHandle.dbPath, 'csv-hook-without');
|
||||
const withCb = await streamAllCSVsToDisk(hookGraph(), repoDir, withDir, () => {});
|
||||
const without = await streamAllCSVsToDisk(hookGraph(), repoDir, withoutDir);
|
||||
|
||||
const manifest = (r: typeof withCb) =>
|
||||
[...r.nodeFiles.entries()].map(([k, v]) => `${k}:${v.rows}`).sort();
|
||||
expect(manifest(withCb)).toEqual(manifest(without));
|
||||
expect(withCb.totalValidRels).toBe(without.totalValidRels);
|
||||
expect(withCb.skippedRels).toBe(without.skippedRels);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
211
gitnexus/test/integration/lbug-load-overlap-errors.test.ts
Normal file
211
gitnexus/test/integration/lbug-load-overlap-errors.test.ts
Normal file
|
|
@ -0,0 +1,211 @@
|
|||
/**
|
||||
* Integration test: loadGraphToLbug overlap ERROR paths (#2226 review F1/F2).
|
||||
*
|
||||
* The node-COPY ‖ relationship-emit overlap (#2203) added two error branches the
|
||||
* happy-path tests don't exercise:
|
||||
* - F1: relationship emit fails while node COPY is in flight — the emit error
|
||||
* must surface from loadGraphToLbug AND the in-flight node-COPY promise must
|
||||
* be settled, so nothing leaks as an unhandled rejection.
|
||||
* - F2: node COPY itself fails — the captured error must be rethrown at the FK
|
||||
* barrier rather than silently swallowed.
|
||||
*
|
||||
* Both are fault-injected by mocking `streamAllCSVsToDisk` so the test controls
|
||||
* when `onNodePhaseComplete` fires (starting node COPY) and whether emit throws,
|
||||
* while a REAL LadybugDB connection (`initLbug`) runs the node COPY. The mock
|
||||
* pattern (vi.hoisted + vi.mock with importOriginal, preserving every other
|
||||
* export) mirrors test/unit/api-graph-streaming.test.ts.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, afterAll, afterEach, vi } from 'vitest';
|
||||
import fs from 'fs/promises';
|
||||
import path from 'path';
|
||||
import os from 'os';
|
||||
import type { StreamedCSVResult } from '../../src/core/lbug/csv-generator.js';
|
||||
import type { NodeTableName } from '../../src/core/lbug/schema.js';
|
||||
import { buildTestGraph } from '../helpers/test-graph.js';
|
||||
|
||||
const { emitMock } = vi.hoisted(() => ({ emitMock: vi.fn() }));
|
||||
|
||||
// Replace ONLY streamAllCSVsToDisk; importOriginal preserves StreamedCSVResult
|
||||
// and every other export lbug-adapter (and its transitive deps) rely on.
|
||||
vi.mock('../../src/core/lbug/csv-generator.js', async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import('../../src/core/lbug/csv-generator.js')>();
|
||||
return { ...actual, streamAllCSVsToDisk: emitMock };
|
||||
});
|
||||
|
||||
type NodeFiles = Map<NodeTableName, { csvPath: string; rows: number }>;
|
||||
|
||||
const REL_HEADER = 'from,to,type,confidence,reason,step';
|
||||
const emptyResult = (): StreamedCSVResult => ({
|
||||
nodeFiles: new Map(),
|
||||
relsByPair: new Map(),
|
||||
relHeader: REL_HEADER,
|
||||
skippedRels: 0,
|
||||
totalValidRels: 0,
|
||||
});
|
||||
|
||||
let tmpBase: string;
|
||||
let storagePath: string;
|
||||
let dbPath: string;
|
||||
|
||||
beforeAll(async () => {
|
||||
tmpBase = await fs.mkdtemp(path.join(os.tmpdir(), 'gitnexus-lbug-overlap-err-'));
|
||||
storagePath = path.join(tmpBase, '.gitnexus');
|
||||
dbPath = path.join(storagePath, 'lbug');
|
||||
await fs.mkdir(dbPath, { recursive: true });
|
||||
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
await adapter.initLbug(dbPath);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
emitMock.mockReset();
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
try {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
await adapter.closeLbug();
|
||||
} catch {
|
||||
/* may already be closed */
|
||||
}
|
||||
if (tmpBase) {
|
||||
for (let attempt = 0; attempt < 5; attempt++) {
|
||||
try {
|
||||
await fs.rm(tmpBase, { recursive: true, force: true });
|
||||
return;
|
||||
} catch {
|
||||
if (attempt < 4) await new Promise((r) => setTimeout(r, 200 * (attempt + 1)));
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
describe('loadGraphToLbug overlap error paths (#2226 F1)', () => {
|
||||
it('relationship-emit failure with node COPY in flight surfaces the emit error and leaks no unhandled rejection', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const graph = buildTestGraph(
|
||||
[{ id: 'File:src/u1.ts', label: 'File', name: 'u1.ts', filePath: 'src/u1.ts' }],
|
||||
[],
|
||||
);
|
||||
|
||||
// emit writes a REAL node CSV (so node COPY actually runs to completion in
|
||||
// flight), fires the hook to start node COPY, then throws as if the
|
||||
// relationship pass blew up (EMFILE / disk-full).
|
||||
emitMock.mockImplementation(
|
||||
async (
|
||||
_g: unknown,
|
||||
_r: unknown,
|
||||
dir: string,
|
||||
onNodePhaseComplete?: (n: NodeFiles) => void,
|
||||
) => {
|
||||
await fs.mkdir(dir, { recursive: true }); // mocked emit skips the real csvDir creation
|
||||
const csvPath = path.join(dir, 'file.csv');
|
||||
await fs.writeFile(
|
||||
csvPath,
|
||||
'id,name,filePath,content\n"File:src/u1.ts","u1.ts","src/u1.ts",""\n',
|
||||
);
|
||||
onNodePhaseComplete?.(new Map([['File', { csvPath, rows: 1 }]]) as NodeFiles);
|
||||
throw new Error('simulated rel emit failure');
|
||||
},
|
||||
);
|
||||
|
||||
const rejections: unknown[] = [];
|
||||
const onUnhandled = (e: unknown): void => {
|
||||
rejections.push(e);
|
||||
};
|
||||
process.on('unhandledRejection', onUnhandled);
|
||||
try {
|
||||
await expect(adapter.loadGraphToLbug(graph, tmpBase, storagePath)).rejects.toThrow(
|
||||
'simulated rel emit failure',
|
||||
);
|
||||
// Flush the macrotask queue so any stray rejection from the in-flight
|
||||
// node-COPY promise actually surfaces BEFORE we assert — otherwise the
|
||||
// check passes vacuously.
|
||||
await new Promise((r) => setImmediate(r));
|
||||
expect(rejections).toEqual([]);
|
||||
} finally {
|
||||
// forks reuse the process across this file's tests — a leaked listener
|
||||
// would corrupt sibling tests.
|
||||
process.off('unhandledRejection', onUnhandled);
|
||||
}
|
||||
});
|
||||
|
||||
it('also logs the node-COPY error when BOTH emit and node COPY fail (emit error still wins)', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const { _captureLogger } = await import('../../src/core/logger.js');
|
||||
const graph = buildTestGraph(
|
||||
[{ id: 'File:src/u1b.ts', label: 'File', name: 'u1b.ts', filePath: 'src/u1b.ts' }],
|
||||
[],
|
||||
);
|
||||
|
||||
// nodeFiles points at a MISSING csv (node COPY fails) AND emit throws.
|
||||
emitMock.mockImplementation(
|
||||
async (
|
||||
_g: unknown,
|
||||
_r: unknown,
|
||||
dir: string,
|
||||
onNodePhaseComplete?: (n: NodeFiles) => void,
|
||||
) => {
|
||||
onNodePhaseComplete?.(
|
||||
new Map([['File', { csvPath: path.join(dir, 'missing-u1b.csv'), rows: 1 }]]) as NodeFiles,
|
||||
);
|
||||
throw new Error('simulated rel emit failure (double)');
|
||||
},
|
||||
);
|
||||
|
||||
const cap = _captureLogger();
|
||||
const rejections: unknown[] = [];
|
||||
const onUnhandled = (e: unknown): void => {
|
||||
rejections.push(e);
|
||||
};
|
||||
process.on('unhandledRejection', onUnhandled);
|
||||
try {
|
||||
await expect(adapter.loadGraphToLbug(graph, tmpBase, storagePath)).rejects.toThrow(
|
||||
'simulated rel emit failure (double)',
|
||||
);
|
||||
await new Promise((r) => setImmediate(r));
|
||||
expect(rejections).toEqual([]);
|
||||
const warned = cap
|
||||
.records()
|
||||
.map((rec) => (typeof rec.msg === 'string' ? rec.msg : ''))
|
||||
.some((m) => m.includes('node COPY also failed'));
|
||||
expect(warned).toBe(true);
|
||||
} finally {
|
||||
cap.restore();
|
||||
process.off('unhandledRejection', onUnhandled);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe('loadGraphToLbug overlap error paths (#2226 F2)', () => {
|
||||
it('a node-COPY hard failure is rethrown at the FK barrier', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const graph = buildTestGraph(
|
||||
[{ id: 'File:src/u2.ts', label: 'File', name: 'u2.ts', filePath: 'src/u2.ts' }],
|
||||
[],
|
||||
);
|
||||
|
||||
// Node COPY targets a MISSING csv → COPY fails at bind time ("No file
|
||||
// found …"), which IGNORE_ERRORS does NOT suppress (it only skips row-level
|
||||
// errors), so copyNodeCSVs throws. Emit otherwise "succeeds" (returns an
|
||||
// empty result), so the only failure is the node COPY captured in
|
||||
// nodeCopyError and rethrown at the FK barrier.
|
||||
emitMock.mockImplementation(
|
||||
async (
|
||||
_g: unknown,
|
||||
_r: unknown,
|
||||
dir: string,
|
||||
onNodePhaseComplete?: (n: NodeFiles) => void,
|
||||
) => {
|
||||
onNodePhaseComplete?.(
|
||||
new Map([['File', { csvPath: path.join(dir, 'missing-u2.csv'), rows: 1 }]]) as NodeFiles,
|
||||
);
|
||||
return emptyResult();
|
||||
},
|
||||
);
|
||||
|
||||
await expect(adapter.loadGraphToLbug(graph, tmpBase, storagePath)).rejects.toThrow(
|
||||
/COPY failed for File/,
|
||||
);
|
||||
});
|
||||
});
|
||||
206
gitnexus/test/integration/lbug-load-overlap.test.ts
Normal file
206
gitnexus/test/integration/lbug-load-overlap.test.ts
Normal file
|
|
@ -0,0 +1,206 @@
|
|||
/**
|
||||
* Integration test: the node-COPY ‖ relationship-emit overlap (#2203) persists
|
||||
* BYTE-IDENTICAL graph content to the legacy strictly-serial path.
|
||||
*
|
||||
* This is the acceptance gate for the #2203 "parallelized emit" follow-up: the
|
||||
* overlap reorders *scheduling* (node COPY runs while relationship CSVs are
|
||||
* still being written), never the data. We prove that by loading one fixture
|
||||
* graph into two fresh DBs — once via the default overlap path, once via the
|
||||
* GITNEXUS_SERIAL_LBUG_LOAD=1 escape hatch — and asserting the two databases are
|
||||
* content-equivalent: same per-table node counts, same per-type edge counts,
|
||||
* same multiline `content`/`text` field bytes, and the same loadGraphToLbug
|
||||
* return accounting.
|
||||
*
|
||||
* Needs a real LadybugDB connection (initLbug), so it lives under integration.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, afterAll } from 'vitest';
|
||||
import fs from 'fs/promises';
|
||||
import path from 'path';
|
||||
import os from 'os';
|
||||
import { buildTestGraph } from '../helpers/test-graph.js';
|
||||
|
||||
let tmpBase: string;
|
||||
let repoDir: string;
|
||||
|
||||
// Source-code-shaped content with embedded newlines AND quotes — the exact
|
||||
// shape PARALLEL=false exists to handle. Read from disk for the File node and
|
||||
// carried inline for the BasicBlock node, so both content sources are checked.
|
||||
const FILE_SRC = 'export function f() {\n const s = "a,b\\"c";\n return s;\n}\n';
|
||||
const BB_TEXT = 'if (cond) {\n log("x = " + x);\n}\nreturn "done";';
|
||||
|
||||
const BB1 = 'BasicBlock:src/a.ts:0';
|
||||
const BB2 = 'BasicBlock:src/a.ts:1';
|
||||
|
||||
const NODE_TABLES_CHECKED = ['File', 'Function', 'Class', 'BasicBlock'] as const;
|
||||
const EDGE_TYPES_CHECKED = ['DEFINES', 'CALLS', 'CFG', 'REACHING_DEF'] as const;
|
||||
|
||||
const buildFixture = () =>
|
||||
buildTestGraph(
|
||||
[
|
||||
{ id: 'File:src/a.ts', label: 'File', name: 'a.ts', filePath: 'src/a.ts' },
|
||||
{
|
||||
id: 'Function:src/a.ts:f:1',
|
||||
label: 'Function',
|
||||
name: 'f',
|
||||
filePath: 'src/a.ts',
|
||||
startLine: 1,
|
||||
endLine: 4,
|
||||
isExported: true,
|
||||
},
|
||||
{
|
||||
id: 'Class:src/a.ts:C:6',
|
||||
label: 'Class',
|
||||
name: 'C',
|
||||
filePath: 'src/a.ts',
|
||||
startLine: 6,
|
||||
endLine: 8,
|
||||
},
|
||||
{
|
||||
id: BB1,
|
||||
label: 'BasicBlock',
|
||||
name: '',
|
||||
filePath: 'src/a.ts',
|
||||
startLine: 1,
|
||||
endLine: 3,
|
||||
extra: { text: BB_TEXT },
|
||||
},
|
||||
{ id: BB2, label: 'BasicBlock', name: '', filePath: 'src/a.ts', startLine: 4, endLine: 4 },
|
||||
],
|
||||
[
|
||||
{ sourceId: 'File:src/a.ts', targetId: 'Function:src/a.ts:f:1', type: 'DEFINES' },
|
||||
{ sourceId: 'File:src/a.ts', targetId: 'Class:src/a.ts:C:6', type: 'DEFINES' },
|
||||
{ sourceId: 'Function:src/a.ts:f:1', targetId: 'Class:src/a.ts:C:6', type: 'CALLS' },
|
||||
{ sourceId: BB1, targetId: BB2, type: 'CFG', reason: 'cfg' },
|
||||
{ sourceId: BB1, targetId: BB2, type: 'REACHING_DEF', reason: 's' },
|
||||
],
|
||||
);
|
||||
|
||||
interface Snapshot {
|
||||
nodeCounts: Record<string, number>;
|
||||
edgeCounts: Record<string, number>;
|
||||
bbText: string | undefined;
|
||||
fileContent: string | undefined;
|
||||
ret: { insertedRels: number; skippedRels: number; warnings: string[] };
|
||||
}
|
||||
|
||||
const snapshotDb = async (
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
adapter: any,
|
||||
ret: Snapshot['ret'],
|
||||
): Promise<Snapshot> => {
|
||||
const nodeCounts: Record<string, number> = {};
|
||||
for (const t of NODE_TABLES_CHECKED) {
|
||||
const rows = await adapter.executeQuery(`MATCH (n:\`${t}\`) RETURN count(n) AS c`);
|
||||
nodeCounts[t] = Number(rows[0].c);
|
||||
}
|
||||
const edgeCounts: Record<string, number> = {};
|
||||
for (const ty of EDGE_TYPES_CHECKED) {
|
||||
const rows = await adapter.executeQuery(
|
||||
`MATCH ()-[r:CodeRelation {type: '${ty}'}]->() RETURN count(r) AS c`,
|
||||
);
|
||||
edgeCounts[ty] = Number(rows[0].c);
|
||||
}
|
||||
const bbRows = await adapter.executeQuery(
|
||||
`MATCH (n:BasicBlock {id: '${BB1}'}) RETURN n.text AS text`,
|
||||
);
|
||||
const fileRows = await adapter.executeQuery(
|
||||
"MATCH (n:File {id: 'File:src/a.ts'}) RETURN n.content AS content",
|
||||
);
|
||||
return {
|
||||
nodeCounts,
|
||||
edgeCounts,
|
||||
bbText: bbRows[0]?.text,
|
||||
fileContent: fileRows[0]?.content,
|
||||
ret,
|
||||
};
|
||||
};
|
||||
|
||||
const loadAndSnapshot = async (label: string, serial: boolean): Promise<Snapshot> => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const storagePath = path.join(tmpBase, label, '.gitnexus');
|
||||
const dbPath = path.join(storagePath, 'lbug');
|
||||
await fs.mkdir(dbPath, { recursive: true });
|
||||
|
||||
if (serial) process.env.GITNEXUS_SERIAL_LBUG_LOAD = '1';
|
||||
else delete process.env.GITNEXUS_SERIAL_LBUG_LOAD;
|
||||
try {
|
||||
await adapter.initLbug(dbPath);
|
||||
const ret = await adapter.loadGraphToLbug(buildFixture(), repoDir, storagePath);
|
||||
const snap = await snapshotDb(adapter, {
|
||||
insertedRels: ret.insertedRels,
|
||||
skippedRels: ret.skippedRels,
|
||||
warnings: ret.warnings,
|
||||
});
|
||||
await adapter.closeLbug();
|
||||
return snap;
|
||||
} finally {
|
||||
delete process.env.GITNEXUS_SERIAL_LBUG_LOAD;
|
||||
}
|
||||
};
|
||||
|
||||
let overlapSnap: Snapshot;
|
||||
let serialSnap: Snapshot;
|
||||
|
||||
beforeAll(async () => {
|
||||
// mkdtemp (not a predictable os.tmpdir join) — secure unique dir, and the
|
||||
// `gitnexus-lbug-` prefix is in TEST_FIXTURE_PREFIXES so the stale-sidecar
|
||||
// sweep recognizes it on Windows (lbug-config.ts).
|
||||
tmpBase = await fs.mkdtemp(path.join(os.tmpdir(), 'gitnexus-lbug-overlap-'));
|
||||
repoDir = path.join(tmpBase, 'repo');
|
||||
await fs.mkdir(path.join(repoDir, 'src'), { recursive: true });
|
||||
await fs.writeFile(path.join(repoDir, 'src', 'a.ts'), FILE_SRC);
|
||||
|
||||
// Default overlap path first, then the serial escape hatch into a fresh DB.
|
||||
overlapSnap = await loadAndSnapshot('overlap', false);
|
||||
serialSnap = await loadAndSnapshot('serial', true);
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
try {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
await adapter.closeLbug();
|
||||
} catch {
|
||||
/* may already be closed */
|
||||
}
|
||||
if (tmpBase) {
|
||||
for (let attempt = 0; attempt < 5; attempt++) {
|
||||
try {
|
||||
await fs.rm(tmpBase, { recursive: true, force: true });
|
||||
return;
|
||||
} catch {
|
||||
if (attempt < 4) await new Promise((r) => setTimeout(r, 200 * (attempt + 1)));
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
describe('node-COPY ‖ rel-emit overlap persists identical content (#2203)', () => {
|
||||
it('per-table node counts are identical between overlap and serial', () => {
|
||||
expect(overlapSnap.nodeCounts).toEqual(serialSnap.nodeCounts);
|
||||
// Sanity: the fixture actually populated every checked table.
|
||||
for (const t of NODE_TABLES_CHECKED) {
|
||||
expect(overlapSnap.nodeCounts[t], `${t} should have rows`).toBeGreaterThan(0);
|
||||
}
|
||||
});
|
||||
|
||||
it('per-type edge counts are identical between overlap and serial', () => {
|
||||
expect(overlapSnap.edgeCounts).toEqual(serialSnap.edgeCounts);
|
||||
for (const ty of EDGE_TYPES_CHECKED) {
|
||||
expect(overlapSnap.edgeCounts[ty], `${ty} should have edges`).toBeGreaterThan(0);
|
||||
}
|
||||
});
|
||||
|
||||
it('multiline content/text fields round-trip identically (byte-for-byte)', () => {
|
||||
expect(overlapSnap.bbText).toBe(serialSnap.bbText);
|
||||
expect(overlapSnap.bbText).toBe(BB_TEXT);
|
||||
expect(overlapSnap.fileContent).toBe(serialSnap.fileContent);
|
||||
expect(overlapSnap.fileContent).toBe(FILE_SRC);
|
||||
});
|
||||
|
||||
it('loadGraphToLbug accounting (insertedRels/skippedRels/warnings) is identical', () => {
|
||||
expect(overlapSnap.ret).toEqual(serialSnap.ret);
|
||||
expect(overlapSnap.ret.insertedRels).toBe(5);
|
||||
expect(overlapSnap.ret.skippedRels).toBe(0);
|
||||
expect(overlapSnap.ret.warnings).toEqual([]);
|
||||
});
|
||||
});
|
||||
|
|
@ -135,6 +135,9 @@ describe('PROF_LBUG_LOAD persistence-path profiling (#2203 U1)', () => {
|
|||
for (const key of ['csv-emit=', 'copy-nodes=', 'copy-rels=', 'fallback=', 'total=']) {
|
||||
expect(line).toContain(key);
|
||||
}
|
||||
// Default load path is the node-COPY ‖ rel-emit overlap (#2203); the prof
|
||||
// line records which path ran. GITNEXUS_SERIAL_LBUG_LOAD is unset here.
|
||||
expect(line).toContain('mode=overlap');
|
||||
// 3 node rows (File, Function, Class), 2 valid rels emitted.
|
||||
expect(line).toContain('(3 nodes, 2 rels)');
|
||||
});
|
||||
|
|
|
|||
|
|
@ -179,3 +179,67 @@ describe('streamed PDG manifest → disjoint-key merge guard (#2202 review #3)',
|
|||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe('streamed PDG manifest → rel-pair collision guard (#2226 F3)', () => {
|
||||
// The rel-pair analogue of the node-CSV guard above, for the case Codex
|
||||
// flagged on PR #2226. The collision check was moved ahead of node COPY so the
|
||||
// serial escape hatch detects it before committing node rows; this asserts the
|
||||
// guard fires in BOTH the overlap (default) and serial paths.
|
||||
|
||||
// A leaky graph carrying a structural BasicBlock→BasicBlock EDGE but NO
|
||||
// BasicBlock nodes: RelPairRouter derives the label from the `BasicBlock:` id
|
||||
// prefix, so the structural relsByPair gets a `BasicBlock|BasicBlock` pair
|
||||
// while nodeFiles stays empty — isolating the REL-pair collision from the
|
||||
// node-CSV one. The manifest declares the same pair via PdgEmitSink.
|
||||
const buildRelCollision = async (label: string) => {
|
||||
const leakyGraph = createKnowledgeGraph();
|
||||
leakyGraph.addRelationship({
|
||||
id: 'CFG:0->1-structural',
|
||||
sourceId: BB(0),
|
||||
targetId: BB(1),
|
||||
type: 'CFG',
|
||||
confidence: 1,
|
||||
reason: 'leak',
|
||||
});
|
||||
const base = await fs.mkdtemp(path.join(os.tmpdir(), `gitnexus-lbug-relcollide-${label}-`));
|
||||
const storage = path.join(base, '.gitnexus');
|
||||
await fs.mkdir(storage, { recursive: true });
|
||||
const sink = new PdgEmitSink(createKnowledgeGraph(), path.join(storage, 'pdg-csv'));
|
||||
sink.addRelationship({
|
||||
id: 'CFG:0->1-manifest',
|
||||
sourceId: BB(0),
|
||||
targetId: BB(1),
|
||||
type: 'CFG',
|
||||
confidence: 1,
|
||||
reason: 'manifest',
|
||||
});
|
||||
const manifest = sink.finalize();
|
||||
return { leakyGraph, base, storage, manifest };
|
||||
};
|
||||
|
||||
it('throws on a rel-pair collision in the overlap (default) path', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const { leakyGraph, base, storage, manifest } = await buildRelCollision('overlap');
|
||||
try {
|
||||
await expect(
|
||||
adapter.loadGraphToLbug(leakyGraph, base, storage, undefined, manifest),
|
||||
).rejects.toThrow(/collides with a structural relationship CSV for pair/);
|
||||
} finally {
|
||||
await fs.rm(base, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
it('throws on a rel-pair collision in the serial (GITNEXUS_SERIAL_LBUG_LOAD=1) path', async () => {
|
||||
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
|
||||
const { leakyGraph, base, storage, manifest } = await buildRelCollision('serial');
|
||||
process.env.GITNEXUS_SERIAL_LBUG_LOAD = '1';
|
||||
try {
|
||||
await expect(
|
||||
adapter.loadGraphToLbug(leakyGraph, base, storage, undefined, manifest),
|
||||
).rejects.toThrow(/collides with a structural relationship CSV for pair/);
|
||||
} finally {
|
||||
delete process.env.GITNEXUS_SERIAL_LBUG_LOAD;
|
||||
await fs.rm(base, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue