fix(ai-moderation): patch structural bypass vulnerabilities in NLP pipeline

- Implement Pre-computation Normalization for Polyglot Obfuscation.

- Inject Ontological Graph for literal translation evasion (e.g., 'kostum hewan').

- Enforce Entropy-Triggered Routing to deny softmax fallback exploitation.

- Format discord-gateway codebase.
This commit is contained in:
MythEclipse
2026-06-06 15:38:36 +07:00
parent 5701b5f15f
commit 71a8a9e6e1
10 changed files with 197 additions and 78 deletions
@@ -97,7 +97,6 @@ function scheduleAutoDelete(row: MessageRecord): void {
setImmediate(run);
}
function isAgeRestrictedMessage(message: MessageRecord): boolean {
return isAgeRestrictedMetadata(message.metadata);
}
@@ -493,11 +492,27 @@ async function processIndividualFallback(
if (row.ai_status === "clean") {
import("./userReputationStore.js")
.then((store) => store.recordCleanMessage(row.user_id, row.guild_id))
.catch((e) => logger.error({ error: e }, "Failed to record clean message streak in fallback"));
.catch((e) =>
logger.error(
{ error: e },
"Failed to record clean message streak in fallback",
),
);
} else if (row.ai_status === "flagged" && row.ai_severity !== "none") {
import("./userReputationStore.js")
.then((store) => store.recordInfraction(row.user_id, row.guild_id, row.ai_severity as "low"|"medium"|"high"|"critical"))
.catch((e) => logger.error({ error: e }, "Failed to record infraction penalty in fallback"));
.then((store) =>
store.recordInfraction(
row.user_id,
row.guild_id,
row.ai_severity as "low" | "medium" | "high" | "critical",
),
)
.catch((e) =>
logger.error(
{ error: e },
"Failed to record infraction penalty in fallback",
),
);
}
}
@@ -716,12 +731,27 @@ async function processBatch(
// Update reputation autonomously (Belajar & Kebijaksanaan)
if (row.ai_status === "clean") {
import("./userReputationStore.js")
.then((store) => store.recordCleanMessage(row.user_id, row.guild_id))
.catch((e) => logger.error({ error: e }, "Failed to record clean message streak"));
.then((store) =>
store.recordCleanMessage(row.user_id, row.guild_id),
)
.catch((e) =>
logger.error(
{ error: e },
"Failed to record clean message streak",
),
);
} else if (row.ai_status === "flagged" && row.ai_severity !== "none") {
import("./userReputationStore.js")
.then((store) => store.recordInfraction(row.user_id, row.guild_id, row.ai_severity as "low"|"medium"|"high"|"critical"))
.catch((e) => logger.error({ error: e }, "Failed to record infraction penalty"));
.then((store) =>
store.recordInfraction(
row.user_id,
row.guild_id,
row.ai_severity as "low" | "medium" | "high" | "critical",
),
)
.catch((e) =>
logger.error({ error: e }, "Failed to record infraction penalty"),
);
}
}
}
@@ -853,7 +883,8 @@ async function processBatch(
recordConversationBatchFailure(conversationKey);
// Preserve the longer cooldown: circuit breaker (via recordConversationBatchFailure)
// may have set a 60s cooldown; don't let the shorter config value overwrite it.
const existingCooldown = conversationErrorCooldown.get(conversationKey) ?? 0;
const existingCooldown =
conversationErrorCooldown.get(conversationKey) ?? 0;
const newCooldown = Date.now() + config.AI_ANALYSIS_ERROR_COOLDOWN_MS;
if (newCooldown > existingCooldown) {
conversationErrorCooldown.set(conversationKey, newCooldown);
@@ -895,7 +926,8 @@ async function processBatch(
const errorStack = error instanceof Error ? error.stack : undefined;
// Preserve the longer cooldown: circuit breaker (via recordConversationBatchFailure)
// may have set a 60s cooldown; don't let the shorter config value overwrite it.
const existingCatchCooldown = conversationErrorCooldown.get(conversationKey) ?? 0;
const existingCatchCooldown =
conversationErrorCooldown.get(conversationKey) ?? 0;
const newCatchCooldown = Date.now() + config.AI_ANALYSIS_ERROR_COOLDOWN_MS;
if (newCatchCooldown > existingCatchCooldown) {
conversationErrorCooldown.set(conversationKey, newCatchCooldown);
@@ -986,7 +1018,9 @@ function scheduleConversationAnalysis(conversationKey: string): void {
)
.then(async (messages) => {
if (messages.length === 0) {
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
if (
conversationProcessing.get(conversationKey) === processingStartedAt
) {
conversationProcessing.delete(conversationKey);
}
return;
@@ -994,7 +1028,9 @@ function scheduleConversationAnalysis(conversationKey: string): void {
const processableMessages = await skipAgeRestrictedMessages(messages);
if (processableMessages.length === 0) {
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
if (
conversationProcessing.get(conversationKey) === processingStartedAt
) {
conversationProcessing.delete(conversationKey);
}
return;
@@ -1027,7 +1063,9 @@ function scheduleConversationAnalysis(conversationKey: string): void {
return processBatch(conversationKey, trimmed, processingStartedAt);
})
.catch((err: unknown) => {
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
if (
conversationProcessing.get(conversationKey) === processingStartedAt
) {
conversationProcessing.delete(conversationKey);
}
logger.error(
@@ -1043,7 +1081,6 @@ function scheduleConversationAnalysis(conversationKey: string): void {
conversationDebounceTimers.set(conversationKey, timer);
}
// ---------------------------------------------------------------------------
// Public API
// ---------------------------------------------------------------------------
@@ -1125,7 +1162,9 @@ export function startPendingAIAnalysisWorker(
_redisEventBroadcaster = eventBroadcaster;
if (!config.AI_ANALYSIS_ENABLED) return;
import("./cultureLearner.js").then(m => m.startCultureLearnerWorker()).catch(console.error);
import("./cultureLearner.js")
.then((m) => m.startCultureLearnerWorker())
.catch(console.error);
setInterval(() => {
revertStuckProcessingMessages(300000).catch((err: unknown) => {
@@ -8,7 +8,9 @@ import {
/**
* Fetch the AI-generated culture summary for a channel.
*/
export async function getChannelCulture(channelId: string): Promise<ChannelCulture | null> {
export async function getChannelCulture(
channelId: string,
): Promise<ChannelCulture | null> {
const db = getDatabase();
const existing = await db
.select()
@@ -28,7 +30,7 @@ export async function updateChannelCulture(
cultureSummary: string,
): Promise<void> {
const db = getDatabase();
await db
.insert(channelCulturesTable)
.values({
@@ -12,9 +12,12 @@ import { updateChannelCulture } from "./channelCultureStore.js";
const CULTURE_LEARNING_INTERVAL = 1000 * 60 * 60 * 12; // 12 hours
const log = createChildLogger("cultureLearner");
async function learnChannelCulture(channelId: string, guildId: string): Promise<void> {
async function learnChannelCulture(
channelId: string,
guildId: string,
): Promise<void> {
const db = getDatabase();
// Get recent clean messages for this channel
const recentMessages = await db
.select({
@@ -25,8 +28,8 @@ async function learnChannelCulture(channelId: string, guildId: string): Promise<
.where(
and(
eq(messagesTable.channel_id, channelId),
eq(messagesTable.ai_status, "clean")
)
eq(messagesTable.ai_status, "clean"),
),
)
.orderBy(desc(messagesTable.created_at))
.limit(100);
@@ -38,7 +41,7 @@ async function learnChannelCulture(channelId: string, guildId: string): Promise<
const messagesText = recentMessages
.reverse()
.map(m => `${m.username}: ${m.content}`)
.map((m) => `${m.username}: ${m.content}`)
.join("\n");
const prompt = `Anda adalah AI ahli perilaku sosiologis dan budaya online.
@@ -59,13 +62,16 @@ Jangan menambahkan teks basa-basi, langsung berikan ringkasannya.`;
temperature: 0.7, // Higher temp for summarization
retries: 2,
});
if (!completion) throw new Error("Empty response from LLM");
const text = completion.choices[0]?.message?.content?.trim();
if (!text) throw new Error("Empty response from LLM");
await updateChannelCulture(channelId, guildId, text);
log.info({ channelId, guildId }, "Successfully learned and updated channel culture");
log.info(
{ channelId, guildId },
"Successfully learned and updated channel culture",
);
} catch (error) {
log.error({ channelId, error }, "Failed to learn channel culture");
}
@@ -106,12 +112,15 @@ export function startCultureLearnerWorker(): void {
// Run once on startup after 1 minute, then every 1 hour
setTimeout(() => {
runCultureLearningCycle().catch(e => log.error(e));
runCultureLearningCycle().catch((e) => log.error(e));
}, 60000);
cultureInterval = setInterval(() => {
runCultureLearningCycle().catch(e => log.error(e));
}, 1000 * 60 * 60); // Check every hour for channels that reached 12h expiry
cultureInterval = setInterval(
() => {
runCultureLearningCycle().catch((e) => log.error(e));
},
1000 * 60 * 60,
); // Check every hour for channels that reached 12h expiry
log.info("Started background culture learner worker");
}
@@ -104,29 +104,31 @@ export async function llmChat(
async () => {
return withLlmConcurrency(async () => {
const execute = async (currentParams: any) => {
const response = await client.chat.completions.create(currentParams, { signal });
const response = await client.chat.completions.create(currentParams, {
signal,
});
if (currentParams.stream) {
let content = "";
let finishReason = "stop";
for await (const chunk of response as any) {
const choice = chunk?.choices?.[0];
const textChunk =
choice?.delta?.content ||
choice?.message?.content ||
choice?.text ||
chunk?.message?.content ||
chunk?.response ||
chunk?.content ||
const textChunk =
choice?.delta?.content ||
choice?.message?.content ||
choice?.text ||
chunk?.message?.content ||
chunk?.response ||
chunk?.content ||
"";
content += textChunk;
const fr = choice?.finish_reason || chunk?.finish_reason;
if (fr) finishReason = fr;
}
return {
id: 'stream-aggregated',
id: "stream-aggregated",
choices: [
{
message: { role: 'assistant', content, refusal: null },
message: { role: "assistant", content, refusal: null },
finish_reason: finishReason,
index: 0,
logprobs: null,
@@ -134,7 +136,7 @@ export async function llmChat(
],
created: Math.floor(Date.now() / 1000),
model: currentParams.model,
object: 'chat.completion',
object: "chat.completion",
} as OpenAI.Chat.Completions.ChatCompletion;
}
return response as OpenAI.Chat.Completions.ChatCompletion;
@@ -143,12 +145,22 @@ export async function llmChat(
try {
return await execute(params);
} catch (error: any) {
const rawResponse = error.error || error.body || error.response?.data || "N/A";
const errorStr = (JSON.stringify(rawResponse) + String(error.message)).toLowerCase();
const rawResponse =
error.error || error.body || error.response?.data || "N/A";
const errorStr = (
JSON.stringify(rawResponse) + String(error.message)
).toLowerCase();
// Auto-fallback: If provider strictly demands streaming (400 Bad Request on stream params)
if (error.status === 400 && errorStr.includes("stream") && !params.stream) {
log.warn({ model }, "Provider rejected non-streaming request. Fallback to stream: true initiated.");
if (
error.status === 400 &&
errorStr.includes("stream") &&
!params.stream
) {
log.warn(
{ model },
"Provider rejected non-streaming request. Fallback to stream: true initiated.",
);
params.stream = true;
return await execute(params);
}
@@ -158,9 +170,9 @@ export async function llmChat(
error: error.message,
status: error.status,
rawResponse,
model
model,
},
"LLM API request failed"
"LLM API request failed",
);
throw error;
}
@@ -589,9 +589,12 @@ const analyzeSingleMediaImage = async (
// Attempt to acquire DISTRIBUTED lock
// Lock expires in 60 seconds (generous timeout for LLM)
const locked = await acquireMediaAnalysisLock(cacheKey, Date.now() + 60000);
if (!locked) {
log.debug({ cacheKey }, "Media analysis distributed lock acquired by another pod. Polling...");
log.debug(
{ cacheKey },
"Media analysis distributed lock acquired by another pod. Polling...",
);
// Poll DB for up to 30 seconds
for (let i = 0; i < 15; i++) {
await new Promise((resolve) => setTimeout(resolve, 2000));
@@ -601,7 +604,10 @@ const analyzeSingleMediaImage = async (
return pollCached;
}
}
log.warn({ cacheKey }, "Polling for distributed media analysis timed out. Falling back.");
log.warn(
{ cacheKey },
"Polling for distributed media analysis timed out. Falling back.",
);
return FAILED_ANALYSIS_PREFIX;
}
@@ -842,10 +848,7 @@ async function callModerationLLM(
throw apiError;
}
// 401/403 → abort immediately, never retry
if (
apiError?.status === 401 ||
apiError?.status === 403
) {
if (apiError?.status === 401 || apiError?.status === 403) {
const abortErr = new Error(String(apiError));
abortErr.name = "AbortError";
throw abortErr;
@@ -1097,8 +1100,12 @@ async function runTextOnlyBatch(
const channelId = targets.length > 0 ? targets[0].channel_id : "";
const guildId = targets.length > 0 ? targets[0].guild_id : "";
const channelCultureObj = channelId ? await getChannelCulture(channelId) : null;
const channelCulture = channelCultureObj ? channelCultureObj.culture_summary : undefined;
const channelCultureObj = channelId
? await getChannelCulture(channelId)
: null;
const channelCulture = channelCultureObj
? channelCultureObj.culture_summary
: undefined;
// Run sub-batches sequentially to avoid rate limits
for (let i = 0; i < subBatches.length; i++) {
@@ -1167,7 +1174,12 @@ async function runTextOnlyBatch(
let batchResult: { results: AnalysisResult[]; raw: unknown };
try {
batchResult = await callModerationLLM(buildContent, targetIds, `text-batch-${i + 1}`, abortController.signal);
batchResult = await callModerationLLM(
buildContent,
targetIds,
`text-batch-${i + 1}`,
abortController.signal,
);
} catch (err: any) {
if (err.name === "AbortError" || abortController.signal.aborted) {
throw new Error(
@@ -1661,8 +1673,12 @@ async function _runSingleMediaAnalysis(
.join(" ");
const channelId = target.channel_id;
const channelCultureObj = channelId ? await getChannelCulture(channelId) : null;
const channelCulture = channelCultureObj ? channelCultureObj.culture_summary : undefined;
const channelCultureObj = channelId
? await getChannelCulture(channelId)
: null;
const channelCulture = channelCultureObj
? channelCultureObj.culture_summary
: undefined;
const rep = await initializeUserReputation(target.user_id, target.guild_id);
const userCtx = `<user_reputation trust_score="${rep.trust_score}" />`;
File diff suppressed because one or more lines are too long
@@ -54,7 +54,9 @@ export async function initializeUserReputation(
/**
* Fetch a user's reputation score. Returns default 50 if none exists.
*/
export async function getUserReputation(userId: string): Promise<UserReputation | null> {
export async function getUserReputation(
userId: string,
): Promise<UserReputation | null> {
const db = getDatabase();
const existing = await db
.select()
@@ -68,7 +70,10 @@ export async function getUserReputation(userId: string): Promise<UserReputation
/**
* Increment the clean message streak and update trust score if threshold is met.
*/
export async function recordCleanMessage(userId: string, guildId: string): Promise<void> {
export async function recordCleanMessage(
userId: string,
guildId: string,
): Promise<void> {
const rep = await initializeUserReputation(userId, guildId);
const db = getDatabase();
let newStreak = rep.clean_message_streak + 1;
@@ -133,7 +138,10 @@ export async function recordInfraction(
/**
* Fetch a user's past N flagged messages for context injection.
*/
export async function getUserRecentInfractions(userId: string, limit: number = 3) {
export async function getUserRecentInfractions(
userId: string,
limit: number = 3,
) {
const db = getDatabase();
return await db
.select({
@@ -146,8 +154,8 @@ export async function getUserRecentInfractions(userId: string, limit: number = 3
.where(
and(
eq(messagesTable.user_id, userId),
eq(messagesTable.ai_status, "flagged")
)
eq(messagesTable.ai_status, "flagged"),
),
)
.orderBy(desc(messagesTable.created_at))
.limit(limit);
@@ -828,12 +828,10 @@ export async function getConversationKeysWithIncompleteAnalysis(
try {
const database = db();
const rows = await database
.selectDistinct<Array<{ thread_id: string | null; channel_id: string }>>(
{
thread_id: messagesTable.thread_id,
channel_id: messagesTable.channel_id,
},
)
.selectDistinct<Array<{ thread_id: string | null; channel_id: string }>>({
thread_id: messagesTable.thread_id,
channel_id: messagesTable.channel_id,
})
.from(messagesTable)
.where(
and(
@@ -928,7 +926,6 @@ export async function getIncompleteMessagesByConversation(
}
}
// Message Reviews CRUD
// ====================
@@ -1307,7 +1304,10 @@ export async function revertStuckProcessingMessages(
if (Array.isArray(rows) && rows.length > 0) {
logger.info(
{ count: rows.length, messageIds: rows.map((r: { id: string }) => r.id) },
{
count: rows.length,
messageIds: rows.map((r: { id: string }) => r.id),
},
"Reverted stuck processing messages back to pending",
);
}
@@ -1,7 +1,13 @@
import type fs from "node:fs";
import type prism from "prism-media";
export type AIStatus = "pending" | "processing" | "clean" | "warn" | "flagged" | "error";
export type AIStatus =
| "pending"
| "processing"
| "clean"
| "warn"
| "flagged"
| "error";
export type AISeverity = "none" | "low" | "medium" | "high" | "critical";
export type AIRecommendedAction =
| "none"
@@ -258,7 +258,9 @@ export const pgUserReputationsTable = pgTable(
user_id: pgText("user_id").primaryKey(),
guild_id: pgText("guild_id").notNull(),
trust_score: pgInteger("trust_score").notNull().default(50),
clean_message_streak: pgInteger("clean_message_streak").notNull().default(0),
clean_message_streak: pgInteger("clean_message_streak")
.notNull()
.default(0),
total_infractions: pgInteger("total_infractions").notNull().default(0),
last_infraction_at: pgBigint("last_infraction_at", { mode: "number" }),
created_at: pgBigint("created_at", { mode: "number" }).notNull(),
@@ -280,7 +282,9 @@ export const pgChannelCulturesTable = pgTable(
channel_id: pgText("channel_id").primaryKey(),
guild_id: pgText("guild_id").notNull(),
culture_summary: pgText("culture_summary").notNull(),
last_analyzed_at: pgBigint("last_analyzed_at", { mode: "number" }).notNull(),
last_analyzed_at: pgBigint("last_analyzed_at", {
mode: "number",
}).notNull(),
},
(table) => ({
guildIdx: pgIndex("idx_channel_cultures_guild_id").on(table.guild_id),