From 10f44138bc194fe3bf29ce7b836001f37c847129 Mon Sep 17 00:00:00 2001 From: MythEclipse Date: Sat, 30 May 2026 21:12:25 +0700 Subject: [PATCH] feat: implement age restriction handling in message analysis and metadata --- src/moderation/aiAnalyzer.ts | 85 +++++++++++++++++++++++++++++-- src/moderation/analyticsStore.ts | 25 ++++++--- src/moderation/messageMetadata.ts | 37 ++++++++++++++ 3 files changed, 136 insertions(+), 11 deletions(-) diff --git a/src/moderation/aiAnalyzer.ts b/src/moderation/aiAnalyzer.ts index 6f43cf7..aaf96e0 100644 --- a/src/moderation/aiAnalyzer.ts +++ b/src/moderation/aiAnalyzer.ts @@ -9,6 +9,7 @@ import { retryWithBackoff } from "../retry.js"; import { attemptAutoDeleteFlaggedMessage } from "./autoDeleteManager.js"; import { buildConversationContext } from "./conversationContext.js"; import { runModerationAnalysis } from "./llmModerationClient.js"; +import { isAgeRestrictedMetadata } from "./messageMetadata.js"; import { getAttachmentsForMessages, getConversationContextBefore, @@ -17,6 +18,7 @@ import { getMessageById, getPendingConversationKeys, getPendingMessagesByConversation, + updateMessageAIAnalysis, updateMessagesAIAnalysisBulk, } from "./messageStore.js"; import type { @@ -56,6 +58,59 @@ function scheduleAutoDelete(row: MessageRecord): void { setImmediate(run); } +function isAgeRestrictedMessage(message: MessageRecord): boolean { + return isAgeRestrictedMetadata(message.metadata); +} + +function buildAgeRestrictedSkipResult(): { + status: "clean"; + flags: string | null; + score: number; + analysis: string; + categories: string[]; + severity: "none"; + confidence: number; + recommendedAction: "none"; + analyzedAt: number; + error: null; +} { + return { + status: "clean", + flags: JSON.stringify(["age_restricted"]), + score: 0, + analysis: "Skipped moderation for age-restricted content.", + categories: ["age_restricted"], + severity: "none", + confidence: 1, + recommendedAction: "none", + analyzedAt: Date.now(), + error: null, + }; +} + +async function skipAgeRestrictedMessages( + messages: MessageRecord[], +): Promise { + const ageRestrictedMessages = messages.filter(isAgeRestrictedMessage); + if (ageRestrictedMessages.length === 0) { + return messages; + } + + const skippedRows = await updateMessagesAIAnalysisBulk( + ageRestrictedMessages.map((message) => ({ + messageId: message.id, + result: buildAgeRestrictedSkipResult(), + })), + ); + + for (const row of skippedRows) { + getModerationBroadcaster()?.messageAnalyzed(row); + } + + const skippedIds = new Set(ageRestrictedMessages.map((message) => message.id)); + return messages.filter((message) => !skippedIds.has(message.id)); +} + // --------------------------------------------------------------------------- // Batch pipeline state // --------------------------------------------------------------------------- @@ -672,13 +727,16 @@ function scheduleConversationAnalysis(conversationKey: string): void { conversationKey, config.AI_ANALYSIS_MAX_BATCH_SIZE, ) - .then((messages) => { + .then(async (messages) => { if (messages.length === 0) return; + const processableMessages = await skipAgeRestrictedMessages(messages); + if (processableMessages.length === 0) return; + // FIX #6: trim to token budget before sending to LLM. // 50 tokens overhead accounts for JSON structure + id/username fields. let trimmed = pickBatchWithinBudget( - messages, + processableMessages, config.AI_ANALYSIS_MAX_TARGET_TOKENS, 50, ); @@ -687,12 +745,12 @@ function scheduleConversationAnalysis(conversationKey: string): void { // pickBatchWithinBudget returns [] — which would leave them permanently // stuck as `pending`. Fall back to the first message alone so at // least one makes progress; the rest will be processed in later ticks. - if (trimmed.length === 0 && messages.length > 0) { - trimmed = messages.slice(0, 1); + if (trimmed.length === 0 && processableMessages.length > 0) { + trimmed = processableMessages.slice(0, 1); logger.warn( { conversationKey, - messageId: messages[0]?.id, + messageId: processableMessages[0]?.id, tokenBudget: config.AI_ANALYSIS_MAX_TARGET_TOKENS, }, "All messages exceed token budget — processing first message alone to avoid stuck-pending deadlock", @@ -731,6 +789,19 @@ export async function queueMessageAnalysis(messageId: string): Promise { logger.warn({ messageId }, "Message not found for analysis queue"); return; } + + if (isAgeRestrictedMessage(message)) { + const updated = await updateMessageAIAnalysis( + message.id, + buildAgeRestrictedSkipResult(), + ); + if (updated) { + getModerationBroadcaster()?.messageAnalyzed(updated); + } + logger.info({ messageId }, "Skipped AI analysis for age-restricted message"); + return; + } + queueConversationAnalysis(getConversationKey(message)); } catch (error) { logger.error( @@ -830,6 +901,10 @@ export function startPendingAIAnalysisWorker(client?: Client): void { key, config.AI_ANALYSIS_INDIVIDUAL_MAX_CONCURRENT, ) + .then(async (msgs) => { + const processableMessages = await skipAgeRestrictedMessages(msgs); + return processableMessages; + }) .then((msgs) => { if (msgs.length > 0) { enqueueIndividualFallbacks(msgs); diff --git a/src/moderation/analyticsStore.ts b/src/moderation/analyticsStore.ts index e30cb26..e171c9e 100644 --- a/src/moderation/analyticsStore.ts +++ b/src/moderation/analyticsStore.ts @@ -1,5 +1,6 @@ import { executeAll, executeGet } from "../database/drizzle.js"; import { createChildLogger } from "../logger.js"; +import { config } from "../config.js"; import type { MessageRecord } from "./types.js"; const logger = createChildLogger("analytics-store"); @@ -106,10 +107,16 @@ export async function getHourlyStats(input: { try { const since = Date.now() - hours * 3600_000; - const sqliteRows = await executeAll( + const isPg = config.DATABASE_TYPE === "postgres"; + + const hourExpr = isPg + ? `to_char(to_timestamp((created_at / 3600000) * 3600), 'YYYY-MM-DD HH24:MI:SS') as hour` + : `datetime((created_at / 3600000) * 3600, 'unixepoch') as hour`; + + const rows = await executeAll( ` SELECT - datetime((created_at / 3600000) * 3600, 'unixepoch') as hour, + ${hourExpr}, count(*) as count, count(case when ai_status = 'clean' then 1 end) as clean, count(case when ai_status = 'warn' then 1 end) as warned, @@ -141,7 +148,7 @@ export async function getHourlyStats(input: { buckets.set(key, { count: 0, clean: 0, warned: 0, flagged: 0, error: 0 }); } - for (const row of sqliteRows) { + for (const row of rows) { const d = new Date(row.hour.replace(" ", "T") + "Z"); const key = d.toISOString().slice(0, 13) + ":00:00Z"; const bucket = buckets.get(key); @@ -344,7 +351,7 @@ export async function getUserLeaderboard(input: { AND created_at >= ? AND deleted_at IS NULL ${channelId ? `AND (channel_id = ? OR thread_id = ?)` : ""} - GROUP BY user_id + GROUP BY user_id, username, avatar_url ORDER BY message_count DESC LIMIT ? `, @@ -379,6 +386,12 @@ export async function getModerationStats(input: { try { const since = Date.now() - hours * 3600_000; + const isPg = config.DATABASE_TYPE === "postgres"; + + const avgScoreExpr = isPg + ? `round(avg(ai_moderation_score)::numeric, 2)` + : `round(avg(ai_moderation_score), 2)`; + const row = await executeGet( ` SELECT @@ -388,7 +401,7 @@ export async function getModerationStats(input: { count(case when ai_status = 'flagged' then 1 end) as flagged, count(case when ai_status = 'error' then 1 end) as error, count(case when ai_status = 'pending' or ai_status IS NULL then 1 end) as pending, - round(avg(ai_moderation_score), 2) as average_score + ${avgScoreExpr} as average_score FROM messages WHERE guild_id = ? AND created_at >= ? @@ -501,7 +514,7 @@ export async function getTopViolators(input: { AND created_at >= ? AND deleted_at IS NULL ${channelId ? `AND (channel_id = ? OR thread_id = ?)` : ""} - GROUP BY user_id + GROUP BY user_id, username, avatar_url HAVING flagged_count > 0 OR warned_count > 0 ORDER BY (flagged_count * 3 + warned_count) DESC LIMIT ? diff --git a/src/moderation/messageMetadata.ts b/src/moderation/messageMetadata.ts index b30e41c..f83d55a 100644 --- a/src/moderation/messageMetadata.ts +++ b/src/moderation/messageMetadata.ts @@ -9,6 +9,9 @@ export interface MessageLocation { threadId: string | null; threadName: string | null; channelName: string | null; + nsfw?: boolean; + nsfwLevel?: string | null; + ageRestricted?: boolean; } export interface StickerEvidence { @@ -74,12 +77,25 @@ export interface RichMessageMetadata { export function getMessageLocation(message: Message): MessageLocation { const channel = message.channel as TextChannel | ThreadChannel; + const safetyChannel = channel as TextChannel & { + nsfw?: boolean; + nsfwLevel?: string | null; + }; if (!channel.isThread?.()) { return { channelId: message.channelId, threadId: null, threadName: null, channelName: "name" in channel ? channel.name : null, + nsfw: typeof safetyChannel.nsfw === "boolean" ? safetyChannel.nsfw : undefined, + nsfwLevel: + typeof safetyChannel.nsfwLevel === "string" + ? safetyChannel.nsfwLevel + : null, + ageRestricted: + typeof safetyChannel.nsfw === "boolean" + ? safetyChannel.nsfw + : undefined, }; } @@ -88,6 +104,13 @@ export function getMessageLocation(message: Message): MessageLocation { threadId: channel.id, threadName: channel.name, channelName: channel.parent?.name ?? null, + nsfw: typeof safetyChannel.nsfw === "boolean" ? safetyChannel.nsfw : undefined, + nsfwLevel: + typeof safetyChannel.nsfwLevel === "string" + ? safetyChannel.nsfwLevel + : null, + ageRestricted: + typeof safetyChannel.nsfw === "boolean" ? safetyChannel.nsfw : undefined, }; } @@ -200,6 +223,20 @@ export function parseRichMessageMetadata( } } +export function isAgeRestrictedMetadata( + metadata: string | null | undefined, +): boolean { + const parsed = parseRichMessageMetadata(metadata); + if (!parsed) return false; + + const nsfwLevel = parsed.channel.nsfwLevel?.toUpperCase(); + return Boolean( + parsed.channel.nsfw || + parsed.channel.ageRestricted || + nsfwLevel === "AGE_RESTRICTED", + ); +} + export function extractMessageMediaEvidence( metadata: string | null | undefined, ): MessageMediaEvidence {