From e8286247d640cbd828c6f81f47a45eb8d789582c Mon Sep 17 00:00:00 2001 From: MythEclipse Date: Mon, 22 Jun 2026 17:42:46 +0700 Subject: [PATCH] fix: parallelize text + media LLM analysis instead of sequential - text-only and media analysis now run concurrently via Promise.all - text no longer blocks on media download + vision analysis - each path independently saves to DB when its own results are ready - same batch still uses single context fetch + attachment lookup --- .../modules/ai-moderation/aiAnalysisWorker.ts | 126 ++++++++++-------- .../ai-moderation/llmModerationClient.ts | 44 +++--- 2 files changed, 94 insertions(+), 76 deletions(-) diff --git a/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts b/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts index 50479e4..18baedf 100644 --- a/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts +++ b/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts @@ -180,65 +180,75 @@ async function processBatch(job: { const allRows: MessageRecord[] = []; - // Phase 1: Text-only → save immediately (fast) - if (textOnly.length > 0) { - const textResult = await runModerationAnalysis({ - targets: textOnly, - contextText: contextLines.join("\n"), - attachments, - }); - const textUpdates = textResult.results.map((analysisResult) => ({ - messageId: analysisResult.messageId, - result: { - status: analysisResult.status, - flags: JSON.stringify(analysisResult.flags), - score: analysisResult.score, - analysis: analysisResult.analysis, - categories: analysisResult.categories, - severity: analysisResult.severity, - confidence: analysisResult.confidence, - recommendedAction: analysisResult.recommendedAction, - analyzedAt: Date.now(), - error: null, - }, - })); - if (textUpdates.length > 0) { - const rows = await updateMessagesAIAnalysisBulk(textUpdates); - allRows.push(...rows); - logger.info( - { count: textUpdates.length, conversationKey }, - "Text-only batch saved — media analysis still in progress", - ); - } - } + // ── Parallel: text-only + media analysis run concurrently ────────── + // Text-only → fast LLM call. Media → download + vision + LLM. + // Running both in parallel means media downloads overlap with text LLM call. + // Each path saves to DB as soon as its own results are ready. + // ──────────────────────────────────────────────────────────────────── + const textPromise = textOnly.length > 0 + ? runModerationAnalysis({ + targets: textOnly, + contextText: contextLines.join("\n"), + attachments, + }).then((result) => { + const updates = result.results.map((analysisResult) => ({ + messageId: analysisResult.messageId, + result: { + status: analysisResult.status, + flags: JSON.stringify(analysisResult.flags), + score: analysisResult.score, + analysis: analysisResult.analysis, + categories: analysisResult.categories, + severity: analysisResult.severity, + confidence: analysisResult.confidence, + recommendedAction: analysisResult.recommendedAction, + analyzedAt: Date.now(), + error: null, + }, + })); + if (updates.length > 0) { + return updateMessagesAIAnalysisBulk(updates).then((rows) => { + allRows.push(...rows); + logger.info( + { count: updates.length, conversationKey }, + "Text-only batch saved — media analysis still in progress", + ); + }); + } + }) + : Promise.resolve(); - // Phase 2: Media → save when done (slow: download + vision) - if (media.length > 0) { - const mediaResult = await runModerationAnalysis({ - targets: media, - contextText: contextLines.join("\n"), - attachments, - }); - const mediaUpdates = mediaResult.results.map((analysisResult) => ({ - messageId: analysisResult.messageId, - result: { - status: analysisResult.status, - flags: JSON.stringify(analysisResult.flags), - score: analysisResult.score, - analysis: analysisResult.analysis, - categories: analysisResult.categories, - severity: analysisResult.severity, - confidence: analysisResult.confidence, - recommendedAction: analysisResult.recommendedAction, - analyzedAt: Date.now(), - error: null, - }, - })); - if (mediaUpdates.length > 0) { - const rows = await updateMessagesAIAnalysisBulk(mediaUpdates); - allRows.push(...rows); - } - } + const mediaPromise = media.length > 0 + ? runModerationAnalysis({ + targets: media, + contextText: contextLines.join("\n"), + attachments, + }).then((result) => { + const updates = result.results.map((analysisResult) => ({ + messageId: analysisResult.messageId, + result: { + status: analysisResult.status, + flags: JSON.stringify(analysisResult.flags), + score: analysisResult.score, + analysis: analysisResult.analysis, + categories: analysisResult.categories, + severity: analysisResult.severity, + confidence: analysisResult.confidence, + recommendedAction: analysisResult.recommendedAction, + analyzedAt: Date.now(), + error: null, + }, + })); + if (updates.length > 0) { + return updateMessagesAIAnalysisBulk(updates).then((rows) => { + allRows.push(...rows); + }); + } + }) + : Promise.resolve(); + + // Wait for both to complete + await Promise.all([textPromise, mediaPromise]); logger.info( { total: messages.length, textOnly: textOnly.length, media: media.length, saved: allRows.length }, diff --git a/services/discord-gateway/src/modules/ai-moderation/llmModerationClient.ts b/services/discord-gateway/src/modules/ai-moderation/llmModerationClient.ts index d68929a..49b4991 100644 --- a/services/discord-gateway/src/modules/ai-moderation/llmModerationClient.ts +++ b/services/discord-gateway/src/modules/ai-moderation/llmModerationClient.ts @@ -707,12 +707,13 @@ async function runTextOnlyBatch( const maxBatchSize = config.AI_LLM_TEXT_BATCH_SIZE ?? 20; const timeoutMs = config.AI_LLM_MEDIA_ANALYSIS_TIMEOUT_MS ?? 60000; + // ── Phase A: Prepare context in parallel ────────────────── + // URL fetching (web_content) and SearXNG (web_searches) are independent — + // both just enrich the LLM prompt. Run them concurrently so SearXNG + // doesn't block on slow URLs (or vice versa). + // // ── Fetch web content from URLs in text-only messages ── - // Prevents LLM from guessing based on domain name alone (e.g., false "scam" flags). - // The fetched text content is injected into the message XML so the LLM can - // analyze the actual page rather than pattern-match the URL string. - const urlFetchMap = new Map(); // url → fetched text content - { + const urlFetchPromise = (async () => { const allUrls = new Set(); for (const msg of targets) { const content = msg.edited_content ?? msg.content; @@ -720,7 +721,7 @@ async function runTextOnlyBatch( allUrls.add(url); } } - const urlArr = Array.from(allUrls).slice(0, 10); // cap to 10 fetches per batch + const urlArr = Array.from(allUrls).slice(0, 10); if (urlArr.length > 0) { log.debug( { urlCount: urlArr.length }, @@ -729,6 +730,7 @@ async function runTextOnlyBatch( const results = await Promise.allSettled( urlArr.map((url) => fetchUrlSafely(url)), ); + const map = new Map(); for (let i = 0; i < urlArr.length; i++) { const r = results[i]; if ( @@ -736,19 +738,16 @@ async function runTextOnlyBatch( r.value.type === "text" && r.value.textContent ) { - urlFetchMap.set(urlArr[i], r.value.textContent); + map.set(urlArr[i], r.value.textContent); } } + return map; } - } + return new Map(); + })(); - // ── SearXNG enrichment for suspicious/ambiguous content ────────── - // Uses SearXNG to look up references mentioned in messages (e.g. anime - // titles, drug names). Results are injected as XML tags - // so the LLM can make informed decisions instead of guessing. - // All messages are searched — Redis cache prevents redundant lookups. - const searxngResults = new Map(); // query → formatted XML - { + // ── SearXNG enrichment ── + const searxngPromise = (async () => { const queries = new Set(); for (const msg of targets) { const content = msg.edited_content ?? msg.content; @@ -757,7 +756,7 @@ async function runTextOnlyBatch( } } if (queries.size > 0) { - const queryArr = Array.from(queries).slice(0, 3); // cap to 3 searches per batch + const queryArr = Array.from(queries).slice(0, 3); log.debug( { searchQueries: queryArr }, "Running SearXNG enrichment for batch", @@ -765,14 +764,23 @@ async function runTextOnlyBatch( const results = await Promise.allSettled( queryArr.map((q) => searchSearxng(q)), ); + const map = new Map(); for (let i = 0; i < queryArr.length; i++) { const r = results[i]; if (r.status === "fulfilled" && r.value.length > 0) { - searxngResults.set(queryArr[i], formatSearchResults(r.value)); + map.set(queryArr[i], formatSearchResults(r.value)); } } + return map; } - } + return new Map(); + })(); + + // Wait for BOTH concurrently + const [urlFetchMap, searxngResults] = await Promise.all([ + urlFetchPromise, + searxngPromise, + ]); // ── Group identical short messages (< 20 chars) to reduce redundant analysis ── // Messages with identical normalized content share a single representative.