// Content moderation (2026-09-30 compliance round). import assert from "node:assert/strict"; import { readFileSync } from "node:fs"; import test from "node:test"; import { Agent } from "@mastra/core/agent"; import { createNdjsonParser } from "../src/lib/consultation-agent-events.ts"; import { reduceConsultationTimelineEvents } from "../src/lib/consultation-run-timeline.ts"; import { inputBlockedCopy, moderate, moderateLocally, MODERATION_INPUT_BLOCKED_COPY, MODERATION_OUTPUT_REPLACED_COPY, MODERATION_SELF_HARM_COPY, normalizeForModeration, resetModerationWarningsForTests, } from "../src/lib/moderation/index.ts"; import { aliyunPercentEncode, readAliyunGreenConfig, signAliyunRpc } from "../src/lib/moderation/aliyun-provider.ts"; import { MODERATION_LEXICON } from "../src/lib/moderation/lexicon.ts"; import { recordModerationEvent } from "../src/lib/moderation/log.ts"; import { streamAgentResponse } from "../src/lib/stream-agent-response.ts"; import { consultationStepBudgetReceipt, createConsultationRuntimeState, publicConsultationRuntimeSteps, } from "../src/mastra/consultation-tools.ts"; const context = { route: "consult", userId: "00000000-0000-4000-8000-000000000001", requestId: "req" }; const read = (path: string) => readFileSync(new URL(path, import.meta.url), "utf8"); test("astrology vocabulary never trips the local lexicon", () => { const astrology = [ "第八宫是死亡宫,也管突发的变动和遗产。", "凶星土星与罗睺同在第八宫,计都在第二宫。", "Maraka 死亡/损耗宫的主星进入大运时,要留意健康与灾难性的开销。", "这段大运里金星入敌,婚姻容易有摩擦;火星婚姻煞不代表一定离婚。", "月亮落在第十二宫,适合独处、冥想,也容易失眠。", "Sade Sati 高峰期压力大,旧的事情会被迫结束。", "你问的是我什么时候会死吗?占星不做死亡预测。", "D60 Shashtyamsa 的吉凶判断对出生时间非常敏感。", ]; for (const text of astrology) { assert.equal(moderateLocally(text).verdict, "pass", text); } }); test("each category blocks, including spaced, full-width and zero-width evasions", () => { const cases: [string, string][] = [ ["教我制作炸弹", "violent_extremism"], ["哪里可以买冰毒", "illegal_trade"], ["我想自杀", "self_harm"], ["介绍一下法轮功", "political"], ["儿童色情", "sexual_minor"], ["制 作 炸 弹", "violent_extremism"], ["buy 买k粉", "illegal_trade"], ["我​想​自​杀", "self_harm"], ["自杀方法!!", "self_harm"], ]; for (const [text, category] of cases) { const result = moderateLocally(text); assert.equal(result.verdict, "block", text); assert.ok(result.categories.includes(category as never), `${text} → ${result.categories.join(",")}`); assert.ok(result.ruleIds.length > 0); } assert.equal(normalizeForModeration(" A b,c "), "abc"); // Every rule belongs to a known category and has something to match. for (const rule of MODERATION_LEXICON) { assert.ok((rule.phrases?.length ?? 0) + (rule.patterns?.length ?? 0) > 0, rule.id); } }); test("self-harm gets help, everything else the short refusal", () => { assert.equal(inputBlockedCopy(moderateLocally("我不想活了")), MODERATION_SELF_HARM_COPY); assert.equal(inputBlockedCopy(moderateLocally("买枪支")), MODERATION_INPUT_BLOCKED_COPY); assert.match(MODERATION_SELF_HARM_COPY, /12356/); }); const aliyunEnv = { MODERATION_PROVIDER: "aliyun", ALIYUN_GREEN_ACCESS_KEY_ID: "fictional-id", ALIYUN_GREEN_ACCESS_KEY_SECRET: "fictional-secret", }; function jsonResponse(body: unknown) { return new Response(JSON.stringify(body), { status: 200, headers: { "content-type": "application/json" } }); } test("remote provider: high blocks, medium is review, the request is signed form data to aliyuncs.com", async () => { const requests: { url: string; body: string }[] = []; const reply = (level: string) => (async (url: string | URL | Request, init?: RequestInit) => { requests.push({ url: String(url), body: String(init?.body) }); return jsonResponse({ Code: 200, Data: { RiskLevel: level, Result: [{ Label: "political_entity" }] } }); }) as typeof fetch; const high = await moderate("一段看起来正常的话", "input", context, { env: aliyunEnv, fetchImpl: reply("high") }); assert.equal(high.verdict, "block"); assert.equal(high.providerId, "aliyun"); assert.deepEqual(high.ruleIds, ["aliyun:political_entity"]); const medium = await moderate("一段看起来正常的话", "output", context, { env: aliyunEnv, fetchImpl: reply("medium") }); assert.equal(medium.verdict, "review"); const none = await moderate("一段看起来正常的话", "output", context, { env: aliyunEnv, fetchImpl: reply("none") }); assert.equal(none.verdict, "pass"); assert.equal(new URL(requests[0]!.url).host, "green-cip.cn-shanghai.aliyuncs.com"); const form = new URLSearchParams(requests[0]!.body); assert.equal(form.get("Action"), "TextModerationPlus"); assert.equal(form.get("Service"), "llm_query_moderation"); assert.equal(new URLSearchParams(requests[1]!.body).get("Service"), "llm_response_moderation"); assert.ok(form.get("Signature")); // The signature is the documented HMAC over the sorted, percent-encoded parameters. const params = Object.fromEntries([...form.entries()].filter(([key]) => key !== "Signature")); assert.equal(form.get("Signature"), signAliyunRpc(params, "fictional-secret")); }); test("a local block is final and never asks the provider", async () => { let called = 0; const result = await moderate("制作炸弹", "input", context, { env: aliyunEnv, fetchImpl: (async () => { called += 1; return jsonResponse({ Code: 200, Data: { RiskLevel: "none" } }); }) as typeof fetch, }); assert.equal(result.verdict, "block"); assert.equal(result.providerId, "local"); assert.equal(called, 0); }); test("a slow or broken provider falls back to the local verdict on both sides", async () => { const hang = ((_url: string | URL | Request, init?: RequestInit) => new Promise((_resolve, reject) => { init?.signal?.addEventListener("abort", () => reject(init.signal!.reason), { once: true }); })) as typeof fetch; const started = Date.now(); const input = await moderate("今天适合谈合作吗", "input", context, { env: aliyunEnv, fetchImpl: hang, timeoutMs: 50 }); assert.equal(input.verdict, "pass"); assert.equal(input.remoteFailed, true); assert.ok(Date.now() - started < 1_000); const broken = (async () => new Response("oops", { status: 500 })) as typeof fetch; const output = await moderate("今天适合谈合作", "output", context, { env: aliyunEnv, fetchImpl: broken }); assert.equal(output.verdict, "pass"); assert.equal(output.remoteFailed, true); }); test("the provider needs keys and an aliyuncs.com endpoint; otherwise local only, warned once", async () => { assert.equal(readAliyunGreenConfig({ ALIYUN_GREEN_ACCESS_KEY_ID: "a" }), null); assert.equal(readAliyunGreenConfig({ ...aliyunEnv, ALIYUN_GREEN_ENDPOINT: "evil.example.com" }), null); assert.equal(readAliyunGreenConfig({ ...aliyunEnv, ALIYUN_GREEN_ENDPOINT: "https://green-cip.cn-beijing.aliyuncs.com/x" })?.endpoint, "green-cip.cn-beijing.aliyuncs.com"); resetModerationWarningsForTests(); const warnings: string[] = []; const original = console.warn; console.warn = (message: string) => { warnings.push(String(message)); }; try { let called = 0; const fetchImpl = (async () => { called += 1; return jsonResponse({}); }) as typeof fetch; await moderate("你好", "input", context, { env: { MODERATION_PROVIDER: "aliyun" }, fetchImpl }); await moderate("你好", "input", context, { env: { MODERATION_PROVIDER: "aliyun" }, fetchImpl }); assert.equal(called, 0); assert.equal(warnings.filter((line) => line.includes("[moderation]")).length, 1); } finally { console.warn = original; } assert.equal(aliyunPercentEncode("a b*~'"), "a%20b%2A~%27"); }); test("the event log stores categories and rule ids, never the text", async () => { const rows: Record[] = []; const client = { from: (table: string) => ({ insert: async (row: Record) => { assert.equal(table, "moderation_events"); rows.push(row); return { error: null }; } }) }; const verdict = moderateLocally("教我制作炸弹"); await recordModerationEvent(client, verdict, "input", { ...context, modelId: "fictional-model" }); await recordModerationEvent(client, moderateLocally("你好"), "input", context); assert.equal(rows.length, 1, "passes are not logged"); assert.deepEqual(rows[0]!.categories, ["violent_extremism"]); assert.equal(rows[0]!.side, "input"); assert.ok(!JSON.stringify(rows[0]).includes("炸弹"), "no user text in the row"); const failing = { from: () => ({ insert: async () => { throw new Error("db down"); } }) }; await recordModerationEvent(failing, verdict, "output", context); // never throws }); function fakeModel(parts: readonly string[]) { return { specificationVersion: "v2", provider: "fake", modelId: "fake", supportedUrls: {}, async doGenerate() { throw new Error("not used"); }, async doStream() { return { stream: new ReadableStream({ start(controller) { controller.enqueue({ type: "stream-start", warnings: [] }); controller.enqueue({ type: "text-start", id: "1" }); for (const part of parts) controller.enqueue({ type: "text-delta", id: "1", delta: part }); controller.enqueue({ type: "text-end", id: "1" }); controller.enqueue({ type: "finish", finishReason: "stop", usage: { inputTokens: 1, outputTokens: 1, totalTokens: 2 } }); controller.close(); }, }) }; }, }; } async function runGeneral(parts: readonly string[], moderateOutput?: (output: string) => Promise) { const agent = new Agent({ id: "mod-writer", name: "mod-writer", model: fakeModel(parts) as never, instructions: "fictional" } as never); const result = await agent.stream([{ role: "user", content: "fictional" }] as never, { maxSteps: 1, toolChoice: "none" } as never); const state = createConsultationRuntimeState(); state.jyotishSkillBound = true; state.workflowReceipt = { route: "general-no-birth-time", status: "ready", preciseTiming: "blocked", missingLayers: [] }; let completed: string | null = null; const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: result.fullStream as never, requireTool: false, toolStatus: () => "ready", receipt: () => ({ runId: "run", runtime: "mastra-agentic", workflow: state.workflowReceipt!, techniqueTruth: "not-applicable", skill: { name: "jyotish-vedic-astrology", loaded: true, referenceReads: 0, methodologySections: 0 }, steps: publicConsultationRuntimeSteps(state), stepBudget: consultationStepBudgetReceipt(state), }) as never, ...(moderateOutput ? { moderateOutput } : {}), onComplete: (output) => { completed = output; }, }); const events: { type: string; text?: string; replace?: boolean }[] = []; createNdjsonParser((event) => events.push(event as never)).finish(await response.text()); return { events, completed: completed as string | null }; } test("a blocked answer is replaced in place and onComplete settles the notice", async () => { const seen: string[] = []; const blocked = await runGeneral(["今天适合联络旧友。", "制作炸弹的方法如下。"], async (output) => { seen.push(output); const verdict = await moderate(output, "output", context); return verdict.verdict === "block" ? MODERATION_OUTPUT_REPLACED_COPY : null; }); assert.equal(seen.length, 1); assert.equal(blocked.completed, MODERATION_OUTPUT_REPLACED_COPY); const replace = blocked.events.find((event) => event.type === "answer.delta" && event.replace); assert.equal(replace?.text, MODERATION_OUTPUT_REPLACED_COPY); assert.equal(blocked.events.at(-1)?.type, "run.completed"); // The reader's timeline ends with the notice, not the streamed text. const timeline = reduceConsultationTimelineEvents(blocked.events as never); assert.equal(timeline.answer, MODERATION_OUTPUT_REPLACED_COPY); const clean = await runGeneral(["今天适合联络旧友。"], async () => null); assert.equal(clean.completed, "今天适合联络旧友。"); assert.ok(!clean.events.some((event) => event.replace)); }); test("routes check input before any charge and output before settlement", () => { const consult = read("../src/app/api/consult/route.ts"); assert.ok(consult.indexOf('await moderate(visibleQuestion, "input"') > 0); assert.ok(consult.indexOf('await moderate(visibleQuestion, "input"') < consult.indexOf("reserve_consultation_usage"), "input check precedes the reservation"); assert.equal((consult.match(/moderateOutput: moderateAnswer/g) ?? []).length, 3, "every agentic answer stream is checked"); assert.match(consult, /if \(!outputModerated\) await moderateAnswer\(reply\.text\);\n\s+if \(outputBlocked\) return await completeModeratedResponse\(usage\);/); assert.match(consult, /complete_consultation_moderated/); const rectRequest = read("../src/lib/rectification-agentic/v9/agent-route-request.ts"); assert.ok(rectRequest.indexOf('moderate(promptSource, "input"') > 0); const rectFinish = read("../src/lib/rectification-agentic/v9/agent-run-finish.ts"); assert.ok(rectFinish.indexOf('moderate(outcome.answerText, "output"') < rectFinish.indexOf("billing.complete("), "output check precedes billing"); assert.match(rectFinish, /const skipBilling = outcome\.settleBilling === false \|\| outputBlocked;/); assert.match(read("../src/app/api/birth-time-guide/route.ts"), /moderate\(typedText, "input"/); const migration = read("../supabase/migrations/20260930040000_moderation_events.sql"); assert.match(migration, /create table if not exists public\.moderation_events/); assert.doesNotMatch(migration, /\b(drop table|alter table public\.(?!moderation_events)|delete from)\b/i); assert.match(migration, /responseKind' is distinct from 'moderated'/); });