diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index 164b29b66..fba92465c 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -13,6 +13,144 @@ import { } from './schema.js'; import { streamAllCSVsToDisk } from './csv-generator.js'; +// --------------------------------------------------------------------------- +// Relationship CSV splitting — extracted for testability (PR #818) +// --------------------------------------------------------------------------- + +/** Factory for creating WriteStreams — injectable for testing. */ +export type WriteStreamFactory = (filePath: string) => import('fs').WriteStream; + +/** Result of splitting the relationship CSV into per-label-pair files. */ +export interface RelCsvSplitResult { + relHeader: string; + relsByPairMeta: Map; + pairWriteStreams: Map; + skippedRels: number; + totalValidRels: number; +} + +/** + * Split a relationship CSV into per-label-pair files on disk. + * + * 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 + * backpressured streams have drained. + * + * @param csvPath Path to the combined relationship CSV + * @param csvDir Directory to write per-pair CSV files + * @param validTables Set of valid node table names + * @param getNodeLabel Function to extract the label from a node ID + * @param wsFactory Optional WriteStream factory (defaults to fs.createWriteStream) + */ +export const splitRelCsvByLabelPair = async ( + csvPath: string, + csvDir: string, + validTables: Set, + getNodeLabel: (id: string) => string, + wsFactory: WriteStreamFactory = (p) => createWriteStream(p, 'utf-8'), +): Promise => { + let relHeader = ''; + const relsByPairMeta = new Map(); + const pairWriteStreams = new Map(); + let skippedRels = 0; + let totalValidRels = 0; + + await new Promise((resolve, reject) => { + const inputStream = createReadStream(csvPath, 'utf-8'); + const rl = createInterface({ + input: inputStream, + crlfDelay: Infinity, + }); + + // Track which streams are already waiting for drain to prevent + // listener accumulation. rl.pause() is not synchronous — buffered + // line events continue firing after pause(), and without this guard + // each line targeting the same pairKey would add another drain listener. + const waitingForDrain = new Set(); + + let settled = false; + const cleanup = (err: Error) => { + if (settled) return; + settled = true; + try { + rl.close(); + } catch {} + try { + inputStream.destroy(); + } catch {} + for (const ws of pairWriteStreams.values()) { + try { + ws.destroy(); + } catch {} + } + reject(err); + }; + + let isFirst = true; + rl.on('line', (line) => { + if (isFirst) { + relHeader = line; + isFirst = false; + return; + } + if (!line.trim()) return; + const match = line.match(/"([^"]*)","([^"]*)"/); + if (!match) { + skippedRels++; + return; + } + const fromLabel = getNodeLabel(match[1]); + const toLabel = getNodeLabel(match[2]); + if (!validTables.has(fromLabel) || !validTables.has(toLabel)) { + skippedRels++; + return; + } + const pairKey = `${fromLabel}|${toLabel}`; + let ws = pairWriteStreams.get(pairKey); + if (!ws) { + const pairCsvPath = path.join(csvDir, `rel_${fromLabel}_${toLabel}.csv`); + ws = wsFactory(pairCsvPath); + // If any per-pair WriteStream errors (disk full, EMFILE, etc.), + // tear down everything and reject the Promise. Without this handler, + // a stream error while rl is paused waiting for drain would cause + // the drain callback to never fire and the Promise to hang forever. + ws.on('error', cleanup); + ws.write(relHeader + '\n'); + pairWriteStreams.set(pairKey, ws); + relsByPairMeta.set(pairKey, { csvPath: pairCsvPath, rows: 0 }); + } + const ok = ws.write(line + '\n'); + relsByPairMeta.get(pairKey)!.rows++; + totalValidRels++; + // Handle backpressure: pause reading when the write buffer is full, + // resume when the stream drains. Prevents unbounded memory growth + // on repos with millions of relationships. + // Guard with waitingForDrain to ensure only one drain listener is + // registered per stream at a time — rl.pause() doesn't stop buffered + // line events immediately. Only resume when ALL streams have drained + // to avoid writing into still-full streams. + if (!ok && !waitingForDrain.has(pairKey)) { + waitingForDrain.add(pairKey); + rl.pause(); + ws.once('drain', () => { + waitingForDrain.delete(pairKey); + if (waitingForDrain.size === 0) rl.resume(); + }); + } + }); + rl.on('close', () => { + if (!settled) { + settled = true; + resolve(); + } + }); + rl.on('error', cleanup); + }); + + return { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels }; +}; + let db: lbug.Database | null = null; let conn: lbug.Connection | null = null; let currentDbPath: string | null = null; @@ -247,105 +385,8 @@ export const loadGraphToLbug = async ( } // Bulk COPY relationships — split by FROM→TO label pair (LadybugDB requires it) - // Stream-read the relation CSV line by line and write directly to per-pair - // temp files on disk. This avoids accumulating potentially millions of CSV - // lines in memory which could exceed V8 Map or array limits on large repos. - let relHeader = ''; - const relsByPairMeta = new Map(); - const pairWriteStreams = new Map(); - let skippedRels = 0; - let totalValidRels = 0; - - await new Promise((resolve, reject) => { - const inputStream = createReadStream(csvResult.relCsvPath, 'utf-8'); - const rl = createInterface({ - input: inputStream, - crlfDelay: Infinity, - }); - - // Track which streams are already waiting for drain to prevent - // listener accumulation. rl.pause() is not synchronous — buffered - // line events continue firing after pause(), and without this guard - // each line targeting the same pairKey would add another drain listener. - const waitingForDrain = new Set(); - - let settled = false; - const cleanup = (err: Error) => { - if (settled) return; - settled = true; - try { - rl.close(); - } catch {} - try { - inputStream.destroy(); - } catch {} - for (const ws of pairWriteStreams.values()) { - try { - ws.destroy(); - } catch {} - } - reject(err); - }; - - let isFirst = true; - rl.on('line', (line) => { - if (isFirst) { - relHeader = line; - isFirst = false; - return; - } - if (!line.trim()) return; - const match = line.match(/"([^"]*)","([^"]*)"/); - if (!match) { - skippedRels++; - return; - } - const fromLabel = getNodeLabel(match[1]); - const toLabel = getNodeLabel(match[2]); - if (!validTables.has(fromLabel) || !validTables.has(toLabel)) { - skippedRels++; - return; - } - const pairKey = `${fromLabel}|${toLabel}`; - let ws = pairWriteStreams.get(pairKey); - if (!ws) { - const pairCsvPath = path.join(csvDir, `rel_${fromLabel}_${toLabel}.csv`); - ws = createWriteStream(pairCsvPath, 'utf-8'); - // If any per-pair WriteStream errors (disk full, EMFILE, etc.), - // tear down everything and reject the Promise. Without this handler, - // a stream error while rl is paused waiting for drain would cause - // the drain callback to never fire and the Promise to hang forever. - ws.on('error', cleanup); - ws.write(relHeader + '\n'); - pairWriteStreams.set(pairKey, ws); - relsByPairMeta.set(pairKey, { csvPath: pairCsvPath, rows: 0 }); - } - const ok = ws.write(line + '\n'); - relsByPairMeta.get(pairKey)!.rows++; - totalValidRels++; - // Handle backpressure: pause reading when the write buffer is full, - // resume when the stream drains. Prevents unbounded memory growth - // on repos with millions of relationships. - // Guard with waitingForDrain to ensure only one drain listener is - // registered per stream at a time — rl.pause() doesn't stop buffered - // line events immediately. - if (!ok && !waitingForDrain.has(pairKey)) { - waitingForDrain.add(pairKey); - rl.pause(); - ws.once('drain', () => { - waitingForDrain.delete(pairKey); - rl.resume(); - }); - } - }); - rl.on('close', () => { - if (!settled) { - settled = true; - resolve(); - } - }); - rl.on('error', cleanup); - }); + const { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels } = + await splitRelCsvByLabelPair(csvResult.relCsvPath, csvDir, validTables, getNodeLabel); // Close all per-pair write streams before COPY await Promise.all( diff --git a/gitnexus/test/unit/rel-csv-split.test.ts b/gitnexus/test/unit/rel-csv-split.test.ts index 4e8eb0654..dab38a913 100644 --- a/gitnexus/test/unit/rel-csv-split.test.ts +++ b/gitnexus/test/unit/rel-csv-split.test.ts @@ -3,13 +3,14 @@ import { EventEmitter } from 'events'; import fs from 'fs'; import path from 'path'; import os from 'os'; +import { splitRelCsvByLabelPair } from '../../src/core/lbug/lbug-adapter.js'; /** - * Regression tests for the relationship CSV splitting logic in lbug-adapter.ts. + * Regression tests for splitRelCsvByLabelPair (PR #818). * - * These tests exercise the backpressure, error handling, and drain-listener - * guard that were fixed in PR #818. They use a mock WriteStream to simulate - * backpressure and error conditions without touching real LadybugDB. + * These tests call the real exported function from lbug-adapter.ts with a + * mock WriteStream factory, exercising the actual backpressure, error + * handling, and drain-listener guard without touching LadybugDB. */ // --------------------------------------------------------------------------- @@ -56,114 +57,6 @@ class MockWriteStream extends EventEmitter { } } -// --------------------------------------------------------------------------- -// Minimal reimplementation of the rel-CSV split logic under test. -// Mirrors the exact pattern in loadGraphToLbug() so these tests break if -// the production code regresses. -// --------------------------------------------------------------------------- -interface SplitResult { - relHeader: string; - pairMeta: Map; - skippedRels: number; - totalValidRels: number; -} - -type WriteStreamFactory = (path: string) => MockWriteStream; - -async function splitRelCsvByLabelPair( - csvPath: string, - csvDir: string, - validTables: Set, - getNodeLabel: (id: string) => string, - wsFactory: WriteStreamFactory, -): Promise { - const { createReadStream } = await import('fs'); - const { createInterface } = await import('readline'); - - let relHeader = ''; - const pairMeta = new Map(); - const pairStreams = new Map(); - let skippedRels = 0; - let totalValidRels = 0; - - await new Promise((resolve, reject) => { - const inputStream = createReadStream(csvPath, 'utf-8'); - const rl = createInterface({ input: inputStream, crlfDelay: Infinity }); - - const waitingForDrain = new Set(); - - let settled = false; - const cleanup = (err: Error) => { - if (settled) return; - settled = true; - try { - rl.close(); - } catch {} - try { - inputStream.destroy(); - } catch {} - for (const ws of pairStreams.values()) { - try { - ws.destroy(); - } catch {} - } - reject(err); - }; - - let isFirst = true; - rl.on('line', (line) => { - if (isFirst) { - relHeader = line; - isFirst = false; - return; - } - if (!line.trim()) return; - const match = line.match(/"([^"]*)","([^"]*)"/); - if (!match) { - skippedRels++; - return; - } - const fromLabel = getNodeLabel(match[1]); - const toLabel = getNodeLabel(match[2]); - if (!validTables.has(fromLabel) || !validTables.has(toLabel)) { - skippedRels++; - return; - } - const pairKey = `${fromLabel}|${toLabel}`; - let ws = pairStreams.get(pairKey); - if (!ws) { - ws = wsFactory(path.join(csvDir, `rel_${fromLabel}_${toLabel}.csv`)); - ws.on('error', cleanup); - ws.write(relHeader + '\n'); - pairStreams.set(pairKey, ws); - pairMeta.set(pairKey, { rows: 0 }); - } - ws.write(line + '\n'); - pairMeta.get(pairKey)!.rows++; - totalValidRels++; - - const ok = !ws.blocked; - if (!ok && !waitingForDrain.has(pairKey)) { - waitingForDrain.add(pairKey); - rl.pause(); - ws.once('drain', () => { - waitingForDrain.delete(pairKey); - rl.resume(); - }); - } - }); - rl.on('close', () => { - if (!settled) { - settled = true; - resolve(); - } - }); - rl.on('error', cleanup); - }); - - return { relHeader, pairMeta, skippedRels, totalValidRels }; -} - // --------------------------------------------------------------------------- // Test helpers // --------------------------------------------------------------------------- @@ -177,6 +70,16 @@ function getNodeLabel(id: string): string { return id.split(':')[0]; } +/** Cast MockWriteStream factory to the real WriteStreamFactory type. */ +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(() => { @@ -208,26 +111,29 @@ describe('splitRelCsvByLabelPair', () => { ]); const streams: MockWriteStream[] = []; - const result = await splitRelCsvByLabelPair(csvPath, tmpDir, validTables, getNodeLabel, () => { - const ws = new MockWriteStream(); - streams.push(ws); - return ws; - }); - - expect(result.totalValidRels).toBe(3); - expect(result.pairMeta.get('Function|Class')?.rows).toBe(2); - expect(result.pairMeta.get('File|Method')?.rows).toBe(1); - }); - - it('captures the CSV header in relHeader', async () => { - const csvPath = writeCsv([HEADER, csvLine('Function:a', 'Class:b')]); - const result = await splitRelCsvByLabelPair( csvPath, tmpDir, validTables, getNodeLabel, - () => new MockWriteStream(), + mockFactory(streams), + ); + + expect(result.totalValidRels).toBe(3); + expect(result.relsByPairMeta.get('Function|Class')?.rows).toBe(2); + expect(result.relsByPairMeta.get('File|Method')?.rows).toBe(1); + }); + + it('captures the CSV header in relHeader', async () => { + const csvPath = writeCsv([HEADER, csvLine('Function:a', 'Class:b')]); + + const streams: MockWriteStream[] = []; + const result = await splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams), ); expect(result.relHeader).toBe(HEADER); @@ -241,12 +147,13 @@ describe('splitRelCsvByLabelPair', () => { csvLine('Function:c', 'Bogus:d'), ]); + const streams: MockWriteStream[] = []; const result = await splitRelCsvByLabelPair( csvPath, tmpDir, validTables, getNodeLabel, - () => new MockWriteStream(), + mockFactory(streams), ); expect(result.totalValidRels).toBe(1); @@ -256,12 +163,13 @@ describe('splitRelCsvByLabelPair', () => { it('ignores blank lines without counting them as skipped', async () => { const csvPath = writeCsv([HEADER, '', csvLine('Function:a', 'Class:b'), '', '']); + const streams: MockWriteStream[] = []; const result = await splitRelCsvByLabelPair( csvPath, tmpDir, validTables, getNodeLabel, - () => new MockWriteStream(), + mockFactory(streams), ); expect(result.totalValidRels).toBe(1); @@ -269,7 +177,6 @@ describe('splitRelCsvByLabelPair', () => { }); it('registers at most 1 drain listener per stream under heavy backpressure', async () => { - // Create many lines targeting the same pair — all will hit backpressure const lines = [HEADER]; for (let i = 0; i < 50; i++) { lines.push(csvLine(`Function:f${i}`, `Class:c${i}`)); @@ -277,15 +184,13 @@ describe('splitRelCsvByLabelPair', () => { const csvPath = writeCsv(lines); const streams: MockWriteStream[] = []; - const factory = () => { - const ws = new MockWriteStream(); - ws.blocked = true; // permanent backpressure - streams.push(ws); - return ws; - }; - - // Start the split — it will pause on first backpressure - const promise = splitRelCsvByLabelPair(csvPath, tmpDir, validTables, getNodeLabel, factory); + const promise = splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams, { blocked: true }), + ); // Give readline time to buffer and fire lines await new Promise((r) => setTimeout(r, 50)); @@ -304,14 +209,13 @@ describe('splitRelCsvByLabelPair', () => { const csvPath = writeCsv([HEADER, csvLine('Function:a', 'Class:b')]); const streams: MockWriteStream[] = []; - const factory = () => { - const ws = new MockWriteStream(); - ws.blocked = true; // hold the promise open via backpressure - streams.push(ws); - return ws; - }; - - const promise = splitRelCsvByLabelPair(csvPath, tmpDir, validTables, getNodeLabel, factory); + const promise = splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams, { blocked: true }), + ); // Wait for readline to process, then error while paused on drain await new Promise((r) => setTimeout(r, 50)); @@ -322,7 +226,6 @@ describe('splitRelCsvByLabelPair', () => { }); it('destroys all streams when one errors (no lingering FDs)', async () => { - // Use enough lines to create two different pair streams const lines = [HEADER]; for (let i = 0; i < 10; i++) { lines.push(csvLine(`Function:f${i}`, `Class:c${i}`)); @@ -331,14 +234,13 @@ describe('splitRelCsvByLabelPair', () => { const csvPath = writeCsv(lines); const streams: MockWriteStream[] = []; - const factory = () => { - const ws = new MockWriteStream(); - ws.blocked = true; // hold all streams open via backpressure - streams.push(ws); - return ws; - }; - - const promise = splitRelCsvByLabelPair(csvPath, tmpDir, validTables, getNodeLabel, factory); + const promise = splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams, { blocked: true }), + ); // Wait for readline to process and create streams await new Promise((r) => setTimeout(r, 50)); @@ -355,12 +257,13 @@ describe('splitRelCsvByLabelPair', () => { it('handles empty CSV (header only) without errors', async () => { const csvPath = writeCsv([HEADER]); + const streams: MockWriteStream[] = []; const result = await splitRelCsvByLabelPair( csvPath, tmpDir, validTables, getNodeLabel, - () => new MockWriteStream(), + mockFactory(streams), ); expect(result.totalValidRels).toBe(0);