From 94d6aa2365d6b7c6fe3d30733698475fbfe9f4cf Mon Sep 17 00:00:00 2001 From: Test Date: Sat, 18 Apr 2026 13:14:21 +0200 Subject: [PATCH] fix: wire queue detection into phase-based pipeline after rebase onto v1.6.2-rc.9 The upstream phase refactor replaced the monolithic pipeline with pipeline-phases/. Add queues.ts phase, queue-extraction.ts inline extractor, and thread allQueuePatterns through ParseOutput so ENQUEUES/PROCESSES edges and Queue nodes are created correctly. Co-Authored-By: Claude Sonnet 4.6 --- .../core/ingestion/pipeline-phases/index.ts | 1 + .../ingestion/pipeline-phases/parse-impl.ts | 9 ++ .../core/ingestion/pipeline-phases/parse.ts | 2 + .../pipeline-phases/queue-extraction.ts | 97 +++++++++++++++++++ .../core/ingestion/pipeline-phases/queues.ts | 92 ++++++++++++++++++ gitnexus/src/core/ingestion/pipeline.ts | 2 + 6 files changed, 203 insertions(+) create mode 100644 gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts create mode 100644 gitnexus/src/core/ingestion/pipeline-phases/queues.ts diff --git a/gitnexus/src/core/ingestion/pipeline-phases/index.ts b/gitnexus/src/core/ingestion/pipeline-phases/index.ts index c05264de1..1912100a2 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/index.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/index.ts @@ -15,6 +15,7 @@ export { parsePhase, type ParseOutput } from './parse.js'; export { routesPhase, type RoutesOutput, type RouteEntry } from './routes.js'; export { toolsPhase, type ToolsOutput, type ToolDef } from './tools.js'; export { ormPhase, type ORMOutput } from './orm.js'; +export { queuesPhase, type QueuesOutput } from './queues.js'; export { crossFilePhase, type CrossFileOutput } from './cross-file.js'; export { mroPhase, type MROOutput } from './mro.js'; export { communitiesPhase, type CommunitiesOutput } from './communities.js'; diff --git a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts index f8e0d3b53..c73e9b4c5 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts @@ -53,6 +53,7 @@ import type { ExtractedDecoratorRoute, ExtractedFetchCall, ExtractedORMQuery, + ExtractedQueuePattern, ExtractedRoute, ExtractedToolDef, FileConstructorBindings, @@ -68,6 +69,7 @@ import { fileURLToPath, pathToFileURL } from 'node:url'; import { isDev } from '../utils/env.js'; import { synthesizeWildcardImportBindings, needsSynthesis } from './wildcard-synthesis.js'; import { extractORMQueriesInline } from './orm-extraction.js'; +import { extractQueuePatternsInline } from './queue-extraction.js'; // ── Constants ────────────────────────────────────────────────────────────── @@ -106,6 +108,7 @@ export async function runChunkedParseAndResolve( allDecoratorRoutes: ExtractedDecoratorRoute[]; allToolDefs: ExtractedToolDef[]; allORMQueries: ExtractedORMQuery[]; + allQueuePatterns: ExtractedQueuePattern[]; bindingAccumulator: BindingAccumulator; resolutionContext: ReturnType; usedWorkerPool: boolean; @@ -248,6 +251,7 @@ export async function runChunkedParseAndResolve( const allDecoratorRoutes: ExtractedDecoratorRoute[] = []; const allToolDefs: ExtractedToolDef[] = []; const allORMQueries: ExtractedORMQuery[] = []; + const allQueuePatterns: ExtractedQueuePattern[] = []; const deferredWorkerCalls: ExtractedCall[] = []; const deferredWorkerHeritage: ExtractedHeritage[] = []; const deferredConstructorBindings: FileConstructorBindings[] = []; @@ -393,6 +397,9 @@ export async function runChunkedParseAndResolve( if (chunkWorkerData.ormQueries?.length) { for (const item of chunkWorkerData.ormQueries) allORMQueries.push(item); } + if (chunkWorkerData.queuePatterns?.length) { + for (const item of chunkWorkerData.queuePatterns) allQueuePatterns.push(item); + } } else { await processImports(graph, chunkFiles, astCache, ctx, undefined, repoPath, allPaths); sequentialChunkPaths.push(chunkPaths); @@ -509,6 +516,7 @@ export async function runChunkedParseAndResolve( } for (const f of chunkFiles) { extractORMQueriesInline(f.path, f.content, allORMQueries); + extractQueuePatternsInline(f.path, f.content, allQueuePatterns); } astCache.clear(); cachedSequentialChunkFiles[chunkIdx] = []; @@ -589,6 +597,7 @@ export async function runChunkedParseAndResolve( allDecoratorRoutes, allToolDefs, allORMQueries, + allQueuePatterns, bindingAccumulator, resolutionContext: ctx, // Whether a worker pool was actually live for this run. False means the diff --git a/gitnexus/src/core/ingestion/pipeline-phases/parse.ts b/gitnexus/src/core/ingestion/pipeline-phases/parse.ts index 6415cb6e5..74b4369f3 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/parse.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/parse.ts @@ -26,6 +26,7 @@ import type { ExtractedDecoratorRoute, ExtractedToolDef, ExtractedORMQuery, + ExtractedQueuePattern, } from '../workers/parse-worker.js'; import type { createResolutionContext } from '../model/resolution-context.js'; import { runChunkedParseAndResolve } from './parse-impl.js'; @@ -47,6 +48,7 @@ export interface ParseOutput { readonly allDecoratorRoutes: readonly ExtractedDecoratorRoute[]; readonly allToolDefs: readonly ExtractedToolDef[]; readonly allORMQueries: readonly ExtractedORMQuery[]; + readonly allQueuePatterns: readonly ExtractedQueuePattern[]; bindingAccumulator: BindingAccumulator; /** Resolution context from the parse phase — carries importMap, namedImportMap, etc. */ resolutionContext: ReturnType; diff --git a/gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts b/gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts new file mode 100644 index 000000000..94475f15f --- /dev/null +++ b/gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts @@ -0,0 +1,97 @@ +/** + * Inline queue pattern extraction (sequential fallback path). + * + * Extracts BullMQ and Temporal queue patterns from source content using + * regex patterns. Used by the sequential parse path when workers are + * not available — the worker path extracts queue patterns via + * extractQueuePatterns in parse-worker.ts instead. + * + * @module + */ + +import type { ExtractedQueuePattern } from '../workers/parse-worker.js'; + +// ── Regex patterns ───────────────────────────────────────────────────────── + +const BULLMQ_ADD_RE = /(\w+)\.(add|addBulk)\s*\(/g; +const BULLMQ_WORKER_RE = /new\s+Worker\s*\(\s*['\"](\w[\w-]*)['\"]/g; +const TEMPORAL_ACTIVITY_RE = /activities\.(\w+)\s*\(/g; +const TEMPORAL_WORKFLOW_START_RE = /client\.workflow\.(start|execute)\s*\(\s*(\w+)/g; + +// ── Extraction function ─────────────────────────────────────────────────── + +/** + * Extract BullMQ and Temporal queue patterns from file content using regex. + * + * Fast-path: skips files that don't contain queue-related markers. + * Results are appended to the `out` array (push pattern avoids allocation). + * + * @param filePath Relative path of the source file + * @param content File content string + * @param out Output array to append extracted patterns to + */ +export function extractQueuePatternsInline( + filePath: string, + content: string, + out: ExtractedQueuePattern[], +): void { + const hasBullMQ = content.includes('new Queue') || content.includes('new Worker'); + const hasTemporal = content.includes('activities.') || content.includes('client.workflow.'); + if (!hasBullMQ && !hasTemporal) return; + + if (hasBullMQ) { + const queueVarMap = new Map(); + const assignRe = /(?:const|let|var)\s+(\w+)\s*=\s*new\s+Queue\s*\(\s*['\"](\w[\w-]*)['\"]/g; + assignRe.lastIndex = 0; + let m; + while ((m = assignRe.exec(content)) !== null) { + queueVarMap.set(m[1], m[2]); + } + BULLMQ_ADD_RE.lastIndex = 0; + while ((m = BULLMQ_ADD_RE.exec(content)) !== null) { + const qn = queueVarMap.get(m[1]); + if (qn) { + out.push({ + filePath, + role: 'producer', + queueName: qn, + method: m[2], + lineNumber: content.substring(0, m.index).split('\n').length - 1, + }); + } + } + BULLMQ_WORKER_RE.lastIndex = 0; + while ((m = BULLMQ_WORKER_RE.exec(content)) !== null) { + out.push({ + filePath, + role: 'consumer', + queueName: m[1], + lineNumber: content.substring(0, m.index).split('\n').length - 1, + }); + } + } + + if (hasTemporal) { + let m; + TEMPORAL_ACTIVITY_RE.lastIndex = 0; + while ((m = TEMPORAL_ACTIVITY_RE.exec(content)) !== null) { + out.push({ + filePath, + role: 'activity', + queueName: m[1], + handlerName: m[1], + lineNumber: content.substring(0, m.index).split('\n').length - 1, + }); + } + TEMPORAL_WORKFLOW_START_RE.lastIndex = 0; + while ((m = TEMPORAL_WORKFLOW_START_RE.exec(content)) !== null) { + out.push({ + filePath, + role: 'workflow', + queueName: m[2], + method: m[1], + lineNumber: content.substring(0, m.index).split('\n').length - 1, + }); + } + } +} diff --git a/gitnexus/src/core/ingestion/pipeline-phases/queues.ts b/gitnexus/src/core/ingestion/pipeline-phases/queues.ts new file mode 100644 index 000000000..45a49e6a1 --- /dev/null +++ b/gitnexus/src/core/ingestion/pipeline-phases/queues.ts @@ -0,0 +1,92 @@ +/** + * Phase: queues + * + * Processes async queue patterns (BullMQ + Temporal) and creates + * ENQUEUES / PROCESSES edges and Queue CodeElement nodes. + * + * @deps parse + * @reads allQueuePatterns (from parse) + * @writes graph (CodeElement nodes, ENQUEUES/PROCESSES edges) + */ + +import type { PipelinePhase, PipelineContext, PhaseResult } from './types.js'; +import { getPhaseOutput } from './types.js'; +import type { ParseOutput } from './parse.js'; +import { generateId } from '../../../lib/utils.js'; +import type { ExtractedQueuePattern } from '../workers/parse-worker.js'; +import type { KnowledgeGraph } from '../../graph/types.js'; +import { isDev } from '../utils/env.js'; + +export interface QueuesOutput { + queuesCreated: number; + edgesCreated: number; +} + +export const queuesPhase: PipelinePhase = { + name: 'queues', + deps: ['parse'], + + async execute( + ctx: PipelineContext, + deps: ReadonlyMap>, + ): Promise { + const { allQueuePatterns } = getPhaseOutput(deps, 'parse'); + + if (allQueuePatterns.length === 0) { + return { queuesCreated: 0, edgesCreated: 0 }; + } + + return processQueuePatterns(ctx.graph, allQueuePatterns); + }, +}; + +function processQueuePatterns( + graph: KnowledgeGraph, + patterns: readonly ExtractedQueuePattern[], +): QueuesOutput { + const queueNodes = new Map(); + const seenEdges = new Set(); + let edgesCreated = 0; + + for (const pt of patterns) { + let queueNodeId = queueNodes.get(pt.queueName); + if (!queueNodeId) { + queueNodeId = generateId('CodeElement', `Queue:${pt.queueName}`); + graph.addNode({ + id: queueNodeId, + label: 'CodeElement', + properties: { + name: pt.queueName, + filePath: '', + description: `Queue: ${pt.queueName}`, + }, + }); + queueNodes.set(pt.queueName, queueNodeId); + } + + const edgeType = + pt.role === 'producer' || pt.role === 'workflow' ? 'ENQUEUES' : 'PROCESSES'; + const fileId = generateId('File', pt.filePath); + const edgeKey = `${fileId}->${queueNodeId}:${edgeType}`; + if (seenEdges.has(edgeKey)) continue; + seenEdges.add(edgeKey); + + graph.addRelationship({ + id: generateId(edgeType, edgeKey), + sourceId: fileId, + targetId: queueNodeId, + type: edgeType, + confidence: 0.9, + reason: `queue-${pt.role}`, + }); + edgesCreated++; + } + + if (isDev) { + console.log( + `Queues: ${edgesCreated} edges (ENQUEUES/PROCESSES), ${queueNodes.size} queue nodes (${patterns.length} total patterns)`, + ); + } + + return { queuesCreated: queueNodes.size, edgesCreated }; +} diff --git a/gitnexus/src/core/ingestion/pipeline.ts b/gitnexus/src/core/ingestion/pipeline.ts index e1f226289..70329f157 100644 --- a/gitnexus/src/core/ingestion/pipeline.ts +++ b/gitnexus/src/core/ingestion/pipeline.ts @@ -29,6 +29,7 @@ import { routesPhase, toolsPhase, ormPhase, + queuesPhase, crossFilePhase, mroPhase, communitiesPhase, @@ -79,6 +80,7 @@ function buildPhaseList(options?: PipelineOptions): PipelinePhase[] { routesPhase, toolsPhase, ormPhase, + queuesPhase, crossFilePhase, ];