From 4f3d8d7c62e016875e506be6507c2235d842689e Mon Sep 17 00:00:00 2001 From: Jesse_Chen Date: Wed, 30 Sep 2026 09:13:40 +0800 Subject: [PATCH] =?UTF-8?q?feat(compliance):=20content=20moderation=20on?= =?UTF-8?q?=20model=20input=20and=20output=20=E2=80=94=20local=20lexicon,?= =?UTF-8?q?=20Aliyun=20adapter,=20free=20completion,=20admin=20log?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01N4f2nya58RoRu4yEmJgRGE --- BLOCKED.md | 7 + CHANGELOG.md | 8 + .../PROGRESS-content-moderation-20260930.md | 40 +++ frontend/DESIGN.md | 6 + frontend/docs/VOICE.md | 7 + .../src/app/admin/moderation-events/page.tsx | 2 + .../app/api/admin/moderation-events/route.ts | 101 +++++++ .../src/app/api/birth-time-guide/route.ts | 17 ++ frontend/src/app/api/consult/route.ts | 50 ++++ frontend/src/components/admin/admin-app.tsx | 1 + .../admin/moderation-events-resource.tsx | 71 +++++ frontend/src/hooks/use-consultation-run.ts | 3 +- frontend/src/lib/admin/providers.ts | 1 + frontend/src/lib/consultation-agent-events.ts | 3 +- frontend/src/lib/consultation-run-timeline.ts | 2 +- .../src/lib/moderation/aliyun-provider.ts | 109 ++++++++ frontend/src/lib/moderation/index.ts | 82 ++++++ frontend/src/lib/moderation/lexicon.ts | 84 ++++++ frontend/src/lib/moderation/local-provider.ts | 27 ++ frontend/src/lib/moderation/log.ts | 39 +++ frontend/src/lib/moderation/types.ts | 29 ++ .../v9/agent-route-request.ts | 15 + .../v9/agent-run-finish.ts | 15 +- frontend/src/lib/stream-agent-response.ts | 12 + .../20260930040000_moderation_events.sql | 152 +++++++++++ frontend/tests/admin-contracts.test.ts | 5 +- frontend/tests/content-moderation.test.ts | 257 ++++++++++++++++++ tests/test_api_server_security.py | 4 + 28 files changed, 1143 insertions(+), 6 deletions(-) create mode 100644 docs/tasks/PROGRESS-content-moderation-20260930.md create mode 100644 frontend/src/app/admin/moderation-events/page.tsx create mode 100644 frontend/src/app/api/admin/moderation-events/route.ts create mode 100644 frontend/src/components/admin/moderation-events-resource.tsx create mode 100644 frontend/src/lib/moderation/aliyun-provider.ts create mode 100644 frontend/src/lib/moderation/index.ts create mode 100644 frontend/src/lib/moderation/lexicon.ts create mode 100644 frontend/src/lib/moderation/local-provider.ts create mode 100644 frontend/src/lib/moderation/log.ts create mode 100644 frontend/src/lib/moderation/types.ts create mode 100644 frontend/supabase/migrations/20260930040000_moderation_events.sql create mode 100644 frontend/tests/content-moderation.test.ts 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',