Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01N4f2nya58RoRu4yEmJgRGE
258 lines
14 KiB
TypeScript
258 lines
14 KiB
TypeScript
// 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<Response>((_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<string, unknown>[] = [];
|
||
const client = { from: (table: string) => ({ insert: async (row: Record<string, unknown>) => { 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<string | null>) {
|
||
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'/);
|
||
});
|