diff --git a/gitnexus/src/core/ingestion/utils/queue-extraction.ts b/gitnexus/src/core/ingestion/utils/queue-extraction.ts index 1fafed729..ce0fbcb7d 100644 --- a/gitnexus/src/core/ingestion/utils/queue-extraction.ts +++ b/gitnexus/src/core/ingestion/utils/queue-extraction.ts @@ -9,13 +9,13 @@ import type { ExtractedQueuePattern } from '../workers/parse-worker.js'; // --------------------------------------------------------------------------- /** 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; +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; +const BULLMQ_WORKER_RE = /new\s+Worker\s*\(\s*['"]([\w][\w:.-]*)['"]/g; // --------------------------------------------------------------------------- // Temporal regexes (more specific to avoid false positives) @@ -45,12 +45,27 @@ 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++; +/** + * Build a sorted array of newline offsets so lineAt lookups are O(log n) + * via binary search instead of O(n) per call. + */ +function buildLineOffsets(content: string): number[] { + const offsets: number[] = []; + for (let i = 0; i < content.length; i++) { + if (content.charCodeAt(i) === 10) offsets.push(i); } - return count; + return offsets; +} + +function lineAt(offsets: number[], index: number): number { + let lo = 0; + let hi = offsets.length; + while (lo < hi) { + const mid = (lo + hi) >>> 1; + if (offsets[mid] < index) lo = mid + 1; + else hi = mid; + } + return lo; } // --------------------------------------------------------------------------- @@ -71,6 +86,8 @@ export function extractQueuePatterns( if (!hasBullMQ && !hasTemporal) return; + const offsets = buildLineOffsets(content); + // --- BullMQ --- if (hasBullMQ) { // Build variable-name -> queue-name map from `new Queue('name')` declarations @@ -91,7 +108,7 @@ export function extractQueuePatterns( role: 'producer', queueName: qn, method: m[2], - lineNumber: lineAt(content, m.index), + lineNumber: lineAt(offsets, m.index), }); } } @@ -103,7 +120,7 @@ export function extractQueuePatterns( filePath, role: 'consumer', queueName: m[1], - lineNumber: lineAt(content, m.index), + lineNumber: lineAt(offsets, m.index), }); } } @@ -131,7 +148,7 @@ export function extractQueuePatterns( queueName: taskQueueName, method: startMethod, handlerName: workflowFnName, - lineNumber: lineAt(content, m.index), + lineNumber: lineAt(offsets, m.index), }); } } @@ -146,7 +163,7 @@ export function extractQueuePatterns( role: 'producer', // activity invocation dispatches work to task queue queueName: m[1], handlerName: m[1], - lineNumber: lineAt(content, m.index), + lineNumber: lineAt(offsets, m.index), }); } } diff --git a/gitnexus/test/fixtures/queue-repo/src/payments-producer.ts b/gitnexus/test/fixtures/queue-repo/src/payments-producer.ts new file mode 100644 index 000000000..bc185a501 --- /dev/null +++ b/gitnexus/test/fixtures/queue-repo/src/payments-producer.ts @@ -0,0 +1,3 @@ +import { Queue } from 'bullmq'; +const paymentQueue = new Queue('payments:high-priority'); +export async function enqueuePayment(paymentId: string) { await paymentQueue.add('process', { paymentId }); } diff --git a/gitnexus/test/fixtures/queue-repo/src/payments-worker.ts b/gitnexus/test/fixtures/queue-repo/src/payments-worker.ts new file mode 100644 index 000000000..5e2b48a65 --- /dev/null +++ b/gitnexus/test/fixtures/queue-repo/src/payments-worker.ts @@ -0,0 +1,2 @@ +import { Worker } from 'bullmq'; +const worker = new Worker('payments:high-priority', async (job) => { console.log('Processing payment: ' + job.data.paymentId); }); diff --git a/gitnexus/test/integration/queue-detection.test.ts b/gitnexus/test/integration/queue-detection.test.ts index a5d31220b..d2b792d27 100644 --- a/gitnexus/test/integration/queue-detection.test.ts +++ b/gitnexus/test/integration/queue-detection.test.ts @@ -46,6 +46,38 @@ describe('Queue Detection', () => { expect(queueNames).toContain('orders'); }); + it('detects colon-namespaced BullMQ queue names', () => { + 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); + expect(queueNames).toContain('payments:high-priority'); + }); + + it('creates ENQUEUES edge for colon-namespaced producer', () => { + const enqueues: any[] = []; + graph.forEachRelationship((r: any) => { if (r.type === 'ENQUEUES') enqueues.push(r); }); + // The target node of the ENQUEUES edge should be the colon-namespaced queue + const targetNames = enqueues.map((r: any) => { + const target = graph.getNode(r.targetId); + return target?.properties?.name; + }); + expect(targetNames).toContain('payments:high-priority'); + }); + + it('creates PROCESSES edge for colon-namespaced consumer', () => { + const processes: any[] = []; + graph.forEachRelationship((r: any) => { if (r.type === 'PROCESSES') processes.push(r); }); + const targetNames = processes.map((r: any) => { + const target = graph.getNode(r.targetId); + return target?.properties?.name; + }); + expect(targetNames).toContain('payments:high-priority'); + }); + 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)