From 2e5ab2be30a363ef1243db2b9382877943a6c5ff Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Mon, 15 Jun 2026 16:16:20 +0000 Subject: [PATCH] perf(lbug): route relationships to per-pair CSVs in the emit pass (#2203 U2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Relationships were written once to a monolithic relations.csv, then re-read line-by-line (regex per edge) and re-split into per-FROM->TO-label-pair files before COPY — writing and reading the entire ~1M-edge set twice. Route each edge to its pair file directly during the single emit pass via a shared RelPairRouter, eliminating the monolithic write + re-read + per-edge regex. The router applies the SAME getNodeLabel + validTables filter as the legacy splitRelCsvByLabelPair, which is retained as a differential oracle. A new differential test asserts the direct-emit per-pair files are byte-for-byte identical to the oracle's, with identical skip/total accounting. The prof line (U1) drops its rel-split stage (routing now folds into csv-emit). Co-Authored-By: Claude Opus 4.8 (1M context) --- gitnexus/src/core/lbug/csv-generator.ts | 70 ++++++--- gitnexus/src/core/lbug/lbug-adapter.ts | 48 ++---- gitnexus/src/core/lbug/rel-pair-routing.ts | 145 ++++++++++++++++++ .../test/integration/csv-pipeline.test.ts | 143 +++++++++++++++-- .../test/integration/lbug-load-prof.test.ts | 11 +- 5 files changed, 341 insertions(+), 76 deletions(-) create mode 100644 gitnexus/src/core/lbug/rel-pair-routing.ts diff --git a/gitnexus/src/core/lbug/csv-generator.ts b/gitnexus/src/core/lbug/csv-generator.ts index c7d8413b5..470cf7ed2 100644 --- a/gitnexus/src/core/lbug/csv-generator.ts +++ b/gitnexus/src/core/lbug/csv-generator.ts @@ -17,7 +17,8 @@ import { createWriteStream, WriteStream } from 'fs'; import path from 'path'; import type { GraphNode, GraphRelationship } from 'gitnexus-shared'; import { KnowledgeGraph } from '../graph/types.js'; -import { NodeTableName } from './schema.js'; +import { NodeTableName, NODE_TABLES } from './schema.js'; +import { RelPairRouter } from './rel-pair-routing.js'; import { parseTruthyEnv } from '../ingestion/utils/env.js'; /** @@ -223,10 +224,33 @@ class BufferedCSVWriter { // STREAMING CSV GENERATION — SINGLE PASS // ============================================================================ +/** Canonical relationship CSV header — shared by the emit pass and the + * `splitRelCsvByLabelPair` differential oracle. */ +export const REL_CSV_HEADER = 'from,to,type,confidence,reason,step'; + +/** Build the escaped CSV row (no trailing newline) for one relationship. + * Single source of the relationship row bytes — used by the emit pass and by + * the byte-identity differential test that feeds the legacy split oracle. */ +export const buildRelRow = (rel: GraphRelationship): string => + [ + escapeCSVField(rel.sourceId), + escapeCSVField(rel.targetId), + escapeCSVField(rel.type), + escapeCSVNumber(rel.confidence, 1.0), + escapeCSVField(rel.reason), + escapeCSVNumber((rel as { step?: number }).step, 0), + ].join(','); + export interface StreamedCSVResult { nodeFiles: Map; - relCsvPath: string; - relRows: number; + /** pairKey (`From|To`) → per-FROM→TO-label-pair CSV file. */ + relsByPair: Map; + /** Header line shared by every per-pair file. */ + relHeader: string; + /** Edges skipped because an endpoint label is not a valid node table. */ + skippedRels: number; + /** Edges routed to a per-pair file. */ + totalValidRels: number; } /** @@ -558,22 +582,24 @@ export const streamAllCSVsToDisk = async ( ]; await Promise.all(allWriters.map((w) => w.finish())); - // --- Stream relationship CSV --- - const relCsvPath = path.join(csvDir, 'relations.csv'); - const relWriter = new BufferedCSVWriter(relCsvPath, 'from,to,type,confidence,reason,step'); - for (const rel of orderedRelationships(graph, sortOutput)) { - await relWriter.addRow( - [ - escapeCSVField(rel.sourceId), - escapeCSVField(rel.targetId), - escapeCSVField(rel.type), - escapeCSVNumber(rel.confidence, 1.0), - escapeCSVField(rel.reason), - escapeCSVNumber((rel as any).step, 0), - ].join(','), - ); + // --- 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(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(); + throw err; } - await relWriter.finish(); // Build result map — only include tables that have rows const nodeFiles = new Map(); @@ -607,5 +633,11 @@ export const streamAllCSVsToDisk = async ( // Restore original process listener limit process.setMaxListeners(prevMax); - return { nodeFiles, relCsvPath, relRows: relWriter.rows }; + return { + nodeFiles, + relsByPair: relRouter.byPair, + relHeader: REL_CSV_HEADER, + skippedRels: relRouter.skipped, + totalValidRels: relRouter.total, + }; }; diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index bc7dc1109..c62466ed0 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -19,6 +19,7 @@ import { NodeTableName, } from './schema.js'; import { streamAllCSVsToDisk } from './csv-generator.js'; +import { getNodeLabel as deriveNodeLabel } from './rel-pair-routing.js'; import type { CachedEmbedding } from '../embeddings/types.js'; import { extensionManager, type ExtensionEnsureOptions } from './extension-loader.js'; import { @@ -902,11 +903,6 @@ export const loadGraphToLbug = async ( const tCsv = mark(); const validTables = new Set(NODE_TABLES as readonly string[]); - const getNodeLabel = (nodeId: string): string => { - if (nodeId.startsWith('comm_')) return 'Community'; - if (nodeId.startsWith('proc_')) return 'Process'; - return nodeId.split(':')[0]; - }; // Bulk COPY all node CSVs (sequential — LadybugDB allows only one write txn at a time) const nodeFiles = [...csvResult.nodeFiles.entries()]; @@ -938,40 +934,30 @@ export const loadGraphToLbug = async ( const tCopyNodes = mark(); - // Bulk COPY relationships — split by FROM→TO label pair (LadybugDB requires it) - const { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels } = - await splitRelCsvByLabelPair(csvResult.relCsvPath, csvDir, validTables, getNodeLabel); - - // Close all per-pair write streams before COPY. `stream/promises.finished` - // resolves on the stream's 'finish' event and rejects on 'error' — replaces - // a hand-rolled promisification with the stdlib primitive. - await Promise.all( - Array.from(pairWriteStreams.values()).map(async (ws) => { - ws.end(); - await finished(ws); - }), - ); - const tSplit = mark(); - let tCopyRels = tSplit; - let tFallback = tSplit; + // Bulk COPY relationships. They were already routed to per-FROM→TO-label-pair + // files during the emit pass (#2203 U2) — there is no monolithic relations.csv + // to re-read/re-split here; we COPY each pair file directly. + const { relsByPair, relHeader, skippedRels, totalValidRels } = csvResult; + let tCopyRels = tCopyNodes; + let tFallback = tCopyNodes; const insertedRels = totalValidRels; const warnings: string[] = []; if (insertedRels > 0) { - log(`Loading edges: ${insertedRels.toLocaleString()} across ${relsByPairMeta.size} types`); + log(`Loading edges: ${insertedRels.toLocaleString()} across ${relsByPair.size} types`); let pairIdx = 0; let failedPairEdges = 0; const failedPairCsvPaths = new Set(); - for (const [pairKey, { csvPath: pairCsvPath, rows }] of relsByPairMeta) { + for (const [pairKey, { csvPath: pairCsvPath, rows }] of relsByPair) { pairIdx++; const [fromLabel, toLabel] = pairKey.split('|'); const normalizedPath = normalizeCopyPath(pairCsvPath); 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}/${relsByPairMeta.size} types (${fromLabel} -> ${toLabel})`); + log(`Loading edges: ${pairIdx}/${relsByPair.size} types (${fromLabel} -> ${toLabel})`); } try { @@ -1017,16 +1003,14 @@ export const loadGraphToLbug = async ( } catch {} } if (allLines.length > 1) { - await fallbackRelationshipInserts(allLines, validTables, getNodeLabel); + await fallbackRelationshipInserts(allLines, validTables, deriveNodeLabel); } } tFallback = mark(); } - // Cleanup all CSVs - try { - await fs.unlink(csvResult.relCsvPath); - } catch {} + // Cleanup all CSVs (per-pair rel files are unlinked in the COPY loop above; + // the remaining sweep below catches node CSVs + any leftover pair files). for (const [, { csvPath }] of csvResult.nodeFiles) { try { await fs.unlink(csvPath); @@ -1050,9 +1034,9 @@ export const loadGraphToLbug = async ( for (const [, { rows }] of csvResult.nodeFiles) totalNodeRows += rows; logger.warn( `[lbug-load prof] csv-emit=${span(tStart, tCsv)}ms ` + - `copy-nodes=${span(tCsv, tCopyNodes)}ms rel-split=${span(tCopyNodes, tSplit)}ms ` + - `copy-rels=${span(tSplit, tCopyRels)}ms fallback=${span(tCopyRels, tFallback)}ms ` + - `total=${span(tStart, tEnd)}ms (${totalNodeRows} nodes, ${insertedRels} rels)`, + `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)`, ); } diff --git a/gitnexus/src/core/lbug/rel-pair-routing.ts b/gitnexus/src/core/lbug/rel-pair-routing.ts new file mode 100644 index 000000000..2ba3ee2a4 --- /dev/null +++ b/gitnexus/src/core/lbug/rel-pair-routing.ts @@ -0,0 +1,145 @@ +/** + * Relationship per-label-pair routing (#2203 U2). + * + * LadybugDB's bulk `COPY` into the single `CodeRelation` rel table requires a + * separate CSV per FROM→TO node-label pair (the `from=`/`to=` COPY params). + * Historically the emit pass wrote one monolithic `relations.csv`, which + * `loadGraphToLbug` then RE-READ line-by-line (regex per edge) and re-split + * into per-pair files — writing and reading the entire ~1M-edge set twice. + * + * This router lets the single emit pass route each edge to its per-pair file + * directly, so the monolithic write + re-read + per-edge regex are all gone. + * The label-derivation + validTables filtering + per-pair-file format here are + * the SAME ones the legacy `splitRelCsvByLabelPair` applies, so the per-pair + * files are byte-identical — see the differential test in + * `test/integration/csv-pipeline.test.ts`. `splitRelCsvByLabelPair` is retained + * as the differential oracle. + * + * Backpressure: at most one stream is awaited at a time (the caller routes + * edges sequentially and awaits the returned drain promise before the next), + * mirroring the legacy split's `for await` invariant. The hot path (existing + * pair, no backpressure) returns `void` — no microtask per edge. + */ +import path from 'path'; +import { createWriteStream, type WriteStream } from 'fs'; +import { once } from 'events'; +import { finished } from 'stream/promises'; + +/** Injectable for tests (backpressure/error simulation), mirroring split. */ +export type WriteStreamFactory = (filePath: string) => WriteStream; + +/** + * Derive a node's table label from its graph id. Matches the legacy + * `getNodeLabel` that lived inline in `loadGraphToLbug`: + * - `comm_*` → Community + * - `proc_*` → Process + * - otherwise the prefix before the first `:` (e.g. `Function:…` → Function) + */ +export const getNodeLabel = (nodeId: string): string => { + if (nodeId.startsWith('comm_')) return 'Community'; + if (nodeId.startsWith('proc_')) return 'Process'; + return nodeId.split(':')[0]; +}; + +export interface RelPairMeta { + csvPath: string; + rows: number; +} + +/** + * Routes already-escaped relationship CSV rows to per-FROM→TO-label-pair + * files. Filters edges whose endpoint labels are not valid node tables + * (counted as `skipped`), exactly as the legacy split did. + */ +export class RelPairRouter { + /** pairKey (`From|To`) → { csvPath, rows } */ + readonly byPair = new Map(); + private readonly streams = new Map(); + skipped = 0; + total = 0; + + private streamError: Error | null = null; + private readonly abort = new AbortController(); + + constructor( + private readonly csvDir: string, + private readonly header: string, + private readonly validTables: Set, + private readonly wsFactory: WriteStreamFactory = (p) => createWriteStream(p, 'utf-8'), + ) {} + + private markError = (err: Error): void => { + this.streamError ??= err; + this.abort.abort(err); + }; + + /** + * Route one already-escaped CSV row (no trailing newline) to its pair file. + * Returns `void` on the synchronous hot path; a `Promise` only when a + * stream signals backpressure (or a new pair's header does) — the caller + * awaits the promise before routing the next edge. + */ + route(fromId: string, toId: string, row: string): void | Promise { + if (this.streamError) throw this.streamError; + + const fromLabel = getNodeLabel(fromId); + const toLabel = getNodeLabel(toId); + if (!this.validTables.has(fromLabel) || !this.validTables.has(toLabel)) { + this.skipped++; + return; + } + + const pairKey = `${fromLabel}|${toLabel}`; + const ws = this.streams.get(pairKey); + if (ws === undefined) { + // First edge for this pair: open the stream, write header + row. + return this.openAndWrite(pairKey, fromLabel, toLabel, row); + } + + this.byPair.get(pairKey)!.rows++; + this.total++; + if (!ws.write(row + '\n')) { + return once(ws, 'drain', { signal: this.abort.signal }).then(() => undefined); + } + } + + private async openAndWrite( + pairKey: string, + fromLabel: string, + toLabel: string, + row: string, + ): Promise { + const csvPath = path.join(this.csvDir, `rel_${fromLabel}_${toLabel}.csv`); + const ws = this.wsFactory(csvPath); + ws.on('error', this.markError); + this.streams.set(pairKey, ws); + this.byPair.set(pairKey, { csvPath, rows: 1 }); + this.total++; + if (!ws.write(this.header + '\n')) { + await once(ws, 'drain', { signal: this.abort.signal }); + } + if (!ws.write(row + '\n')) { + await once(ws, 'drain', { signal: this.abort.signal }); + } + } + + /** Flush + close every pair stream. Rejects if any stream errored. */ + async close(): Promise { + if (this.streamError) { + this.destroy(); + throw this.streamError; + } + await Promise.all( + Array.from(this.streams.values()).map(async (ws) => { + ws.end(); + await finished(ws); + }), + ); + if (this.streamError) throw this.streamError; + } + + /** Tear down all streams (no flush) — used on the error path. */ + destroy(): void { + for (const ws of this.streams.values()) ws.destroy(); + } +} diff --git a/gitnexus/test/integration/csv-pipeline.test.ts b/gitnexus/test/integration/csv-pipeline.test.ts index cf44c85b2..28690587b 100644 --- a/gitnexus/test/integration/csv-pipeline.test.ts +++ b/gitnexus/test/integration/csv-pipeline.test.ts @@ -6,15 +6,43 @@ */ import { describe, it, expect, beforeAll, afterAll } from 'vitest'; import fs from 'fs/promises'; +import { finished } from 'stream/promises'; import path from 'path'; import { createTempDir, type TestDBHandle } from '../helpers/test-db.js'; import { buildTestGraph, type TestNodeInput, type TestRelInput } from '../helpers/test-graph.js'; -import { streamAllCSVsToDisk } from '../../src/core/lbug/csv-generator.js'; +import { + streamAllCSVsToDisk, + buildRelRow, + REL_CSV_HEADER, +} from '../../src/core/lbug/csv-generator.js'; +import { splitRelCsvByLabelPair } from '../../src/core/lbug/lbug-adapter.js'; +import { getNodeLabel } from '../../src/core/lbug/rel-pair-routing.js'; +import { NODE_TABLES } from '../../src/core/lbug/schema.js'; let tmpHandle: TestDBHandle; let csvDir: string; let repoDir: string; +/** Data rows (header dropped) of one CSV file's text. */ +const dataRowsOf = (csv: string): string[] => + csv + .trim() + .split('\n') + .slice(1) + .filter((l) => l.length > 0); + +/** Concatenate data rows from every per-pair rel file (#2203 U2), pair keys + * sorted so the concatenation order is deterministic regardless of map order. */ +const readAllRelRows = async ( + relsByPair: Map, +): Promise => { + const rows: string[] = []; + for (const key of [...relsByPair.keys()].sort()) { + rows.push(...dataRowsOf(await fs.readFile(relsByPair.get(key)!.csvPath, 'utf-8'))); + } + return rows; +}; + beforeAll(async () => { tmpHandle = await createTempDir('csv-pipeline-test-'); csvDir = path.join(tmpHandle.dbPath, 'csv'); @@ -76,9 +104,9 @@ describe('streamAllCSVsToDisk', () => { { id: 'folder:src', label: 'Folder', name: 'src', filePath: 'src' }, ], [ - { sourceId: 'func:main', targetId: 'func:helper', type: 'CALLS' }, - { sourceId: 'file:src/index.ts', targetId: 'func:main', type: 'CONTAINS' }, - { sourceId: 'file:src/utils.ts', targetId: 'func:helper', type: 'CONTAINS' }, + { sourceId: 'Function:main', targetId: 'Function:helper', type: 'CALLS' }, + { sourceId: 'File:src/index.ts', targetId: 'Function:main', type: 'CONTAINS' }, + { sourceId: 'File:src/utils.ts', targetId: 'Function:helper', type: 'CONTAINS' }, ], ); @@ -86,7 +114,8 @@ describe('streamAllCSVsToDisk', () => { // Check that CSV files were created expect(result.nodeFiles.size).toBeGreaterThan(0); - expect(result.relRows).toBe(3); + expect(result.totalValidRels).toBe(3); + expect(result.skippedRels).toBe(0); // Verify File CSV const fileCsv = result.nodeFiles.get('File'); @@ -108,10 +137,12 @@ describe('streamAllCSVsToDisk', () => { expect(folderCsv).toBeDefined(); expect(folderCsv!.rows).toBe(1); - // Verify relations CSV exists - const relContent = await fs.readFile(result.relCsvPath, 'utf-8'); - const relLines = relContent.trim().split('\n'); - expect(relLines.length).toBe(4); // header + 3 relationships + // Relationships are routed to per-FROM→TO-label-pair files (#2203 U2): + // Function→Function (CALLS) + File→Function (2× CONTAINS). + expect(result.relsByPair.has('Function|Function')).toBe(true); + expect(result.relsByPair.has('File|Function')).toBe(true); + expect(result.relsByPair.get('File|Function')!.rows).toBe(2); + expect(await readAllRelRows(result.relsByPair)).toHaveLength(3); }); it('CSV content is properly escaped', async () => { @@ -210,7 +241,8 @@ describe('streamAllCSVsToDisk', () => { const graph = buildTestGraph([], []); const result = await streamAllCSVsToDisk(graph, repoDir, csvDir); expect(result.nodeFiles.size).toBe(0); - expect(result.relRows).toBe(0); + expect(result.totalValidRels).toBe(0); + expect(result.relsByPair.size).toBe(0); }); it('handles node with empty string properties', async () => { @@ -234,15 +266,18 @@ describe('streamAllCSVsToDisk — deterministic output ordering', () => { // Folder nodes: single-line CSV rows (no multi-line `content` column), so the // id is the first comma-separated field and split('\n') is safe. ids are // deliberately NOT in insertion order (c, a, b). + // ids use the `Folder:` prefix so getNodeLabel derives the valid `Folder` + // table — edges route to rel_Folder_Folder.csv (#2203 U2). Deliberately NOT + // in insertion order (c, a, b). const NODES: TestNodeInput[] = [ - { id: 'folder:c', label: 'Folder', name: 'c', filePath: 'c' }, - { id: 'folder:a', label: 'Folder', name: 'a', filePath: 'a' }, - { id: 'folder:b', label: 'Folder', name: 'b', filePath: 'b' }, + { id: 'Folder:c', label: 'Folder', name: 'c', filePath: 'c' }, + { id: 'Folder:a', label: 'Folder', name: 'a', filePath: 'a' }, + { id: 'Folder:b', label: 'Folder', name: 'b', filePath: 'b' }, ]; const RELS: TestRelInput[] = [ - { sourceId: 'folder:c', targetId: 'folder:a', type: 'CONTAINS' }, - { sourceId: 'folder:a', targetId: 'folder:b', type: 'CONTAINS' }, - { sourceId: 'folder:b', targetId: 'folder:c', type: 'CONTAINS' }, + { sourceId: 'Folder:c', targetId: 'Folder:a', type: 'CONTAINS' }, + { sourceId: 'Folder:a', targetId: 'Folder:b', type: 'CONTAINS' }, + { sourceId: 'Folder:b', targetId: 'Folder:c', type: 'CONTAINS' }, ]; const dataRows = (csv: string): string[] => csv @@ -270,7 +305,7 @@ describe('streamAllCSVsToDisk — deterministic output ordering', () => { const folderIds = folderCsv ? dataRows(await fs.readFile(folderCsv.csvPath, 'utf-8')).map(firstCol) : []; - const relRows = dataRows(await fs.readFile(result.relCsvPath, 'utf-8')); + const relRows = await readAllRelRows(result.relsByPair); return { folderIds, relRows }; } finally { delete process.env.GITNEXUS_SORT_GRAPH_OUTPUT; @@ -307,3 +342,77 @@ describe('streamAllCSVsToDisk — deterministic output ordering', () => { expect([...onFwd.relRows].sort()).toEqual([...offFwd.relRows].sort()); }); }); + +/** + * #2203 U2 byte-identity: the direct per-pair emit must produce per-pair files + * byte-for-byte identical to the legacy splitRelCsvByLabelPair oracle run over + * an equivalent monolithic relations.csv from the same graph. This is the + * load-bearing guard for "byte-identical graph content" (issue acceptance). + */ +describe('streamAllCSVsToDisk — direct per-pair emit matches the split oracle', () => { + it('produces byte-identical per-pair files + identical skip/total accounting', async () => { + // Multiple valid pairs, getNodeLabel special prefixes (comm_/proc_), and one + // invalid-label edge that BOTH paths must skip identically. + const graph = buildTestGraph( + [ + { id: 'File:a.ts', label: 'File', name: 'a.ts', filePath: 'a.ts' }, + { id: 'Function:a.ts:f:1', label: 'Function', name: 'f', filePath: 'a.ts' }, + { id: 'Function:a.ts:g:5', label: 'Function', name: 'g', filePath: 'a.ts' }, + { id: 'comm_1', label: 'Community' as never, name: 'c1', filePath: '' }, + { id: 'comm_2', label: 'Community' as never, name: 'c2', filePath: '' }, + ], + [ + { sourceId: 'File:a.ts', targetId: 'Function:a.ts:f:1', type: 'CONTAINS' }, + { sourceId: 'File:a.ts', targetId: 'Function:a.ts:g:5', type: 'CONTAINS' }, + { sourceId: 'Function:a.ts:f:1', targetId: 'Function:a.ts:g:5', type: 'CALLS' }, + { sourceId: 'comm_1', targetId: 'comm_2', type: 'CONTAINS' }, + // Invalid FROM label ('Bogus' ∉ NODE_TABLES) — skipped by both paths. + { sourceId: 'Bogus:x', targetId: 'File:a.ts', type: 'CONTAINS' }, + ], + ); + + const directDir = path.join(csvDir, 'diff-direct'); + const oracleDir = path.join(csvDir, 'diff-oracle'); + await fs.mkdir(oracleDir, { recursive: true }); + + // Direct emit (production path). + const direct = await streamAllCSVsToDisk(graph, repoDir, directDir); + + // Oracle: build the monolithic relations.csv this graph would have produced + // (same insertion order, same row bytes via buildRelRow), then split it. + const relCsv = path.join(oracleDir, 'relations.csv'); + const lines = [REL_CSV_HEADER]; + for (const rel of graph.iterRelationships()) lines.push(buildRelRow(rel)); + await fs.writeFile(relCsv, lines.join('\n') + '\n', 'utf-8'); + + const split = await splitRelCsvByLabelPair( + relCsv, + oracleDir, + new Set(NODE_TABLES), + getNodeLabel, + ); + await Promise.all( + Array.from(split.pairWriteStreams.values()).map(async (ws) => { + ws.end(); + await finished(ws); + }), + ); + + // Identical accounting. + expect(direct.totalValidRels).toBe(split.totalValidRels); + expect(direct.totalValidRels).toBe(4); + expect(direct.skippedRels).toBe(split.skippedRels); + expect(direct.skippedRels).toBe(1); + expect(direct.relHeader).toBe(split.relHeader); + + // Identical pair set. + expect([...direct.relsByPair.keys()].sort()).toEqual([...split.relsByPairMeta.keys()].sort()); + + // Byte-identical per-pair file contents. + for (const key of direct.relsByPair.keys()) { + const directContent = await fs.readFile(direct.relsByPair.get(key)!.csvPath, 'utf-8'); + const oracleContent = await fs.readFile(split.relsByPairMeta.get(key)!.csvPath, 'utf-8'); + expect(directContent, `pair ${key}`).toBe(oracleContent); + } + }); +}); diff --git a/gitnexus/test/integration/lbug-load-prof.test.ts b/gitnexus/test/integration/lbug-load-prof.test.ts index 8484e0e2c..f83cd75c2 100644 --- a/gitnexus/test/integration/lbug-load-prof.test.ts +++ b/gitnexus/test/integration/lbug-load-prof.test.ts @@ -130,14 +130,9 @@ describe('PROF_LBUG_LOAD persistence-path profiling (#2203 U1)', () => { expect(lines).toHaveLength(1); const line = lines[0]; - for (const key of [ - 'csv-emit=', - 'copy-nodes=', - 'rel-split=', - 'copy-rels=', - 'fallback=', - 'total=', - ]) { + // Relationships are routed to per-pair files during csv-emit (#2203 U2), + // so there is no separate rel-split stage. + for (const key of ['csv-emit=', 'copy-nodes=', 'copy-rels=', 'fallback=', 'total=']) { expect(line).toContain(key); } // 3 node rows (File, Function, Class), 2 valid rels emitted.