diff --git a/gitnexus/src/core/ingestion/pipeline-phases/processes.ts b/gitnexus/src/core/ingestion/pipeline-phases/processes.ts index 462d1523f..835dce7b3 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/processes.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/processes.ts @@ -38,11 +38,19 @@ export function computeDynamicMaxProcesses(symbolCount: number): number { export const processesPhase: PipelinePhase = { name: 'processes', - // `structure` supplies `totalFiles` (progress counter) without the spurious - // structural data dependency on `parse`. `pruneLocalSymbols` is declared + // `structure` supplies `totalFiles` (progress counter), which is why this + // phase historically avoided depending on `parse` at all — that dependency + // was spurious for a progress number. + // + // It is no longer spurious. R3-6 reads `allFetchCalls` / `allORMQueries` from + // the parse output to learn WHERE the program reaches outward, which is what + // lets a trace end at a sink instead of only at a leaf. That is a real data + // dependency, so it is declared rather than reached for implicitly — and the + // read below still fails open, so a pipeline without that output detects no + // sinks rather than failing the phase. `pruneLocalSymbols` is declared // explicitly so process extraction always reads the trimmed graph even if a // future option drops the intervening `mro`/`communities` phases. - deps: ['communities', 'routes', 'tools', 'pruneLocalSymbols', 'structure'], + deps: ['communities', 'routes', 'tools', 'pruneLocalSymbols', 'structure', 'parse'], async execute( ctx: PipelineContext, @@ -66,6 +74,29 @@ export const processesPhase: PipelinePhase = { }); const dynamicMaxProcesses = computeDynamicMaxProcesses(symbolCount); + // R3-6: where the program reaches outward. Already collected by the parse + // phase for FILE-level FETCHES/QUERIES edges; reused here at function + // granularity so a trace can end somewhere meaningful instead of only at a + // leaf. Absent (or an older parse output) simply yields no sinks and the + // previous behaviour. + let parseOutput: + | { + allFetchCalls?: { filePath: string; lineNumber: number }[]; + allORMQueries?: { filePath: string; lineNumber: number }[]; + } + | undefined; + try { + parseOutput = getPhaseOutput(deps, 'parse'); + } catch { + // Fail open: no sinks, previous behaviour. A missing parse output is a + // pipeline-composition question, not a reason to lose every process. + parseOutput = undefined; + } + const outwardActionSites = [ + ...(parseOutput?.allFetchCalls ?? []), + ...(parseOutput?.allORMQueries ?? []), + ].filter((s) => typeof s?.filePath === 'string' && typeof s?.lineNumber === 'number'); + const processResult = await processProcesses( ctx.graph, communityResult.memberships, @@ -79,6 +110,7 @@ export const processesPhase: PipelinePhase = { }); }, { maxProcesses: dynamicMaxProcesses, minSteps: 3 }, + outwardActionSites, ); if (isDev) { diff --git a/gitnexus/src/core/ingestion/process-processor.ts b/gitnexus/src/core/ingestion/process-processor.ts index 18f99a5f3..aa0f51ac4 100644 --- a/gitnexus/src/core/ingestion/process-processor.ts +++ b/gitnexus/src/core/ingestion/process-processor.ts @@ -83,8 +83,16 @@ export const processProcesses = async ( memberships: CommunityMembership[], onProgress?: (message: string, progress: number) => void, config: Partial = {}, + /** + * Places the program reaches outward — fetch calls and ORM queries, each with + * a file and a line (R3-6). Attributed to their enclosing function to form the + * sink set; omitted, behaviour is exactly as before. + */ + outwardActionSites: readonly OutwardActionSite[] = [], ): Promise => { const cfg = { ...DEFAULT_CONFIG, ...config }; + const sinkFunctions = buildSinkFunctionSet(knowledgeGraph, outwardActionSites); + const isSink = (nodeId: string): boolean => sinkFunctions.has(nodeId); onProgress?.('Finding entry points...', 0); @@ -109,7 +117,7 @@ export const processProcesses = async ( for (let i = 0; i < entryPoints.length && allTraces.length < cfg.maxProcesses * 2; i++) { const entryId = entryPoints[i]; - const traces = traceFromEntryPoint(entryId, callsEdges, cfg); + const traces = traceFromEntryPoint(entryId, callsEdges, cfg, isSink); // Filter out traces that are too short traces.filter((t) => t.length >= cfg.minSteps).forEach((t) => allTraces.push(t)); @@ -125,7 +133,7 @@ export const processProcesses = async ( onProgress?.(`Found ${allTraces.length} traces, deduplicating...`, 60); // Step 3: Deduplicate similar traces (subset removal) - const uniqueTraces = deduplicateTraces(allTraces); + const uniqueTraces = deduplicateTraces(allTraces, isSink); // Step 3b: Deduplicate by entry+terminal pair (keep longest path per pair) const endpointDeduped = deduplicateByEndpoints(uniqueTraces); @@ -169,8 +177,19 @@ export const processProcesses = async ( // whatever leaf it happens to bottom out in. Fixing that means teaching the // walk what a sink is (I/O, route handler, external call), which is a change // to what a process IS rather than to how processes are ranked. + // + // R3-6 adds one rule ahead of depth: a SINK-terminated trace outranks a + // leaf-terminated one. A flow that ends where the program does something — + // places an order, writes a row — is what a reader came for; a chain that + // ends in a date helper is where control happened to stop. Depth still orders + // within each group. const tracesByTerminal = new Map(); - for (const trace of [...endpointDeduped].sort((a, b) => b.length - a.length)) { + const rankedByInterest = [...endpointDeduped].sort((a, b) => { + const aSink = Number(isSink(a[a.length - 1] ?? '')); + const bSink = Number(isSink(b[b.length - 1] ?? '')); + return bSink - aSink || b.length - a.length; + }); + for (const trace of rankedByInterest) { const terminalId = trace[trace.length - 1]; if (terminalId === undefined) continue; const existing = tracesByTerminal.get(terminalId); @@ -402,6 +421,13 @@ export const traceFromEntryPoint = ( entryId: string, callsEdges: AdjacencyList, config: ProcessDetectionConfig, + /** + * Functions that reach outward (R3-6). A trace also ENDS here, even though + * the walk continues past it: a flow whose meaningful endpoint calls onward + * was otherwise never a candidate, only ever surviving as whatever leaf it + * bottomed out in. + */ + isSink: (nodeId: string) => boolean = () => false, ): string[][] => { const traces: string[][] = []; @@ -443,6 +469,13 @@ export const traceFromEntryPoint = ( traces.push([...path]); } } else { + // A SINK ends a trace without ending the walk (R3-6). Emitting here is + // what lets `placeOrder` be an endpoint while `placeOrder -> formatDate` + // still exists as its own longer trace; the two answer different + // questions and neither should suppress the other. + if (isSink(currentId) && path.length >= config.minSteps) { + traces.push([...path]); + } // Continue tracing - limit branching const limitedCallees = callees.slice(0, config.maxBranching); let addedBranch = false; @@ -484,6 +517,75 @@ export const traceFromEntryPoint = ( return traces; }; +// ============================================================================ +// HELPER: Function-level sink set +// ============================================================================ + +/** A place in the source where the program reaches outward. */ +export interface OutwardActionSite { + readonly filePath: string; + readonly lineNumber: number; +} + +const CALLABLE_SINK_LABELS: ReadonlySet = new Set([ + 'Function', + 'Method', + 'Constructor', +]); + +/** + * Functions that DO something outward — issue a request, run a query (R3-6). + * + * The missing layer behind "business flows are never processes". A trace is + * only emitted at a node with NO outgoing calls, so a real flow — scan, score, + * arm, place the order — is always a PREFIX of some longer chain that runs on + * into date helpers and formatters, and can never be a process in its own + * right. Ending a trace somewhere meaningful needs a notion of an endpoint that + * is not a leaf, and the walk had none. + * + * GitNexus already knew where the program reaches outward: the parse phase + * collects fetch calls and ORM queries carrying `filePath` + `lineNumber`. Those + * facts only ever produced FILE-level edges (`File -[FETCHES]-> Route`), which + * is too coarse to end a trace on — every function in a file containing one + * would qualify. Attributing each site to the function whose range CONTAINS it + * turns the same facts into the function-level signal the walk needs, with no + * new extraction, no new relation pair, and no schema change. + * + * Innermost wins: a nested closure that performs the call is the sink, not the + * outer function that merely spans it. + */ +export function buildSinkFunctionSet( + graph: KnowledgeGraph, + sites: readonly OutwardActionSite[], +): ReadonlySet { + const sinks = new Set(); + if (sites.length === 0) return sinks; + + const byFile = new Map(); + for (const node of graph.iterNodes()) { + if (!CALLABLE_SINK_LABELS.has(node.label)) continue; + const props = node.properties as { filePath?: string; startLine?: number; endLine?: number }; + if (typeof props.filePath !== 'string') continue; + if (typeof props.startLine !== 'number' || typeof props.endLine !== 'number') continue; + const entry = { id: node.id, start: props.startLine, end: props.endLine }; + const list = byFile.get(props.filePath); + if (list === undefined) byFile.set(props.filePath, [entry]); + else list.push(entry); + } + + for (const site of sites) { + const candidates = byFile.get(site.filePath); + if (candidates === undefined) continue; + let best: { id: string; start: number; end: number } | undefined; + for (const c of candidates) { + if (site.lineNumber < c.start || site.lineNumber > c.end) continue; + if (best === undefined || c.end - c.start < best.end - best.start) best = c; + } + if (best !== undefined) sinks.add(best.id); + } + return sinks; +} + // ============================================================================ // HELPER: Deduplicate traces // ============================================================================ @@ -492,7 +594,11 @@ export const traceFromEntryPoint = ( * Merge traces that are subsets of other traces. * Keep longer traces, remove redundant shorter ones. */ -const deduplicateTraces = (traces: string[][]): string[][] => { +const deduplicateTraces = ( + traces: string[][], + /** See `buildSinkFunctionSet` — a sink-terminated trace survives subsumption. */ + isSink: (nodeId: string) => boolean = () => false, +): string[][] => { if (traces.length === 0) return []; // Sort by length descending @@ -500,6 +606,17 @@ const deduplicateTraces = (traces: string[][]): string[][] => { const unique: string[][] = []; for (const trace of sorted) { + // A SINK-TERMINATED trace is never redundant (R3-6), even though it is by + // definition a prefix of the longer chain that runs on past the sink into + // helpers. That is the whole shape of a business flow — scan, score, arm, + // PLACE THE ORDER — so subsuming it here is exactly what kept such flows + // from ever being processes. Emitting one at the walk and deleting it one + // step later would have been a no-op fix. + const terminal = trace[trace.length - 1]; + if (terminal !== undefined && isSink(terminal)) { + unique.push(trace); + continue; + } // Check if this trace is a subset of any already-added trace const traceKey = trace.join('->'); const isSubset = unique.some((existing) => { diff --git a/gitnexus/test/unit/process-processor.test.ts b/gitnexus/test/unit/process-processor.test.ts index 7aa047089..6265a99db 100644 --- a/gitnexus/test/unit/process-processor.test.ts +++ b/gitnexus/test/unit/process-processor.test.ts @@ -2,6 +2,7 @@ import { describe, it, expect, vi } from 'vitest'; import { processProcesses, traceFromEntryPoint, + buildSinkFunctionSet, type ProcessDetectionConfig, } from '../../src/core/ingestion/process-processor.js'; import { computeDynamicMaxProcesses } from '../../src/core/ingestion/pipeline-phases/processes.js'; @@ -635,6 +636,109 @@ describe('process depth (D1/D2)', () => { // asserted against `traceFromEntryPoint` directly, above. }); +describe('sink-terminated flows (R3-6)', () => { + const addFn = ( + graph: ReturnType, + id: string, + line: number, + ): void => { + graph.addNode({ + id, + label: 'Function', + properties: { + name: id.split(':')[1], + filePath: 'src/flow.ts', + startLine: line, + endLine: line + 2, + }, + }); + }; + const addCall = ( + graph: ReturnType, + from: string, + to: string, + ): void => { + graph.addRelationship({ + id: `rel:${from}->${to}`, + sourceId: from, + targetId: to, + type: 'CALLS', + confidence: 1, + reason: 'test', + }); + }; + + /** + * The shape the whole item is about: a business flow whose meaningful + * endpoint CALLS ONWARD into helpers. `placeOrder` is where the program does + * something; `formatDate` is merely where control stops. + */ + const flowGraph = (): ReturnType => { + const graph = createKnowledgeGraph(); + addFn(graph, 'func:scan', 1); + addFn(graph, 'func:score', 10); + addFn(graph, 'func:placeOrder', 20); + addFn(graph, 'func:formatDate', 30); + addFn(graph, 'func:pad', 40); + addCall(graph, 'func:scan', 'func:score'); + addCall(graph, 'func:score', 'func:placeOrder'); + addCall(graph, 'func:placeOrder', 'func:formatDate'); + addCall(graph, 'func:formatDate', 'func:pad'); + return graph; + }; + + // `placeOrder` spans lines 20-22, so an outward action on line 21 belongs to + // it — the attribution the file-level FETCHES edge could not express. + const ORDER_SITE = [{ filePath: 'src/flow.ts', lineNumber: 21 }]; + + it('ends a trace where the program reaches outward', async () => { + const result = await processProcesses(flowGraph(), [], undefined, {}, ORDER_SITE); + expect(result.processes.map((p) => p.terminalId)).toContain('func:placeOrder'); + }); + + // The half a naive implementation gets wrong: emitting the sink trace at the + // walk and then letting subset-removal delete it one step later is a no-op, + // because a sink-terminated flow is BY DEFINITION a prefix of the longer + // chain that runs on past it. + it('keeps the sink flow even though it is a prefix of a longer chain', async () => { + const result = await processProcesses(flowGraph(), [], undefined, {}, ORDER_SITE); + const terminals = result.processes.map((p) => p.terminalId); + expect(terminals).toContain('func:placeOrder'); + // The longer chain still exists — the two answer different questions. + expect(terminals).toContain('func:pad'); + }); + + it('ranks the sink flow above the leaf chain', async () => { + const result = await processProcesses(flowGraph(), [], undefined, {}, ORDER_SITE); + const first = result.processes[0]?.terminalId; + expect(first).toBe('func:placeOrder'); + }); + + // Without sites, behaviour must be exactly what it was. + it('changes nothing when no outward action is known', async () => { + const result = await processProcesses(flowGraph(), [], undefined, {}, []); + expect(result.processes.map((p) => p.terminalId)).not.toContain('func:placeOrder'); + }); + + it('attributes a site to the INNERMOST enclosing function', () => { + const graph = createKnowledgeGraph(); + // An outer function spanning the inner one; the inner performs the call. + graph.addNode({ + id: 'func:outer', + label: 'Function', + properties: { name: 'outer', filePath: 'src/a.ts', startLine: 1, endLine: 50 }, + }); + graph.addNode({ + id: 'func:inner', + label: 'Function', + properties: { name: 'inner', filePath: 'src/a.ts', startLine: 10, endLine: 20 }, + }); + const sinks = buildSinkFunctionSet(graph, [{ filePath: 'src/a.ts', lineNumber: 15 }]); + expect(sinks.has('func:inner')).toBe(true); + expect(sinks.has('func:outer')).toBe(false); + }); +}); + describe('process selection diversity (R2-3)', () => { const addFn = (graph: ReturnType, id: string): void => { graph.addNode({