From 760f3ce86e0fc124704e33644662c44634a54b27 Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Mon, 15 Jun 2026 17:02:06 +0000 Subject: [PATCH] fix(review): apply autofix feedback (#2203) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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) --- gitnexus/src/core/lbug/csv-generator.ts | 671 +++++++++--------- gitnexus/src/core/lbug/lbug-adapter.ts | 17 +- gitnexus/src/core/lbug/rel-pair-routing.ts | 10 + .../test/integration/csv-pipeline.test.ts | 19 +- gitnexus/test/unit/rel-pair-routing.test.ts | 175 +++++ 5 files changed, 554 insertions(+), 338 deletions(-) create mode 100644 gitnexus/test/unit/rel-pair-routing.test.ts diff --git a/gitnexus/src/core/lbug/csv-generator.ts b/gitnexus/src/core/lbug/csv-generator.ts index aaee437df..f0315f051 100644 --- a/gitnexus/src/core/lbug/csv-generator.ts +++ b/gitnexus/src/core/lbug/csv-generator.ts @@ -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(); - 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 = { - 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(); + // 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 | 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(); + 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 = { + 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(); + + // --- 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 | 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(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(); - 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(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(); + 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, - }; }; diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index c62466ed0..6b5d9c8fa 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -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 diff --git a/gitnexus/src/core/lbug/rel-pair-routing.ts b/gitnexus/src/core/lbug/rel-pair-routing.ts index 2ba3ee2a4..453ffe0b1 100644 --- a/gitnexus/src/core/lbug/rel-pair-routing.ts +++ b/gitnexus/src/core/lbug/rel-pair-routing.ts @@ -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` only when a diff --git a/gitnexus/test/integration/csv-pipeline.test.ts b/gitnexus/test/integration/csv-pipeline.test.ts index 28690587b..74b1abf1f 100644 --- a/gitnexus/test/integration/csv-pipeline.test.ts +++ b/gitnexus/test/integration/csv-pipeline.test.ts @@ -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); diff --git a/gitnexus/test/unit/rel-pair-routing.test.ts b/gitnexus/test/unit/rel-pair-routing.test.ts new file mode 100644 index 000000000..0bf68d84e --- /dev/null +++ b/gitnexus/test/unit/rel-pair-routing.test.ts @@ -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(['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); + }); +});