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 <noreply@anthropic.com>
This commit is contained in:
Test 2026-04-18 13:14:21 +02:00
parent af492be01b
commit 94d6aa2365
6 changed files with 203 additions and 0 deletions

View file

@ -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';

View file

@ -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<typeof createResolutionContext>;
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

View file

@ -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<typeof createResolutionContext>;

View file

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

View file

@ -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<QueuesOutput> = {
name: 'queues',
deps: ['parse'],
async execute(
ctx: PipelineContext,
deps: ReadonlyMap<string, PhaseResult<unknown>>,
): Promise<QueuesOutput> {
const { allQueuePatterns } = getPhaseOutput<ParseOutput>(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<string, string>();
const seenEdges = new Set<string>();
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 };
}

View file

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