diff --git a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts index 00390d183..99ef89ec1 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts @@ -338,6 +338,34 @@ export async function runChunkedParseAndResolve( let chunkCacheMisses = 0; try { + // U1 — bounded chunk concurrency (B1 from PR #1693 review): pre-fetch + // chunk file contents up to `parseChunkConcurrency` chunks ahead of the + // dispatch cursor so file I/O overlaps with worker compute. Worker + // dispatch itself stays serial because `WorkerPool.dispatch` is not + // reentrant (concurrent calls would race on the shared per-slot + // busy/in-flight state). With concurrency=1 behavior is identical to + // the pure-serial loop. F4: deferred-state aggregation still happens + // in chunkIdx order (the for-loop below iterates sequentially), so + // cross-chunk processors see deterministic input regardless of + // file-read completion order. Honors options.parseChunkConcurrency + // (threaded from the CLI), then GITNEXUS_PARSE_CHUNK_CONCURRENCY env + // (default 2 — matches the help text the CLI advertises). + const parseChunkConcurrency = ((): number => { + const opt = options?.parseChunkConcurrency; + if (typeof opt === 'number' && Number.isInteger(opt) && opt >= 1) return opt; + const env = Number(process.env.GITNEXUS_PARSE_CHUNK_CONCURRENCY); + if (Number.isInteger(env) && env >= 1) return env; + return 2; + })(); + const chunkContentPromises = new Array> | undefined>(numChunks); + const startChunkPrefetch = (i: number): void => { + if (i >= numChunks || chunkContentPromises[i] !== undefined) return; + chunkContentPromises[i] = readFileContents(repoPath, chunks[i]); + }; + for (let i = 0; i < Math.min(parseChunkConcurrency, numChunks); i++) { + startChunkPrefetch(i); + } + for (let chunkIdx = 0; chunkIdx < numChunks; chunkIdx++) { const chunkPaths = chunks[chunkIdx]; // Start wall-clock for the per-chunk throughput log emitted at end @@ -350,7 +378,9 @@ export async function runChunkedParseAndResolve( const verboseThroughputLog = isDev || isVerboseIngestionEnabled(); const chunkStartMs: number | null = verboseThroughputLog ? Date.now() : null; - const chunkContents = await readFileContents(repoPath, chunkPaths); + const chunkContents = await chunkContentPromises[chunkIdx]!; + chunkContentPromises[chunkIdx] = undefined; // release the in-memory copy + startChunkPrefetch(chunkIdx + parseChunkConcurrency); const chunkFiles = chunkPaths .filter((p) => chunkContents.has(p)) .map((p) => ({ path: p, content: chunkContents.get(p)! })); diff --git a/gitnexus/src/core/ingestion/pipeline.ts b/gitnexus/src/core/ingestion/pipeline.ts index 65ef049f5..eefbd87db 100644 --- a/gitnexus/src/core/ingestion/pipeline.ts +++ b/gitnexus/src/core/ingestion/pipeline.ts @@ -81,6 +81,20 @@ export interface PipelineOptions { * `process.env` state across analyze invocations. */ workerPoolSize?: number; + /** + * Number of chunks whose file contents may be read into memory in + * parallel while the worker pool is busy dispatching the current + * chunk. Pre-fetching overlaps disk I/O for chunk N+1..N+K with the + * worker compute on chunk N — modest but real wall-clock win on + * repos large enough to chunk. Worker dispatch itself remains serial + * because `WorkerPool.dispatch` is not reentrant (concurrent calls + * would race on the shared per-slot busy/in-flight state). + * + * `1` matches today's pure-serial behavior; `2` is the documented + * default (`GITNEXUS_PARSE_CHUNK_CONCURRENCY`). Falls back to the + * env var when undefined; defaults to 2 when neither is set. + */ + parseChunkConcurrency?: number; } // ── Phase registry ───────────────────────────────────────────────────────── diff --git a/gitnexus/test/unit/parse-impl-chunk-concurrency.test.ts b/gitnexus/test/unit/parse-impl-chunk-concurrency.test.ts new file mode 100644 index 000000000..22223e544 --- /dev/null +++ b/gitnexus/test/unit/parse-impl-chunk-concurrency.test.ts @@ -0,0 +1,131 @@ +/** + * U1 — Bounded chunk concurrency (B1 from PR #1693 review). + * + * Verifies that the new `parseChunkConcurrency` PipelineOption (and the + * paired `GITNEXUS_PARSE_CHUNK_CONCURRENCY` env-var fallback) flow through + * `runChunkedParseAndResolve` without changing graph output. Pre-fetching + * chunk file contents up to N chunks ahead of the worker-dispatch cursor + * is a wall-clock optimization (file I/O overlaps with worker compute), + * not a graph-semantics change — the deferred-state aggregation still + * runs in `chunkIdx` order so cross-chunk processors see deterministic + * input regardless of file-read completion order. + */ +import { describe, it, expect, beforeEach, afterEach } from 'vitest'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; + +import { runChunkedParseAndResolve } from '../../src/core/ingestion/pipeline-phases/parse-impl.js'; +import { createKnowledgeGraph } from '../../src/core/graph/graph.js'; + +function scanned(repo: string, files: string[]) { + return files.map((rel) => ({ + path: rel, + size: fs.statSync(path.join(repo, rel)).size, + })); +} + +describe('parse-impl chunk concurrency (U1)', () => { + let repoPath = ''; + + beforeEach(() => { + repoPath = fs.mkdtempSync(path.join(os.tmpdir(), 'parse-impl-chunk-concurrency-')); + fs.writeFileSync(path.join(repoPath, 'a.ts'), 'export function foo() { return 1; }\n'); + fs.writeFileSync( + path.join(repoPath, 'b.ts'), + 'import { foo } from "./a";\nexport function bar() { return foo(); }\n', + ); + fs.writeFileSync( + path.join(repoPath, 'c.ts'), + 'import { bar } from "./b";\nexport class Baz { run() { return bar(); } }\n', + ); + }); + + afterEach(() => { + if (repoPath && fs.existsSync(repoPath)) { + fs.rmSync(repoPath, { recursive: true, force: true }); + } + }); + + it('produces identical graph output across parseChunkConcurrency values', async () => { + const files = ['a.ts', 'b.ts', 'c.ts']; + const scan = scanned(repoPath, files); + + const g1 = createKnowledgeGraph(); + await runChunkedParseAndResolve(g1, scan, files, files.length, repoPath, Date.now(), () => {}, { + skipWorkers: true, + parseChunkConcurrency: 1, + }); + + const g2 = createKnowledgeGraph(); + await runChunkedParseAndResolve(g2, scan, files, files.length, repoPath, Date.now(), () => {}, { + skipWorkers: true, + parseChunkConcurrency: 2, + }); + + // Same fixture under different concurrency values must produce the + // same graph — F4 (wildcard-synthesis ordering): per-chunk results + // merge in chunkIdx order regardless of file-read completion order, + // so cross-chunk processors see deterministic input. + expect(g2.nodeCount).toBe(g1.nodeCount); + expect(g2.relationshipCount).toBe(g1.relationshipCount); + }); + + it('accepts parseChunkConcurrency=1 (serial-equivalent) and produces the expected fixture symbols', async () => { + const files = ['a.ts', 'b.ts', 'c.ts']; + const graph = createKnowledgeGraph(); + await runChunkedParseAndResolve( + graph, + scanned(repoPath, files), + files, + files.length, + repoPath, + Date.now(), + () => {}, + { skipWorkers: true, parseChunkConcurrency: 1 }, + ); + // Exact assertions per DoD §2.7: pin specific symbols from the fixture + // so a regression in either the chunk loop or the resolver surfaces + // here instead of being masked by a bounds-only nodeCount check. + const symbolNames = Array.from(graph.nodes.values()).map( + (n) => (n.properties as { name?: string } | undefined)?.name, + ); + expect(symbolNames.includes('foo')).toBe(true); + expect(symbolNames.includes('bar')).toBe(true); + expect(symbolNames.includes('Baz')).toBe(true); + }); + + it('falls back to GITNEXUS_PARSE_CHUNK_CONCURRENCY env when option is undefined', async () => { + const original = process.env.GITNEXUS_PARSE_CHUNK_CONCURRENCY; + process.env.GITNEXUS_PARSE_CHUNK_CONCURRENCY = '3'; + try { + const files = ['a.ts', 'b.ts', 'c.ts']; + const graph = createKnowledgeGraph(); + await runChunkedParseAndResolve( + graph, + scanned(repoPath, files), + files, + files.length, + repoPath, + Date.now(), + () => {}, + { skipWorkers: true }, + ); + // Resolver reads the env when options.parseChunkConcurrency is + // undefined. The env value (3) must produce the same fixture + // symbols on this fixture as the other concurrency values do. + const symbolNames = Array.from(graph.nodes.values()).map( + (n) => (n.properties as { name?: string } | undefined)?.name, + ); + expect(symbolNames.includes('foo')).toBe(true); + expect(symbolNames.includes('bar')).toBe(true); + expect(symbolNames.includes('Baz')).toBe(true); + } finally { + if (original === undefined) { + delete process.env.GITNEXUS_PARSE_CHUNK_CONCURRENCY; + } else { + process.env.GITNEXUS_PARSE_CHUNK_CONCURRENCY = original; + } + } + }); +});