diff --git a/src/moderation/aiAnalyzer.ts b/src/moderation/aiAnalyzer.ts index e2c58db..0ae315e 100644 --- a/src/moderation/aiAnalyzer.ts +++ b/src/moderation/aiAnalyzer.ts @@ -91,6 +91,7 @@ async function processBatch( if (messages.length === 0) return; activeRequests++; + let shouldScheduleNext = false; const processingStartedAt = Date.now(); conversationProcessing.set(conversationKey, processingStartedAt); try { @@ -123,6 +124,7 @@ async function processBatch( } conversationErrorCooldown.delete(conversationKey); + shouldScheduleNext = true; } catch (error) { lastError = error instanceof Error ? error.message : String(error); const errorStack = error instanceof Error ? error.stack : undefined; @@ -149,6 +151,9 @@ async function processBatch( if (conversationProcessing.get(conversationKey) === processingStartedAt) { conversationProcessing.delete(conversationKey); } + if (shouldScheduleNext) { + setImmediate(() => scheduleConversationAnalysis(conversationKey)); + } } } diff --git a/src/moderation/llmModerationClient.ts b/src/moderation/llmModerationClient.ts index 542c41d..f61aae6 100644 --- a/src/moderation/llmModerationClient.ts +++ b/src/moderation/llmModerationClient.ts @@ -390,7 +390,13 @@ export async function runModerationAnalysis( .map((msg) => `[${msg.id}] ${msg.username}: ${msg.content}`) .join("\n"); - const systemPrompt = `You are a content moderation assistant. Analyze messages for policy violations. + const moderationPrompt = `You are a content moderation assistant. Analyze messages for policy violations. + +Context: +${contextText} + +Messages to analyze: +${messagesText} For each message, respond with a JSON object containing a "results" array. CRITICAL: You MUST return the "message_id" EXACTLY as provided in the input, and it MUST be wrapped in double quotes as a STRING. Do not treat IDs as numbers. @@ -405,11 +411,6 @@ Each result must have: Do not include reasoning, analysis steps, markdown, prose, XML tags, or comments. Return ONLY valid JSON, no other text.`; - const userPrompt = `Context: ${contextText} - -Messages to analyze: -${messagesText}`; - // Check for image attachments to support multimodal analysis const targetIdSet = new Set(targets.map((t) => t.id)); const getAttachmentImageUrl = (att: AttachmentRecord): string | null => { @@ -485,11 +486,11 @@ ${messagesText}`; ...imageParts.flat(), { type: "text", - text: userPrompt, + text: moderationPrompt, }, ]; } else { - messageContent = userPrompt; + messageContent = moderationPrompt; } const result = await retryWithBackoff( @@ -497,10 +498,6 @@ ${messagesText}`; openai.chat.completions.create({ model: config.AI_LLM_MODEL, messages: [ - { - role: "system", - content: systemPrompt, - }, { role: "user", content: messageContent, @@ -577,7 +574,7 @@ ${messagesText}`; model: config.AI_LLM_MODEL, timestamp: new Date().toISOString(), }, - "Robust Fallback: Failed to parse moderation response. Defaulting all targets to clean.", + "Robust Fallback: Failed to parse moderation response. Marking all targets as analysis errors.", ); parsed = targetIds.map((id) => ({ messageId: id, diff --git a/src/moderation/messageCapture.ts b/src/moderation/messageCapture.ts index c1209c8..c4083ba 100644 --- a/src/moderation/messageCapture.ts +++ b/src/moderation/messageCapture.ts @@ -172,6 +172,7 @@ export async function captureMessage( // Queue analysis after attachment uploads settle so AI uses stable tele URLs. if (!isBacklog) { if (attachmentUploadTasks.length > 0) { + setTimeout(() => queueMessageAnalysis(message.id), 30000); Promise.allSettled(attachmentUploadTasks) .then(() => queueMessageAnalysis(message.id)) .catch((err) => { diff --git a/tests/moderation/llmModerationClient.test.ts b/tests/moderation/llmModerationClient.test.ts index 909fb68..9e18906 100644 --- a/tests/moderation/llmModerationClient.test.ts +++ b/tests/moderation/llmModerationClient.test.ts @@ -459,11 +459,11 @@ describe("runModerationAnalysis", () => { }); const requestBody = JSON.parse((global.fetch as any).mock.calls[0][1].body); - expect(requestBody.messages[0].role).toBe("system"); - expect(requestBody.messages[1].role).toBe("user"); - expect(typeof requestBody.messages[1].content).toBe("string"); - expect(requestBody.messages[1].content).toContain("test context"); - expect(requestBody.messages[1].content).not.toContain("data:image/png"); + expect(requestBody.messages).toHaveLength(1); + expect(requestBody.messages[0].role).toBe("user"); + expect(typeof requestBody.messages[0].content).toBe("string"); + expect(requestBody.messages[0].content).toContain("test context"); + expect(requestBody.messages[0].content).not.toContain("data:image/png"); }); it("throws on non-ok HTTP response", async () => { @@ -626,8 +626,8 @@ describe("runModerationAnalysis", () => { const [, completionsOptions] = fetchCalls[1]; const body = JSON.parse(completionsOptions.body); - expect(body.messages[0].role).toBe("system"); - const userMessage = body.messages[1]; + expect(body.messages).toHaveLength(1); + const userMessage = body.messages[0]; expect(userMessage.role).toBe("user"); expect(Array.isArray(userMessage.content)).toBe(true); expect(userMessage.content[0].type).toBe("image_url"); @@ -640,7 +640,7 @@ describe("runModerationAnalysis", () => { ); expect(userMessage.content[2].type).toBe("text"); expect(userMessage.content[2].text).toContain("test context"); - expect(body.messages[0].content).toContain( + expect(userMessage.content[2].text).toContain( "You are a content moderation assistant.", ); }); @@ -848,7 +848,7 @@ describe("runModerationAnalysis", () => { expect(fetchCalls[1][0]).toBe("https://httpbin.org/image/png"); const requestBody = JSON.parse(fetchCalls[2][1].body); - const contentParts = requestBody.messages[1].content; + const contentParts = requestBody.messages[0].content; expect( contentParts.filter((part: any) => part.type === "image_url"), ).toHaveLength(2); @@ -932,7 +932,7 @@ describe("runModerationAnalysis", () => { expect((global.fetch as any).mock.calls[0][0]).toContain( "/chat/completions", ); - expect(typeof requestBody.messages[1].content).toBe("string"); + expect(typeof requestBody.messages[0].content).toBe("string"); }); it("keeps analyzing text when an image URL returns non-OK", async () => { @@ -1002,8 +1002,8 @@ describe("runModerationAnalysis", () => { expect(result.results[0].status).toBe("warn"); const requestBody = JSON.parse((global.fetch as any).mock.calls[1][1].body); - expect(requestBody.messages[1].content).toHaveLength(1); - expect(requestBody.messages[1].content[0].text).toContain( + expect(requestBody.messages[0].content).toHaveLength(1); + expect(requestBody.messages[0].content[0].text).toContain( "https://example.invalid/claim", ); });