fix: address review findings — dedupe extraction, fix Temporal modeling

1. Move extractQueuePatterns() to shared utility module
   (utils/queue-extraction.ts) eliminating copy-paste between
   parse-worker.ts and pipeline.ts.

2. Fix Temporal semantic modeling:
   - Extract actual taskQueue name from workflow.start() options
     instead of using the workflow function name
   - Activity calls now create ENQUEUES edges (producer), not
     PROCESSES edges — they dispatch work, not consume it
   - Add proxyActivities guard to avoid false positives on generic
     activities.foo() calls in non-Temporal code
   - Simplify role union to 'producer' | 'consumer'

3. Add integration tests verifying taskQueue extraction and
   proxyActivities guard behavior.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Test 2026-03-26 11:43:19 +01:00
parent d53e790395
commit 149c4222d4
4 changed files with 210 additions and 30 deletions

View file

@ -3311,7 +3311,7 @@ export const processQueuePatterns = (graph: KnowledgeGraph, patterns: ExtractedQ
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';
const fid = generateId('File', pt.filePath), rt = pt.role === 'producer' ? '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++;
}

View file

@ -0,0 +1,154 @@
/**
* Shared queue pattern extraction for BullMQ + Temporal.
* Used by both parse-worker (worker threads) and pipeline (sequential fallback).
*/
import type { ExtractedQueuePattern } from '../workers/parse-worker.js';
// ---------------------------------------------------------------------------
// BullMQ regexes
// ---------------------------------------------------------------------------
/** Matches `const q = new Queue('orders')` -- captures var name + queue name */
const BULLMQ_QUEUE_DECL_RE = /(?:const|let|var)\s+(\w+)\s*=\s*new\s+Queue\s*\(\s*['"](\w[\w-]*)['"]/g;
/** Matches `q.add(...)` or `q.addBulk(...)` -- captures var name + method */
const BULLMQ_ADD_RE = /(\w+)\.(add|addBulk)\s*\(/g;
/** Matches `new Worker('orders', ...)` -- captures queue name */
const BULLMQ_WORKER_RE = /new\s+Worker\s*\(\s*['"](\w[\w-]*)['"]/g;
// ---------------------------------------------------------------------------
// Temporal regexes (more specific to avoid false positives)
// ---------------------------------------------------------------------------
/**
* Matches Temporal workflow.start/execute with taskQueue option.
* Pattern: `client.workflow.start(workflowFn, { taskQueue: 'orders' })`
* Captures: [1] = start|execute, [2] = workflow function name
*/
const TEMPORAL_WORKFLOW_START_RE = /client\.workflow\.(start|execute)\s*\(\s*(\w+)/g;
/**
* Matches Temporal activity invocations -- but ONLY when preceded by a
* `proxyActivities` import/call to reduce false positives on generic
* `activities.foo()` calls in non-Temporal code.
*/
const TEMPORAL_PROXY_ACTIVITIES_RE = /proxyActivities/;
/**
* Matches `activities.methodName(...)` -- captures method name.
* Only used when TEMPORAL_PROXY_ACTIVITIES_RE confirms Temporal context.
*/
const TEMPORAL_ACTIVITY_CALL_RE = /activities\.(\w+)\s*\(/g;
// ---------------------------------------------------------------------------
// Line number helper
// ---------------------------------------------------------------------------
function lineAt(content: string, index: number): number {
let count = 0;
for (let i = 0; i < index; i++) {
if (content.charCodeAt(i) === 10) count++;
}
return count;
}
// ---------------------------------------------------------------------------
// Main extraction
// ---------------------------------------------------------------------------
/**
* Extract BullMQ and Temporal queue patterns from file content.
* Appends results to `out` array (avoids allocation when no patterns found).
*/
export function extractQueuePatterns(
filePath: string,
content: string,
out: ExtractedQueuePattern[],
): void {
const hasBullMQ = content.includes('new Queue') || content.includes('new Worker');
const hasTemporal = content.includes('proxyActivities') || content.includes('client.workflow.');
if (!hasBullMQ && !hasTemporal) return;
// --- BullMQ ---
if (hasBullMQ) {
// Build variable-name -> queue-name map from `new Queue('name')` declarations
const queueVarMap = new Map<string, string>();
BULLMQ_QUEUE_DECL_RE.lastIndex = 0;
let m: RegExpExecArray | null;
while ((m = BULLMQ_QUEUE_DECL_RE.exec(content)) !== null) {
queueVarMap.set(m[1], m[2]);
}
// Producer: q.add() / q.addBulk()
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: lineAt(content, m.index),
});
}
}
// Consumer: new Worker('name', ...)
BULLMQ_WORKER_RE.lastIndex = 0;
while ((m = BULLMQ_WORKER_RE.exec(content)) !== null) {
out.push({
filePath,
role: 'consumer',
queueName: m[1],
lineNumber: lineAt(content, m.index),
});
}
}
// --- Temporal ---
if (hasTemporal) {
let m: RegExpExecArray | null;
// Workflow starter: client.workflow.start(workflowFn, { taskQueue: 'orders' })
// Extract the actual taskQueue name from the options object
if (content.includes('client.workflow.')) {
TEMPORAL_WORKFLOW_START_RE.lastIndex = 0;
while ((m = TEMPORAL_WORKFLOW_START_RE.exec(content)) !== null) {
const workflowFnName = m[2];
const startMethod = m[1];
// Look ahead in the next ~500 chars for taskQueue option
const lookAhead = content.substring(m.index, m.index + 500);
const tqMatch = lookAhead.match(/taskQueue\s*:\s*['"]([^'"]+)['"]/);
const taskQueueName = tqMatch ? tqMatch[1] : workflowFnName;
out.push({
filePath,
role: 'producer', // workflow.start is a producer (enqueues work)
queueName: taskQueueName,
method: startMethod,
handlerName: workflowFnName,
lineNumber: lineAt(content, m.index),
});
}
}
// Activity calls: activities.methodName() -- these are ENQUEUES (dispatching work)
// Only match when file uses proxyActivities (Temporal-specific)
if (TEMPORAL_PROXY_ACTIVITIES_RE.test(content)) {
TEMPORAL_ACTIVITY_CALL_RE.lastIndex = 0;
while ((m = TEMPORAL_ACTIVITY_CALL_RE.exec(content)) !== null) {
out.push({
filePath,
role: 'producer', // activity invocation dispatches work to task queue
queueName: m[1],
handlerName: m[1],
lineNumber: lineAt(content, m.index),
});
}
}
}
}

View file

@ -77,6 +77,7 @@ import {
buildCollisionGroups,
} from '../utils/method-props.js';
import type { LanguageProvider } from '../language-provider.js';
import { extractQueuePatterns } from '../utils/queue-extraction.js';
// ============================================================================
// Types for serializable results
@ -219,7 +220,7 @@ export interface ExtractedORMQuery {
export interface ExtractedQueuePattern {
filePath: string;
role: 'producer' | 'consumer' | 'workflow' | 'activity';
role: 'producer' | 'consumer';
queueName: string;
method?: string;
handlerName?: string;
@ -699,28 +700,7 @@ 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 }); }
}
}
// extractQueuePatterns is imported from ../utils/queue-extraction.js (shared with pipeline.ts)
// ============================================================================
// Process a batch of files

View file

@ -1,12 +1,58 @@
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); });
let graph: any;
beforeAll(async () => {
const result = await runPipelineFromRepo(QUEUE_REPO, () => {});
graph = result.graph;
}, 60_000);
it('creates ENQUEUES edges for BullMQ producers and Temporal starters', () => {
const enqueues: any[] = [];
graph.forEachRelationship((r: any) => { if (r.type === 'ENQUEUES') enqueues.push(r); });
expect(enqueues.length).toBeGreaterThanOrEqual(1);
});
it('creates PROCESSES edges for BullMQ consumers', () => {
const processes: any[] = [];
graph.forEachRelationship((r: any) => { if (r.type === 'PROCESSES') processes.push(r); });
expect(processes.length).toBeGreaterThanOrEqual(1);
});
it('creates CodeElement queue nodes', () => {
const queueNodes: any[] = [];
graph.forEachNode((n: any) => {
if (n.label === 'CodeElement' && n.properties?.description?.startsWith('Queue:')) {
queueNodes.push(n);
}
});
expect(queueNodes.length).toBeGreaterThanOrEqual(1);
});
it('extracts actual taskQueue name from Temporal workflow.start', () => {
const queueNodes: any[] = [];
graph.forEachNode((n: any) => {
if (n.label === 'CodeElement' && n.properties?.description?.startsWith('Queue:')) {
queueNodes.push(n);
}
});
const queueNames = queueNodes.map((n: any) => n.properties.name);
// starter.ts has taskQueue: 'orders' -- should use that, not the workflow function name
expect(queueNames).toContain('orders');
});
it('uses proxyActivities guard for Temporal activity detection', () => {
// workflow.ts has proxyActivities + activities.validateOrder etc.
// These should create ENQUEUES edges (producer role, not 'activity' role)
const enqueues: any[] = [];
graph.forEachRelationship((r: any) => { if (r.type === 'ENQUEUES') enqueues.push(r); });
const reasons = enqueues.map((r: any) => r.reason);
// Activity calls should have producer in reason
expect(reasons.some((r: string) => r.includes('producer'))).toBe(true);
});
});