mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-07 02:58:02 +00:00
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:
parent
972af286f6
commit
cacc07a7fa
6 changed files with 203 additions and 0 deletions
|
|
@ -15,6 +15,7 @@ export { parsePhase, type ParseOutput } from './parse.js';
|
||||||
export { routesPhase, type RoutesOutput, type RouteEntry } from './routes.js';
|
export { routesPhase, type RoutesOutput, type RouteEntry } from './routes.js';
|
||||||
export { toolsPhase, type ToolsOutput, type ToolDef } from './tools.js';
|
export { toolsPhase, type ToolsOutput, type ToolDef } from './tools.js';
|
||||||
export { ormPhase, type ORMOutput } from './orm.js';
|
export { ormPhase, type ORMOutput } from './orm.js';
|
||||||
|
export { queuesPhase, type QueuesOutput } from './queues.js';
|
||||||
export { crossFilePhase, type CrossFileOutput } from './cross-file.js';
|
export { crossFilePhase, type CrossFileOutput } from './cross-file.js';
|
||||||
export { mroPhase, type MROOutput } from './mro.js';
|
export { mroPhase, type MROOutput } from './mro.js';
|
||||||
export { communitiesPhase, type CommunitiesOutput } from './communities.js';
|
export { communitiesPhase, type CommunitiesOutput } from './communities.js';
|
||||||
|
|
|
||||||
|
|
@ -53,6 +53,7 @@ import type {
|
||||||
ExtractedDecoratorRoute,
|
ExtractedDecoratorRoute,
|
||||||
ExtractedFetchCall,
|
ExtractedFetchCall,
|
||||||
ExtractedORMQuery,
|
ExtractedORMQuery,
|
||||||
|
ExtractedQueuePattern,
|
||||||
ExtractedRoute,
|
ExtractedRoute,
|
||||||
ExtractedToolDef,
|
ExtractedToolDef,
|
||||||
FileConstructorBindings,
|
FileConstructorBindings,
|
||||||
|
|
@ -68,6 +69,7 @@ import { fileURLToPath, pathToFileURL } from 'node:url';
|
||||||
import { isDev } from '../utils/env.js';
|
import { isDev } from '../utils/env.js';
|
||||||
import { synthesizeWildcardImportBindings, needsSynthesis } from './wildcard-synthesis.js';
|
import { synthesizeWildcardImportBindings, needsSynthesis } from './wildcard-synthesis.js';
|
||||||
import { extractORMQueriesInline } from './orm-extraction.js';
|
import { extractORMQueriesInline } from './orm-extraction.js';
|
||||||
|
import { extractQueuePatternsInline } from './queue-extraction.js';
|
||||||
|
|
||||||
// ── Constants ──────────────────────────────────────────────────────────────
|
// ── Constants ──────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
@ -106,6 +108,7 @@ export async function runChunkedParseAndResolve(
|
||||||
allDecoratorRoutes: ExtractedDecoratorRoute[];
|
allDecoratorRoutes: ExtractedDecoratorRoute[];
|
||||||
allToolDefs: ExtractedToolDef[];
|
allToolDefs: ExtractedToolDef[];
|
||||||
allORMQueries: ExtractedORMQuery[];
|
allORMQueries: ExtractedORMQuery[];
|
||||||
|
allQueuePatterns: ExtractedQueuePattern[];
|
||||||
bindingAccumulator: BindingAccumulator;
|
bindingAccumulator: BindingAccumulator;
|
||||||
resolutionContext: ReturnType<typeof createResolutionContext>;
|
resolutionContext: ReturnType<typeof createResolutionContext>;
|
||||||
usedWorkerPool: boolean;
|
usedWorkerPool: boolean;
|
||||||
|
|
@ -248,6 +251,7 @@ export async function runChunkedParseAndResolve(
|
||||||
const allDecoratorRoutes: ExtractedDecoratorRoute[] = [];
|
const allDecoratorRoutes: ExtractedDecoratorRoute[] = [];
|
||||||
const allToolDefs: ExtractedToolDef[] = [];
|
const allToolDefs: ExtractedToolDef[] = [];
|
||||||
const allORMQueries: ExtractedORMQuery[] = [];
|
const allORMQueries: ExtractedORMQuery[] = [];
|
||||||
|
const allQueuePatterns: ExtractedQueuePattern[] = [];
|
||||||
const deferredWorkerCalls: ExtractedCall[] = [];
|
const deferredWorkerCalls: ExtractedCall[] = [];
|
||||||
const deferredWorkerHeritage: ExtractedHeritage[] = [];
|
const deferredWorkerHeritage: ExtractedHeritage[] = [];
|
||||||
const deferredConstructorBindings: FileConstructorBindings[] = [];
|
const deferredConstructorBindings: FileConstructorBindings[] = [];
|
||||||
|
|
@ -393,6 +397,9 @@ export async function runChunkedParseAndResolve(
|
||||||
if (chunkWorkerData.ormQueries?.length) {
|
if (chunkWorkerData.ormQueries?.length) {
|
||||||
for (const item of chunkWorkerData.ormQueries) allORMQueries.push(item);
|
for (const item of chunkWorkerData.ormQueries) allORMQueries.push(item);
|
||||||
}
|
}
|
||||||
|
if (chunkWorkerData.queuePatterns?.length) {
|
||||||
|
for (const item of chunkWorkerData.queuePatterns) allQueuePatterns.push(item);
|
||||||
|
}
|
||||||
} else {
|
} else {
|
||||||
await processImports(graph, chunkFiles, astCache, ctx, undefined, repoPath, allPaths);
|
await processImports(graph, chunkFiles, astCache, ctx, undefined, repoPath, allPaths);
|
||||||
sequentialChunkPaths.push(chunkPaths);
|
sequentialChunkPaths.push(chunkPaths);
|
||||||
|
|
@ -509,6 +516,7 @@ export async function runChunkedParseAndResolve(
|
||||||
}
|
}
|
||||||
for (const f of chunkFiles) {
|
for (const f of chunkFiles) {
|
||||||
extractORMQueriesInline(f.path, f.content, allORMQueries);
|
extractORMQueriesInline(f.path, f.content, allORMQueries);
|
||||||
|
extractQueuePatternsInline(f.path, f.content, allQueuePatterns);
|
||||||
}
|
}
|
||||||
astCache.clear();
|
astCache.clear();
|
||||||
cachedSequentialChunkFiles[chunkIdx] = [];
|
cachedSequentialChunkFiles[chunkIdx] = [];
|
||||||
|
|
@ -589,6 +597,7 @@ export async function runChunkedParseAndResolve(
|
||||||
allDecoratorRoutes,
|
allDecoratorRoutes,
|
||||||
allToolDefs,
|
allToolDefs,
|
||||||
allORMQueries,
|
allORMQueries,
|
||||||
|
allQueuePatterns,
|
||||||
bindingAccumulator,
|
bindingAccumulator,
|
||||||
resolutionContext: ctx,
|
resolutionContext: ctx,
|
||||||
// Whether a worker pool was actually live for this run. False means the
|
// Whether a worker pool was actually live for this run. False means the
|
||||||
|
|
|
||||||
|
|
@ -26,6 +26,7 @@ import type {
|
||||||
ExtractedDecoratorRoute,
|
ExtractedDecoratorRoute,
|
||||||
ExtractedToolDef,
|
ExtractedToolDef,
|
||||||
ExtractedORMQuery,
|
ExtractedORMQuery,
|
||||||
|
ExtractedQueuePattern,
|
||||||
} from '../workers/parse-worker.js';
|
} from '../workers/parse-worker.js';
|
||||||
import type { createResolutionContext } from '../model/resolution-context.js';
|
import type { createResolutionContext } from '../model/resolution-context.js';
|
||||||
import { runChunkedParseAndResolve } from './parse-impl.js';
|
import { runChunkedParseAndResolve } from './parse-impl.js';
|
||||||
|
|
@ -47,6 +48,7 @@ export interface ParseOutput {
|
||||||
readonly allDecoratorRoutes: readonly ExtractedDecoratorRoute[];
|
readonly allDecoratorRoutes: readonly ExtractedDecoratorRoute[];
|
||||||
readonly allToolDefs: readonly ExtractedToolDef[];
|
readonly allToolDefs: readonly ExtractedToolDef[];
|
||||||
readonly allORMQueries: readonly ExtractedORMQuery[];
|
readonly allORMQueries: readonly ExtractedORMQuery[];
|
||||||
|
readonly allQueuePatterns: readonly ExtractedQueuePattern[];
|
||||||
bindingAccumulator: BindingAccumulator;
|
bindingAccumulator: BindingAccumulator;
|
||||||
/** Resolution context from the parse phase — carries importMap, namedImportMap, etc. */
|
/** Resolution context from the parse phase — carries importMap, namedImportMap, etc. */
|
||||||
resolutionContext: ReturnType<typeof createResolutionContext>;
|
resolutionContext: ReturnType<typeof createResolutionContext>;
|
||||||
|
|
|
||||||
|
|
@ -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,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
92
gitnexus/src/core/ingestion/pipeline-phases/queues.ts
Normal file
92
gitnexus/src/core/ingestion/pipeline-phases/queues.ts
Normal 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 };
|
||||||
|
}
|
||||||
|
|
@ -29,6 +29,7 @@ import {
|
||||||
routesPhase,
|
routesPhase,
|
||||||
toolsPhase,
|
toolsPhase,
|
||||||
ormPhase,
|
ormPhase,
|
||||||
|
queuesPhase,
|
||||||
crossFilePhase,
|
crossFilePhase,
|
||||||
mroPhase,
|
mroPhase,
|
||||||
communitiesPhase,
|
communitiesPhase,
|
||||||
|
|
@ -79,6 +80,7 @@ function buildPhaseList(options?: PipelineOptions): PipelinePhase[] {
|
||||||
routesPhase,
|
routesPhase,
|
||||||
toolsPhase,
|
toolsPhase,
|
||||||
ormPhase,
|
ormPhase,
|
||||||
|
queuesPhase,
|
||||||
crossFilePhase,
|
crossFilePhase,
|
||||||
];
|
];
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue