fix(ingestion): address #2038 tri-review findings (parse-phase memory)

Resolves the confirmed review findings on PR #2038:

- P1: thread exportedTypeMap through the sequential parse path
  (processParsingSequential) so a no-worker run over a partially-warm
  cache no longer silently drops the sequential-miss files' exported
  types. Cache hits made exportedTypeMap.size > 0, suppressing the
  end-of-loop buildExportedTypeMapFromGraph rebuild, but the sequential
  path never populated the map. Regression test added (fails on the
  pre-fix tree, passes after) plus a fully-sequential differential oracle.
- P2: saveParseCache builds its on-disk index from hashes actually
  written/copied (writtenKeys), never a usedKeys hash whose shard write
  or copy was skipped — no more phantom index entries.
- P2: add a unit test asserting SCOPE_RESOLUTION_LANGUAGES stays in sync
  with SCOPE_RESOLVERS (asymmetric drift would lose a language's ParsedFile).
- Backfill cache coverage: loadParseCacheChunk missing/corrupt -> undefined,
  pruneCache onDiskKeys branch, slim preserves nodes, saveParseCache
  copy-evicted-shard round-trip.
- Cleanups: single-source heap-probe gating via isDebugHeapEnabled();
  hoist the per-chunk mkdir in persistParseCacheChunk behind a
  process-scoped Set; gate COBOL's unused worker-side ParsedFile
  extraction (graph nodes still come from cobolPhase) while keeping
  fileCount/progress unconditional.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Gergo Magyar 2026-06-04 18:17:35 +00:00
parent 3e5ba36bd7
commit 9b6650822f
7 changed files with 284 additions and 24 deletions

View file

@ -396,6 +396,7 @@ const processParsingSequential = async (
astCache: ASTCache,
scopeTreeCache: ASTCache | undefined,
onFileProgress?: FileProgressCallback,
exportedTypeMap?: ExportedTypeMap,
) => {
const parser = await loadParser();
const total = files.length;
@ -937,6 +938,21 @@ const processParsingSequential = async (
qualifiedName: qualifiedTypeName,
});
// #1983: populate the incremental ExportedTypeMap on the sequential path
// too. Without this, a no-worker run over a partially-warm cache silently
// drops the sequential-miss files' exported types: cache-hit chunks make
// `exportedTypeMap.size > 0`, which suppresses the end-of-loop
// `buildExportedTypeMapFromGraph` rebuild, but the sequential path never
// populated the map. Placed right after `symbolTable.add` so accumulate's
// `lookupExactAll` finds this node's own def — the node depends only on
// its own symbol, so per-node placement is order-safe (unlike the worker
// path's two-pass shape in `mergeChunkResults`). `graph.addNode` above is
// the sole node-creation site in this path, so this covers exactly the
// node set the full-graph rebuild would scan.
if (exportedTypeMap) {
accumulateExportedTypesFromParsedNode(exportedTypeMap, node, symbolTable);
}
// ── HAS_METHOD / HAS_PROPERTY: link member to enclosing class ──
const ownerIdForMemberEdge = enclosingClassId ?? objectLiteralOwnerInfo?.ownerId ?? null;
@ -1111,7 +1127,9 @@ export const processParsing = async (
return data;
}
// Fallback: sequential parsing (no pre-extracted data)
// Fallback: sequential parsing (no pre-extracted data). Thread the
// exportedTypeMap so no-worker runs populate it during the chunk loop (#1983)
// — see the accumulate call in processParsingSequential.
await processParsingSequential(
graph,
files,
@ -1119,6 +1137,7 @@ export const processParsing = async (
astCache,
scopeTreeCache,
reportProgress,
exportedTypeMap,
);
return null;
};

View file

@ -64,7 +64,7 @@ import fs from 'node:fs';
import path from 'node:path';
import { fileURLToPath, pathToFileURL } from 'node:url';
import { isDev, parseTruthyEnv } from '../utils/env.js';
import { isDev } from '../utils/env.js';
import { isVerboseIngestionEnabled } from '../utils/verbose.js';
import {
endTimer,
@ -72,7 +72,7 @@ import {
logDeferredProfile,
startTimer,
} from '../utils/deferred-resolution-profile.js';
import { logHeapProbe } from '../utils/heap-probe.js';
import { isDebugHeapEnabled, logHeapProbe } from '../utils/heap-probe.js';
import { extractORMQueriesInline } from './orm-extraction.js';
import { logger } from '../../logger.js';
@ -476,10 +476,7 @@ export async function runChunkedParseAndResolve(
// body, which re-read process.env on every iteration even though
// the env can't change mid-run.
const verboseThroughputLog = isDev || isVerboseIngestionEnabled();
const heapProbeEveryN =
parseTruthyEnv(process.env.GITNEXUS_DEBUG_HEAP) || isDeferredResolutionProfileEnabled()
? 25
: 0;
const heapProbeEveryN = isDebugHeapEnabled() ? 25 : 0;
for (let chunkIdx = 0; chunkIdx < numChunks; chunkIdx++) {
if (heapProbeEveryN > 0 && chunkIdx > 0 && chunkIdx % heapProbeEveryN === 0) {

View file

@ -839,24 +839,36 @@ const processBatch = (
// extractParsedFile directly — no tree-sitter involved.
if (provider.emitScopeCaptures) {
for (const file of langFiles) {
const parsedFile = extractParsedFile(
provider,
file.content,
file.path,
(message) => {
if (parentPort) {
parentPort.postMessage({ type: 'warning', message });
} else {
logger.warn(message);
}
},
undefined, // no cachedTree for standalone providers
);
// #1983: skip building the ParsedFile for registry-primary
// (scope-resolution) languages — the scope-resolution phase
// re-extracts from source on the main thread, so the worker copy is
// unused work + retained RAM. COBOL is a standalone provider that is
// in SCOPE_RESOLUTION_LANGUAGES; its graph nodes come from cobolPhase,
// not from this ParsedFile, so gating here is safe.
const parsedFile = isScopeResolutionLanguage(language)
? undefined
: extractParsedFile(
provider,
file.content,
file.path,
(message) => {
if (parentPort) {
parentPort.postMessage({ type: 'warning', message });
} else {
logger.warn(message);
}
},
undefined, // no cachedTree for standalone providers
);
if (parsedFile !== undefined) {
result.parsedFiles.push(parsedFile);
result.fileCount++;
onFileProcessed?.();
}
// fileCount / progress fire per file regardless of whether a
// ParsedFile was produced — matching the tree-sitter branch (which
// increments fileCount outside the parsedFile gate). Otherwise gated
// COBOL files would vanish from worker progress counts.
result.fileCount++;
onFileProcessed?.();
}
}
continue;

View file

@ -227,6 +227,14 @@ export const loadParseCacheChunk = async (
return undefined;
};
/**
* Cache directories already created this process. `persistParseCacheChunk` runs
* once per cache-miss chunk; without this guard every miss re-issues a redundant
* `mkdir` syscall (hundreds on a large cold repo) (#1983). Storage paths are
* process-scoped, so the Set stays bounded.
*/
const createdCacheDirs = new Set<string>();
/**
* Persist one chunk shard and avoid retaining it in RAM for the rest of the
* run. Falls back to `cache.entries` when `storagePath` is unset (unit tests).
@ -238,7 +246,11 @@ export const persistParseCacheChunk = async (
): Promise<void> => {
const slim = slimParseWorkerResultsForCache(chunkResults);
if (cache.storagePath) {
await fs.mkdir(getCacheDirPath(cache.storagePath), { recursive: true });
const cacheDir = getCacheDirPath(cache.storagePath);
if (!createdCacheDirs.has(cacheDir)) {
await fs.mkdir(cacheDir, { recursive: true });
createdCacheDirs.add(cacheDir);
}
const payload = JSON.stringify(slim, mapReplacer);
await fs.writeFile(getCacheChunkPath(cache.storagePath, chunkHash), payload, 'utf-8');
cache.onDiskKeys ??= new Set<string>();
@ -335,6 +347,13 @@ export const saveParseCache = async (storagePath: string, cache: ParseCache): Pr
await fs.mkdir(tmpDir, { recursive: true });
const keys = [...cache.usedKeys].filter(isValidChunkCacheKey).sort();
// Track hashes whose shard was actually written/copied this save. A hash can
// be in `usedKeys` without a backing shard — its in-memory serialize threw, or
// its on-disk copy failed/was-absent (e.g. a worker-quarantined chunk added to
// usedKeys but never persisted). Writing such a hash into `index.keys` would
// make the next load reference a shard that doesn't exist (#1983). Build the
// index from what we persisted, not from the raw usedKeys snapshot.
const writtenKeys: string[] = [];
for (const chunkHash of keys) {
const chunkPath = path.join(tmpDir, `${chunkHash}.json`);
const inMemory = cache.entries.get(chunkHash);
@ -346,11 +365,13 @@ export const saveParseCache = async (storagePath: string, cache: ParseCache): Pr
continue;
}
await fs.writeFile(chunkPath, payload, 'utf-8');
writtenKeys.push(chunkHash);
continue;
}
const existingPath = getCacheChunkPath(storagePath, chunkHash);
try {
await fs.copyFile(existingPath, chunkPath);
writtenKeys.push(chunkHash);
} catch {
/* shard missing — skip; next run treats as cache miss */
}
@ -358,7 +379,7 @@ export const saveParseCache = async (storagePath: string, cache: ParseCache): Pr
const index: ShardedParseCacheIndex = {
version: cache.version,
keys,
keys: writtenKeys,
};
await fs.writeFile(path.join(tmpDir, CACHE_INDEX_FILENAME), JSON.stringify(index), 'utf-8');

View file

@ -135,6 +135,18 @@ describe('pruneCache', () => {
expect(pruneCache(cache, cache.usedKeys)).toBe(0);
expect(cache.entries.size).toBe(2);
});
it('drops onDiskKeys entries not in the used-set and counts them', () => {
const cache: ParseCache = {
version: PARSE_CACHE_VERSION,
entries: new Map<string, ParseWorkerResult[]>(),
usedKeys: new Set<string>(['disk-A']),
onDiskKeys: new Set<string>(['disk-A', 'disk-B', 'disk-C']),
};
const removed = pruneCache(cache, new Set(['disk-A']));
expect(removed).toBe(2);
expect([...(cache.onDiskKeys ?? [])].sort()).toEqual(['disk-A']);
});
});
describe('loadParseCache / saveParseCache (round-trip)', () => {
@ -249,6 +261,10 @@ describe('loadParseCache / saveParseCache (round-trip)', () => {
expect(loaded.onDiskKeys?.size).toBe(3);
const chunk = await loadParseCacheChunk(loaded, goodKey);
expect(chunk?.[0]?.fileCount).toBe(3);
// A shard listed in the index but absent on disk, and a corrupt-JSON
// shard, both resolve to undefined (graceful cache miss) — not a throw.
expect(await loadParseCacheChunk(loaded, missingKey)).toBeUndefined();
expect(await loadParseCacheChunk(loaded, badKey)).toBeUndefined();
} finally {
await rm(dir, { recursive: true, force: true });
}
@ -476,6 +492,23 @@ describe('loadParseCache / saveParseCache (round-trip)', () => {
expect(slim.fileCount).toBe(raw.fileCount);
});
it('slimParseWorkerResultsForCache preserves nodes (incremental exportedTypeMap depends on them)', () => {
const raw = minimalResult({
nodes: [
{
id: 'Function:a.ts:foo',
label: 'Function',
properties: { name: 'foo', filePath: 'a.ts', isExported: true },
},
] as ParseWorkerResult['nodes'],
});
const slim = slimParseWorkerResultsForCache([raw])[0];
// `nodes` (and `symbols`) must survive slimming — on a warm cache hit they
// are what mergeChunkResults replays to rebuild the ExportedTypeMap.
expect(slim.nodes).toEqual(raw.nodes);
expect(slim.nodes).toHaveLength(1);
});
it('persistParseCacheChunk writes to disk without retaining in-memory entries', async () => {
const dir = await mkdtemp(path.join(tmpdir(), 'gnx-pc-'));
try {
@ -496,4 +529,50 @@ describe('loadParseCache / saveParseCache (round-trip)', () => {
await rm(dir, { recursive: true, force: true });
}
});
it('saveParseCache excludes a usedKeys hash whose shard was never persisted (no phantom index key)', async () => {
const dir = await mkdtemp(path.join(tmpdir(), 'gnx-pc-'));
try {
const realKey = 'a'.repeat(64);
const phantomKey = 'b'.repeat(64); // in usedKeys but has no entry and no on-disk shard
const cache: ParseCache = {
version: PARSE_CACHE_VERSION,
entries: new Map([[realKey, [minimalResult({ fileCount: 3 })]]]),
usedKeys: new Set([realKey, phantomKey]),
};
await saveParseCache(dir, cache);
const loaded = await loadParseCache(dir);
expect(loaded.onDiskKeys?.has(realKey)).toBe(true);
// The phantom key was never written, so it must not appear in the index.
expect(loaded.onDiskKeys?.has(phantomKey)).toBe(false);
expect((await loadParseCacheChunk(loaded, realKey))?.[0]?.fileCount).toBe(3);
expect(await loadParseCacheChunk(loaded, phantomKey)).toBeUndefined();
} finally {
await rm(dir, { recursive: true, force: true });
}
});
it('saveParseCache copies a persisted-but-evicted shard (copyFile branch) and round-trips', async () => {
const dir = await mkdtemp(path.join(tmpdir(), 'gnx-pc-'));
try {
const key = 'c'.repeat(64);
const cache: ParseCache = {
version: PARSE_CACHE_VERSION,
entries: new Map(),
usedKeys: new Set([key]),
storagePath: dir,
onDiskKeys: new Set(),
};
// persist writes the shard to the live dir and evicts it from `entries`,
// so saveParseCache must hit the copyFile branch to carry it forward.
await persistParseCacheChunk(cache, key, [minimalResult({ fileCount: 42 })]);
expect(cache.entries.has(key)).toBe(false);
await saveParseCache(dir, cache);
const loaded = await loadParseCache(dir);
expect(loaded.onDiskKeys?.has(key)).toBe(true);
expect((await loadParseCacheChunk(loaded, key))?.[0]?.fileCount).toBe(42);
} finally {
await rm(dir, { recursive: true, force: true });
}
});
});

View file

@ -14,6 +14,7 @@ import { pathToFileURL } from 'node:url';
import { createKnowledgeGraph } from '../../src/core/graph/graph.js';
import { runChunkedParseAndResolve } from '../../src/core/ingestion/pipeline-phases/parse-impl.js';
import { buildExportedTypeMapFromGraph } from '../../src/core/ingestion/call-processor.js';
import { computeChunkHash, fileContentHash } from '../../src/storage/parse-cache.js';
import type { ParseWorkerResult } from '../../src/core/ingestion/workers/parse-worker.js';
@ -52,6 +53,37 @@ const emptyWorkerResult = (filePath: string, name: string): ParseWorkerResult =>
fileCount: 1,
});
// A cached chunk result carrying an exported, typed symbol — enough to make a
// cache-hit replay push exportedTypeMap.size > 0, the precondition that
// suppresses the full-graph rebuild and exposes the sequential-miss gap (#2038).
const exportedTypedResult = (
filePath: string,
name: string,
returnType: string,
): ParseWorkerResult => {
const id = `Function:${filePath}:${name}`;
return {
...emptyWorkerResult(filePath, name),
nodes: [
{
id,
label: 'Function',
properties: {
name,
filePath,
startLine: 1,
endLine: 1,
language: 'typescript',
isExported: true,
},
},
],
symbols: [
{ filePath, name, nodeId: id, type: 'Function', returnType },
] as ParseWorkerResult['symbols'],
};
};
const writeReadyWorker = (workerPath: string, markerPath: string): void => {
fs.writeFileSync(
workerPath,
@ -322,4 +354,76 @@ describe('parse-impl worker pool lazy startup', () => {
else process.env.GITNEXUS_WORKER_POOL_SIZE = saved;
}
});
it('threads exportedTypeMap through the sequential path: a no-worker run over a partially-warm cache keeps the sequential-miss chunk exported types (#2038)', async () => {
const saved = process.env.GITNEXUS_WORKER_POOL_SIZE;
process.env.GITNEXUS_WORKER_POOL_SIZE = '0'; // force the no-worker (sequential) path
try {
// Chunk A — cache HIT, pre-seeded with an exported typed symbol so the
// replay makes exportedTypeMap.size > 0 and the size===0 rebuild is skipped.
const relA = 'src/a_hit.ts';
const contentA = 'export function aWidget(): number { return 2; }\n';
const fullA = path.join(repoDir, relA);
fs.mkdirSync(path.dirname(fullA), { recursive: true });
fs.writeFileSync(fullA, contentA);
// Chunk B — cache MISS, parsed sequentially for real; exported + typed.
const relB = 'src/b_miss.ts';
const contentB = 'export function bWidget(): number { return 1; }\n';
const fullB = path.join(repoDir, relB);
fs.writeFileSync(fullB, contentB);
const chunkHashA = computeChunkHash([
{ filePath: relA, contentHash: fileContentHash(contentA) },
]);
const parseCache = {
version: 'test',
entries: new Map<string, ParseWorkerResult[]>([
[chunkHashA, [exportedTypedResult(relA, 'aWidget', 'number')]],
]),
usedKeys: new Set<string>(),
};
const graph = createKnowledgeGraph();
const result = await runChunkedParseAndResolve(
graph,
[
{ path: relA, size: fs.statSync(fullA).size },
{ path: relB, size: fs.statSync(fullB).size },
],
[relA, relB],
2,
repoDir,
Date.now(),
() => {},
{
// 1-byte budget → each file is its own chunk, so A hits while B misses.
chunkByteBudget: 1,
parseCache,
},
);
expect(result.usedWorkerPool).toBe(false);
// Sanity: the cache-hit chunk populated the map — this is what makes
// size > 0 and suppresses the full-graph rebuild on the size===0 guard.
expect(result.exportedTypeMap.get(relA)?.get('aWidget')).toBe('number');
// Regression (#2038): the sequential-miss chunk's exported type must
// survive. Fails on pre-fix HEAD — the sequential path never populated
// exportedTypeMap, and the rebuild was skipped because the hit made size > 0.
expect(result.exportedTypeMap.get(relB)?.get('bWidget')).toBe('number');
// Differential oracle: the threaded map must match a fresh full-graph build
// (both directions for the entries under test) on the actual mixed path.
const oracle = buildExportedTypeMapFromGraph(graph, result.model.symbols);
expect(result.exportedTypeMap.get(relB)?.get('bWidget')).toBe(
oracle.get(relB)?.get('bWidget'),
);
expect(result.exportedTypeMap.get(relA)?.get('aWidget')).toBe(
oracle.get(relA)?.get('aWidget'),
);
} finally {
if (saved === undefined) delete process.env.GITNEXUS_WORKER_POOL_SIZE;
else process.env.GITNEXUS_WORKER_POOL_SIZE = saved;
}
});
});

View file

@ -0,0 +1,28 @@
import { describe, it, expect } from 'vitest';
import { SCOPE_RESOLVERS } from '../../../src/core/ingestion/scope-resolution/pipeline/registry.js';
import {
SCOPE_RESOLUTION_LANGUAGES,
isScopeResolutionLanguage,
} from '../../../src/core/ingestion/scope-resolution/pipeline/migrated-languages.js';
describe('SCOPE_RESOLUTION_LANGUAGES drift guard', () => {
// The parse worker gates ParsedFile emission on `isScopeResolutionLanguage`,
// which is derived from SCOPE_RESOLUTION_LANGUAGES — a hand-maintained
// duplicate of the SCOPE_RESOLVERS key set (kept resolver-import-free so the
// worker bundle stays light). The dangerous drift is asymmetric: a language
// in this Set but missing from SCOPE_RESOLVERS would have its ParsedFile
// skipped in the worker AND never re-extracted by scope resolution →
// permanent loss. This test fails if the two ever diverge (#1983).
it('covers exactly the languages registered in SCOPE_RESOLVERS', () => {
const resolverLangs = [...SCOPE_RESOLVERS.keys()].sort();
const skipSetLangs = [...SCOPE_RESOLUTION_LANGUAGES].sort();
expect(skipSetLangs).toEqual(resolverLangs);
});
it('isScopeResolutionLanguage returns true for every registered resolver language and false for null', () => {
for (const lang of SCOPE_RESOLVERS.keys()) {
expect(isScopeResolutionLanguage(lang)).toBe(true);
}
expect(isScopeResolutionLanguage(null)).toBe(false);
});
});