1. 移除多个导出接口的显式response_model声明 2. 调整access_rule和test_case的创建接口位置,修复静态路径冲突 3. 优化适配器配置校验的异常处理逻辑 4. 重构集成路由的查询逻辑,统一使用get_integration_or_raise 5. 新增channels路由组下的capability、reports、dashboard、webhook、wizard、doctor、directory、session共8个子路由模块 6. 注册channels_router到全局路由列表
298 lines
12 KiB
Python
298 lines
12 KiB
Python
"""会话资源域 Router。
|
||
|
||
提供渠道会话管理的 6 个 HTTP 端点,全部由 ``get_admin_user`` 守门,仅管理员
|
||
可访问。Router 不含任何业务逻辑,仅做协议翻译:将 HTTP 请求参数组装为契约层
|
||
命令/参数,调用 ``use_cases.session_management`` 的类型化方法走控制面管道,
|
||
再通过 ``raiseOnControlFailure`` 转译失败、``serialize_control_data`` 序列化
|
||
成功结果(模板 A)。
|
||
|
||
端点清单:
|
||
- GET /sessions SES-QUERY-01 列出渠道会话
|
||
- POST /sessions/batch-close SES-BATCH-CLOSE-01 批量关闭会话
|
||
- GET /sessions/{session_id} SES-QUERY-02 查询会话详情
|
||
- POST /sessions/{session_id}/merge SES-05 合并会话(FR-07)
|
||
- POST /sessions/{session_id}/transfer SES-13 转移会话所有者(FR-26)
|
||
- POST /sessions/{session_id}/close SES-CLOSE-01 关闭会话
|
||
- GET /sessions/{session_id}/messages SES-MSG-01 查询会话消息流
|
||
- GET /sessions/{session_id}/stats SES-STATS-01 查询会话统计指标
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from typing import Any
|
||
|
||
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.conversation import MergeConversationCmd
|
||
from yuxi.channels.contract.dtos.session import (
|
||
BatchCloseSessionsCmd,
|
||
CloseSessionCmd,
|
||
OwnerTransferCmd,
|
||
)
|
||
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
|
||
|
||
session_router = APIRouter(tags=["channels-session"])
|
||
|
||
|
||
class MergeSessionRequest(BaseModel):
|
||
"""合并会话请求体。``target_conversation_id`` 由路径参数提供,不在 body 中。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
source_conversation_id: str = Field(..., min_length=1, description="源会话 ID(被合并方)")
|
||
reason: str = Field(default="", max_length=512, description="合并原因(审计用)")
|
||
|
||
|
||
class TransferOwnerRequest(BaseModel):
|
||
"""转移所有者请求体。``conversation_id`` 由路径参数提供,不在 body 中。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
new_owner_id: str = Field(..., min_length=1, description="新所有者对端 ID")
|
||
|
||
|
||
class CloseSessionRequest(BaseModel):
|
||
"""关闭会话请求体。``session_id`` 由路径参数提供,不在 body 中。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
reason: str | None = Field(default=None, max_length=512, description="关闭原因(审计用)")
|
||
|
||
|
||
class BatchCloseSessionsRequest(BaseModel):
|
||
"""批量关闭会话请求体(SES-BATCH-CLOSE-01)。
|
||
|
||
``session_ids`` 与 ``filter`` 二者不可同时为空(由 dispatch handler 校验)。
|
||
"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
session_ids: list[str] | None = Field(
|
||
default=None,
|
||
description="显式会话 ID 列表(与 filter 二选一)",
|
||
)
|
||
filter: dict[str, Any] | None = Field(
|
||
default=None,
|
||
description="筛选条件(含 channel_type / inactive_before / status)",
|
||
)
|
||
max_count: int = Field(default=100, ge=1, le=1000, description="单次最大关闭数")
|
||
reason: str | None = Field(default=None, max_length=512, description="关闭原因(审计用)")
|
||
|
||
|
||
@session_router.get("/sessions", response_model=dict)
|
||
async def list_sessions(
|
||
request: Request,
|
||
channel_type: ChannelType | None = Query(default=None, description="按渠道类型过滤,留空跨渠道查询"),
|
||
peer_id: str | None = Query(default=None, description="按对端 ID 模糊匹配"),
|
||
limit: int = Query(default=100, 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]:
|
||
"""列出渠道会话(运维排查用例,SES-QUERY-01)。
|
||
|
||
供管理员排查"用户消息未到达"等运维问题。``channel_type`` 留空时跨渠道查询;
|
||
``peer_id`` 提供时按对端 ID 模糊匹配。对应控制面操作 ``session/list``。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
result = await use_cases.session_management.listSessions(
|
||
channel_type=channel_type,
|
||
limit=limit,
|
||
operator=operator,
|
||
offset=offset,
|
||
peer_id=peer_id,
|
||
)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@session_router.post("/sessions/batch-close", response_model=dict)
|
||
async def batch_close_sessions(
|
||
payload: BatchCloseSessionsRequest,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""批量关闭渠道会话(SES-BATCH-CLOSE-01)。
|
||
|
||
支持 ``session_ids`` 显式列表与 ``filter`` 筛选条件两种模式(二者不可
|
||
同时空,由 dispatch handler 校验)。逐条独立事务复用单条 ``closeSession``
|
||
逻辑(模式 D,部分成功语义),单条失败不回滚已成功条目。对应控制面
|
||
操作 ``session/batch_close``。
|
||
|
||
静态路径先于 ``/sessions/{session_id}`` 声明,避免被动态参数捕获。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
cmd = BatchCloseSessionsCmd(
|
||
operator=operator,
|
||
session_ids=tuple(payload.session_ids) if payload.session_ids else (),
|
||
filter=payload.filter,
|
||
max_count=payload.max_count,
|
||
reason=payload.reason,
|
||
)
|
||
result = await use_cases.session_management.batchCloseSessions(cmd)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@session_router.get("/sessions/{session_id}", response_model=dict)
|
||
async def get_session(
|
||
session_id: str,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""按 session_id 查询渠道会话详情(SES-QUERY-02)。
|
||
|
||
供管理员审批 mergeSession 前核对源/目标会话信息。软删除的会话返回 404。
|
||
对应控制面操作 ``session/get``。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
result = await use_cases.session_management.getSession(
|
||
session_id=session_id,
|
||
operator=operator,
|
||
)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@session_router.post("/sessions/{session_id}/merge", response_model=dict)
|
||
async def merge_session(
|
||
session_id: str,
|
||
payload: MergeSessionRequest,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""合并会话(FR-07,SES-05)。
|
||
|
||
将源会话合并至目标会话(路径 ``session_id``),迁移消息并软删除源会话。
|
||
受 ``merge_strategy_enabled`` 配置策略门控制,falsy 时返回 422。
|
||
对应控制面操作 ``session/merge``。
|
||
|
||
协议翻译:路径参数 ``session_id`` → ``MergeConversationCmd.target_conversation_id``。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
cmd = MergeConversationCmd(
|
||
source_conversation_id=payload.source_conversation_id,
|
||
target_conversation_id=session_id,
|
||
operator=operator,
|
||
reason=payload.reason,
|
||
)
|
||
result = await use_cases.session_management.mergeSession(cmd)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@session_router.post("/sessions/{session_id}/transfer", response_model=dict)
|
||
async def transfer_session_owner(
|
||
session_id: str,
|
||
payload: TransferOwnerRequest,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""转移会话所有者(FR-26,SES-13)。
|
||
|
||
将会话(路径 ``session_id``)所有者从当前对端转移至新对端,纯 DB 写入,
|
||
加入控制面管道事务。审计日志由 AuditStage 在 SHARED 事务中统一写入(fail-closed)。
|
||
对应控制面操作 ``session/transfer_owner``。
|
||
|
||
协议翻译:路径参数 ``session_id`` → ``OwnerTransferCmd.conversation_id``。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
cmd = OwnerTransferCmd(
|
||
conversation_id=session_id,
|
||
new_owner_id=payload.new_owner_id,
|
||
operator=operator,
|
||
)
|
||
result = await use_cases.session_management.transferSessionOwner(cmd)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@session_router.post("/sessions/{session_id}/close", response_model=dict)
|
||
async def close_session(
|
||
session_id: str,
|
||
payload: CloseSessionRequest,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""关闭渠道会话(SES-CLOSE-01)。
|
||
|
||
主动关闭指定会话,停止接收新消息。已关闭会话不可重新激活。
|
||
对应控制面操作 ``session/close``。
|
||
|
||
协议翻译:路径参数 ``session_id`` → ``CloseSessionCmd.session_id``。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
cmd = CloseSessionCmd(
|
||
session_id=session_id,
|
||
reason=payload.reason,
|
||
operator=operator,
|
||
)
|
||
result = await use_cases.session_management.closeSession(cmd)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@session_router.get("/sessions/{session_id}/messages", response_model=dict)
|
||
async def list_session_messages(
|
||
session_id: str,
|
||
request: Request,
|
||
limit: int = Query(default=50, ge=1, le=200, 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]:
|
||
"""查询会话消息流(SES-MSG-01)。
|
||
|
||
按时间顺序列出会话内的消息,支持分页。
|
||
对应控制面操作 ``session/list_messages``。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
result = await use_cases.session_management.listSessionMessages(
|
||
session_id=session_id,
|
||
limit=limit,
|
||
offset=offset,
|
||
operator=operator,
|
||
)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@session_router.get("/sessions/{session_id}/stats", response_model=dict)
|
||
async def get_session_stats(
|
||
session_id: str,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""查询单会话统计指标(SES-STATS-01)。
|
||
|
||
按 ``session_id`` 查询会话的统计指标(消息计数、首响时间、平均响应、
|
||
会话时长等),由 dispatch handler 组合 ConversationPort 消息查询计算。
|
||
对应控制面操作 ``session/stats``。
|
||
|
||
协议翻译:路径参数 ``session_id`` 直接传入
|
||
``getSessionStats(session_id=..., operator=...)``。
|
||
"""
|
||
operator = build_operator(current_user, request)
|
||
result = await use_cases.session_management.getSessionStats(
|
||
session_id=session_id,
|
||
operator=operator,
|
||
)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|