From 2c2ab32b212324f2241523dd1d8f82dcca845506 Mon Sep 17 00:00:00 2001 From: Test Date: Sat, 18 Apr 2026 13:37:28 +0200 Subject: [PATCH] fix: wire queue detection into pipeline-phases correctly Previous cherry-pick brought #521's duplicate queue-extraction.ts with wrong types ('activity'/'workflow' instead of 'consumer'/'producer'). Use existing utils/queue-extraction.ts which has correct types and full BullMQ + Temporal extraction logic. - Delete duplicate pipeline-phases/queue-extraction.ts - parse-impl.ts: import extractQueuePatterns from utils/ - queues.ts: simplify role check to 'producer' - parsing-processor.ts: drop typeEnvBindings ref (not on this branch) Co-Authored-By: Claude Opus 4.7 (1M context) --- .../src/core/ingestion/parsing-processor.ts | 1 - .../ingestion/pipeline-phases/parse-impl.ts | 4 +- .../pipeline-phases/queue-extraction.ts | 97 ------------------- .../core/ingestion/pipeline-phases/queues.ts | 3 +- 4 files changed, 3 insertions(+), 102 deletions(-) delete mode 100644 gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts diff --git a/gitnexus/src/core/ingestion/parsing-processor.ts b/gitnexus/src/core/ingestion/parsing-processor.ts index 716ffca07..190538869 100644 --- a/gitnexus/src/core/ingestion/parsing-processor.ts +++ b/gitnexus/src/core/ingestion/parsing-processor.ts @@ -160,7 +160,6 @@ const processParsingWithWorkers = async ( if (result.ormQueries) for (const item of result.ormQueries) allORMQueries.push(item); if (result.queuePatterns) for (const item of result.queuePatterns) allQueuePatterns.push(item); for (const item of result.constructorBindings) allConstructorBindings.push(item); - if (result.typeEnvBindings) for (const item of result.typeEnvBindings) allTypeEnvBindings.push(item); if (result.fileScopeBindings) for (const item of result.fileScopeBindings) fileScopeBindingsByFile.push(item); } diff --git a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts index c73e9b4c5..4a0a1b828 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts @@ -69,7 +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'; +import { extractQueuePatterns } from '../utils/queue-extraction.js'; // ── Constants ────────────────────────────────────────────────────────────── @@ -516,7 +516,7 @@ export async function runChunkedParseAndResolve( } for (const f of chunkFiles) { extractORMQueriesInline(f.path, f.content, allORMQueries); - extractQueuePatternsInline(f.path, f.content, allQueuePatterns); + extractQueuePatterns(f.path, f.content, allQueuePatterns); } astCache.clear(); cachedSequentialChunkFiles[chunkIdx] = []; diff --git a/gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts b/gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts deleted file mode 100644 index 94475f15f..000000000 --- a/gitnexus/src/core/ingestion/pipeline-phases/queue-extraction.ts +++ /dev/null @@ -1,97 +0,0 @@ -/** - * 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 index 45a49e6a1..15fd7504e 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/queues.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/queues.ts @@ -64,8 +64,7 @@ function processQueuePatterns( queueNodes.set(pt.queueName, queueNodeId); } - const edgeType = - pt.role === 'producer' || pt.role === 'workflow' ? 'ENQUEUES' : 'PROCESSES'; + const edgeType = pt.role === 'producer' ? 'ENQUEUES' : 'PROCESSES'; const fileId = generateId('File', pt.filePath); const edgeKey = `${fileId}->${queueNodeId}:${edgeType}`; if (seenEdges.has(edgeKey)) continue;