diff --git a/gitnexus/src/core/ingestion/languages/typescript/captures.ts b/gitnexus/src/core/ingestion/languages/typescript/captures.ts index b82edcff8..06ff691a0 100644 --- a/gitnexus/src/core/ingestion/languages/typescript/captures.ts +++ b/gitnexus/src/core/ingestion/languages/typescript/captures.ts @@ -64,7 +64,10 @@ const CALL_TAGS = [ '@reference.call.constructor', ] as const; -function pickFirstDefined(grouped: CaptureMatch, tags: readonly string[]): Capture | undefined { +function pickFirstDefined( + grouped: Record, + tags: readonly string[], +): T | undefined { for (const tag of tags) { const cap = grouped[tag]; if (cap !== undefined) return cap; @@ -113,6 +116,27 @@ function shouldEmitReadMember(memberNode: SyntaxNode): boolean { } } +function findSelfOrAncestorOfType(node: SyntaxNode | undefined, type: string): SyntaxNode | null { + let current: SyntaxNode | null | undefined = node; + while (current !== undefined && current !== null) { + if (current.type === type) return current; + current = current.parent; + } + return null; +} + +function findSelfOrAncestorOfTypes( + node: SyntaxNode | undefined, + types: readonly string[], +): SyntaxNode | null { + let current: SyntaxNode | null | undefined = node; + while (current !== undefined && current !== null) { + if (types.includes(current.type)) return current; + current = current.parent; + } + return null; +} + export function emitTsScopeCaptures( sourceText: string, filePath: string, @@ -151,9 +175,11 @@ export function emitTsScopeCaptures( // `@`; we put it back so the central extractor's prefix lookups // (`@scope.`, `@declaration.`, …) work. const grouped: Record = {}; + const groupedNodes: Record = {}; for (const c of m.captures) { const tag = '@' + c.name; grouped[tag] = nodeToCapture(tag, c.node); + groupedNodes[tag] = c.node; } if (Object.keys(grouped).length === 0) continue; @@ -165,6 +191,10 @@ export function emitTsScopeCaptures( if (grouped['@import.statement'] !== undefined) { const stmtCapture = grouped['@import.statement']; const stmtNode = + findSelfOrAncestorOfTypes(groupedNodes['@import.statement'], [ + 'import_statement', + 'export_statement', + ]) ?? findNodeAtRange(tree.rootNode, stmtCapture.range, 'import_statement') ?? findNodeAtRange(tree.rootNode, stmtCapture.range, 'export_statement'); if (stmtNode !== null) { @@ -183,7 +213,9 @@ export function emitTsScopeCaptures( // `splitDynamicImport` branch consumes. if (grouped['@import.dynamic'] !== undefined) { const dynCapture = grouped['@import.dynamic']; - const callNode = findNodeAtRange(tree.rootNode, dynCapture.range, 'call_expression'); + const callNode = + findSelfOrAncestorOfType(groupedNodes['@import.dynamic'], 'call_expression') ?? + findNodeAtRange(tree.rootNode, dynCapture.range, 'call_expression'); if (callNode !== null) { const decomposed = splitImportStatement(callNode); for (const d of decomposed) out.push(d); @@ -197,7 +229,9 @@ export function emitTsScopeCaptures( // we rely on this emit-side filter so the query stays simple. if (grouped['@reference.read.member'] !== undefined) { const anchor = grouped['@reference.read.member']; - const memberNode = findNodeAtRange(tree.rootNode, anchor.range, 'member_expression'); + const memberNode = + findSelfOrAncestorOfType(groupedNodes['@reference.read.member'], 'member_expression') ?? + findNodeAtRange(tree.rootNode, anchor.range, 'member_expression'); if (memberNode === null || !shouldEmitReadMember(memberNode)) { continue; } @@ -209,8 +243,9 @@ export function emitTsScopeCaptures( // function_signature, so `parameterTypes` is populated when // available. const declAnchor = pickFirstDefined(grouped, FUNCTION_DECL_TAGS); + const declAnchorNode = pickFirstDefined(groupedNodes, FUNCTION_DECL_TAGS); if (declAnchor !== undefined) { - const fnNode = findFunctionNode(tree.rootNode, declAnchor.range); + const fnNode = findFunctionNode(tree.rootNode, declAnchor.range, declAnchorNode); if (fnNode !== null) { const arity = computeTsArityMetadata(fnNode); if (arity.parameterCount !== undefined) { @@ -256,8 +291,10 @@ export function emitTsScopeCaptures( // synthesizer would need to count `jsx_attribute` children of the // opening tag instead of `arguments`. const callAnchor = pickFirstDefined(grouped, CALL_TAGS); + const callAnchorNode = pickFirstDefined(groupedNodes, CALL_TAGS); if (callAnchor !== undefined && grouped['@reference.arity'] === undefined) { const callNode = + findSelfOrAncestorOfTypes(callAnchorNode, ['call_expression', 'new_expression']) ?? findNodeAtRange(tree.rootNode, callAnchor.range, 'call_expression') ?? findNodeAtRange(tree.rootNode, callAnchor.range, 'new_expression'); if (callNode !== null) { @@ -293,7 +330,11 @@ export function emitTsScopeCaptures( // lookup instead of synthesis — covered by `tsReceiverBinding`. const scopeFnAnchor = grouped['@scope.function']; if (scopeFnAnchor !== undefined) { - const fnNode = findFunctionNode(tree.rootNode, scopeFnAnchor.range); + const fnNode = findFunctionNode( + tree.rootNode, + scopeFnAnchor.range, + groupedNodes['@scope.function'], + ); if (fnNode !== null) { const synth = synthesizeTsReceiverBinding(fnNode); if (synth !== null) out.push(synth); @@ -518,7 +559,13 @@ function inferArgType(argNode: SyntaxNode): string { * The `@scope.function` anchor range covers the whole node, but the * tag alone doesn't identify which node type among the many TS * function-likes. */ -function findFunctionNode(rootNode: SyntaxNode, range: Capture['range']): SyntaxNode | null { +function findFunctionNode( + rootNode: SyntaxNode, + range: Capture['range'], + anchorNode?: SyntaxNode, +): SyntaxNode | null { + const fromAnchor = findSelfOrAncestorOfTypes(anchorNode, FUNCTION_NODE_TYPES); + if (fromAnchor !== null) return fromAnchor; for (const nodeType of FUNCTION_NODE_TYPES) { const n = findNodeAtRange(rootNode, range, nodeType); if (n !== null) return n; diff --git a/gitnexus/src/core/ingestion/parsing-processor.ts b/gitnexus/src/core/ingestion/parsing-processor.ts index 5cf398bee..7a669504e 100644 --- a/gitnexus/src/core/ingestion/parsing-processor.ts +++ b/gitnexus/src/core/ingestion/parsing-processor.ts @@ -37,7 +37,7 @@ import { } from './utils/template-arguments.js'; import type { LanguageProvider } from './language-provider.js'; import type { ParsedFile } from 'gitnexus-shared'; -import { WorkerPool } from './workers/worker-pool.js'; +import { WorkerPool, WorkerPoolDispatchError } from './workers/worker-pool.js'; import { logger } from '../logger.js'; import type { ParseWorkerResult, @@ -886,12 +886,37 @@ export const processParsing = async ( ); } catch (err) { const message = err instanceof Error ? err.message : String(err); + let fallbackFiles = files; + if (err instanceof WorkerPoolDispatchError && err.fallbackExcludePaths.length > 0) { + const excluded = new Set(err.fallbackExcludePaths); + fallbackFiles = files.filter((file) => !excluded.has(file.path)); + logger.warn( + { + skippedPaths: err.fallbackExcludePaths, + }, + 'Skipping worker-timeout files in sequential fallback:', + ); + reportProgress?.( + lastProgress, + files.length, + `Skipping ${files.length - fallbackFiles.length} worker-timeout file(s) in sequential fallback`, + ); + } logger.warn({ message }, 'Worker pool parsing stopped; continuing with sequential parser:'); reportProgress?.( lastProgress, files.length, `Sequential fallback after worker issue: ${message}`, ); + await processParsingSequential( + graph, + fallbackFiles, + symbolTable, + astCache, + scopeTreeCache, + reportProgress, + ); + return null; } } diff --git a/gitnexus/src/core/ingestion/workers/worker-pool.ts b/gitnexus/src/core/ingestion/workers/worker-pool.ts index 211368567..0e5c3e75b 100644 --- a/gitnexus/src/core/ingestion/workers/worker-pool.ts +++ b/gitnexus/src/core/ingestion/workers/worker-pool.ts @@ -29,6 +29,16 @@ export interface WorkerPoolOptions { timeoutBackoffFactor?: number; } +export class WorkerPoolDispatchError extends Error { + readonly fallbackExcludePaths: readonly string[]; + + constructor(message: string, fallbackExcludePaths: readonly string[] = []) { + super(message); + this.name = 'WorkerPoolDispatchError'; + this.fallbackExcludePaths = fallbackExcludePaths; + } +} + /** Message shapes sent back by worker threads. */ type WorkerOutgoingMessage = | { type: 'progress'; filesProcessed: number } @@ -336,13 +346,15 @@ export const createWorkerPool = ( return true; } + const stalledPath = itemPath(job.items[0]); void fail( - new Error( + new WorkerPoolDispatchError( `Worker ${workerIndex} parse job idle timeout after ${job.timeoutMs / 1000}s ` + - `(single item${itemPath(job.items[0]) ? `: ${itemPath(job.items[0])}` : ''}, ` + + `(single item${stalledPath ? `: ${stalledPath}` : ''}, ` + `${job.estimatedBytes} bytes, last progress: ${lastProgress}). ` + `Analyze will retry through sequential fallback. Increase with ` + `--worker-timeout or GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS.`, + stalledPath ? [stalledPath] : [], ), ); return false; diff --git a/gitnexus/test/unit/parsing-worker-fallback.test.ts b/gitnexus/test/unit/parsing-worker-fallback.test.ts index d25fe938b..ab434eb75 100644 --- a/gitnexus/test/unit/parsing-worker-fallback.test.ts +++ b/gitnexus/test/unit/parsing-worker-fallback.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it, vi } from 'vitest'; import { createASTCache } from '../../src/core/ingestion/ast-cache.js'; import { processParsing } from '../../src/core/ingestion/parsing-processor.js'; import type { WorkerPool } from '../../src/core/ingestion/workers/worker-pool.js'; +import { WorkerPoolDispatchError } from '../../src/core/ingestion/workers/worker-pool.js'; import { createKnowledgeGraph } from '../../src/core/graph/graph.js'; import { createSymbolTable } from '../../src/core/ingestion/model/symbol-table.js'; @@ -41,4 +42,40 @@ describe('processParsing worker fallback', () => { graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'a'), ).toBe(true); }); + + it('skips worker-timeout singleton files during sequential fallback', async () => { + const graph = createKnowledgeGraph(); + const progressDetails: string[] = []; + const workerPool: WorkerPool = { + size: 1, + dispatch: vi.fn(async () => { + throw new WorkerPoolDispatchError('injected worker idle timeout', ['src/stuck.ts']); + }), + terminate: vi.fn(async () => undefined), + }; + + const result = await processParsing( + graph, + [ + { path: 'src/stuck.ts', content: 'export function stuck() { return 0; }\n' }, + { path: 'src/a.ts', content: 'export function a() { return 1; }\n' }, + ], + createSymbolTable(), + createASTCache(), + createASTCache(), + (_current, _total, detail) => { + progressDetails.push(detail); + }, + workerPool, + ); + + expect(result).toBeNull(); + expect(progressDetails).toContain('Skipping 1 worker-timeout file(s) in sequential fallback'); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'a'), + ).toBe(true); + expect( + graph.nodes.some((node) => node.label === 'Function' && node.properties.name === 'stuck'), + ).toBe(false); + }); });