diff --git a/BLOCKED.md b/BLOCKED.md
index e81b10f8..ee39a666 100644
--- a/BLOCKED.md
+++ b/BLOCKED.md
@@ -32,6 +32,13 @@
- 运营主体与联系邮箱仍是 `lib/legal-entity.ts` 占位值。
- 未推送、未部署。
+## 合规轮 · 敏感内容审核(2026-09-30)
+
+- **阿里云内容安全账号与密钥**:需产品开通「内容安全 · 大模型输入/输出审核(TextModerationPlus,llm_query_moderation / llm_response_moderation)」并提供 `ALIYUN_GREEN_ACCESS_KEY_ID` / `ALIYUN_GREEN_ACCESS_KEY_SECRET`(可选 `ALIYUN_GREEN_ENDPOINT`),再设 `MODERATION_PROVIDER=aliyun`。适配器按公开 API 文档写成、签名有单测,但**未对真实服务调用过**;接入后需在 staging 用受控账号发一条命中样例验证。缺钥匙时只用本地词表(启动日志警告一次)。
+- **本地词表需法务审定**:`frontend/src/lib/moderation/lexicon.ts`(五类、保守短语),以及自伤求助句里的热线号码 12356。
+- **数据库测试**:新表 `moderation_events` 与函数 `complete_consultation_moderated`(`20260930040000_moderation_events.sql`)未跑 `npm run test:db`(本机无 Docker);函数由 `complete_consultation_free` 逐行复制,仅改 responseKind、长度上限与释放原因。部署到 staging 时的迁移步骤是第一次真实执行。
+- **真机**:命中后的前端表现(输入被拦的提示条、回答被整段替换)未在浏览器实测。
+
## TASK-chart-surface-polish:受控登录、实体手机与基线全量失败(2026-09-29)
- 本轮 Chrome、Edge、Docker 均可用,不套用历史“无 Chrome / 无 Docker”。真实 Chrome + golden 的本地组件验收及骨架高度修复后复验已完成(70+8 项通过);没有受控线上登录账号、实体手机和读屏实测;不能代替完整账户/人物/请求链路。清单见 `docs/testing/chart-surface-polish-20260929.md`。
diff --git a/CHANGELOG.md b/CHANGELOG.md
index e711ffe8..2c88e333 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -24,6 +24,14 @@
- 后台新增只读的「注销申请」列表(不显示邮箱)。
- 新增迁移 `20260930020000_account_deletion_requests.sql`(只加表、函数、触发器)。Skill 版本不 bump。
+## 2026-09-30 — 敏感内容审核(待验收)
+
+- 普通咨询、生时校正、出生时间引导里用户打的字,发给模型之前先过一遍审核。命中的话不发给模型、不扣点,回一句「这个话题我们不能讨论,换个问题试试。」;涉及自伤的,改为给出心理援助热线 12356 的求助提示。
+- 模型写完的回答也会审核;命中时整段换成「这段回答里有不适合展示的内容,已经撤下,本次不扣点。换个问法再试试。」,这一轮不扣点。
+- 默认用本地敏感词表(政治、暴恐、涉未成年色情、自伤、违法交易五类,占星用语不会误拦);配好阿里云内容安全的账号后可切换为「本地词表 + 阿里云」双层。审核服务超时或出错时只用本地词表,不会因此拦截正常内容。
+- 后台新增「内容审核」只读页:只记类别和命中规则编号,不保存用户或模型原文。
+- 数据库新增一张 `moderation_events` 表和一个免扣点结算函数(只加不改)。Skill 版本不 bump。
+
## 2026-09-30 — 首页去掉校正提示、星盘类型文字完整显示、加载改为轨道环动画、「那一刻的天空」换小星座图标(待验收)
- 首页不再显示「上次那次校正还没完成,可以在历史对话里接着做。」(BUG-1111)。
diff --git a/docs/tasks/PROGRESS-content-moderation-20260930.md b/docs/tasks/PROGRESS-content-moderation-20260930.md
new file mode 100644
index 00000000..ffe05623
--- /dev/null
+++ b/docs/tasks/PROGRESS-content-moderation-20260930.md
@@ -0,0 +1,40 @@
+# PROGRESS · 敏感内容审核 · 2026-09-30
+
+> 合规轮四项之一(产品 2026-09-30「合规与法律的前四个你可以帮我做吗」)。执行:Claude fork D,直接执行。分支 `codex/content-moderation-20260930`,基于 `codex/compliance-base-20260930`(= origin/staging `5208f19b` + `lib/legal-entity.ts`)。
+> 产品决定:输入命中 → 不发模型、不扣点、固定一句;输出命中 → 整段替换并退点;命中留脱敏日志。
+
+## 做了什么
+
+| 部分 | 位置 | 说明 |
+|---|---|---|
+| 审核接口 | `src/lib/moderation/index.ts` | `moderate(text, side, context)` → `{verdict, categories, ruleIds, providerId}`;本地词表先跑,命中即终判;配置了阿里云再跑远端(≤1.5 s);远端失败/超时 → 沿用本地判定(输入:本地兜底不放空;输出:不因服务不可用而扣下回答) |
+| 本地词表 | `src/lib/moderation/lexicon.ts` | 12 条规则、5 类(政治 3 / 暴恐 3 / 涉未成年色情 1 / 自伤 2 / 违法交易 3),均为多字短语;归一化去空格标点零宽、全角转半角 |
+| 阿里云适配 | `src/lib/moderation/aliyun-provider.ts` | TextModerationPlus(llm_query_moderation / llm_response_moderation),RPC 签名 v1;high → 拦,medium → 待复核(放行并记录);endpoint 限定 `*.aliyuncs.com`;超 2000 字分段,最多 4 段 |
+| 日志 | `src/lib/moderation/log.ts` + 迁移 `20260930040000_moderation_events.sql` | 只记类别、规则编号/服务标签、方向、入口、模型、request、user_id;不存原文;仅 service_role |
+| 普通咨询 | `app/api/consult/route.ts` | 输入:在 prompt-extraction 守卫之后、预留点数之前;输出:三条 agentic 流 `moderateOutput` + `completeResponse` 兜底(覆盖旧 text/plain 运行时);命中走新函数 `complete_consultation_moderated`(复制自免费结算,释放预留) |
+| 流式替换 | `lib/stream-agent-response.ts`、`lib/consultation-agent-events.ts`、`lib/consultation-run-timeline.ts`、`hooks/use-consultation-run.ts` | `answer.delta` 增加可选 `replace: true`;客户端时间线与发送钩子遇到即整段替换 |
+| 生时校正 | `v9/agent-route-request.ts`(输入)、`v9/agent-run-finish.ts`(输出) | 输出命中:`skipBilling` → `billing.release()`,spoken 换成撤下提示,走既有 replace |
+| 出生时间引导 | `app/api/birth-time-guide/route.ts` | `draft_evidence.message` / `reframe_unmatched.note` 输入审核 |
+| 后台 | `app/api/admin/moderation-events`、`components/admin/moderation-events-resource.tsx`、菜单「内容审核」 | 只读、`audit.read`、14 天分类计数 |
+
+未接入(无用户自由文本进模型):报告生成、今日星语、合盘、星盘 / 星历。Python API 未发现把用户自由文本送进模型的端点(模型调用都在 Next.js 侧)。
+
+## 既有断言改动(原值 / 新值 / 原因)
+
+| 文件 | 原值 | 新值 | 原因 |
+|---|---|---|---|
+| `tests/test_api_server_security.py` app_routes 列表 | admin/models 后直接 admin/orders | 中间加 admin/moderation-events | 新增只读后台页 |
+| `frontend/tests/admin-contracts.test.ts` 只读资源 | 四个 | 加 moderation-events | 同一只读约束覆盖新页 |
+
+## 验证(Claude fork 实跑,Linux,Node 22.14)
+
+| 项 | 结果 |
+|---|---|
+| `tsc --noEmit` | 0 错 |
+| `npm run lint` | 0 error,126 warning(同基线) |
+| 新增 `content-moderation.test.ts` | 10/10:占星用语 8 句全部放行;五类各自拦截含空格/全角/零宽规避;自伤给求助句;阿里云 high→拦、medium→待复核、none→放行,签名与表单校验;本地拦截不调远端;远端超时(50 ms)/500 → 沿用本地;无钥匙/非 aliyuncs.com 只用本地且只警告一次;日志不含原文、写失败不抛;流式回答被整段替换且 onComplete 收到提示句、时间线最终文本为提示句;路由里输入审核在预留点数之前、校正输出审核在结算之前 |
+| `npm test` 全量 | 4,428 / fail 26。与基线(24 条环境失败)相比多 2 条:`a lookup does not reset the answer clock…` 与 `a refresh inside a rectification conversation stays in it…`,均为并发负载下超时(后者等待 16.8 s 超时),单独各跑两次全部通过,且不涉及本次改动文件。新增 10 条,消失 0 条 |
+| `next build` | `/` 仍 `○ Static`;首屏 gzip-9 646,505 B(基线 646,480 B,+25 B) |
+| Python | `tests/test_api_server_security.py::test_capability_audit_scans_registry_and_local_sources` 通过;快速门 Python 步 1000 passed / 2 skipped(门内 `npm test` 步走系统 Node 20,属已知环境失败) |
+
+未验证:阿里云真实调用;迁移在真实 PostgreSQL 上执行(无 Docker);浏览器中被拦与被替换的表现;遗留 text/plain 运行时(无法就地替换,只保证存档与结算为提示句)。
diff --git a/frontend/DESIGN.md b/frontend/DESIGN.md
index ff75d57e..af19f669 100644
--- a/frontend/DESIGN.md
+++ b/frontend/DESIGN.md
@@ -574,6 +574,12 @@ The chart page, the ephemeris and the report list wait with one motion: the home
- **/terms and /privacy:** public, static, one 720px column on `--color-canvas-soft`; a top row with 「返回 Jyotisha」 and a link to the other document; while the text is a draft, a warning-tinted note 「草稿 · 版本 …,待律师审定」 above the title; h1 in the display face, h2 per section, bullet lists for enumerations.
- **Re-consent:** when a final version replaces the one a signed-in user accepted, one modal 「协议已更新」 with the same checkbox and 「同意并继续」; it cannot be dismissed without agreeing. Off while the documents are a draft.
+### Content moderation (2026-09-30, compliance round)
+
+- **Blocked input** comes back before the turn starts, as the same error notice a rejected send already uses (HTTP 400, `code: "content_blocked"`); nothing is added to the conversation and nothing is charged.
+- **Blocked output** in 普通咨询 arrives as one `answer.delta` with `replace: true` at the end of the stream: the streamed answer is swapped for the one-line notice in place (no second bubble, no error styling), and the turn completes free through `complete_consultation_moderated`. 生时校正 uses its existing replace path and releases the turn's charge. The legacy text/plain runtime cannot replace in place; its stored reply is the notice.
+- **Admin:** 「内容审核」 is a read-only list (time, side, verdict tag, route, categories, rule ids, provider, user id, request id) with a 14-day per-category count strip; it never shows user or model text.
+
### Personal report centre
- **Structure:** inside the app shell, not a page of its own. The name 「我的报告」 sits alone in the 46px header. The body opens with the **generate card** (`.report-center-create`, 2026-09-29, product sketch): a centred `--color-canvas` sheet, `--space-6` below the header (BUG-1103: it used to sit flush under it), with a file icon, the title 「完整本命报告」 (what you get), one line 「星盘、力量、大运、年运、瑜伽共 6 章,中英两版;生成后可离开,完成时下方自动出现。」 and the filled 「生成报告」 button (what it does) — title and button no longer say the same thing — the page's only generate entry; the header button was deleted rather than kept as a second one. Below it, 「过往的报告」 heads the **row list** of reports; with none yet, a single quiet line 「还没有个人报告。」 replaces the old second big empty card. The supporting paragraph (「有填报到分钟的出生时间即可生成……」) and the 「共 N 份 · 已完成 N 份」 overview line above it were removed on 2026-09-28 (TASK-self-edit-avatar-menu-20260928 S3, product: 「这里的提示去掉」); the minute requirement still lives in the 生成 button's hover title.
diff --git a/frontend/docs/VOICE.md b/frontend/docs/VOICE.md
index d8b4a651..1eaf3542 100644
--- a/frontend/docs/VOICE.md
+++ b/frontend/docs/VOICE.md
@@ -22,6 +22,13 @@ Jyotisha 的可见文案是产品的一部分。正确性红线(真实性、
首页开场语下面可以有一行今日趋势(今日星语卡片的 trend,每日生成);还没生成时写「今天的星语还没写出来。」,没有出生分钟时不写。入口按钮下面平时不写字,只在用户需要动手时写一句:有没做完的校正写「上次那次校正还没完成,可以在历史对话里接着做。」;当前人物不是本人写「生时校正暂时只支持本人。」。不再写「不确定出生时间时,用记得住的经历一步步缩小范围」「上次已经校正完,可以拿最新资料再来一次」,也不要用「·」把两句不相干的提示拼成一行。
+## 内容审核的三句话(2026-09-30,合规轮)
+
+- 用户输入被拦:「这个话题我们不能讨论,换个问题试试。」——一句,不说教、不解释规则、不说「违规」。
+- 涉及自伤:「听起来你现在很难受。请马上联系身边信任的人,或拨打心理援助热线 12356;如果有紧急危险,请拨打 110 或 120。」——先共情、给出路,不拒绝、不评判(热线号码上线前请法务/产品确认)。
+- 模型回答被撤下:「这段回答里有不适合展示的内容,已经撤下,本次不扣点。换个问法再试试。」——说清发生了什么、钱没扣、下一步。
+- 三句都在 `src/lib/moderation/index.ts`,改文案只改那里。
+
## 性别(选填)(2026-09-27,TASK-consult-gender-optional)
- 标题「性别(选填)」,三个选项「女」「男」「不填」;下面只有一行说明:「用于婚恋解读里判断夫星 / 妻星,不填也能用。」不解释为什么只有两项,不劝用户填。
diff --git a/frontend/src/app/admin/moderation-events/page.tsx b/frontend/src/app/admin/moderation-events/page.tsx
new file mode 100644
index 00000000..43bccb38
--- /dev/null
+++ b/frontend/src/app/admin/moderation-events/page.tsx
@@ -0,0 +1,2 @@
+import ModerationEventsResource from "@/components/admin/moderation-events-resource";
+export default function Page() { return ; }
diff --git a/frontend/src/app/api/admin/moderation-events/route.ts b/frontend/src/app/api/admin/moderation-events/route.ts
new file mode 100644
index 00000000..d3b3e1d8
--- /dev/null
+++ b/frontend/src/app/api/admin/moderation-events/route.ts
@@ -0,0 +1,101 @@
+import { NextResponse } from "next/server";
+
+import { requirePermission } from "@/lib/admin/auth";
+import { pageOffset, queryAdminRows } from "@/lib/admin/database";
+import {
+ adminErrorResponse,
+ invalidQueryResponse,
+ parseListQuery,
+ readonlyAdminMutation,
+} from "@/lib/admin/http";
+
+export const runtime = "nodejs";
+
+type ModerationRow = {
+ id: string;
+ user_id: string | null;
+ side: string;
+ verdict: string;
+ route: string;
+ categories: string[];
+ rule_ids: string[];
+ provider: string;
+ model_id: string | null;
+ request_id: string | null;
+ created_at: Date;
+ total_count: string;
+};
+
+type SummaryRow = { day: string; category: string; hits: string };
+
+const sortColumns = new Map([
+ ["createdAt", "m.created_at"],
+ ["route", "m.route"],
+]);
+
+export const POST = readonlyAdminMutation;
+export const PUT = readonlyAdminMutation;
+export const PATCH = readonlyAdminMutation;
+export const DELETE = readonlyAdminMutation;
+
+/**
+ * Content moderation hits (2026-09-30). Read-only, audit permission. Rows
+ * carry categories and rule ids only — no user or model text is stored.
+ * `?summary=1` returns hit counts per category per day for the last 14 days.
+ */
+export async function GET(request: Request) {
+ try {
+ await requirePermission("audit.read");
+ if (new URL(request.url).searchParams.get("summary") === "1") {
+ const rows = await queryAdminRows(`
+ select to_char(date_trunc('day', m.created_at), 'YYYY-MM-DD') as day, c.category, count(*)::text as hits
+ from public.moderation_events m, unnest(m.categories) as c(category)
+ where m.created_at >= now() - interval '14 days'
+ group by 1, 2
+ order by 1 desc, 3 desc
+ `, []);
+ return NextResponse.json({ data: rows.map((row) => ({ day: row.day, category: row.category, hits: Number(row.hits) })) });
+ }
+ const parsed = parseListQuery(request);
+ if (!parsed.success) return invalidQueryResponse(parsed.error.flatten());
+ const { page, pageSize, sort, order, q, status } = parsed.data;
+ const values: unknown[] = [];
+ const conditions: string[] = [];
+ if (q) {
+ values.push(`%${q}%`);
+ conditions.push(`(m.route ilike $${values.length} or m.request_id ilike $${values.length} or m.user_id::text ilike $${values.length})`);
+ }
+ if (status === "input" || status === "output") {
+ values.push(status);
+ conditions.push(`m.side = $${values.length}`);
+ }
+ values.push(pageSize, pageOffset(page, pageSize));
+ const sortColumn = sortColumns.get(sort ?? "createdAt") ?? "m.created_at";
+ const rows = await queryAdminRows(`
+ select m.id::text as id, m.user_id, m.side, m.verdict, m.route, m.categories, m.rule_ids, m.provider,
+ m.model_id, m.request_id, m.created_at, count(*) over()::text as total_count
+ from public.moderation_events m
+ ${conditions.length ? `where ${conditions.join(" and ")}` : ""}
+ order by ${sortColumn} ${order === "asc" ? "asc" : "desc"}, m.id desc
+ limit $${values.length - 1} offset $${values.length}
+ `, values);
+ return NextResponse.json({
+ data: rows.map((row) => ({
+ id: row.id,
+ userId: row.user_id,
+ side: row.side,
+ verdict: row.verdict,
+ route: row.route,
+ categories: row.categories,
+ ruleIds: row.rule_ids,
+ provider: row.provider,
+ modelId: row.model_id,
+ requestId: row.request_id,
+ createdAt: row.created_at.toISOString(),
+ })),
+ total: Number(rows[0]?.total_count ?? 0),
+ });
+ } catch (error) {
+ return adminErrorResponse(error);
+ }
+}
diff --git a/frontend/src/app/api/birth-time-guide/route.ts b/frontend/src/app/api/birth-time-guide/route.ts
index ed93eb10..3ae79700 100644
--- a/frontend/src/app/api/birth-time-guide/route.ts
+++ b/frontend/src/app/api/birth-time-guide/route.ts
@@ -16,6 +16,8 @@ import {
import { StaleJourneyTurnError } from "@/lib/birth-time-journey-turn-persistence";
import { jsonForSupabaseSetupFailure } from "@/lib/api/service-unavailable";
import { createAdminSupabaseClient } from "@/lib/supabase/admin";
+import { inputBlockedCopy, moderate } from "@/lib/moderation";
+import { recordModerationEvent } from "@/lib/moderation/log";
import { createServerSupabaseClient } from "@/lib/supabase/server";
import { getBirthTimeGuideAgent } from "@/mastra";
import { loadLanguageModelCatalog } from "@/lib/model-catalog";
@@ -59,6 +61,21 @@ export async function POST(request: Request) {
);
}
+ // Content moderation, input side (2026-09-30): typed text never reaches the guide model if it fails.
+ const typedText = parsed.data.type === "draft_evidence" ? parsed.data.message
+ : parsed.data.type === "reframe_unmatched" ? parsed.data.note : "";
+ if (typedText.trim()) {
+ const context = { route: "birth-time-guide", userId: user.id };
+ const verdict = await moderate(typedText, "input", context);
+ if (verdict.verdict !== "pass") await recordModerationEvent(createAdminSupabaseClient() as never, verdict, "input", context);
+ if (verdict.verdict === "block") {
+ return NextResponse.json(
+ { error: "无法处理该请求", code: "content_blocked", message: inputBlockedCopy(verdict) },
+ { status: 400 },
+ );
+ }
+ }
+
const store = createSupabaseBirthTimeJourneyStore(createAdminSupabaseClient());
const actions = createJourneyTurnActions({ store });
const journey = createBirthTimeJourneyService({
diff --git a/frontend/src/app/api/consult/route.ts b/frontend/src/app/api/consult/route.ts
index d0fae84e..d3efea85 100644
--- a/frontend/src/app/api/consult/route.ts
+++ b/frontend/src/app/api/consult/route.ts
@@ -10,6 +10,8 @@ import {
runConsultationWorkflow,
} from "@/mastra";
import { blocksPromptExtraction } from "@/lib/consult-safety";
+import { inputBlockedCopy, moderate, MODERATION_OUTPUT_REPLACED_COPY } from "@/lib/moderation";
+import { recordModerationEvent } from "@/lib/moderation/log";
import { consultationDomainSchema } from "@/lib/consultation-domain-registry";
import { parseAgentReply } from "@/lib/agent-reply";
import { createConsultationReplyMetadata } from "@/lib/consultation-reply-metadata";
@@ -409,6 +411,17 @@ export async function POST(request: Request) {
{ status: 400 },
);
}
+ // Content moderation, input side (2026-09-30): before any reservation, so a
+ // blocked question is never sent to a model and never charged.
+ const moderationContext = { route: "consult", userId: user.id, requestId: parsed.data.requestId, modelId: chatSession.model_id };
+ const inputModeration = await moderate(visibleQuestion, "input", moderationContext);
+ if (inputModeration.verdict !== "pass") await recordModerationEvent(accounting, inputModeration, "input", moderationContext);
+ if (inputModeration.verdict === "block") {
+ return NextResponse.json(
+ { error: "无法处理该请求", code: "content_blocked", message: inputBlockedCopy(inputModeration) },
+ { status: 400 },
+ );
+ }
type SelectedModel = NonNullable>>;
type ReservationResult = { success: boolean; credits: number | null; error_code: string | null };
type ModelSelection = Awaited>>;
@@ -785,6 +798,38 @@ export async function POST(request: Request) {
};
}
+ // Content moderation, output side (2026-09-30). The finished answer is
+ // checked once; a hit replaces it with a notice and the turn completes free.
+ let outputModerated = false;
+ let outputBlocked = false;
+ async function moderateAnswer(text: string): Promise {
+ outputModerated = true;
+ const verdict = await moderate(text, "output", moderationContext);
+ if (verdict.verdict !== "pass") await recordModerationEvent(accounting, verdict, "output", moderationContext);
+ if (verdict.verdict !== "block") return null;
+ outputBlocked = true;
+ return MODERATION_OUTPUT_REPLACED_COPY;
+ }
+
+ async function completeModeratedResponse(usage: Promise<{ inputTokens?: number; outputTokens?: number }>): Promise {
+ const actualUsage = await usagePayload(usage);
+ await retryDetachedSettlement(async () => {
+ const { data, error } = await accounting.rpc("complete_consultation_moderated", {
+ p_user_id: userId,
+ p_request_id: requestId,
+ p_session_id: sessionId,
+ p_response_message: { role: "assistant", text: MODERATION_OUTPUT_REPLACED_COPY, responseKind: "moderated" },
+ p_actual_usage: actualUsage,
+ });
+ const result = consultationCompletionSchema.safeParse(first(data ?? []));
+ if (error || !result.success || !result.data.success) {
+ throw new CreditRpcError(error?.message || (result.success ? result.data.error_code : "invalid_moderated_completion") || "moderated_completion_rejected");
+ }
+ return result.data;
+ });
+ return "completed";
+ }
+
async function completeResponse(
rawTransformedText: string,
usage: Promise<{ inputTokens?: number; outputTokens?: number }>,
@@ -800,6 +845,8 @@ export async function POST(request: Request) {
createConsultationReplyMetadata({ question: visibleQuestion }),
);
if (!reply.text) throw new Error("empty_agent_reply");
+ if (!outputModerated) await moderateAnswer(reply.text);
+ if (outputBlocked) return await completeModeratedResponse(usage);
const persistedThinking = thinkingText?.trim().slice(0, 4_000);
const responseMessage = {
role: "assistant" as const,
@@ -1164,6 +1211,7 @@ export async function POST(request: Request) {
: {}),
});
return streamAgentResponse({
+ moderateOutput: moderateAnswer,
runId: requestId,
requestId,
sideEvent: titleSideEvent,
@@ -1266,6 +1314,7 @@ export async function POST(request: Request) {
: {}),
});
return streamAgentResponse({
+ moderateOutput: moderateAnswer,
runId: requestId,
requestId,
sideEvent: titleSideEvent,
@@ -1402,6 +1451,7 @@ export async function POST(request: Request) {
: {}),
});
return streamAgentResponse({
+ moderateOutput: moderateAnswer,
runId: requestId,
requestId,
sideEvent: titleSideEvent,
diff --git a/frontend/src/components/admin/admin-app.tsx b/frontend/src/components/admin/admin-app.tsx
index 9b977d77..1e2d6b6b 100644
--- a/frontend/src/components/admin/admin-app.tsx
+++ b/frontend/src/components/admin/admin-app.tsx
@@ -121,6 +121,7 @@ export function AdminApp({ children }: { children: ReactNode }) {
{ name: "feature-flags", list: "/admin/feature-flags", meta: { label: "功能开关", icon: } },
{ name: "security", list: "/admin/security", meta: { label: "安全验证", icon: } },
{ name: "audit-logs", list: "/admin/audit-logs", meta: { label: "审计日志", icon: } },
+ { name: "moderation-events", list: "/admin/moderation-events", meta: { label: "内容审核", icon: } },
]}
options={{
syncWithLocation: true,
diff --git a/frontend/src/components/admin/moderation-events-resource.tsx b/frontend/src/components/admin/moderation-events-resource.tsx
new file mode 100644
index 00000000..3c7ded1b
--- /dev/null
+++ b/frontend/src/components/admin/moderation-events-resource.tsx
@@ -0,0 +1,71 @@
+"use client";
+
+import { useEffect, useState } from "react";
+import { Descriptions, Tag, type TableColumnsType } from "antd";
+
+import { formatAdminDate, ResourceTable } from "@/components/admin/resource-table";
+
+type ModerationRecord = {
+ id: string;
+ userId: string | null;
+ side: "input" | "output";
+ verdict: "block" | "review";
+ route: string;
+ categories: string[];
+ ruleIds: string[];
+ provider: string;
+ requestId: string | null;
+ createdAt: string;
+};
+
+const CATEGORY_LABELS: Record = {
+ political: "政治",
+ violent_extremism: "暴恐",
+ sexual_minor: "涉未成年色情",
+ self_harm: "自伤",
+ illegal_trade: "违法交易",
+ remote: "审核服务判定",
+};
+
+const columns: TableColumnsType = [
+ { title: "时间", dataIndex: "createdAt", sorter: true, render: formatAdminDate },
+ { title: "方向", dataIndex: "side", render: (side: string) => (side === "input" ? "用户输入" : "模型输出") },
+ { title: "处理", dataIndex: "verdict", render: (verdict: string) => {verdict === "block" ? "拦截" : "待复核"} },
+ { title: "入口", dataIndex: "route", sorter: true },
+ { title: "类别", dataIndex: "categories", render: (list: string[]) => list.map((c) => CATEGORY_LABELS[c] ?? c).join("、") },
+ { title: "命中规则", dataIndex: "ruleIds", render: (list: string[]) => list.join(", ") },
+ { title: "来源", dataIndex: "provider" },
+ { title: "用户 ID", dataIndex: "userId" },
+ { title: "Request ID", dataIndex: "requestId" },
+];
+
+function Summary() {
+ const [rows, setRows] = useState<{ day: string; category: string; hits: number }[] | null>(null);
+ useEffect(() => {
+ let cancelled = false;
+ fetch("/api/admin/moderation-events?summary=1", { credentials: "same-origin" })
+ .then((response) => (response.ok ? response.json() : { data: [] }))
+ .then((json: { data?: { day: string; category: string; hits: number }[] }) => { if (!cancelled) setRows(json.data ?? []); })
+ .catch(() => { if (!cancelled) setRows([]); });
+ return () => { cancelled = true; };
+ }, []);
+ const totals = new Map();
+ for (const row of rows ?? []) totals.set(row.category, (totals.get(row.category) ?? 0) + row.hits);
+ return ({ key: category, label: `近 14 天 · ${CATEGORY_LABELS[category] ?? category}`, children: String(hits) })),
+ ]} />;
+}
+
+export default function ModerationEventsPage() {
+ return
+ resource="moderation-events"
+ title="内容审核记录(只读)"
+ columns={columns}
+ statusOptions={[
+ { label: "用户输入", value: "input" },
+ { label: "模型输出", value: "output" },
+ ]}
+ extra={}
+ />;
+}
diff --git a/frontend/src/hooks/use-consultation-run.ts b/frontend/src/hooks/use-consultation-run.ts
index 7252ab58..fdfef76e 100644
--- a/frontend/src/hooks/use-consultation-run.ts
+++ b/frontend/src/hooks/use-consultation-run.ts
@@ -938,7 +938,8 @@ export function useConsultationRun(params: ConsultationRunParams) {
if (responseKind !== "smalltalk") timelineState = reduceConsultationTimeline(timelineState, event);
frames.touch();
if (event.type === "answer.delta") {
- answer += event.text;
+ // A moderated answer arrives as one replacing delta (2026-09-30).
+ answer = event.replace ? event.text : answer + event.text;
frames.setAnswer(answer);
}
if (event.type === "thinking.delta" && typeof event.text === "string") {
diff --git a/frontend/src/lib/admin/providers.ts b/frontend/src/lib/admin/providers.ts
index 73341cd2..c0112c6f 100644
--- a/frontend/src/lib/admin/providers.ts
+++ b/frontend/src/lib/admin/providers.ts
@@ -192,6 +192,7 @@ const resourcePermissions: Record = {
"credit-transactions": { read: "billing.orders.read" },
consultations: { read: "billing.orders.read" },
"audit-logs": { read: "audit.read" },
+ "moderation-events": { read: "audit.read" },
payments: { read: "billing.orders.read" },
packages: { read: "billing.products.read", write: "billing.products.write" },
products: { read: "billing.products.read", write: "billing.products.write" },
diff --git a/frontend/src/lib/consultation-agent-events.ts b/frontend/src/lib/consultation-agent-events.ts
index 0441b361..79cc8581 100644
--- a/frontend/src/lib/consultation-agent-events.ts
+++ b/frontend/src/lib/consultation-agent-events.ts
@@ -101,7 +101,8 @@ const toolFailedSchema = z.object({
type: z.literal("tool.failed"), callId: z.string(), tool: z.literal("run-jyotish-consultation"),
code: z.enum(["calculation_failed", "timeout", "cancelled"]),
}).strict();
-const answerDeltaSchema = z.object({ type: z.literal("answer.delta"), text: z.string() }).strict();
+/** `replace: true` swaps the whole answer for `text` (a moderated answer, 2026-09-30). */
+const answerDeltaSchema = z.object({ type: z.literal("answer.delta"), text: z.string(), replace: z.literal(true).optional() }).strict();
const thinkingDeltaSchema = z.object({ type: z.literal("thinking.delta"), text: z.string() }).strict();
const thinkingSectionEventSchema = publicThinkingSectionSchema.extend({
type: z.literal("thinking.section"),
diff --git a/frontend/src/lib/consultation-run-timeline.ts b/frontend/src/lib/consultation-run-timeline.ts
index f50c3665..beaedd59 100644
--- a/frontend/src/lib/consultation-run-timeline.ts
+++ b/frontend/src/lib/consultation-run-timeline.ts
@@ -152,7 +152,7 @@ export function reduceConsultationTimeline(
});
}
if (event.type === "answer.delta") {
- const next = { ...completeLiveThink(state), answer: `${state.answer}${event.text}` };
+ const next = { ...completeLiveThink(state), answer: event.replace ? event.text : `${state.answer}${event.text}` };
return syncWriteRows(next, false);
}
if (event.type === "run.completed" || event.type === "run.failed") {
diff --git a/frontend/src/lib/moderation/aliyun-provider.ts b/frontend/src/lib/moderation/aliyun-provider.ts
new file mode 100644
index 00000000..a32cd61a
--- /dev/null
+++ b/frontend/src/lib/moderation/aliyun-provider.ts
@@ -0,0 +1,109 @@
+import { createHmac, randomUUID } from "node:crypto";
+
+import type { ModerationResult, ModerationSide } from "./types";
+
+/**
+ * Aliyun 内容安全 (Green 2.0) text moderation for large-model input and output
+ * (`TextModerationPlus`, services `llm_query_moderation` / `llm_response_moderation`),
+ * signed with the public RPC signature v1 (HMAC-SHA1).
+ *
+ * Written from Aliyun's public API documentation without credentials, so it
+ * has not been exercised against the live service (BLOCKED.md). The endpoint
+ * is restricted to `*.aliyuncs.com` so a mis-set env cannot point the server at
+ * an arbitrary host (AGENTS §8-3).
+ */
+export type AliyunGreenConfig = Readonly<{
+ accessKeyId: string;
+ accessKeySecret: string;
+ endpoint: string;
+}>;
+
+const API_VERSION = "2022-03-02";
+/** Service limit per request for the LLM moderation services. */
+export const ALIYUN_CHUNK_CHARS = 2_000;
+const MAX_CHUNKS = 4;
+
+export function readAliyunGreenConfig(env: Readonly> = process.env): AliyunGreenConfig | null {
+ const accessKeyId = env.ALIYUN_GREEN_ACCESS_KEY_ID?.trim();
+ const accessKeySecret = env.ALIYUN_GREEN_ACCESS_KEY_SECRET?.trim();
+ if (!accessKeyId || !accessKeySecret) return null;
+ const endpoint = (env.ALIYUN_GREEN_ENDPOINT?.trim() || "green-cip.cn-shanghai.aliyuncs.com").replace(/^https?:\/\//, "").replace(/\/.*$/, "");
+ if (!/^[a-z0-9.-]+\.aliyuncs\.com$/.test(endpoint)) return null;
+ return { accessKeyId, accessKeySecret, endpoint };
+}
+
+/** RFC 3986 encoding as Aliyun's RPC signature requires. */
+export function aliyunPercentEncode(value: string): string {
+ return encodeURIComponent(value)
+ .replace(/[!'()*]/g, (char) => `%${char.charCodeAt(0).toString(16).toUpperCase()}`);
+}
+
+export function signAliyunRpc(params: Readonly>, secret: string, method = "POST"): string {
+ const canonical = Object.keys(params)
+ .sort()
+ .map((key) => `${aliyunPercentEncode(key)}=${aliyunPercentEncode(params[key]!)}`)
+ .join("&");
+ const stringToSign = `${method}&${aliyunPercentEncode("/")}&${aliyunPercentEncode(canonical)}`;
+ return createHmac("sha1", `${secret}&`).update(stringToSign).digest("base64");
+}
+
+function service(side: ModerationSide): string {
+ return side === "input" ? "llm_query_moderation" : "llm_response_moderation";
+}
+
+async function moderateChunk(
+ config: AliyunGreenConfig,
+ content: string,
+ side: ModerationSide,
+ fetchImpl: typeof fetch,
+ signal: AbortSignal,
+): Promise<{ riskLevel: string; labels: string[] }> {
+ const params: Record = {
+ Format: "JSON",
+ Version: API_VERSION,
+ AccessKeyId: config.accessKeyId,
+ SignatureMethod: "HMAC-SHA1",
+ SignatureVersion: "1.0",
+ SignatureNonce: randomUUID(),
+ Timestamp: new Date().toISOString().replace(/\.\d{3}Z$/, "Z"),
+ Action: "TextModerationPlus",
+ Service: service(side),
+ ServiceParameters: JSON.stringify({ content }),
+ };
+ params.Signature = signAliyunRpc(params, config.accessKeySecret);
+ const response = await fetchImpl(`https://${config.endpoint}/`, {
+ method: "POST",
+ headers: { "content-type": "application/x-www-form-urlencoded" },
+ body: new URLSearchParams(params).toString(),
+ signal,
+ cache: "no-store",
+ });
+ const json = await response.json().catch(() => null) as { Code?: number; Data?: { RiskLevel?: string; Result?: { Label?: string }[] } } | null;
+ if (!response.ok || json?.Code !== 200 || !json.Data) throw new Error("aliyun_green_rejected");
+ return {
+ riskLevel: String(json.Data.RiskLevel ?? "none"),
+ labels: (json.Data.Result ?? []).map((row) => String(row.Label ?? "")).filter((label) => label && label !== "nonLabel"),
+ };
+}
+
+export async function moderateWithAliyun(
+ config: AliyunGreenConfig,
+ text: string,
+ side: ModerationSide,
+ options: Readonly<{ signal: AbortSignal; fetchImpl?: typeof fetch }>,
+): Promise {
+ const chunks: string[] = [];
+ for (let start = 0; start < text.length && chunks.length < MAX_CHUNKS; start += ALIYUN_CHUNK_CHARS) {
+ chunks.push(text.slice(start, start + ALIYUN_CHUNK_CHARS));
+ }
+ const results = await Promise.all(chunks.map((chunk) => moderateChunk(config, chunk, side, options.fetchImpl ?? fetch, options.signal)));
+ const high = results.some((result) => result.riskLevel === "high");
+ const medium = results.some((result) => result.riskLevel === "medium");
+ const labels = [...new Set(results.flatMap((result) => result.labels))].slice(0, 8);
+ return {
+ verdict: high ? "block" : medium ? "review" : "pass",
+ categories: high || medium ? ["remote"] : [],
+ ruleIds: labels.map((label) => `aliyun:${label}`.slice(0, 64)),
+ providerId: "aliyun",
+ };
+}
diff --git a/frontend/src/lib/moderation/index.ts b/frontend/src/lib/moderation/index.ts
new file mode 100644
index 00000000..b0d2b59b
--- /dev/null
+++ b/frontend/src/lib/moderation/index.ts
@@ -0,0 +1,82 @@
+import { moderateWithAliyun, readAliyunGreenConfig, type AliyunGreenConfig } from "./aliyun-provider";
+import { moderateLocally } from "./local-provider";
+import type { ModerationContext, ModerationResult, ModerationSide } from "./types";
+
+export type { ModerationCategory, ModerationContext, ModerationResult, ModerationSide, ModerationVerdict } from "./types";
+export { normalizeForModeration, moderateLocally } from "./local-provider";
+
+/** The remote provider gets at most this long; a slow provider never holds a turn. */
+export const MODERATION_REMOTE_TIMEOUT_MS = 1_500;
+
+/** What the user sees when their message is not sent (not charged). */
+export const MODERATION_INPUT_BLOCKED_COPY = "这个话题我们不能讨论,换个问题试试。";
+/**
+ * Self-harm is answered with help, not with a refusal. The national
+ * psychological assistance hotline 12356 is a public service; legal/product
+ * review this line before launch.
+ */
+export const MODERATION_SELF_HARM_COPY = "听起来你现在很难受。请马上联系身边信任的人,或拨打心理援助热线 12356;如果有紧急危险,请拨打 110 或 120。";
+/** Replaces a whole answer whose text failed moderation; the turn is not charged. */
+export const MODERATION_OUTPUT_REPLACED_COPY = "这段回答里有不适合展示的内容,已经撤下,本次不扣点。换个问法再试试。";
+
+export type ModerationOptions = Readonly<{
+ env?: Readonly>;
+ fetchImpl?: typeof fetch;
+ timeoutMs?: number;
+}>;
+
+let warnedMissingKeys = false;
+
+function remoteConfig(env: Readonly>): AliyunGreenConfig | null {
+ if ((env.MODERATION_PROVIDER ?? "local").trim().toLowerCase() !== "aliyun") return null;
+ const config = readAliyunGreenConfig(env);
+ if (!config && !warnedMissingKeys) {
+ warnedMissingKeys = true;
+ console.warn("[moderation] MODERATION_PROVIDER=aliyun but keys or endpoint are missing or invalid; using the local lexicon only");
+ }
+ return config;
+}
+
+/**
+ * One moderation decision for one text.
+ *
+ * The local lexicon always runs and a local block is final. When the remote
+ * provider is configured it runs next, within `MODERATION_REMOTE_TIMEOUT_MS`.
+ * If it fails or times out the local verdict stands, on both sides: input is
+ * therefore never let through *without* the local check (fail-closed to the
+ * local floor), and output is never withheld because the provider is down
+ * (fail-open beyond the local floor). Nothing is blocked only because a
+ * provider was unreachable.
+ */
+export async function moderate(
+ text: string,
+ side: ModerationSide,
+ _context: ModerationContext,
+ options: ModerationOptions = {},
+): Promise {
+ const local = moderateLocally(text);
+ if (local.verdict === "block") return local;
+ const config = remoteConfig(options.env ?? process.env);
+ if (!config || !text.trim()) return local;
+ const controller = new AbortController();
+ const timer = setTimeout(
+ () => controller.abort(new DOMException("moderation provider timed out", "TimeoutError")),
+ options.timeoutMs ?? MODERATION_REMOTE_TIMEOUT_MS,
+ );
+ try {
+ return await moderateWithAliyun(config, text, side, { signal: controller.signal, fetchImpl: options.fetchImpl });
+ } catch {
+ return { ...local, remoteFailed: true };
+ } finally {
+ clearTimeout(timer);
+ }
+}
+
+/** The line shown instead of sending a blocked message. */
+export function inputBlockedCopy(result: ModerationResult): string {
+ return result.categories.includes("self_harm") ? MODERATION_SELF_HARM_COPY : MODERATION_INPUT_BLOCKED_COPY;
+}
+
+export function resetModerationWarningsForTests(): void {
+ warnedMissingKeys = false;
+}
diff --git a/frontend/src/lib/moderation/lexicon.ts b/frontend/src/lib/moderation/lexicon.ts
new file mode 100644
index 00000000..dd744bf8
--- /dev/null
+++ b/frontend/src/lib/moderation/lexicon.ts
@@ -0,0 +1,84 @@
+import type { ModerationCategory } from "./types";
+
+/**
+ * Local sensitive-content lexicon (2026-09-30, compliance round).
+ *
+ * Reviewable data, not logic. Rules match against normalized text (see
+ * `normalizeForModeration`: lower case, half-width, whitespace / punctuation /
+ * zero-width removed). Entries are deliberately multi-character phrases so the
+ * astrology vocabulary this product lives on — 死亡宫、第八宫、凶星、灾难、
+ * Maraka、罗睺、计都 — never trips them (see moderation tests). The list is a
+ * floor, not a filter of record: legal must review it before launch, and the
+ * remote provider (MODERATION_PROVIDER=aliyun) is the production judge.
+ */
+export type LexiconRule = Readonly<{
+ id: string;
+ category: Exclude;
+ /** Plain phrases (matched as substrings of the normalized text). */
+ phrases?: readonly string[];
+ /** Patterns run on the normalized text. */
+ patterns?: readonly RegExp[];
+}>;
+
+export const MODERATION_LEXICON: readonly LexiconRule[] = [
+ {
+ id: "political.subversion",
+ category: "political",
+ phrases: ["颠覆国家政权", "推翻共产党", "推翻中共", "分裂国家", "煽动颠覆"],
+ },
+ {
+ id: "political.separatism",
+ category: "political",
+ phrases: ["台独", "藏独", "疆独", "港独", "东突厥斯坦"],
+ },
+ {
+ id: "political.banned-org",
+ category: "political",
+ phrases: ["法轮功", "法轮大法", "全能神"],
+ },
+ {
+ id: "extremism.weapons-howto",
+ category: "violent_extremism",
+ phrases: ["炸弹制作", "制作炸弹", "自制炸药", "炸药配方", "土制炸弹", "怎么做炸弹"],
+ },
+ {
+ id: "extremism.terror",
+ category: "violent_extremism",
+ phrases: ["恐怖袭击计划", "策划恐袭", "圣战招募", "加入isis", "加入伊斯兰国"],
+ },
+ {
+ id: "extremism.harm-others",
+ category: "violent_extremism",
+ phrases: ["怎么杀人不被发现", "如何杀人不被发现", "怎么投毒", "如何投毒"],
+ },
+ {
+ id: "sexual.minor",
+ category: "sexual_minor",
+ phrases: ["儿童色情", "幼女色情", "萝莉色情", "未成年色情", "恋童"],
+ },
+ {
+ id: "selfharm.method",
+ category: "self_harm",
+ phrases: ["自杀方法", "怎么自杀", "如何自杀", "怎样自杀", "怎么割腕", "吃多少安眠药会死", "无痛自杀"],
+ },
+ {
+ id: "selfharm.intent",
+ category: "self_harm",
+ phrases: ["我想自杀", "我要自杀", "我不想活了", "想结束自己的生命"],
+ },
+ {
+ id: "trade.drugs",
+ category: "illegal_trade",
+ phrases: ["买冰毒", "卖冰毒", "买海洛因", "买大麻", "代购毒品", "买k粉"],
+ },
+ {
+ id: "trade.weapons",
+ category: "illegal_trade",
+ phrases: ["买枪支", "卖枪支", "网购枪支", "仿真枪出售"],
+ },
+ {
+ id: "trade.documents",
+ category: "illegal_trade",
+ phrases: ["办假证", "代开发票", "买卖身份证", "出售银行卡"],
+ },
+];
diff --git a/frontend/src/lib/moderation/local-provider.ts b/frontend/src/lib/moderation/local-provider.ts
new file mode 100644
index 00000000..30348dac
--- /dev/null
+++ b/frontend/src/lib/moderation/local-provider.ts
@@ -0,0 +1,27 @@
+import { MODERATION_LEXICON, type LexiconRule } from "./lexicon";
+import type { ModerationCategory, ModerationResult } from "./types";
+
+const ZERO_WIDTH = /[-]/g;
+const SEPARATORS = /[\s\p{P}\p{S}_]+/gu;
+
+/** Lower case, half-width, no whitespace / punctuation / zero-width characters. */
+export function normalizeForModeration(text: string): string {
+ return text
+ .normalize("NFKC")
+ .replace(ZERO_WIDTH, "")
+ .toLowerCase()
+ .replace(SEPARATORS, "");
+}
+
+function matches(rule: LexiconRule, normalized: string): boolean {
+ if (rule.phrases?.some((phrase) => normalized.includes(normalizeForModeration(phrase)))) return true;
+ return rule.patterns?.some((pattern) => pattern.test(normalized)) ?? false;
+}
+
+export function moderateLocally(text: string, lexicon: readonly LexiconRule[] = MODERATION_LEXICON): ModerationResult {
+ const normalized = normalizeForModeration(text);
+ const hits = normalized ? lexicon.filter((rule) => matches(rule, normalized)) : [];
+ if (!hits.length) return { verdict: "pass", categories: [], ruleIds: [], providerId: "local" };
+ const categories = [...new Set(hits.map((rule) => rule.category))] as ModerationCategory[];
+ return { verdict: "block", categories, ruleIds: hits.map((rule) => rule.id), providerId: "local" };
+}
diff --git a/frontend/src/lib/moderation/log.ts b/frontend/src/lib/moderation/log.ts
new file mode 100644
index 00000000..caf493c1
--- /dev/null
+++ b/frontend/src/lib/moderation/log.ts
@@ -0,0 +1,39 @@
+import type { ModerationContext, ModerationResult, ModerationSide } from "./types";
+
+type InsertClient = {
+ from(table: string): { insert(row: Record): PromiseLike<{ error: unknown }> };
+};
+
+export const MODERATION_EVENTS_TABLE = "moderation_events";
+
+/**
+ * One row per hit. It records what fired (categories, lexicon rule ids or the
+ * remote provider's labels), where (route, side), the model and the request —
+ * never the user's text or the model's text (AGENTS §8). `user_id` is kept so
+ * support can answer a complaint; the table is service-role only.
+ * Logging never blocks or fails the turn.
+ */
+export async function recordModerationEvent(
+ client: InsertClient,
+ result: ModerationResult,
+ side: ModerationSide,
+ context: ModerationContext,
+): Promise {
+ if (result.verdict === "pass") return;
+ try {
+ const { error } = await client.from(MODERATION_EVENTS_TABLE).insert({
+ user_id: context.userId ?? null,
+ side,
+ verdict: result.verdict,
+ route: context.route.slice(0, 64),
+ categories: [...result.categories],
+ rule_ids: result.ruleIds.map((id) => id.slice(0, 64)).slice(0, 16),
+ provider: result.providerId,
+ model_id: context.modelId?.slice(0, 120) ?? null,
+ request_id: context.requestId?.slice(0, 120) ?? null,
+ });
+ if (error) console.warn("[moderation] event log write failed");
+ } catch {
+ console.warn("[moderation] event log write failed");
+ }
+}
diff --git a/frontend/src/lib/moderation/types.ts b/frontend/src/lib/moderation/types.ts
new file mode 100644
index 00000000..d591b63c
--- /dev/null
+++ b/frontend/src/lib/moderation/types.ts
@@ -0,0 +1,29 @@
+/** Which side of a model call the text is on. */
+export type ModerationSide = "input" | "output";
+
+export type ModerationVerdict = "pass" | "block" | "review";
+
+export type ModerationCategory =
+ | "political"
+ | "violent_extremism"
+ | "sexual_minor"
+ | "self_harm"
+ | "illegal_trade"
+ | "remote";
+
+export type ModerationResult = Readonly<{
+ verdict: ModerationVerdict;
+ categories: readonly ModerationCategory[];
+ /** Local lexicon rule ids, or the remote provider's labels. Never the user's text. */
+ ruleIds: readonly string[];
+ providerId: "local" | "aliyun";
+ /** True when the remote provider was configured but could not answer in time. */
+ remoteFailed?: boolean;
+}>;
+
+export type ModerationContext = Readonly<{
+ route: string;
+ userId?: string | null;
+ requestId?: string | null;
+ modelId?: string | null;
+}>;
diff --git a/frontend/src/lib/rectification-agentic/v9/agent-route-request.ts b/frontend/src/lib/rectification-agentic/v9/agent-route-request.ts
index 9b1577a5..be849bb0 100644
--- a/frontend/src/lib/rectification-agentic/v9/agent-route-request.ts
+++ b/frontend/src/lib/rectification-agentic/v9/agent-route-request.ts
@@ -9,6 +9,8 @@
*/
import { NextResponse } from "next/server";
import { blocksPromptExtraction } from "@/lib/consult-safety";
+import { inputBlockedCopy, moderate } from "@/lib/moderation";
+import { recordModerationEvent } from "@/lib/moderation/log";
import { loadRuntimeFeatureFlags } from "@/lib/feature-flags";
import { isProductEnabled } from "@/lib/product-access";
import {
@@ -58,6 +60,19 @@ export async function resolveRectificationAgentRequest(input: Readonly<{
{ status: 400 },
);
}
+ // Content moderation, input side (2026-09-30): a typed message is checked
+ // before the turn opens, so a blocked one reaches no model and is not charged.
+ if (promptSource.trim()) {
+ const context = { route: "rectification", userId: user.id, requestId: parsed.data.requestId };
+ const verdict = await moderate(promptSource, "input", context);
+ if (verdict.verdict !== "pass") await recordModerationEvent(accounting as never, verdict, "input", context);
+ if (verdict.verdict === "block") {
+ return NextResponse.json(
+ { error: "无法处理该请求", code: "content_blocked", message: inputBlockedCopy(verdict) },
+ { status: 400 },
+ );
+ }
+ }
const userId = user.id;
const { caseId, sessionId, requestId, action } = parsed.data;
diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run-finish.ts b/frontend/src/lib/rectification-agentic/v9/agent-run-finish.ts
index a2206cc8..132ae90b 100644
--- a/frontend/src/lib/rectification-agentic/v9/agent-run-finish.ts
+++ b/frontend/src/lib/rectification-agentic/v9/agent-run-finish.ts
@@ -22,6 +22,8 @@ import {
} from "./agent-run-support";
import { dropUngroundedFactSentences } from "./spoken-grounding";
import { RECTIFICATION_USER_COPY } from "../user-copy";
+import { moderate, MODERATION_OUTPUT_REPLACED_COPY } from "@/lib/moderation";
+import { recordModerationEvent } from "@/lib/moderation/log";
import type { V9AgentRunResult } from "./agent-run";
import type { V9TimedTurn } from "./agent-run-prepare";
@@ -65,7 +67,17 @@ export async function finishV9AgentTurn(
};
}
- const skipBilling = outcome.settleBilling === false;
+ // Content moderation, output side (2026-09-30): a model answer that fails is
+ // replaced by a notice below and the turn's charge is released, not settled.
+ const moderationContext = { route: "rectification", userId, requestId: turnId };
+ const outputVerdict = outcome.answerText.trim()
+ ? await moderate(outcome.answerText, "output", moderationContext)
+ : null;
+ if (outputVerdict && outputVerdict.verdict !== "pass") {
+ await recordModerationEvent(accounting as never, outputVerdict, "output", moderationContext);
+ }
+ const outputBlocked = outputVerdict?.verdict === "block";
+ const skipBilling = outcome.settleBilling === false || outputBlocked;
if (!skipBilling) {
const durationMs = Date.now() - startedAt;
const completed = await billing.complete({ ...outcome.usage, durationMs });
@@ -186,6 +198,7 @@ export async function finishV9AgentTurn(
if (action === "evidence" && interviewIdle?.terminalNote && interviewIdle.hostNarration) {
spokenAnswer = composeIdleGapIntoSpoken(spokenAnswer, interviewIdle.hostNarration);
}
+ if (outputBlocked) spokenAnswer = MODERATION_OUTPUT_REPLACED_COPY;
if (spokenAnswer !== (outcome.streamedText ?? outcome.answerText)) {
await emit({ type: "answer.delta", text: spokenAnswer, replace: true });
}
diff --git a/frontend/src/lib/stream-agent-response.ts b/frontend/src/lib/stream-agent-response.ts
index adbf6d7b..60102364 100644
--- a/frontend/src/lib/stream-agent-response.ts
+++ b/frontend/src/lib/stream-agent-response.ts
@@ -353,6 +353,13 @@ type StreamAgentResponseOptions = EventOptions & {
headers?: HeadersInit;
onFirstActivity?: () => void | Promise;
onFirstOutput?: () => void | Promise;
+ /**
+ * Content moderation of the finished answer (2026-09-30). Returns the text
+ * that must replace the whole answer, or null to keep it. A replacement is
+ * sent as one `answer.delta` with `replace: true` and is what `onComplete`
+ * receives.
+ */
+ moderateOutput?: (output: string) => Promise;
onComplete?: (
output: string,
receipt: AgentExecutionReceipt,
@@ -1056,6 +1063,11 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
throw new Error("answer_truncated");
}
settling = true;
+ const replacement = options.moderateOutput ? await options.moderateOutput(fullOutput) : null;
+ if (replacement !== null) {
+ fullOutput = replacement;
+ send(controller, { type: "answer.delta", text: replacement, replace: true });
+ }
const receipt = agentExecutionReceiptSchema.parse(options.receipt());
const thinkingSections = applyThinkingSectionProgress(options.state.thinkingPlan ?? [], fullOutput);
// Public thinking text came from the removed interpret pass, which
diff --git a/frontend/supabase/migrations/20260930040000_moderation_events.sql b/frontend/supabase/migrations/20260930040000_moderation_events.sql
new file mode 100644
index 00000000..dc9fc50e
--- /dev/null
+++ b/frontend/supabase/migrations/20260930040000_moderation_events.sql
@@ -0,0 +1,152 @@
+-- Content moderation hits (2026-09-30, compliance round). Add-only.
+-- One row per blocked or flagged input/output. No user or model text is
+-- stored: only categories, rule ids / provider labels, route, side, model and
+-- request ids. Service role only; operators read it through /admin.
+create table if not exists public.moderation_events (
+ id bigint generated always as identity primary key,
+ created_at timestamptz not null default now(),
+ user_id uuid references auth.users(id) on delete set null,
+ side text not null check (side in ('input', 'output')),
+ verdict text not null check (verdict in ('block', 'review')),
+ route text not null check (char_length(route) between 1 and 64),
+ categories text[] not null default '{}',
+ rule_ids text[] not null default '{}',
+ provider text not null check (provider in ('local', 'aliyun')),
+ model_id text,
+ request_id text
+);
+
+create index if not exists moderation_events_created_idx on public.moderation_events (created_at desc);
+create index if not exists moderation_events_user_idx on public.moderation_events (user_id, created_at desc);
+
+alter table public.moderation_events enable row level security;
+revoke all on public.moderation_events from public, anon, authenticated;
+grant select, insert on public.moderation_events to service_role;
+
+-- Output moderation (same round): an answer whose text failed moderation is
+-- replaced by a short notice and the reservation is released, exactly like
+-- the small-talk free completion above it in history, but for the notice
+-- (responseKind 'moderated', ≤ 200 chars). Service role only.
+begin;
+create or replace function public.complete_consultation_moderated(
+ p_user_id uuid,
+ p_request_id text,
+ p_session_id uuid,
+ p_response_message jsonb,
+ p_actual_usage jsonb
+)
+returns table(success boolean, credits integer, error_code text)
+language plpgsql
+security definer
+set search_path = ''
+as $$
+declare
+ v_request public.consultation_requests%rowtype;
+ v_res public.usage_reservations%rowtype;
+ v_settlement record;
+ v_balance integer;
+begin
+ if p_user_id is null or p_session_id is null or btrim(coalesce(p_request_id, '')) = '' then
+ return query select false, null::integer, 'invalid_request'::text;
+ return;
+ end if;
+ if jsonb_typeof(p_response_message) is distinct from 'object'
+ or p_response_message->>'role' is distinct from 'assistant'
+ or p_response_message->>'responseKind' is distinct from 'moderated'
+ or jsonb_typeof(p_response_message->'text') is distinct from 'string'
+ or btrim(coalesce(p_response_message->>'text', '')) = ''
+ or length(p_response_message->>'text') > 200
+ or (p_response_message - array['role', 'text', 'responseKind']) <> '{}'::jsonb then
+ return query select false, null::integer, 'invalid_response_message'::text;
+ return;
+ end if;
+ if jsonb_typeof(p_actual_usage) is distinct from 'object' then
+ return query select false, null::integer, 'invalid_actual_usage'::text;
+ return;
+ end if;
+
+ perform pg_advisory_xact_lock(hashtextextended(p_user_id::text || ':' || btrim(p_request_id), 0));
+ select request.* into v_request
+ from public.consultation_requests as request
+ where request.user_id = p_user_id
+ and request.request_id = btrim(p_request_id)
+ and request.session_id = p_session_id
+ for update;
+ if not found then
+ return query select false, null::integer, 'request_missing'::text;
+ return;
+ end if;
+ if v_request.status = 'cancelled' then
+ return query select false, null::integer, 'request_cancelled'::text;
+ return;
+ end if;
+ if v_request.status = 'completed' then
+ select profile.credits into v_balance from public.profiles as profile where profile.id = p_user_id;
+ return query select coalesce(v_request.response_message = p_response_message, false), v_balance,
+ case when v_request.response_message = p_response_message then null::text else 'response_conflict'::text end;
+ return; -- Never refund a previously completed paid consultation.
+ end if;
+ if v_request.status <> 'reserved' then
+ return query select false, null::integer, 'invalid_request_status'::text;
+ return;
+ end if;
+
+ select reservation.* into v_res from public.usage_reservations as reservation
+ where reservation.user_id = p_user_id and reservation.request_id = btrim(p_request_id)
+ for update;
+ if not found or v_res.status <> 'reserved' or v_res.feature_key <> 'chat.standard' then
+ return query select false, null::integer, 'invalid_reservation'::text;
+ return;
+ end if;
+ select profile.credits into v_balance from public.profiles as profile
+ where profile.id = p_user_id for update;
+ if not found then
+ return query select false, null::integer, 'profile_missing'::text;
+ return;
+ end if;
+
+ update public.chat_sessions as session
+ set messages = session.messages || jsonb_build_array(p_response_message), updated_at = clock_timestamp()
+ where session.id = p_session_id and session.user_id = p_user_id and session.session_type = 'consultation';
+ if not found then
+ return query select false, null::integer, 'session_missing'::text;
+ return;
+ end if;
+
+ -- complete_usage is the sole cost ledger writer; it does not debit credits.
+ -- Keep cost even though the reservation is released below (also frees subscription quota).
+ select * into v_settlement from public.complete_usage(
+ p_user_id, btrim(p_request_id),
+ p_actual_usage || jsonb_build_object('metadata', coalesce(p_actual_usage->'metadata', '{}'::jsonb)
+ || jsonb_build_object('responseKind', 'moderated', 'freeCompletion', true))
+ );
+ if not coalesce(v_settlement.success, false) then
+ raise exception 'consultation_moderated_usage_settlement_failed:%', coalesce(v_settlement.error_code, 'unknown');
+ end if;
+
+ -- Same refund amount, balance lock and transaction identity as release_usage.
+ -- Only reachable from a reserved request and reserved usage row under the shared lock.
+ if v_res.source = 'credits' and v_res.credit_amount > 0 then
+ update public.profiles as profile
+ set credits = profile.credits + v_res.credit_amount, updated_at = clock_timestamp()
+ where profile.id = p_user_id returning profile.credits into v_balance;
+ insert into public.credit_transactions(user_id, transaction_type, amount, balance_after, request_id, model)
+ values(p_user_id, 'refund', v_res.credit_amount, v_balance, v_res.request_id, v_res.requested_model_id);
+ -- No ON CONFLICT: an inconsistent prior refund must roll back everything, never double-credit.
+ end if;
+ update public.usage_reservations
+ set status = 'released', released_at = clock_timestamp(), release_reason = 'consultation_output_moderated'
+ where id = v_res.id;
+ update public.consultation_requests
+ set status = 'completed', response_message = p_response_message, updated_at = clock_timestamp()
+ where user_id = p_user_id and request_id = btrim(p_request_id);
+
+ return query select true, v_balance, null::text;
+end;
+$$;
+
+revoke all on function public.complete_consultation_moderated(uuid, text, uuid, jsonb, jsonb)
+ from public, anon, authenticated;
+grant execute on function public.complete_consultation_moderated(uuid, text, uuid, jsonb, jsonb)
+ to service_role;
+commit;
diff --git a/frontend/tests/admin-contracts.test.ts b/frontend/tests/admin-contracts.test.ts
index 1cc1fa43..4f80e4c3 100644
--- a/frontend/tests/admin-contracts.test.ts
+++ b/frontend/tests/admin-contracts.test.ts
@@ -36,7 +36,7 @@ const adminRootRoute = readFileSync(new URL("../src/app/admin/route.ts", import.
const forbiddenPage = readFileSync(new URL("../src/app/forbidden.tsx", import.meta.url), "utf8");
const nextConfig = readFileSync(new URL("../next.config.ts", import.meta.url), "utf8");
const stagingCaddy = readFileSync(new URL("../../deploy/Caddyfile.staging", import.meta.url), "utf8");
-const readonlyRoutes = ["customers", "credit-transactions", "consultations", "audit-logs"].map((resource) =>
+const readonlyRoutes = ["customers", "credit-transactions", "consultations", "audit-logs", "moderation-events"].map((resource) =>
readFileSync(new URL(`../src/app/api/admin/${resource}/route.ts`, import.meta.url), "utf8"),
);
const packageJson = JSON.parse(readFileSync(new URL("../package.json", import.meta.url), "utf8"));
@@ -187,7 +187,8 @@ test("initial Owner recovery is single-candidate, fail-closed, and independent o
});
test("readonly resources cannot be mutated through Refine access control", () => {
- for (const resource of ["customers", "credit-transactions", "consultations", "audit-logs"]) {
+ // 原值:四个只读资源;新值:加入 moderation-events;原因:2026-09-30 内容审核记录只读页,同一约束覆盖它
+ for (const resource of ["customers", "credit-transactions", "consultations", "audit-logs", "moderation-events"]) {
assert.match(providers, new RegExp(resource.includes("-") ? `"${resource}"` : `${resource}:`));
}
assert.match(providers, /const resourcePermissions/);
diff --git a/frontend/tests/content-moderation.test.ts b/frontend/tests/content-moderation.test.ts
new file mode 100644
index 00000000..97699890
--- /dev/null
+++ b/frontend/tests/content-moderation.test.ts
@@ -0,0 +1,257 @@
+// 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'/);
+});
diff --git a/tests/test_api_server_security.py b/tests/test_api_server_security.py
index b69fcaff..cd41e5ca 100644
--- a/tests/test_api_server_security.py
+++ b/tests/test_api_server_security.py
@@ -1581,6 +1581,10 @@ def test_capability_audit_scans_registry_and_local_sources() -> None:
'admin/feature-pricing',
'admin/model-releases',
'admin/models',
+ # 原值: admin/models 后面直接是 admin/orders
+ # 新值: 中间加入 admin/moderation-events
+ # 原因: 内容审核命中记录的只读后台页(2026-09-30 合规轮)是新的 app 路由
+ 'admin/moderation-events',
'admin/orders',
'admin/packages',
'admin/payments',