refactor(channels): 替换报告路由为路由绑定管理路由

删除原有的reports_router.py报告路由文件,新增route_binding_router.py实现渠道-Agent路由绑定的CRUD和测试功能,统一使用模板A处理鉴权、调用和响应序列化。
This commit is contained in:
Kris 2026-07-06 20:50:23 +08:00
parent bba1775220
commit 627915f114
2 changed files with 263 additions and 195 deletions

View File

@ -1,195 +0,0 @@
"""Reports 聚合视图域 RouterRPT-05 / RPT-06 / RPT-ONEOFF-01 / RPT-ONEOFF-02
统一采用模板 A控制面端口路由鉴权依赖 get_admin_user调用
use_cases.report_management.<method>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非法格式抛
# ValidationError400避免原生异常在控制面管道内被兜底为 500。
# - 过去时间校验run_at 非空时不得为过去时间Port 契约 @pre
# 避免创建永远不会触发的 scheduler 任务。
# - datetime 直接透传至 UseCasePort 契约要求 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)}

View File

@ -0,0 +1,263 @@
"""路由绑定管理 RouterRB-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`` tiermatch_value 必须为空
- ``chat_type`` tiermatch_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)}