mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-06 02:49:56 +00:00
perf(lbug): route relationships to per-pair CSVs in the emit pass (#2203 U2)
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) <noreply@anthropic.com>
This commit is contained in:
parent
8681c6a6d0
commit
2e5ab2be30
5 changed files with 341 additions and 76 deletions
|
|
@ -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<NodeTableName, { csvPath: string; rows: number }>;
|
||||
relCsvPath: string;
|
||||
relRows: number;
|
||||
/** pairKey (`From|To`) → per-FROM→TO-label-pair CSV file. */
|
||||
relsByPair: Map<string, { csvPath: string; rows: number }>;
|
||||
/** 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<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();
|
||||
throw err;
|
||||
}
|
||||
await relWriter.finish();
|
||||
|
||||
// Build result map — only include tables that have rows
|
||||
const nodeFiles = new Map<NodeTableName, { csvPath: string; rows: number }>();
|
||||
|
|
@ -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,
|
||||
};
|
||||
};
|
||||
|
|
|
|||
|
|
@ -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<string>(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<string>();
|
||||
|
||||
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)`,
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
145
gitnexus/src/core/lbug/rel-pair-routing.ts
Normal file
145
gitnexus/src/core/lbug/rel-pair-routing.ts
Normal file
|
|
@ -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<string, RelPairMeta>();
|
||||
private readonly streams = new Map<string, WriteStream>();
|
||||
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<string>,
|
||||
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<void>` 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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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<string, { csvPath: string; rows: number }>,
|
||||
): Promise<string[]> => {
|
||||
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<string>(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);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue