feat(moderation): improve message analysis scheduling and update moderation prompt structure
This commit is contained in:
@@ -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));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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) => {
|
||||
|
||||
@@ -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",
|
||||
);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user