fix(review): apply autofix feedback (#2203)

- P1: router backpressure drain-await rejected with a generic AbortError,
  masking the real EMFILE/disk-full error. Expose RelPairRouter.lastError and
  rethrow it in the emit catch — mirrors the oracle's throw streamError ?? err.
- P1: cover RelPairRouter error + backpressure + teardown paths with a new
  unit test (test/unit/rel-pair-routing.test.ts) using an injected mock stream.
- P2: wrap streamAllCSVsToDisk body in try/finally so the setMaxListeners bump
  is always restored (the U2 rel-routing throw path could leak it).
- P2: dedup WriteStreamFactory — re-export the canonical type from
  rel-pair-routing instead of a second identical declaration.
- P2: annotate splitRelCsvByLabelPair @internal as the retained differential
  oracle so a future dead-code sweep doesn't delete the byte-identity guard.
- P3: differential test now covers the proc_ prefix + clears
  GITNEXUS_SORT_GRAPH_OUTPUT to prevent env-leak desync.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Gergo Magyar 2026-06-15 17:02:06 +00:00
parent e305561d34
commit 760f3ce86e
5 changed files with 554 additions and 338 deletions

View file

@ -284,259 +284,185 @@ export const streamAllCSVsToDisk = async (
const prevMax = process.getMaxListeners();
process.setMaxListeners(prevMax + 40);
const contentCache = new FileContentCache(repoPath);
// try/finally so the listener bump is ALWAYS restored — including the
// rel-routing throw path (#2203 U2) and any node-writer finish() rejection,
// not just the success path (avoids leaking +40 listeners across failed runs
// in long-lived hosts / the test suite).
try {
const contentCache = new FileContentCache(repoPath);
// Create writers for every node type up-front
const fileWriter = new BufferedCSVWriter(
path.join(csvDir, 'file.csv'),
'id,name,filePath,content',
);
const folderWriter = new BufferedCSVWriter(path.join(csvDir, 'folder.csv'), 'id,name,filePath');
const codeElementHeader = 'id,name,filePath,startLine,endLine,isExported,content,description';
const functionWriter = new BufferedCSVWriter(
path.join(csvDir, 'function.csv'),
codeElementHeader,
);
const classWriter = new BufferedCSVWriter(path.join(csvDir, 'class.csv'), codeElementHeader);
const interfaceWriter = new BufferedCSVWriter(
path.join(csvDir, 'interface.csv'),
codeElementHeader,
);
const methodHeader =
'id,name,filePath,startLine,endLine,isExported,content,description,parameterCount,returnType';
const methodWriter = new BufferedCSVWriter(path.join(csvDir, 'method.csv'), methodHeader);
const codeElemWriter = new BufferedCSVWriter(
path.join(csvDir, 'codeelement.csv'),
codeElementHeader,
);
const communityWriter = new BufferedCSVWriter(
path.join(csvDir, 'community.csv'),
'id,label,heuristicLabel,keywords,description,enrichedBy,cohesion,symbolCount',
);
const processWriter = new BufferedCSVWriter(
path.join(csvDir, 'process.csv'),
'id,label,heuristicLabel,processType,stepCount,communities,entryPointId,terminalId',
);
// Section nodes have an extra 'level' column
const sectionWriter = new BufferedCSVWriter(
path.join(csvDir, 'section.csv'),
'id,name,filePath,startLine,endLine,level,content,description',
);
// Route nodes for API endpoint mapping
const routeWriter = new BufferedCSVWriter(
path.join(csvDir, 'route.csv'),
'id,name,filePath,responseKeys,errorKeys,middleware',
);
// Tool nodes for MCP tool definitions
const toolWriter = new BufferedCSVWriter(
path.join(csvDir, 'tool.csv'),
'id,name,filePath,description',
);
// BasicBlock nodes — taint/PDG substrate (issue #2080). No `name` column;
// blocks are identified by id + source span. Emitted by no phase yet.
const basicBlockWriter = new BufferedCSVWriter(
path.join(csvDir, 'basicblock.csv'),
'id,filePath,startLine,endLine,text',
);
// Multi-language node types share the same CSV shape (no isExported column)
const multiLangHeader = 'id,name,filePath,startLine,endLine,content,description';
const MULTI_LANG_TYPES = [
'Struct',
'Enum',
'Macro',
'Typedef',
'Union',
'Namespace',
'Trait',
'Impl',
'TypeAlias',
'Const',
'Static',
'Variable',
'Property',
'Record',
'Delegate',
'Annotation',
'Constructor',
'Template',
'Module',
] as const;
const propertyHeader = 'id,name,filePath,startLine,endLine,content,description,declaredType';
const multiLangWriters = new Map<string, BufferedCSVWriter>();
for (const t of MULTI_LANG_TYPES) {
multiLangWriters.set(
t,
new BufferedCSVWriter(
path.join(csvDir, `${t.toLowerCase()}.csv`),
t === 'Property' ? propertyHeader : multiLangHeader,
),
// Create writers for every node type up-front
const fileWriter = new BufferedCSVWriter(
path.join(csvDir, 'file.csv'),
'id,name,filePath,content',
);
const folderWriter = new BufferedCSVWriter(path.join(csvDir, 'folder.csv'), 'id,name,filePath');
const codeElementHeader = 'id,name,filePath,startLine,endLine,isExported,content,description';
const functionWriter = new BufferedCSVWriter(
path.join(csvDir, 'function.csv'),
codeElementHeader,
);
const classWriter = new BufferedCSVWriter(path.join(csvDir, 'class.csv'), codeElementHeader);
const interfaceWriter = new BufferedCSVWriter(
path.join(csvDir, 'interface.csv'),
codeElementHeader,
);
const methodHeader =
'id,name,filePath,startLine,endLine,isExported,content,description,parameterCount,returnType';
const methodWriter = new BufferedCSVWriter(path.join(csvDir, 'method.csv'), methodHeader);
const codeElemWriter = new BufferedCSVWriter(
path.join(csvDir, 'codeelement.csv'),
codeElementHeader,
);
const communityWriter = new BufferedCSVWriter(
path.join(csvDir, 'community.csv'),
'id,label,heuristicLabel,keywords,description,enrichedBy,cohesion,symbolCount',
);
const processWriter = new BufferedCSVWriter(
path.join(csvDir, 'process.csv'),
'id,label,heuristicLabel,processType,stepCount,communities,entryPointId,terminalId',
);
}
const codeWriterMap: Record<string, BufferedCSVWriter> = {
Function: functionWriter,
Class: classWriter,
Interface: interfaceWriter,
CodeElement: codeElemWriter,
};
// Section nodes have an extra 'level' column
const sectionWriter = new BufferedCSVWriter(
path.join(csvDir, 'section.csv'),
'id,name,filePath,startLine,endLine,level,content,description',
);
// Deduplicate all node types — the pipeline can produce duplicate IDs across
// all symbol types (Class, Method, Function, etc.), not just File nodes.
// A single Set covering every label prevents PK violations on COPY.
const seenNodeIds = new Set<string>();
// Route nodes for API endpoint mapping
const routeWriter = new BufferedCSVWriter(
path.join(csvDir, 'route.csv'),
'id,name,filePath,responseKeys,errorKeys,middleware',
);
// --- SINGLE PASS over all nodes ---
for (const node of orderedNodes(graph, sortOutput)) {
if (seenNodeIds.has(node.id)) continue;
seenNodeIds.add(node.id);
// Tool nodes for MCP tool definitions
const toolWriter = new BufferedCSVWriter(
path.join(csvDir, 'tool.csv'),
'id,name,filePath,description',
);
// addRow returns a promise only when it flushes; awaiting it once after the
// switch (instead of `await`-ing every addRow) skips a per-row microtask
// tick on the ~FLUSH_EVERY-1 buffered rows between flushes (#2203 U3).
let pending: Promise<void> | undefined;
switch (node.label) {
case 'File': {
const content = await extractContent(node, contentCache);
pending = fileWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVField(content),
].join(','),
);
break;
}
case 'Folder':
pending = folderWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
].join(','),
);
break;
case 'Community': {
const keywords = node.properties.keywords || [];
const keywordsStr = `[${keywords.map((k: string) => `'${k.replace(/\\/g, '\\\\').replace(/'/g, "''").replace(/,/g, '\\,')}'`).join(',')}]`;
pending = communityWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.heuristicLabel || ''),
keywordsStr,
escapeCSVField(node.properties.description || ''),
escapeCSVField(node.properties.enrichedBy || 'heuristic'),
escapeCSVNumber(node.properties.cohesion, 0),
escapeCSVNumber(node.properties.symbolCount, 0),
].join(','),
);
break;
}
case 'Process': {
const communities = node.properties.communities || [];
const communitiesStr = `[${communities.map((c: string) => `'${c.replace(/'/g, "''")}'`).join(',')}]`;
pending = processWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.heuristicLabel || ''),
escapeCSVField(node.properties.processType || ''),
escapeCSVNumber(node.properties.stepCount, 0),
escapeCSVField(communitiesStr),
escapeCSVField(node.properties.entryPointId || ''),
escapeCSVField(node.properties.terminalId || ''),
].join(','),
);
break;
}
case 'Method': {
const content = await extractContent(node, contentCache);
pending = methodWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVNumber(node.properties.startLine, -1),
escapeCSVNumber(node.properties.endLine, -1),
node.properties.isExported ? 'true' : 'false',
escapeCSVField(content),
escapeCSVField(node.properties.description || ''),
escapeCSVNumber(node.properties.parameterCount, 0),
escapeCSVField(node.properties.returnType || ''),
].join(','),
);
break;
}
case 'Section': {
const content = await extractContent(node, contentCache);
pending = sectionWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVNumber(node.properties.startLine, -1),
escapeCSVNumber(node.properties.endLine, -1),
escapeCSVNumber(node.properties.level, 1),
escapeCSVField(content),
escapeCSVField(node.properties.description || ''),
].join(','),
);
break;
}
case 'Route': {
const responseKeys = node.properties.responseKeys || [];
// LadybugDB array literal inside a quoted CSV field: escapeCSVField wraps in "..."
// and the array uses single-quoted elements
const keysStr = `[${responseKeys.map((k: string) => `'${k.replace(/'/g, "''")}'`).join(',')}]`;
const errorKeys = node.properties.errorKeys || [];
const errorKeysStr = `[${errorKeys.map((k: string) => `'${k.replace(/'/g, "''")}'`).join(',')}]`;
const middleware = node.properties.middleware || [];
const middlewareStr = `[${middleware.map((m: string) => `'${m.replace(/'/g, "''")}'`).join(',')}]`;
pending = routeWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVField(keysStr),
escapeCSVField(errorKeysStr),
escapeCSVField(middlewareStr),
].join(','),
);
break;
}
case 'Tool':
pending = toolWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVField(node.properties.description || ''),
].join(','),
);
break;
case 'BasicBlock':
pending = basicBlockWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.filePath || ''),
escapeCSVNumber(node.properties.startLine, -1),
escapeCSVNumber(node.properties.endLine, -1),
escapeCSVField(node.properties.text || ''),
].join(','),
);
break;
default: {
// Code element nodes (Function, Class, Interface, CodeElement)
const writer = codeWriterMap[node.label];
if (writer) {
// BasicBlock nodes — taint/PDG substrate (issue #2080). No `name` column;
// blocks are identified by id + source span. Emitted by no phase yet.
const basicBlockWriter = new BufferedCSVWriter(
path.join(csvDir, 'basicblock.csv'),
'id,filePath,startLine,endLine,text',
);
// Multi-language node types share the same CSV shape (no isExported column)
const multiLangHeader = 'id,name,filePath,startLine,endLine,content,description';
const MULTI_LANG_TYPES = [
'Struct',
'Enum',
'Macro',
'Typedef',
'Union',
'Namespace',
'Trait',
'Impl',
'TypeAlias',
'Const',
'Static',
'Variable',
'Property',
'Record',
'Delegate',
'Annotation',
'Constructor',
'Template',
'Module',
] as const;
const propertyHeader = 'id,name,filePath,startLine,endLine,content,description,declaredType';
const multiLangWriters = new Map<string, BufferedCSVWriter>();
for (const t of MULTI_LANG_TYPES) {
multiLangWriters.set(
t,
new BufferedCSVWriter(
path.join(csvDir, `${t.toLowerCase()}.csv`),
t === 'Property' ? propertyHeader : multiLangHeader,
),
);
}
const codeWriterMap: Record<string, BufferedCSVWriter> = {
Function: functionWriter,
Class: classWriter,
Interface: interfaceWriter,
CodeElement: codeElemWriter,
};
// Deduplicate all node types — the pipeline can produce duplicate IDs across
// all symbol types (Class, Method, Function, etc.), not just File nodes.
// A single Set covering every label prevents PK violations on COPY.
const seenNodeIds = new Set<string>();
// --- SINGLE PASS over all nodes ---
for (const node of orderedNodes(graph, sortOutput)) {
if (seenNodeIds.has(node.id)) continue;
seenNodeIds.add(node.id);
// addRow returns a promise only when it flushes; awaiting it once after the
// switch (instead of `await`-ing every addRow) skips a per-row microtask
// tick on the ~FLUSH_EVERY-1 buffered rows between flushes (#2203 U3).
let pending: Promise<void> | undefined;
switch (node.label) {
case 'File': {
const content = await extractContent(node, contentCache);
pending = writer.addRow(
pending = fileWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVField(content),
].join(','),
);
break;
}
case 'Folder':
pending = folderWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
].join(','),
);
break;
case 'Community': {
const keywords = node.properties.keywords || [];
const keywordsStr = `[${keywords.map((k: string) => `'${k.replace(/\\/g, '\\\\').replace(/'/g, "''").replace(/,/g, '\\,')}'`).join(',')}]`;
pending = communityWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.heuristicLabel || ''),
keywordsStr,
escapeCSVField(node.properties.description || ''),
escapeCSVField(node.properties.enrichedBy || 'heuristic'),
escapeCSVNumber(node.properties.cohesion, 0),
escapeCSVNumber(node.properties.symbolCount, 0),
].join(','),
);
break;
}
case 'Process': {
const communities = node.properties.communities || [];
const communitiesStr = `[${communities.map((c: string) => `'${c.replace(/'/g, "''")}'`).join(',')}]`;
pending = processWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.heuristicLabel || ''),
escapeCSVField(node.properties.processType || ''),
escapeCSVNumber(node.properties.stepCount, 0),
escapeCSVField(communitiesStr),
escapeCSVField(node.properties.entryPointId || ''),
escapeCSVField(node.properties.terminalId || ''),
].join(','),
);
break;
}
case 'Method': {
const content = await extractContent(node, contentCache);
pending = methodWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
@ -546,110 +472,193 @@ export const streamAllCSVsToDisk = async (
node.properties.isExported ? 'true' : 'false',
escapeCSVField(content),
escapeCSVField(node.properties.description || ''),
escapeCSVNumber(node.properties.parameterCount, 0),
escapeCSVField(node.properties.returnType || ''),
].join(','),
);
} else {
// Multi-language node types (Struct, Impl, Trait, Macro, etc.)
const mlWriter = multiLangWriters.get(node.label);
if (mlWriter) {
break;
}
case 'Section': {
const content = await extractContent(node, contentCache);
pending = sectionWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVNumber(node.properties.startLine, -1),
escapeCSVNumber(node.properties.endLine, -1),
escapeCSVNumber(node.properties.level, 1),
escapeCSVField(content),
escapeCSVField(node.properties.description || ''),
].join(','),
);
break;
}
case 'Route': {
const responseKeys = node.properties.responseKeys || [];
// LadybugDB array literal inside a quoted CSV field: escapeCSVField wraps in "..."
// and the array uses single-quoted elements
const keysStr = `[${responseKeys.map((k: string) => `'${k.replace(/'/g, "''")}'`).join(',')}]`;
const errorKeys = node.properties.errorKeys || [];
const errorKeysStr = `[${errorKeys.map((k: string) => `'${k.replace(/'/g, "''")}'`).join(',')}]`;
const middleware = node.properties.middleware || [];
const middlewareStr = `[${middleware.map((m: string) => `'${m.replace(/'/g, "''")}'`).join(',')}]`;
pending = routeWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVField(keysStr),
escapeCSVField(errorKeysStr),
escapeCSVField(middlewareStr),
].join(','),
);
break;
}
case 'Tool':
pending = toolWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVField(node.properties.description || ''),
].join(','),
);
break;
case 'BasicBlock':
pending = basicBlockWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.filePath || ''),
escapeCSVNumber(node.properties.startLine, -1),
escapeCSVNumber(node.properties.endLine, -1),
escapeCSVField(node.properties.text || ''),
].join(','),
);
break;
default: {
// Code element nodes (Function, Class, Interface, CodeElement)
const writer = codeWriterMap[node.label];
if (writer) {
const content = await extractContent(node, contentCache);
pending = mlWriter.addRow(
pending = writer.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVNumber(node.properties.startLine, -1),
escapeCSVNumber(node.properties.endLine, -1),
node.properties.isExported ? 'true' : 'false',
escapeCSVField(content),
escapeCSVField(node.properties.description || ''),
...(node.label === 'Property'
? [escapeCSVField(node.properties.declaredType || '')]
: []),
].join(','),
);
} else {
// Multi-language node types (Struct, Impl, Trait, Macro, etc.)
const mlWriter = multiLangWriters.get(node.label);
if (mlWriter) {
const content = await extractContent(node, contentCache);
pending = mlWriter.addRow(
[
escapeCSVField(node.id),
escapeCSVField(node.properties.name || ''),
escapeCSVField(node.properties.filePath || ''),
escapeCSVNumber(node.properties.startLine, -1),
escapeCSVNumber(node.properties.endLine, -1),
escapeCSVField(content),
escapeCSVField(node.properties.description || ''),
...(node.label === 'Property'
? [escapeCSVField(node.properties.declaredType || '')]
: []),
].join(','),
);
}
}
break;
}
break;
}
}
if (pending) await pending;
}
// Finish all node writers
const allWriters = [
fileWriter,
folderWriter,
functionWriter,
classWriter,
interfaceWriter,
methodWriter,
codeElemWriter,
communityWriter,
processWriter,
sectionWriter,
routeWriter,
toolWriter,
basicBlockWriter,
...multiLangWriters.values(),
];
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();
throw err;
}
// Build result map — only include tables that have rows
const nodeFiles = new Map<NodeTableName, { csvPath: string; rows: number }>();
const tableMap: [NodeTableName, BufferedCSVWriter][] = [
['File', fileWriter],
['Folder', folderWriter],
['Function', functionWriter],
['Class', classWriter],
['Interface', interfaceWriter],
['Method', methodWriter],
['CodeElement', codeElemWriter],
['Community', communityWriter],
['Process', processWriter],
['Section' as NodeTableName, sectionWriter],
['Route' as NodeTableName, routeWriter],
['Tool' as NodeTableName, toolWriter],
['BasicBlock' as NodeTableName, basicBlockWriter],
...Array.from(multiLangWriters.entries()).map(
([name, w]) => [name as NodeTableName, w] as [NodeTableName, BufferedCSVWriter],
),
];
for (const [name, writer] of tableMap) {
if (writer.rows > 0) {
nodeFiles.set(name, {
csvPath: path.join(csvDir, `${name.toLowerCase()}.csv`),
rows: writer.rows,
});
// Finish all node writers
const allWriters = [
fileWriter,
folderWriter,
functionWriter,
classWriter,
interfaceWriter,
methodWriter,
codeElemWriter,
communityWriter,
processWriter,
sectionWriter,
routeWriter,
toolWriter,
basicBlockWriter,
...multiLangWriters.values(),
];
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
const nodeFiles = new Map<NodeTableName, { csvPath: string; rows: number }>();
const tableMap: [NodeTableName, BufferedCSVWriter][] = [
['File', fileWriter],
['Folder', folderWriter],
['Function', functionWriter],
['Class', classWriter],
['Interface', interfaceWriter],
['Method', methodWriter],
['CodeElement', codeElemWriter],
['Community', communityWriter],
['Process', processWriter],
['Section' as NodeTableName, sectionWriter],
['Route' as NodeTableName, routeWriter],
['Tool' as NodeTableName, toolWriter],
['BasicBlock' as NodeTableName, basicBlockWriter],
...Array.from(multiLangWriters.entries()).map(
([name, w]) => [name as NodeTableName, w] as [NodeTableName, BufferedCSVWriter],
),
];
for (const [name, writer] of tableMap) {
if (writer.rows > 0) {
nodeFiles.set(name, {
csvPath: path.join(csvDir, `${name.toLowerCase()}.csv`),
rows: writer.rows,
});
}
}
return {
nodeFiles,
relsByPair: relRouter.byPair,
relHeader: REL_CSV_HEADER,
skippedRels: relRouter.skipped,
totalValidRels: relRouter.total,
};
} finally {
// Restore original process listener limit on every path (success or throw).
process.setMaxListeners(prevMax);
}
// Restore original process listener limit
process.setMaxListeners(prevMax);
return {
nodeFiles,
relsByPair: relRouter.byPair,
relHeader: REL_CSV_HEADER,
skippedRels: relRouter.skipped,
totalValidRels: relRouter.total,
};
};

View file

@ -19,7 +19,7 @@ import {
NodeTableName,
} from './schema.js';
import { streamAllCSVsToDisk } from './csv-generator.js';
import { getNodeLabel as deriveNodeLabel } from './rel-pair-routing.js';
import { getNodeLabel as deriveNodeLabel, type WriteStreamFactory } from './rel-pair-routing.js';
import type { CachedEmbedding } from '../embeddings/types.js';
import { extensionManager, type ExtensionEnsureOptions } from './extension-loader.js';
import {
@ -50,8 +50,10 @@ import { logger } from '../logger.js';
// Relationship CSV splitting — extracted for testability (PR #818)
// ---------------------------------------------------------------------------
/** Factory for creating WriteStreams — injectable for testing. */
export type WriteStreamFactory = (filePath: string) => import('fs').WriteStream;
/** Factory for creating WriteStreams — injectable for testing. Canonical
* definition lives in rel-pair-routing.ts (imported above for local use by
* splitRelCsvByLabelPair); re-exported here to preserve this module's surface. */
export type { WriteStreamFactory };
/** Result of splitting the relationship CSV into per-label-pair files. */
export interface RelCsvSplitResult {
@ -65,6 +67,15 @@ export interface RelCsvSplitResult {
/**
* Split a relationship CSV into per-label-pair files on disk.
*
* @internal RETAINED AS A DIFFERENTIAL ORACLE. As of #2203 U2, production emit
* routes relationships to per-pair files directly during the single pass (see
* RelPairRouter in `rel-pair-routing.ts`), so this function has NO production
* callers — it is kept ONLY so the byte-identity test in
* `test/integration/csv-pipeline.test.ts` ("direct per-pair emit matches the
* split oracle") can diff the direct-emit output against this proven path. Do
* NOT delete it as dead code without also removing that test and accepting the
* loss of the byte-identity guard (and likewise `test/unit/rel-csv-split.test.ts`).
*
* Streams the CSV line-by-line, routing each relationship to a file named
* `rel_{fromLabel}_{toLabel}.csv`. Handles backpressure correctly: only one
* drain listener per stream at a time, and readline resumes only when ALL

View file

@ -73,6 +73,16 @@ export class RelPairRouter {
this.abort.abort(err);
};
/**
* The first stream error observed, if any. Lets the emit caller rethrow the
* real error (EMFILE / disk-full) instead of the generic `AbortError` that a
* pending `once(ws,'drain',{signal})` rejects with when the abort fires —
* mirroring the retained `splitRelCsvByLabelPair`'s `throw streamError ?? err`.
*/
get lastError(): Error | null {
return this.streamError;
}
/**
* 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

View file

@ -4,7 +4,7 @@
* Tests: streamAllCSVsToDisk with real graph data.
* Covers hardening fixes: LRU cache (#24), BufferedCSVWriter flush
*/
import { describe, it, expect, beforeAll, afterAll } from 'vitest';
import { describe, it, expect, beforeAll, beforeEach, afterAll } from 'vitest';
import fs from 'fs/promises';
import { finished } from 'stream/promises';
import path from 'path';
@ -350,9 +350,16 @@ describe('streamAllCSVsToDisk — deterministic output ordering', () => {
* load-bearing guard for "byte-identical graph content" (issue acceptance).
*/
describe('streamAllCSVsToDisk — direct per-pair emit matches the split oracle', () => {
// The oracle always emits in graph.iterRelationships() (unsorted) order; the
// production path honours GITNEXUS_SORT_GRAPH_OUTPUT. Clear it so a value
// leaked from a prior test can't desync the two and produce a spurious diff.
beforeEach(() => {
delete process.env.GITNEXUS_SORT_GRAPH_OUTPUT;
});
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.
// Multiple valid pairs, getNodeLabel special prefixes (comm_ AND 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' },
@ -360,12 +367,16 @@ describe('streamAllCSVsToDisk — direct per-pair emit matches the split oracle'
{ 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: '' },
{ id: 'proc_1', label: 'Process' as never, name: 'p1', filePath: '' },
{ id: 'proc_2', label: 'Process' as never, name: 'p2', 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' },
// proc_ prefix → Process label (getNodeLabel special case).
{ sourceId: 'proc_1', targetId: 'proc_2', type: 'CONTAINS' },
// Invalid FROM label ('Bogus' ∉ NODE_TABLES) — skipped by both paths.
{ sourceId: 'Bogus:x', targetId: 'File:a.ts', type: 'CONTAINS' },
],
@ -400,7 +411,7 @@ describe('streamAllCSVsToDisk — direct per-pair emit matches the split oracle'
// Identical accounting.
expect(direct.totalValidRels).toBe(split.totalValidRels);
expect(direct.totalValidRels).toBe(4);
expect(direct.totalValidRels).toBe(5);
expect(direct.skippedRels).toBe(split.skippedRels);
expect(direct.skippedRels).toBe(1);
expect(direct.relHeader).toBe(split.relHeader);

View file

@ -0,0 +1,175 @@
import { describe, it, expect, beforeEach, afterEach } from 'vitest';
import { EventEmitter } from 'events';
import fs from 'fs';
import path from 'path';
import os from 'os';
import { RelPairRouter, getNodeLabel } from '../../src/core/lbug/rel-pair-routing.js';
/**
* Unit tests for RelPairRouter (#2203 U2) — the production per-pair emit path.
*
* Mirrors test/unit/rel-csv-split.test.ts: drives the router with an injected
* mock WriteStream factory so the error, backpressure, and teardown paths are
* exercised without LadybugDB or real disk streams. These paths are otherwise
* unreachable in the integration suite (which only hits the no-backpressure
* happy path), so this is the coverage for the router's failure modes.
*/
// Controllable backpressure + error injection (same shape as the split oracle's mock).
class MockWriteStream extends EventEmitter {
public chunks: string[] = [];
public destroyed = false;
public ended = false;
public blocked = false;
public maxDrainListenersSeen = 0;
// State flags + events so `stream/promises.finished(ws)` (used by the
// router's close()) resolves against this mock instead of hanging.
public writable = true;
public writableEnded = false;
public writableFinished = false;
write(chunk: string): boolean {
this.chunks.push(chunk);
const count = this.listenerCount('drain');
if (count > this.maxDrainListenersSeen) this.maxDrainListenersSeen = count;
return !this.blocked;
}
end(cb?: (err?: Error) => void): this {
this.ended = true;
this.writableEnded = true;
this.writableFinished = true;
this.writable = false;
if (cb) cb();
queueMicrotask(() => {
this.emit('finish');
this.emit('close');
});
return this;
}
destroy(): this {
this.destroyed = true;
return this;
}
unblock(): void {
this.blocked = false;
this.emit('drain');
}
triggerError(err: Error): void {
this.emit('error', err);
}
}
const HEADER = '"from","to","type","confidence","reason","step"';
const VALID = new Set<string>(['File', 'Function', 'Community', 'Process']);
const row = (from: string, to: string, type = 'CALLS'): string =>
`"${from}","${to}","${type}",1.0,"auto",0`;
function mockFactory(streams: MockWriteStream[], opts?: { blocked?: boolean }) {
return (() => {
const ws = new MockWriteStream();
if (opts?.blocked) ws.blocked = true;
streams.push(ws);
return ws;
}) as unknown as (filePath: string) => import('fs').WriteStream;
}
let tmpDir: string;
beforeEach(() => {
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'rel-pair-routing-test-'));
});
afterEach(() => {
fs.rmSync(tmpDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 50 });
});
describe('getNodeLabel', () => {
it('maps comm_/proc_ prefixes and otherwise splits on the first colon', () => {
expect(getNodeLabel('comm_42')).toBe('Community');
expect(getNodeLabel('proc_7')).toBe('Process');
expect(getNodeLabel('Function:src/a.ts:f:1')).toBe('Function');
expect(getNodeLabel('File:src/a.ts')).toBe('File');
});
});
describe('RelPairRouter', () => {
it('routes valid edges to per-pair files (header first) and skips invalid-label edges', async () => {
const streams: MockWriteStream[] = [];
const router = new RelPairRouter(tmpDir, HEADER, VALID, mockFactory(streams));
const route = async (from: string, to: string) => {
const p = router.route(from, to, row(from, to));
if (p) await p;
};
await route('File:a', 'Function:a:f:1');
await route('File:a', 'Function:a:g:2'); // same pair
await route('Function:a:f:1', 'Function:a:g:2'); // different pair
await route('Bogus:x', 'File:a'); // invalid FROM label → skipped
await router.close();
expect(router.skipped).toBe(1);
expect(router.total).toBe(3);
expect([...router.byPair.keys()].sort()).toEqual(['File|Function', 'Function|Function']);
expect(router.byPair.get('File|Function')!.rows).toBe(2);
// Header is the first chunk written to each pair stream.
expect(streams[0].chunks[0]).toBe(HEADER + '\n');
expect(streams.every((s) => s.ended)).toBe(true);
});
it('returns a drain promise under backpressure and completes once unblocked', async () => {
const streams: MockWriteStream[] = [];
const router = new RelPairRouter(
tmpDir,
HEADER,
VALID,
mockFactory(streams, { blocked: true }),
);
const pending = router.route('File:a', 'Function:a:f:1', row('File:a', 'Function:a:f:1'));
expect(pending).toBeInstanceOf(Promise); // header write hit backpressure
streams[0].unblock();
await pending;
expect(streams[0].maxDrainListenersSeen).toBeLessThanOrEqual(1);
expect(streams[0].chunks[0]).toBe(HEADER + '\n');
expect(router.total).toBe(1);
});
it('on a stream error: route() throws the real error, lastError exposes it, close() rejects + destroys', async () => {
const streams: MockWriteStream[] = [];
const router = new RelPairRouter(tmpDir, HEADER, VALID, mockFactory(streams));
const first = router.route('File:a', 'Function:a:f:1', row('File:a', 'Function:a:f:1'));
if (first) await first;
const err = new Error('EMFILE: too many open files');
streams[0].triggerError(err);
// The next route surfaces the REAL error, not a generic AbortError.
expect(() => router.route('File:a', 'Function:a:g:2', row('File:a', 'Function:a:g:2'))).toThrow(
'EMFILE',
);
expect(router.lastError).toBe(err);
await expect(router.close()).rejects.toThrow('EMFILE');
expect(streams[0].destroyed).toBe(true);
});
it('destroy() tears down every open pair stream', async () => {
const streams: MockWriteStream[] = [];
const router = new RelPairRouter(tmpDir, HEADER, VALID, mockFactory(streams));
const a = router.route('File:a', 'Function:a:f:1', row('File:a', 'Function:a:f:1'));
if (a) await a;
const b = router.route('Community:1', 'Community:2', row('Community:1', 'Community:2'));
if (b) await b;
router.destroy();
expect(streams.length).toBe(2);
expect(streams.every((s) => s.destroyed)).toBe(true);
});
});