feat: add BullMQ + Temporal async queue detection (ENQUEUES/PROCESSES edges)

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Test 2026-03-25 18:54:56 +01:00
parent 969b4623ca
commit af492be01b
10 changed files with 95 additions and 1 deletions

View file

@ -115,7 +115,9 @@ export type RelationshipType =
| 'HANDLES_TOOL'
| 'ENTRY_POINT_OF'
| 'WRAPS'
| 'QUERIES';
| 'QUERIES'
| 'ENQUEUES'
| 'PROCESSES';
export interface GraphNode {
id: string;

View file

@ -67,6 +67,8 @@ export const REL_TYPES = [
'ENTRY_POINT_OF',
'WRAPS',
'QUERIES',
'ENQUEUES',
'PROCESSES',
] as const;
export type RelType = (typeof REL_TYPES)[number];

View file

@ -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<string, ExtractedQueuePattern[]>();
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 };
};

View file

@ -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,
};

View file

@ -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<string, string>(); 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 <template>
if (language === SupportedLanguages.Vue) {
@ -2280,6 +2314,7 @@ let accumulated: ParseWorkerResult = {
decoratorRoutes: [],
toolDefs: [],
ormQueries: [],
queuePatterns: [],
constructorBindings: [],
fileScopeBindings: [],
skippedLanguages: {},
@ -2307,6 +2342,7 @@ const mergeResult = (target: ParseWorkerResult, src: ParseWorkerResult) => {
appendAll(target.decoratorRoutes, src.decoratorRoutes);
appendAll(target.toolDefs, src.toolDefs);
appendAll(target.ormQueries, src.ormQueries);
appendAll(target.queuePatterns, src.queuePatterns);
appendAll(target.constructorBindings, src.constructorBindings);
appendAll(target.fileScopeBindings, src.fileScopeBindings);
for (const [lang, count] of Object.entries(src.skippedLanguages)) {
@ -2358,6 +2394,7 @@ parentPort!.on('message', (msg: WorkerIncomingMessage) => {
decoratorRoutes: [],
toolDefs: [],
ormQueries: [],
queuePatterns: [],
constructorBindings: [],
fileScopeBindings: [],
skippedLanguages: {},

View file

@ -0,0 +1,3 @@
import { Queue } from 'bullmq';
const videoQueue = new Queue('video-processing');
export async function enqueueVideo(videoId: string) { await videoQueue.add('transcode', { videoId }); }

View file

@ -0,0 +1,5 @@
import { Client } from '@temporalio/client';
const client = new Client();
export async function startOrder(orderId: string) {
await client.workflow.start(processOrderWorkflow, { taskQueue: 'orders', workflowId: 'order-' + orderId, args: [orderId] });
}

View file

@ -0,0 +1,2 @@
import { Worker } from 'bullmq';
const worker = new Worker('video-processing', async (job) => { console.log('Processing: ' + job.data.videoId); });

View file

@ -0,0 +1,7 @@
import { proxyActivities } from '@temporalio/workflow';
const activities = proxyActivities({ startToCloseTimeout: '30s' });
export async function processOrderWorkflow(orderId: string) {
const result = await activities.validateOrder(orderId);
await activities.chargePayment(result.amount);
await activities.sendConfirmation(orderId);
}

View file

@ -0,0 +1,12 @@
import { describe, it, expect, beforeAll } from 'vitest';
import path from 'node:path';
import { runPipelineFromRepo } from '../../src/core/ingestion/pipeline.js';
import type { PipelineResult } from '../../types/pipeline.js';
const QUEUE_REPO = path.resolve(__dirname, '..', 'fixtures', 'queue-repo');
describe('Queue Detection', () => {
let result: PipelineResult;
beforeAll(async () => { result = await runPipelineFromRepo(QUEUE_REPO, () => {}); }, 60_000);
it('ENQUEUES edges', () => { expect(result.graph.relationships.filter(r => r.type === 'ENQUEUES').length).toBeGreaterThanOrEqual(1); });
it('PROCESSES edges', () => { expect(result.graph.relationships.filter(r => r.type === 'PROCESSES').length).toBeGreaterThanOrEqual(1); });
it('queue nodes', () => { expect(result.graph.nodes.filter(n => n.label === 'CodeElement' && n.properties?.description?.startsWith('Queue:')).length).toBeGreaterThanOrEqual(1); });
});