From 627915f1147b15f922d5ef2a407ea2e23257ec92 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Mon, 6 Jul 2026 20:50:23 +0800 Subject: [PATCH] =?UTF-8?q?refactor(channels):=20=E6=9B=BF=E6=8D=A2?= =?UTF-8?q?=E6=8A=A5=E5=91=8A=E8=B7=AF=E7=94=B1=E4=B8=BA=E8=B7=AF=E7=94=B1?= =?UTF-8?q?=E7=BB=91=E5=AE=9A=E7=AE=A1=E7=90=86=E8=B7=AF=E7=94=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 删除原有的reports_router.py报告路由文件,新增route_binding_router.py实现渠道-Agent路由绑定的CRUD和测试功能,统一使用模板A处理鉴权、调用和响应序列化。 --- .../server/routers/channels/reports_router.py | 195 ------------- .../routers/channels/route_binding_router.py | 263 ++++++++++++++++++ 2 files changed, 263 insertions(+), 195 deletions(-) delete mode 100644 backend/server/routers/channels/reports_router.py create mode 100644 backend/server/routers/channels/route_binding_router.py diff --git a/backend/server/routers/channels/reports_router.py b/backend/server/routers/channels/reports_router.py deleted file mode 100644 index 5ee68aa6..00000000 --- a/backend/server/routers/channels/reports_router.py +++ /dev/null @@ -1,195 +0,0 @@ -"""Reports 聚合视图域 Router(RPT-05 / RPT-06 / RPT-ONEOFF-01 / RPT-ONEOFF-02)。 - -统一采用模板 A(控制面端口路由),鉴权依赖 get_admin_user,调用 -use_cases.report_management.,raiseOnControlFailure 转译失败, -serialize_control_data 序列化响应。 - -路径设计(静态先于动态,避免 /{report_id} 误捕获 /oneoff): - - POST /reports/oneoff RPT-05 createOneoffReport - - GET /reports/oneoff RPT-ONEOFF-01 listOneoffReports - - GET /reports/oneoff/{task_id}/download RPT-ONEOFF-02 downloadReport - - POST /reports/oneoff/{task_id}/retry RPT-ONEOFF-RETRY retryReport - - GET /reports/{report_id} RPT-06 getReport - -子 router 自身不设置 prefix,根前缀 ``/channels`` 由 ``channels_router`` -聚合 router 统一追加。 -""" - -from __future__ import annotations - -from datetime import datetime, timezone -from typing import Any, Literal - -from fastapi import APIRouter, Depends, Query, Request -from pydantic import BaseModel, ConfigDict, Field -from yuxi.channels.contract.dtos.report import ( - DEFAULT_LIST_LIMIT, - MAX_LIST_LIMIT, - MIN_LIST_LIMIT, -) -from yuxi.channels.contract.errors import ValidationError -from yuxi.storage.postgres.models_business import User - -from server.routers.channels import ( - build_operator, - get_channel_use_cases, - parse_datetime, - raiseOnControlFailure, - serialize_control_data, -) -from server.utils.auth_middleware import get_admin_user - -reports_router = APIRouter(tags=["channels-reports"]) - - -class CreateOneoffReportRequest(BaseModel): - """创建一次性报告请求体(RPT-05)。 - - 字段: - report_type: 报告类型枚举(5 个值)。 - params: 报告生成参数(时间范围 / 过滤条件等)。 - run_at: 计划执行时间(ISO 8601 字符串),为 None 时立即执行。 - """ - - model_config = ConfigDict(frozen=True) - - report_type: Literal[ - "message_stats", - "session_stats", - "account_stats", - "delivery_stats", - "dashboard_overview", - ] = Field(..., description="报告类型") - params: dict[str, Any] = Field(default_factory=dict, description="报告生成参数") - run_at: str | None = Field(default=None, description="计划执行时间(ISO 8601),为空时立即执行") - - -@reports_router.post("/reports/oneoff", response_model=dict) -async def create_oneoff_report( - request: Request, - body: CreateOneoffReportRequest, - use_cases=Depends(get_channel_use_cases), - current_user: User = Depends(get_admin_user), -) -> dict[str, Any]: - """创建一次性报告(RPT-05)。对应控制面操作 reports/create_oneoff。 - - 异步生成报告:创建 pending 状态报告记录 + scheduler 一次性任务, - worker 进程通过 ChannelReportHandler 生成内容后状态变 ready。 - """ - operator = build_operator(current_user, request) - # 时间参数统一在 Router 层校验格式与业务约束,与 list_oneoff_reports - # 的 start_time / end_time 处理一致(FR-34 / RPT-05 @pre): - # - 格式校验:parse_datetime 翻译为 aware UTC datetime,非法格式抛 - # ValidationError(400),避免原生异常在控制面管道内被兜底为 500。 - # - 过去时间校验:run_at 非空时不得为过去时间(Port 契约 @pre), - # 避免创建永远不会触发的 scheduler 任务。 - # - datetime 直接透传至 UseCase(Port 契约要求 datetime | None), - # 消除 str → datetime → str → datetime 的冗余往返转换。 - run_at_dt = parse_datetime("run_at", body.run_at) - if run_at_dt is not None and run_at_dt <= datetime.now(timezone.utc): - raise ValidationError( - field="run_at", - message="run_at must be a future time", - ) - result = await use_cases.report_management.createOneoffReport( - report_type=body.report_type, - params=body.params, - run_at=run_at_dt, - operator=operator, - ) - raiseOnControlFailure(result) - return {"success": True, "data": serialize_control_data(result.data)} - - -@reports_router.get("/reports/oneoff", response_model=dict) -async def list_oneoff_reports( - request: Request, - status: Literal["pending", "generating", "ready", "failed"] | None = Query(default=None, description="按状态过滤"), - report_type: Literal["message_stats", "session_stats", "account_stats", "delivery_stats", "dashboard_overview"] | None = Query(default=None, description="按报告类型过滤"), - start_time: str | None = Query(default=None, description="起始时间(ISO 8601)"), - end_time: str | None = Query(default=None, description="截止时间(ISO 8601)"), - limit: int = Query(default=DEFAULT_LIST_LIMIT, ge=MIN_LIST_LIMIT, le=MAX_LIST_LIMIT, description="分页大小(1-200)"), - 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]: - """列出一次性报告(RPT-ONEOFF-01)。对应控制面操作 reports/list_oneoff。 - - 分页查询报告列表,支持按状态 / 类型 / 时间范围过滤。 - """ - operator = build_operator(current_user, request) - start_time_dt = parse_datetime("start_time", start_time) - end_time_dt = parse_datetime("end_time", end_time) - result = await use_cases.report_management.listOneoffReports( - operator=operator, - status=status, - report_type=report_type, - start_time=start_time_dt, - end_time=end_time_dt, - limit=limit, - offset=offset, - ) - raiseOnControlFailure(result) - return {"success": True, "data": serialize_control_data(result.data)} - - -@reports_router.get("/reports/oneoff/{task_id}/download", response_model=dict) -async def download_oneoff_report( - request: Request, - task_id: str, - use_cases=Depends(get_channel_use_cases), - current_user: User = Depends(get_admin_user), -) -> dict[str, Any]: - """下载一次性报告(RPT-ONEOFF-02)。对应控制面操作 reports/download。 - - 根据 task_id 下载已就绪报告内容。未就绪(pending/generating)返回 409, - 失败(failed)返回 409。 - """ - operator = build_operator(current_user, request) - result = await use_cases.report_management.downloadReport( - task_id=task_id, - operator=operator, - ) - raiseOnControlFailure(result) - return {"success": True, "data": serialize_control_data(result.data)} - - -@reports_router.post("/reports/oneoff/{task_id}/retry", response_model=dict) -async def retry_oneoff_report( - request: Request, - task_id: str, - use_cases=Depends(get_channel_use_cases), - current_user: User = Depends(get_admin_user), -) -> dict[str, Any]: - """重试失败报告(RPT-ONEOFF-RETRY)。对应控制面操作 reports/oneoff_retry。 - - 校验原报告 status=="failed" 后创建新 pending 报告并入队 scheduler - 任务,同时标记原报告 ``retried_at``。返回新旧 task_id。 - """ - operator = build_operator(current_user, request) - result = await use_cases.report_management.retryReport( - task_id=task_id, - operator=operator, - ) - raiseOnControlFailure(result) - return {"success": True, "data": serialize_control_data(result.data)} - - -@reports_router.get("/reports/{report_id}", response_model=dict) -async def get_report( - request: Request, - report_id: str, - use_cases=Depends(get_channel_use_cases), - current_user: User = Depends(get_admin_user), -) -> dict[str, Any]: - """查询报告详情(RPT-06)。对应控制面操作 reports/get。 - - 返回报告完整信息(含 content / status / download_url)。 - """ - operator = build_operator(current_user, request) - result = await use_cases.report_management.getReport( - report_id=report_id, - operator=operator, - ) - raiseOnControlFailure(result) - return {"success": True, "data": serialize_control_data(result.data)} \ No newline at end of file diff --git a/backend/server/routers/channels/route_binding_router.py b/backend/server/routers/channels/route_binding_router.py new file mode 100644 index 00000000..72b1a978 --- /dev/null +++ b/backend/server/routers/channels/route_binding_router.py @@ -0,0 +1,263 @@ +"""路由绑定管理 Router(RB-01~RB-05)。 + +提供渠道-Agent 路由绑定规则的 5 个 HTTP 端点,全部由 ``get_admin_user`` 守门, +仅管理员可访问。Router 不含任何业务逻辑,仅做协议翻译:将 HTTP 请求参数组装 +为契约层命令/参数,调用 ``use_cases.route_binding_management`` 的类型化方法走 +控制面管道,再通过 ``raiseOnControlFailure`` 转译失败、``serialize_control_data`` +序列化成功结果(模板 A)。 + +鉴权策略:全部端点使用 ``get_admin_user`` 依赖,要求管理员或超级管理员角色。 +角色校验由依赖函数完成,router 内不做角色判断(规范 §4)。 + +路径设计:``/route-bindings`` 与 ``/route-bindings/resolve`` 静态路径先于 +``/route-bindings/{binding_id}`` 动态路径声明(规范 §6.5),避免 ``resolve`` +被捕获为 binding_id。子 router 不自行设置 prefix,根前缀 ``/channels`` 由 +``channels_router`` 聚合 router 统一追加。 + +端点清单: + - GET /route-bindings RB-01 list 列出路由绑定规则 + - POST /route-bindings RB-02 create 创建路由绑定规则 + - POST /route-bindings/resolve RB-05 resolve 测试路由解析 + - PUT /route-bindings/{binding_id} RB-03 update 更新路由绑定规则 + - DELETE /route-bindings/{binding_id} RB-04 delete 删除路由绑定规则 +""" + +from __future__ import annotations + +from typing import Any + +from fastapi import APIRouter, Depends, Query, Request +from pydantic import BaseModel, ConfigDict, Field, model_validator +from yuxi.channels.contract.dtos.channel import ChannelType +from yuxi.channels.contract.dtos.route import ( + BindingContext, + RouteBindingFilter, + SaveRouteBindingCmd, + UpdateRouteBindingCmd, +) +from yuxi.channels.contract.errors import ValidationError +from yuxi.storage.postgres.models_business import User + +from server.routers.channels import ( + build_operator, + get_channel_use_cases, + raiseOnControlFailure, + serialize_control_data, +) +from server.utils.auth_middleware import get_admin_user + +route_binding_router = APIRouter(tags=["channels-route-binding"]) + + +# ---------------- Request Schemas ---------------- + + +class CreateRouteBindingRequest(BaseModel): + """创建路由绑定规则请求。字段对齐 ``SaveRouteBindingCmd``。 + + ``match_value`` 与 ``match_source`` 的组合约束由 ``_validate_match_value`` + 校验:``account`` / ``default`` tier 不允许携带 match_value;``chat_type`` + tier 的 match_value 必须为 ``p2p`` 或 ``group``。同时拒绝当前运行时尚未 + 实现 matcher 的实验性 tier(``channel_session`` / ``channel_type``)。 + """ + + model_config = ConfigDict(frozen=True) + + channel_type: ChannelType = Field(..., description="渠道类型(必填)") + account_id: str = Field(..., min_length=1, description="渠道账户 ID") + match_source: str = Field(..., min_length=1, description="匹配来源(tier 名)") + match_value: str | None = Field(default=None, description="匹配值(按 tier 约束)") + agent_binding: str = Field(..., min_length=1, description="Agent 绑定(agent slug)") + description: str | None = Field(default=None, description="描述") + + #: 当前运行时已实现 matcher 的 tier 集合;其余 tier 创建规则会被拒绝。 + #: 与 factory._registerDefaultMatchers / route_match_registry.applyRule 保持一致。 + _SUPPORTED_MATCH_SOURCES: frozenset[str] = frozenset( + {"session_key", "identity_id", "peer_id", "chat_type", "account", "default"} + ) + + @model_validator(mode="after") + def _validate_match_value(self) -> CreateRouteBindingRequest: + """校验 match_source 与 match_value 的组合约束。 + + - 不支持 ``channel_session`` / ``channel_type`` 等实验性 tier, + 运行时尚未注册 matcher,创建规则会导致永久失效。 + - ``account`` / ``default`` tier:match_value 必须为空。 + - ``chat_type`` tier:match_value 必须为 ``p2p`` 或 ``group``。 + """ + if self.match_source not in self._SUPPORTED_MATCH_SOURCES: + raise ValidationError( + "match_source", + f"match_source '{self.match_source}' is not supported by runtime matcher", + ) + if self.match_source in ("account", "default"): + if self.match_value: + raise ValidationError( + "match_value", + f"match_value must be empty for match_source={self.match_source}", + ) + elif self.match_source == "chat_type": + if self.match_value not in ("p2p", "group"): + raise ValidationError( + "match_value", + "match_value must be 'p2p' or 'group' for match_source=chat_type", + ) + return self + + +class UpdateRouteBindingRequest(BaseModel): + """更新路由绑定规则请求。所有字段可选,仅传入字段被更新。 + + 字段语义:未提供(不在请求体中)表示不修改,由 ``UpdateRouteBindingCmd`` + 的 ``None`` 默认值表示"不更新"。``binding_id`` 由路径参数提供,不在请求体中。 + """ + + model_config = ConfigDict(frozen=True) + + match_value: str | None = Field(default=None, description="匹配值") + agent_binding: str | None = Field(default=None, description="Agent 绑定(agent slug)") + enabled: bool | None = Field(default=None, description="是否启用") + description: str | None = Field(default=None, description="描述") + + +class ResolveRouteBindingRequest(BaseModel): + """测试路由解析请求。字段对齐 ``BindingContext``。""" + + model_config = ConfigDict(frozen=True) + + session_key: str = Field(..., min_length=1, description="会话键") + channel_type: ChannelType = Field(..., description="渠道类型") + account_id: str = Field(..., min_length=1, description="渠道账户 ID") + chat_type: str = Field(..., description="会话类型(p2p | group)") + peer_id: str = Field(..., min_length=1, description="对端 ID") + unified_identity_id: str | None = Field(default=None, description="统一身份 ID") + conversation_id: str | None = Field(default=None, description="内部会话 ID") + + +# ---------------- Endpoints ---------------- + + +@route_binding_router.get("/route-bindings", response_model=dict) +async def list_route_bindings( + request: Request, + channel_type: ChannelType | None = Query(default=None, description="按渠道类型过滤"), + account_id: str | None = Query(default=None, description="按账户 ID 过滤"), + match_source: str | None = Query(default=None, description="按匹配来源过滤"), + enabled: bool | None = Query(default=None, description="按启用状态过滤"), + agent_binding: str | None = Query(default=None, description="按绑定 Agent 过滤"), + limit: int = Query(default=20, ge=1, le=1000, description="分页大小"), + 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]: + """列出路由绑定规则(RB-01)。 + + 支持按渠道类型 / 账户 / 匹配来源 / 启用状态 / 绑定 Agent 过滤,分页返回。 + 对应控制面操作 ``route_binding/list``。 + """ + operator = build_operator(current_user, request) + filter = RouteBindingFilter( + channel_type=channel_type, + account_id=account_id, + match_source=match_source, + enabled=enabled, + agent_binding=agent_binding, + ) + result = await use_cases.route_binding_management.listRouteBindings(operator, filter, limit=limit, offset=offset) + raiseOnControlFailure(result) + return {"success": True, "data": serialize_control_data(result.data)} + + +@route_binding_router.post("/route-bindings", response_model=dict) +async def create_route_binding( + payload: CreateRouteBindingRequest, + request: Request, + use_cases=Depends(get_channel_use_cases), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """创建路由绑定规则(RB-02)。 + + 对应控制面操作 ``route_binding/create``。 + """ + operator = build_operator(current_user, request) + cmd = SaveRouteBindingCmd( + channel_type=payload.channel_type, + account_id=payload.account_id, + match_source=payload.match_source, + match_value=payload.match_value, + agent_binding=payload.agent_binding, + description=payload.description, + ) + result = await use_cases.route_binding_management.createRouteBinding(operator, cmd) + raiseOnControlFailure(result) + return {"success": True, "data": serialize_control_data(result.data)} + + +@route_binding_router.post("/route-bindings/resolve", response_model=dict) +async def resolve_route_binding( + payload: ResolveRouteBindingRequest, + request: Request, + use_cases=Depends(get_channel_use_cases), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """测试路由解析(RB-05)。 + + 传入 ``BindingContext`` 字段,返回匹配到的路由绑定规则(含匹配来源、 + 命中层级等诊断元数据)。对应控制面操作 ``route_binding/resolve``。 + """ + operator = build_operator(current_user, request) + ctx = BindingContext( + session_key=payload.session_key, + channel_type=payload.channel_type, + account_id=payload.account_id, + chat_type=payload.chat_type, + peer_id=payload.peer_id, + unified_identity_id=payload.unified_identity_id, + conversation_id=payload.conversation_id, + ) + result = await use_cases.route_binding_management.resolveRouteBinding(operator, ctx) + raiseOnControlFailure(result) + return {"success": True, "data": serialize_control_data(result.data)} + + +@route_binding_router.put("/route-bindings/{binding_id}", response_model=dict) +async def update_route_binding( + binding_id: str, + payload: UpdateRouteBindingRequest, + request: Request, + use_cases=Depends(get_channel_use_cases), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """更新路由绑定规则(RB-03)。 + + ``binding_id`` 由路径参数提供,请求体字段全部可选,仅传入字段被更新。 + 对应控制面操作 ``route_binding/update``。 + """ + operator = build_operator(current_user, request) + cmd = UpdateRouteBindingCmd( + binding_id=binding_id, + match_value=payload.match_value, + agent_binding=payload.agent_binding, + enabled=payload.enabled, + description=payload.description, + ) + result = await use_cases.route_binding_management.updateRouteBinding(operator, cmd) + raiseOnControlFailure(result) + return {"success": True, "data": serialize_control_data(result.data)} + + +@route_binding_router.delete("/route-bindings/{binding_id}", response_model=dict) +async def delete_route_binding( + binding_id: str, + request: Request, + use_cases=Depends(get_channel_use_cases), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """删除路由绑定规则(RB-04)。 + + 对应控制面操作 ``route_binding/delete``。 + """ + operator = build_operator(current_user, request) + result = await use_cases.route_binding_management.deleteRouteBinding(operator, binding_id) + raiseOnControlFailure(result) + return {"success": True, "data": serialize_control_data(result.data)}