ForcePilot/backend/server/routers/channels/content_review_router.py
Kris b7ceeb11f3 feat(webhook/content-review): 新增功能与查询过滤能力
1. 为webhook入站请求合并URL query参数到headers,支持微信公众号等渠道的验签
2. 为内容审核列表接口新增resource_type和trace_id查询过滤支持
2026-07-09 04:21:51 +08:00

407 lines
18 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""内容审核域 RouterCR-01 / CR-02 / CR-03 / CR-STATS-01 / CR-DECISION-BATCH
本 router 实现内容审核域的全部 HTTP 端点,覆盖管理员主动触发的 dry-run
预审核、审核历史查询、审核统计与批量人工决定共 5 个操作。内容审核是
跨渠道的聚合能力,路径中 **不** 包含 ``{channel_type}`` 段;子 router
不自行设置 prefix根前缀 ``/channels`` 由 ``channels_router`` 聚合
router 统一追加。
模板选型:
- CR-01 走模板 A控制面端口路由——端点函数体仅做
``operator = build_operator`` → ``use_cases.content_review.previewContent``
→ ``raiseOnControlFailure`` →
``return {"success": True, "data": serialize_control_data(result.data)}``。
``previewContent`` 内部由 ``ChannelControlService`` 走控制面管道
auth → permission → rate_limit → dispatch → auditdispatch 阶段
调度 ``ContentModerationAdapter.review`` 执行审核并写入审核历史存储
``source=manual_preview``),审计日志经 ``audit_stage`` 写入
``tx_strategy=COOPERATIVE``)。
- CR-STATS-01 走模板 A控制面端口路由——与 CR-01 同模板,经
``ChannelControlService.getReviewStats`` 走控制面管道dispatch 阶段
聚合审核记录计算统计指标,审计日志记录 ``admin_query``(谁查看了
统计指标)。``start_time`` / ``end_time`` 由 DTO ``__post_init__``
校验 ``start_time < end_time``。
- CR-DECISION-BATCH 走模板 A控制面端口路由——经
``ChannelControlService.batchReviewDecision`` 走控制面管道dispatch
阶段使用 ``_executeBatch`` 逐条独立事务Pattern D执行人工决定
覆盖,审计日志记录 ``content_review_decided``,返回部分成功结果
``total`` / ``succeeded`` / ``failed``)。``apply_to_pending_messages=True``
时显式抛 ``NotImplementedError``501不静默忽略。
- CR-02 / CR-03 走模板 B数据面端口路由——端点函数体仅做
构造查询命令 → ``use_cases.content_review_query.<method>(cmd)``
→ ``return {"success": True, "data": dataclass_to_dict(result)}``。
CR-03 额外构造 ``operator``(供 ``NotFoundError.trace_id`` 透传),
CR-02 无此需求(列表查询不抛 ``NotFoundError``)。数据面用例直接
读取 ``ContentReviewRepositoryPort``**不** 经控制面管道,**不** 调
``raiseOnControlFailure``(数据面异常由全局 ``unified_error_handler``
统一映射为 HTTP 响应),**不** 写审计日志(查询操作无副作用)。
鉴权策略:全部端点使用 ``get_admin_user`` 依赖,要求管理员或超级管理员角色。
角色校验由依赖函数完成router 内不做角色判断(规范 §4
路径设计:静态后缀(``/preview`` / ``/history`` / ``/stats`` /
``/history/batch-decision``)先于动态单段路径(``/history/{review_id}``
声明,避免 ``preview`` / ``history`` / ``stats`` / ``batch-decision`` 被
捕获为 ``review_id``(规范 §6.5)。
协议翻译CR-02 的 ``verdict`` 查询参数为 string 类型,需在 router 层
翻译为枚举CR-02 与 CR-STATS-01 的 ``start_time`` / ``end_time`` 查询
参数为 string 类型,需在 router 层翻译为 datetime失败显式抛
``ValidationError``,由全局异常处理器映射为 400 ``VALIDATION_ERROR``。
时间范围 ``start_time < end_time`` 由各 DTO ``__post_init__`` 统一校验。
端点清单对应《内容审核域设计方案》§2.1
- POST /content-review/preview CR-01 preview_content
- GET /content-review/history CR-02 list_review_history
- GET /content-review/history/{review_id} CR-03 get_review_detail
- GET /content-review/stats CR-STATS-01 get_review_stats
- POST /content-review/history/batch-decision CR-DECISION-BATCH batch_review_decision
"""
from __future__ import annotations
from typing import Any, Literal
from fastapi import APIRouter, Depends, Query, Request
from pydantic import BaseModel, ConfigDict, Field
from yuxi.channels.contract.dtos.channel import ChannelType
from yuxi.channels.contract.dtos.content_review import (
BatchReviewDecisionCmd,
ContentReviewDetailQueryCmd,
ContentReviewHistoryQueryCmd,
ContentReviewResourceType,
ContentReviewStatsQuery,
ContentReviewVerdict,
ReviewDecisionItem,
)
from yuxi.channels.contract.errors import ValidationError
from yuxi.storage.postgres.models_business import User
from server.routers.channels import (
Granularity,
build_operator,
dataclass_to_dict,
get_channel_use_cases,
parse_datetime,
raiseOnControlFailure,
serialize_control_data,
)
from server.utils.auth_middleware import get_admin_user
content_review_router = APIRouter(tags=["channels-content-review"])
class ContentReviewPreviewRequest(BaseModel):
"""内容预审核请求体CR-01
frozen=True 保证请求体在端点函数内不被意外修改;字段约束与
``ContentReviewResourceType`` 枚举值由 dispatch 阶段二次校验router
层仅做基本的非空与长度约束。
"""
model_config = ConfigDict(frozen=True)
channel_type: ChannelType = Field(..., description="渠道类型")
account_id: str = Field(..., min_length=1, description="账户 ID")
resource_type: ContentReviewResourceType = Field(
...,
description="资源类型message_text / message_attachment / user_profile",
)
content: str = Field(..., min_length=1, max_length=10000, description="待审核文本")
peer_id: str | None = Field(default=None, description="对端 ID可选用于上下文")
class ReviewDecisionItemRequest(BaseModel):
"""单条审核决定请求体CR-DECISION-BATCH
``decision`` 取值 ``pass`` / ``block``,由 Pydantic ``Literal`` 约束
在入参边界拦截非法值(与 DTO ``ReviewDecisionItem.__post_init__``
二次校验形成纵深防御);``reason`` / ``categories`` 为可选字段,
``categories`` 默认空列表。
"""
model_config = ConfigDict(frozen=True)
review_id: str = Field(..., min_length=1, description="审核记录 ID")
decision: Literal["pass", "block"] = Field(..., description="决定结果pass / block")
reason: str | None = Field(default=None, description="决定原因(可选)")
categories: list[str] = Field(default_factory=list, description="命中分类(可选)")
class BatchReviewDecisionRequest(BaseModel):
"""批量审核决定请求体CR-DECISION-BATCH
``decisions`` 数量限制 1-100由 dispatch 阶段二次校验;
``apply_to_pending_messages`` 为预留字段,本期不实现关联处理。
"""
model_config = ConfigDict(frozen=True)
decisions: list[ReviewDecisionItemRequest] = Field(
..., min_length=1, max_length=100, description="决定条目列表1-100"
)
apply_to_pending_messages: bool = Field(default=False, description="是否关联处理待审消息(本期预留)")
def _parse_verdict(raw: str) -> ContentReviewVerdict:
"""string → ContentReviewVerdict失败抛 ValidationError。
将 CR-02 的 ``verdict`` 查询参数string翻译为 ``ContentReviewVerdict``
枚举;非法值由全局异常处理器映射为 400 ``VALIDATION_ERROR``。
"""
try:
return ContentReviewVerdict(raw)
except ValueError as exc:
raise ValidationError(
"verdict",
f"unsupported verdict: {raw} (expected: pass | review | block)",
) from exc
def _parse_resource_type(raw: str) -> ContentReviewResourceType:
"""string → ContentReviewResourceType失败抛 ValidationError。
将 CR-02 的 ``resource_type`` 查询参数string翻译为
``ContentReviewResourceType`` 枚举;非法值由全局异常处理器映射为
400 ``VALIDATION_ERROR``。
"""
try:
return ContentReviewResourceType(raw)
except ValueError as exc:
raise ValidationError(
"resource_type",
f"unsupported resource_type: {raw} "
"(expected: message_text | message_attachment | user_profile)",
) from exc
# ---------------- 静态路径端点(须先于动态路径声明) ----------------
@content_review_router.post(
"/content-review/preview",
response_model=dict,
)
async def preview_content(
payload: ContentReviewPreviewRequest,
request: Request,
use_cases=Depends(get_channel_use_cases),
current_user: User = Depends(get_admin_user),
) -> dict[str, Any]:
"""内容预审核dry-runFR-30CR-01
对应控制面操作 ``content_review/preview``(由
``ChannelControlService.previewContent`` 内部构造 ``ControlCmd`` 并委托
``_executeControl`` 执行控制面管道。dry-run 语义:不触发消息投递、
不触发告警、不修改消息状态、不触发任何 webhook 回调;仅记录审核结果
``source=manual_preview``)与审计日志(``audit_type=
content_review_previewed``)。
编排链路HTTP 入参 → ``build_operator`` 构造操作人 →
``use_cases.content_review.previewContent`` 调用端口方法 → 控制面管道
调度 ``ContentModerationAdapter.review`` 获取审核结论 → 生成
``review_id`` 并写入 ``ContentReviewRepositoryPort`` → ``ControlResult``
经 ``raiseOnControlFailure`` 转译失败 → ``serialize_control_data``
序列化审核结果返回。
"""
operator = build_operator(current_user, request)
result = await use_cases.content_review.previewContent(
channel_type=payload.channel_type,
account_id=payload.account_id,
resource_type=payload.resource_type,
content=payload.content,
peer_id=payload.peer_id,
operator=operator,
)
raiseOnControlFailure(result)
return {"success": True, "data": serialize_control_data(result.data)}
@content_review_router.get(
"/content-review/history",
response_model=dict,
)
async def list_review_history(
channel_type: ChannelType | None = Query(default=None, description="按渠道类型过滤"),
account_id: str | None = Query(default=None, description="按账户 ID 过滤"),
verdict: str | None = Query(
default=None,
description="按审核结论过滤pass / review / block",
),
resource_type: str | None = Query(
default=None,
description="按资源类型过滤message_text / message_attachment / user_profile",
),
trace_id: str | None = Query(
default=None,
description="按链路追踪 ID 过滤",
),
start_time: str | None = Query(
default=None,
description="起始时间过滤ISO 8601",
),
end_time: str | None = Query(
default=None,
description="结束时间过滤ISO 8601",
),
limit: int = Query(default=20, ge=1, le=100, description="每页数量1-100"),
offset: int = Query(default=0, ge=0, description="分页偏移量"),
use_cases=Depends(get_channel_use_cases),
current_user: User = Depends(get_admin_user),
) -> dict[str, Any]:
"""查询审核历史列表FR-30CR-02
非控制面管道路径(模板 B由 ``ContentReviewQueryService`` 直接读取
``ContentReviewRepositoryPort``,不经控制面管道,不写审计日志(查询
操作无副作用)。``verdict`` / ``resource_type`` / ``start_time`` /
``end_time`` 查询参数经协议翻译辅助函数转换为枚举 / datetime失败抛
``ValidationError``。
编排链路HTTP 入参 → 协议翻译(``_parse_verdict`` /
``_parse_resource_type`` / ``parse_datetime``)→ 构造
``ContentReviewHistoryQueryCmd``
DTO ``__post_init__`` 校验 limit/offset/time_range
``use_cases.content_review_query.listHistory`` →
``dataclass_to_dict`` 序列化 ``ContentReviewHistoryList`` 返回。
数据面查询无副作用,不写审计日志,故不构造 ``operator``。
"""
cmd = ContentReviewHistoryQueryCmd(
channel_type=channel_type,
account_id=account_id,
verdict=_parse_verdict(verdict) if verdict is not None else None,
resource_type=_parse_resource_type(resource_type) if resource_type is not None else None,
trace_id=trace_id,
start_time=parse_datetime("start_time", start_time),
end_time=parse_datetime("end_time", end_time),
limit=limit,
offset=offset,
)
result = await use_cases.content_review_query.listHistory(cmd)
return {"success": True, "data": dataclass_to_dict(result)}
@content_review_router.get(
"/content-review/stats",
response_model=dict,
)
async def get_review_stats(
request: Request,
channel_type: ChannelType | None = Query(default=None, description="按渠道类型过滤"),
account_id: str | None = Query(default=None, description="按账户 ID 过滤"),
start_time: str | None = Query(
default=None,
description="起始时间过滤ISO 8601",
),
end_time: str | None = Query(
default=None,
description="结束时间过滤ISO 8601",
),
granularity: Granularity = Query(default="day", description="时间粒度"),
use_cases=Depends(get_channel_use_cases),
current_user: User = Depends(get_admin_user),
) -> dict[str, Any]:
"""审核统计FR-30CR-STATS-01
对应控制面操作 ``content_review/stats``(由
``ChannelControlService.getReviewStats`` 内部构造 ``ControlCmd`` 并委托
``_executeControl`` 执行控制面管道)。聚合审核记录,计算通过/拦截/待复核
计数、通过率、拦截率、人工介入率、平均决策时长与按分类切片的统计及
时间桶趋势。``start_time`` / ``end_time`` 可选,同时提供时需满足
``start_time < end_time``。
编排链路HTTP 入参 → ``build_operator`` 构造操作人 → 协议翻译
``parse_datetime``)→ 构造 ``ContentReviewStatsQuery`` →
``use_cases.content_review.getReviewStats`` → ``raiseOnControlFailure``
转译失败 → ``serialize_control_data`` 序列化统计结果返回。
"""
operator = build_operator(current_user, request)
query = ContentReviewStatsQuery(
channel_type=channel_type,
account_id=account_id,
start_time=parse_datetime("start_time", start_time),
end_time=parse_datetime("end_time", end_time),
granularity=granularity,
)
result = await use_cases.content_review.getReviewStats(query=query, operator=operator)
raiseOnControlFailure(result)
return {"success": True, "data": serialize_control_data(result.data)}
@content_review_router.post(
"/content-review/history/batch-decision",
response_model=dict,
)
async def batch_review_decision(
payload: BatchReviewDecisionRequest,
request: Request,
use_cases=Depends(get_channel_use_cases),
current_user: User = Depends(get_admin_user),
) -> dict[str, Any]:
"""批量审核决定FR-30CR-DECISION-BATCH
对应控制面操作 ``content_review/batch_decision``(由
``ChannelControlService.batchReviewDecision`` 内部构造 ``ControlCmd`` 并委托
``_executeControl`` 执行控制面管道。dispatch 阶段使用 ``_executeBatch``
逐条独立事务Pattern D对审核记录进行人工决定覆盖``reviewer`` 记录
为 ``operator.user_id``,单条失败不回滚已成功条目,返回部分成功结果
``total`` / ``succeeded`` / ``failed``)。
编排链路HTTP 入参 → ``build_operator`` 构造操作人 → 构造
``BatchReviewDecisionCmd``DTO 含 ``ReviewDecisionItem`` 元组)→
``use_cases.content_review.batchReviewDecision`` →
``raiseOnControlFailure`` 转译失败 → ``serialize_control_data`` 序列化
部分成功结果返回。
"""
operator = build_operator(current_user, request)
cmd = BatchReviewDecisionCmd(
decisions=tuple(
ReviewDecisionItem(
review_id=item.review_id,
decision=item.decision,
reason=item.reason,
categories=tuple(item.categories),
)
for item in payload.decisions
),
operator=operator,
apply_to_pending_messages=payload.apply_to_pending_messages,
)
result = await use_cases.content_review.batchReviewDecision(cmd=cmd)
raiseOnControlFailure(result)
return {"success": True, "data": serialize_control_data(result.data)}
# ---------------- 动态路径端点 ----------------
@content_review_router.get(
"/content-review/history/{review_id}",
response_model=dict,
)
async def get_review_detail(
review_id: str,
request: Request,
use_cases=Depends(get_channel_use_cases),
current_user: User = Depends(get_admin_user),
) -> dict[str, Any]:
"""查询单次审核详情FR-30CR-03
非控制面管道路径(模板 B由 ``ContentReviewQueryService`` 直接读取
``ContentReviewRepositoryPort``,返回完整审核记录含命中片段。
``review_id`` 不存在时由 ``ContentReviewQueryService.getDetail`` 显式
转换为 ``NotFoundError``INV-7 错误显式化),由全局异常处理器映射为
404 ``NOT_FOUND``。
编排链路HTTP 入参 → ``build_operator`` 构造操作人 → 构造
``ContentReviewDetailQueryCmd``DTO ``__post_init__`` 校验 review_id
非空)→ ``use_cases.content_review_query.getDetail`` →
``dataclass_to_dict`` 序列化 ``ContentReviewDetail`` 返回。
"""
operator = build_operator(current_user, request)
cmd = ContentReviewDetailQueryCmd(
review_id=review_id,
operator=operator,
)
result = await use_cases.content_review_query.getDetail(cmd)
return {"success": True, "data": dataclass_to_dict(result)}