From 4996acfacf23c348ce05785ac2bdae07b5cd1f6e Mon Sep 17 00:00:00 2001 From: MythEclipse Date: Thu, 4 Jun 2026 17:39:40 +0700 Subject: [PATCH] refactor(ai-moderation): unify worker pool entry points with discriminated union routing --- .../modules/ai-moderation/aiAnalysisWorker.ts | 311 ++++++++---------- .../src/modules/ai-moderation/aiAnalyzer.ts | 3 +- 2 files changed, 138 insertions(+), 176 deletions(-) diff --git a/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts b/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts index 2421271..a784327 100644 --- a/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts +++ b/services/discord-gateway/src/modules/ai-moderation/aiAnalysisWorker.ts @@ -23,38 +23,31 @@ async function ensureDb() { } // --------------------------------------------------------------------------- -// Batch analysis (existing) +// Job types — the default export routes on `type` // --------------------------------------------------------------------------- -export interface AnalysisWorkerRequest { - conversationKey: string; - messages: MessageRecord[]; -} +type WorkerJob = + | { type: "batch"; conversationKey: string; messages: MessageRecord[] } + | { type: "individual"; message: MessageRecord; skipNormalAnalysis: boolean }; -export type AnalysisWorkerResponse = - | { - ok: true; - conversationKey: string; - rows: MessageRecord[]; - } - | { - ok: false; - conversationKey: string; - rows: MessageRecord[]; - error: string; - }; +type BatchOkResponse = { ok: true; conversationKey: string; rows: MessageRecord[] }; +type BatchErrorResponse = { ok: false; conversationKey: string; rows: MessageRecord[]; error: string }; +type IndividualOkResponse = { ok: true; results: AnalysisResult[] }; +type IndividualErrorResponse = { ok: false; results: AnalysisResult[]; error: string }; -export default async function processAnalysisRequest({ - conversationKey, - messages, -}: AnalysisWorkerRequest): Promise { +type WorkerResponse = BatchOkResponse | BatchErrorResponse | IndividualOkResponse | IndividualErrorResponse; + +/** + * Default export — Piscina worker entry point. + * Routes to the correct handler based on `type` field. + */ +export default async function workerRouter(job: WorkerJob): Promise { if (!config.AI_LLM_API_KEY) { console.error( JSON.stringify({ level: "FATAL", context: "aiAnalysisWorker", - error: - "AI_LLM_API_KEY is missing from environment. Force closing worker operation.", + error: "AI_LLM_API_KEY is missing from environment. Force closing worker operation.", timestamp: new Date().toISOString(), }), ); @@ -62,165 +55,133 @@ export default async function processAnalysisRequest({ } try { - try { - await ensureDb(); - } catch (dbError) { - const msg = dbError instanceof Error ? dbError.message : String(dbError); - return { - ok: false, - conversationKey, - rows: [], - error: `Database init failed: ${msg}`, - }; + await ensureDb(); + } catch (dbError) { + const msg = dbError instanceof Error ? dbError.message : String(dbError); + if (job.type === "batch") { + return { ok: false, conversationKey: job.conversationKey, rows: [], error: `Database init failed: ${msg}` }; } - - const firstMessage = messages[0]; - if (!firstMessage) return { ok: true, conversationKey, rows: [] }; - - const contextBefore = await getConversationContextBefore({ - channelId: firstMessage.channel_id, - threadId: firstMessage.thread_id, - beforeCreatedAt: firstMessage.created_at, - limit: config.AI_ANALYSIS_CONTEXT_MESSAGE_LIMIT, - }); - - const contextLines = buildConversationContext({ - contextBefore, - targets: messages, - maxTokens: config.AI_ANALYSIS_MAX_CONTEXT_TOKENS, - }); - - const targetIds = messages.map((m) => m.id); - const contextIds = contextBefore.map((m) => m.id); - const allMessageIds = [...targetIds, ...contextIds]; - const attachments = await getAttachmentsForMessages(allMessageIds); - - const result = await runModerationAnalysis({ - targets: messages, - contextText: contextLines.join("\n"), - attachments, - }); - - 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, - }, - })); - - try { - const rows = await updateMessagesAIAnalysisBulk(updates); - return { ok: true, conversationKey, rows }; - } catch (dbErr) { - throw new Error( - `Failed to update DB: ${dbErr instanceof Error ? dbErr.message : String(dbErr)}`, - ); - } - } catch (error) { - const errorMessage = error instanceof Error ? error.message : String(error); - const errorStack = error instanceof Error ? error.stack : undefined; - const rows: MessageRecord[] = []; - - console.error( - JSON.stringify({ - level: "ERROR", - context: "aiAnalysisWorker", - conversationKey, - messageCount: messages.length, - error: errorMessage, - stack: errorStack, - timestamp: new Date().toISOString(), - }), - ); - - return { ok: false, conversationKey, rows, error: errorMessage }; - } -} - -// --------------------------------------------------------------------------- -// Individual fallback analysis (offloaded from main thread) -// --------------------------------------------------------------------------- - -export interface IndividualWorkerRequest { - message: MessageRecord; - /** Optional — if true, skip normal analysis and go straight to simple fallback */ - skipNormalAnalysis: boolean; -} - -export type IndividualWorkerResponse = - | { - ok: true; - results: AnalysisResult[]; - } - | { - ok: false; - results: AnalysisResult[]; - error: string; - }; - -/** - * Processes a single message analysis in the worker thread. - * Fetches context, attachments, runs LLM analysis (or simple fallback), - * and returns the result — does NOT update DB or broadcast. - * - * The caller (main thread) handles DB writes, broadcasting, and auto-delete - * scheduling. - */ -export async function processIndividualAnalysis({ - message, - skipNormalAnalysis, -}: IndividualWorkerRequest): Promise { - if (!config.AI_LLM_API_KEY) { - return { ok: false, results: [], error: "AI_LLM_API_KEY is missing" }; + return { ok: false, results: [], error: `Database init failed: ${msg}` }; } try { - await ensureDb(); - - const contextBefore = await getConversationContextBefore({ - channelId: message.channel_id, - threadId: message.thread_id, - beforeCreatedAt: message.created_at, - limit: config.AI_ANALYSIS_CONTEXT_MESSAGE_LIMIT, - }); - - const contextLines = buildConversationContext({ - contextBefore, - targets: [message], - maxTokens: config.AI_ANALYSIS_MAX_CONTEXT_TOKENS, - }); - - const contextIds = contextBefore.map((m) => m.id); - const attachments = await getAttachmentsForMessages([message.id, ...contextIds]); - - let results: AnalysisResult[]; - - if (skipNormalAnalysis) { - // Go straight to simple text fallback (no JSON, no complex prompt) - const simpleResult = await runSimpleTextFallback(message); - results = [simpleResult]; - } else { - // Try normal analysis first - const moderationResult = await runModerationAnalysis({ - targets: [message], - contextText: contextLines.join("\n"), - attachments, - }); - results = moderationResult.results; + if (job.type === "batch") { + return await processBatch(job); } - - return { ok: true, results }; + return await processIndividual(job); } catch (error) { const errorMessage = error instanceof Error ? error.message : String(error); + const errorStack = error instanceof Error ? error.stack : undefined; + console.error(JSON.stringify({ + level: "ERROR", + context: "aiAnalysisWorker", + type: job.type, + error: errorMessage, + stack: errorStack, + timestamp: new Date().toISOString(), + })); + if (job.type === "batch") { + return { ok: false, conversationKey: job.conversationKey, rows: [], error: errorMessage }; + } return { ok: false, results: [], error: errorMessage }; } } + +// --------------------------------------------------------------------------- +// Batch handler +// --------------------------------------------------------------------------- + +async function processBatch(job: { type: "batch"; conversationKey: string; messages: MessageRecord[] }): Promise { + const { conversationKey, messages } = job; + const firstMessage = messages[0]; + if (!firstMessage) return { ok: true, conversationKey, rows: [] }; + + const contextBefore = await getConversationContextBefore({ + channelId: firstMessage.channel_id, + threadId: firstMessage.thread_id, + beforeCreatedAt: firstMessage.created_at, + limit: config.AI_ANALYSIS_CONTEXT_MESSAGE_LIMIT, + }); + + const contextLines = buildConversationContext({ + contextBefore, + targets: messages, + maxTokens: config.AI_ANALYSIS_MAX_CONTEXT_TOKENS, + }); + + const targetIds = messages.map((m) => m.id); + const contextIds = contextBefore.map((m) => m.id); + const allMessageIds = [...targetIds, ...contextIds]; + const attachments = await getAttachmentsForMessages(allMessageIds); + + const result = await runModerationAnalysis({ + targets: messages, + contextText: contextLines.join("\n"), + attachments, + }); + + 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, + }, + })); + + try { + const rows = await updateMessagesAIAnalysisBulk(updates); + return { ok: true, conversationKey, rows }; + } catch (dbErr) { + throw new Error( + `Failed to update DB: ${dbErr instanceof Error ? dbErr.message : String(dbErr)}`, + ); + } +} + +// --------------------------------------------------------------------------- +// Individual fallback handler (offloaded from main thread) +// --------------------------------------------------------------------------- + +async function processIndividual(job: { type: "individual"; message: MessageRecord; skipNormalAnalysis: boolean }): Promise { + const { message, skipNormalAnalysis } = job; + + const contextBefore = await getConversationContextBefore({ + channelId: message.channel_id, + threadId: message.thread_id, + beforeCreatedAt: message.created_at, + limit: config.AI_ANALYSIS_CONTEXT_MESSAGE_LIMIT, + }); + + const contextLines = buildConversationContext({ + contextBefore, + targets: [message], + maxTokens: config.AI_ANALYSIS_MAX_CONTEXT_TOKENS, + }); + + const contextIds = contextBefore.map((m) => m.id); + const attachments = await getAttachmentsForMessages([message.id, ...contextIds]); + + let results: AnalysisResult[]; + + if (skipNormalAnalysis) { + const simpleResult = await runSimpleTextFallback(message); + results = [simpleResult]; + } else { + const moderationResult = await runModerationAnalysis({ + targets: [message], + contextText: contextLines.join("\n"), + attachments, + }); + results = moderationResult.results; + } + + return { ok: true, results }; +} diff --git a/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts b/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts index dc47a13..be6e397 100644 --- a/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts +++ b/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts @@ -412,7 +412,7 @@ async function processIndividualFallback( ); const simpleResult = await workerPool.run({ - type: "individual_simple", + type: "individual", message, skipNormalAnalysis: true, } as any) as @@ -654,6 +654,7 @@ async function processBatch( conversationProcessing.set(conversationKey, processingStartedAt); try { const result = (await workerPool.run({ + type: "batch", conversationKey, messages, })) as AnalysisWorkerResponse;