From af492be01b6d2d4345a61baba6816018fa5d1693 Mon Sep 17 00:00:00 2001 From: Test Date: Wed, 25 Mar 2026 18:54:56 +0100 Subject: [PATCH] feat: add BullMQ + Temporal async queue detection (ENQUEUES/PROCESSES edges) Co-Authored-By: Claude Opus 4.6 (1M context) --- gitnexus-shared/src/graph/types.ts | 4 +- gitnexus-shared/src/lbug/schema-constants.ts | 2 + gitnexus/src/core/ingestion/call-processor.ts | 18 +++++++++ .../src/core/ingestion/parsing-processor.ts | 6 +++ .../core/ingestion/workers/parse-worker.ts | 37 +++++++++++++++++++ .../test/fixtures/queue-repo/src/queue.ts | 3 ++ .../test/fixtures/queue-repo/src/starter.ts | 5 +++ .../fixtures/queue-repo/src/video-worker.ts | 2 + .../test/fixtures/queue-repo/src/workflow.ts | 7 ++++ .../test/integration/queue-detection.test.ts | 12 ++++++ 10 files changed, 95 insertions(+), 1 deletion(-) create mode 100644 gitnexus/test/fixtures/queue-repo/src/queue.ts create mode 100644 gitnexus/test/fixtures/queue-repo/src/starter.ts create mode 100644 gitnexus/test/fixtures/queue-repo/src/video-worker.ts create mode 100644 gitnexus/test/fixtures/queue-repo/src/workflow.ts create mode 100644 gitnexus/test/integration/queue-detection.test.ts diff --git a/gitnexus-shared/src/graph/types.ts b/gitnexus-shared/src/graph/types.ts index 49762d145..555db98ac 100644 --- a/gitnexus-shared/src/graph/types.ts +++ b/gitnexus-shared/src/graph/types.ts @@ -115,7 +115,9 @@ export type RelationshipType = | 'HANDLES_TOOL' | 'ENTRY_POINT_OF' | 'WRAPS' - | 'QUERIES'; + | 'QUERIES' + | 'ENQUEUES' + | 'PROCESSES'; export interface GraphNode { id: string; diff --git a/gitnexus-shared/src/lbug/schema-constants.ts b/gitnexus-shared/src/lbug/schema-constants.ts index 656ffe552..af1ab22d5 100644 --- a/gitnexus-shared/src/lbug/schema-constants.ts +++ b/gitnexus-shared/src/lbug/schema-constants.ts @@ -67,6 +67,8 @@ export const REL_TYPES = [ 'ENTRY_POINT_OF', 'WRAPS', 'QUERIES', + 'ENQUEUES', + 'PROCESSES', ] as const; export type RelType = (typeof REL_TYPES)[number]; diff --git a/gitnexus/src/core/ingestion/call-processor.ts b/gitnexus/src/core/ingestion/call-processor.ts index c30d83848..4435d6be2 100644 --- a/gitnexus/src/core/ingestion/call-processor.ts +++ b/gitnexus/src/core/ingestion/call-processor.ts @@ -70,6 +70,7 @@ import type { ExtractedAssignment, ExtractedRoute, ExtractedFetchCall, + ExtractedQueuePattern, FileConstructorBindings, } from './workers/parse-worker.js'; import { normalizeFetchURL, routeMatches } from './route-extractors/nextjs.js'; @@ -3300,3 +3301,20 @@ export const extractFetchCallsFromFiles = async ( return result; }; + +export const processQueuePatterns = (graph: KnowledgeGraph, patterns: ExtractedQueuePattern[]): { queuesCreated: number; edgesCreated: number } => { + if (patterns.length === 0) return { queuesCreated: 0, edgesCreated: 0 }; + const byQueue = new Map(); + for (const p of patterns) { const e = byQueue.get(p.queueName); if (e) e.push(p); else byQueue.set(p.queueName, [p]); } + let queuesCreated = 0, edgesCreated = 0; + for (const [qn, qp] of byQueue) { + const qid = generateId('CodeElement', 'queue:' + qn); + if (!graph.getNode(qid)) { graph.addNode({ id: qid, label: 'CodeElement', properties: { name: qn, filePath: qp[0].filePath, description: 'Queue: ' + qn } }); queuesCreated++; } + for (const pt of qp) { + const fid = generateId('File', pt.filePath), rt = (pt.role === 'producer' || pt.role === 'workflow') ? 'ENQUEUES' : 'PROCESSES'; + graph.addRelationship({ id: generateId(rt, fid+'->'+qid+':'+pt.lineNumber), sourceId: fid, targetId: qid, type: rt, confidence: 0.9, reason: pt.role+'-'+(pt.method ?? 'handler') }); + edgesCreated++; + } + } + return { queuesCreated, edgesCreated }; +}; diff --git a/gitnexus/src/core/ingestion/parsing-processor.ts b/gitnexus/src/core/ingestion/parsing-processor.ts index 90186c659..377a05e53 100644 --- a/gitnexus/src/core/ingestion/parsing-processor.ts +++ b/gitnexus/src/core/ingestion/parsing-processor.ts @@ -45,6 +45,7 @@ import type { FileConstructorBindings, FileScopeBindings, ExtractedORMQuery, + ExtractedQueuePattern, } from './workers/parse-worker.js'; import { getTreeSitterBufferSize, TREE_SITTER_MAX_BUFFER } from './constants.js'; @@ -60,6 +61,7 @@ export interface WorkerExtractedData { decoratorRoutes: ExtractedDecoratorRoute[]; toolDefs: ExtractedToolDef[]; ormQueries: ExtractedORMQuery[]; + queuePatterns: ExtractedQueuePattern[]; constructorBindings: FileConstructorBindings[]; fileScopeBindings: FileScopeBindings[]; } @@ -94,6 +96,7 @@ const processParsingWithWorkers = async ( decoratorRoutes: [], toolDefs: [], ormQueries: [], + queuePatterns: [], constructorBindings: [], fileScopeBindings: [], }; @@ -118,6 +121,7 @@ const processParsingWithWorkers = async ( const allDecoratorRoutes: ExtractedDecoratorRoute[] = []; const allToolDefs: ExtractedToolDef[] = []; const allORMQueries: ExtractedORMQuery[] = []; + const allQueuePatterns: ExtractedQueuePattern[] = []; const allConstructorBindings: FileConstructorBindings[] = []; const fileScopeBindingsByFile: FileScopeBindings[] = []; for (const result of chunkResults) { @@ -154,6 +158,7 @@ const processParsingWithWorkers = async ( for (const item of result.decoratorRoutes) allDecoratorRoutes.push(item); for (const item of result.toolDefs) allToolDefs.push(item); 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.fileScopeBindings) for (const item of result.fileScopeBindings) fileScopeBindingsByFile.push(item); @@ -185,6 +190,7 @@ const processParsingWithWorkers = async ( decoratorRoutes: allDecoratorRoutes, toolDefs: allToolDefs, ormQueries: allORMQueries, + queuePatterns: allQueuePatterns, constructorBindings: allConstructorBindings, fileScopeBindings: fileScopeBindingsByFile, }; diff --git a/gitnexus/src/core/ingestion/workers/parse-worker.ts b/gitnexus/src/core/ingestion/workers/parse-worker.ts index 203e027fa..cda935f2c 100644 --- a/gitnexus/src/core/ingestion/workers/parse-worker.ts +++ b/gitnexus/src/core/ingestion/workers/parse-worker.ts @@ -217,6 +217,15 @@ export interface ExtractedORMQuery { lineNumber: number; } +export interface ExtractedQueuePattern { + filePath: string; + role: 'producer' | 'consumer' | 'workflow' | 'activity'; + queueName: string; + method?: string; + handlerName?: string; + lineNumber: number; +} + /** Constructor bindings keyed by filePath for cross-file type resolution */ export interface FileConstructorBindings { filePath: string; @@ -266,6 +275,7 @@ export interface ParseWorkerResult { decoratorRoutes: ExtractedDecoratorRoute[]; toolDefs: ExtractedToolDef[]; ormQueries: ExtractedORMQuery[]; + queuePatterns: ExtractedQueuePattern[]; constructorBindings: FileConstructorBindings[]; /** All-scope type bindings from TypeEnv for BindingAccumulator (includes function-local). */ fileScopeBindings: FileScopeBindings[]; @@ -688,6 +698,28 @@ const cachedExportCheck = ( // DEFINITION_CAPTURE_KEYS and getDefinitionNodeFromCaptures imported from ../utils.js + +// BullMQ + Temporal Queue Pattern Extraction +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; +export function extractQueuePatterns(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 }); } + } +} + // ============================================================================ // Process a batch of files // ============================================================================ @@ -709,6 +741,7 @@ const processBatch = ( decoratorRoutes: [], toolDefs: [], ormQueries: [], + queuePatterns: [], constructorBindings: [], fileScopeBindings: [], skippedLanguages: {}, @@ -2246,6 +2279,7 @@ const processFileGroup = ( // Extract ORM queries (Prisma, Supabase) extractORMQueries(file.path, parseContent, result.ormQueries); + extractQueuePatterns(file.path, file.content, result.queuePatterns); // Vue: emit CALLS edges for components used in