ForcePilot/backend/server/routers/external_systems/execution_router.py
Kris 002e6a356b chore: 批量整理代码变更,修复多类细节问题
1.  清理测试文件中未使用的导入与冗余代码
2.  修复审计日志与批量操作的空值约束,统一填充"global"作为默认渠道
3.  调整批量消息撤回的响应语义,对齐其他端点的部分成功契约
4.  修复访问规则批量克隆的唯一约束问题,新增后缀自动处理逻辑
5.  替换anyio为asyncio并行调用,修正时间UTC导入路径
6.  优化前端外部系统概览页的刷新状态提示与缓存逻辑
7.  修复测试用例中的断言与请求方式问题,适配httpx删除请求特性
8.  重构后端路由的依赖注入,移除冗余的数据库会话依赖
9.  调整测试用例的权限校验逻辑,修正强制登出的权限判断
10. 修复语义分块测试的numpy依赖问题,清理冗余导入
2026-07-13 20:48:29 +08:00

334 lines
14 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.

"""Execution 子域 Router。
提供执行记录的查询 / 详情 / trace / 重试 / 清理 / 统计聚合 /
慢执行 / 重试链 / 执行上下文 / 按 trace 查询 API。挂载到
``/system/external-systems/executions`` 前缀下(根前缀由聚合 router 追加)。
路径顺序约束:静态路径 ``/cleanup`` / ``/stats`` / ``/stats/grouped`` /
``/stats/error-analysis`` / ``/stats/trend`` / ``/slow`` / ``/by-trace/{trace_id}``
必须在 ``/{execution_id}`` 前声明,避免被路径参数捕获。
设计依据docs/vibe/v1.1/设计方案/2026-06-17-Router-API-设计方案.md
扩展依据docs/vibe/v1.1/设计方案/RouterAPI扩展设计/10-execution_router扩展设计方案.md
"""
from __future__ import annotations
from datetime import datetime
from typing import Any, Literal
from fastapi import APIRouter, Depends, Path, Query
from pydantic import BaseModel, ConfigDict, Field
from yuxi.external_systems.infrastructure.container import UseCases
from yuxi.external_systems.use_cases.dto.execution import (
CleanupInput,
ExecutionStatusLiteral,
GetExecutionContextInput,
GetExecutionDetailInput,
GetExecutionStatsInput,
GetExecutionTraceInput,
GetExecutionTrendInput,
GetSlowExecutionsInput,
GetTraceExecutionsInput,
ListRetriesInput,
PaginateExecutionsInput,
RetryExecutionInput,
)
from yuxi.storage.postgres.models_business import User
from server.routers.external_systems import get_use_cases
from server.utils.auth_middleware import get_admin_user, get_required_user
execution_router = APIRouter(prefix="/executions", tags=["external-systems-execution"])
# =============================================================================
# === Request Schemas ===
# =============================================================================
class CleanupExecutionsRequest(BaseModel):
"""清理历史执行记录请求体。字段对齐 ``CleanupInput``。"""
model_config = ConfigDict(frozen=True)
system_id: int | None = Field(default=None, ge=1)
before_at: datetime | None = None
status: ExecutionStatusLiteral | None = None
batch_size: int = Field(default=1000, ge=1, le=10000)
# =============================================================================
# === 静态路径端点(必须在 /{execution_id} 之前声明) ===
# =============================================================================
@execution_router.post("/cleanup", response_model=dict)
async def cleanup_old_executions(
payload: CleanupExecutionsRequest,
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_admin_user),
) -> dict[str, Any]:
"""清理历史执行记录(管理员)。"""
input_dto = CleanupInput(
system_id=payload.system_id,
before_at=payload.before_at,
status=payload.status,
batch_size=payload.batch_size,
)
output = await use_cases.execution_service.cleanup_old_executions(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("", response_model=dict)
async def paginate_executions(
limit: int = Query(100, ge=1, le=500),
offset: int = Query(0, ge=0),
system_id: int | None = Query(None, ge=1),
env_key: str | None = Query(None, max_length=32),
tool_slug: str | None = Query(None, max_length=128),
status: ExecutionStatusLiteral | None = Query(
None, description="执行状态pending/running/success/failed/timeout/throttled/cancelled"
),
caller: str | None = Query(None, max_length=32),
caller_id: str | None = Query(None, max_length=64, description="按调用方实体过滤"),
operation: str | None = Query(None, max_length=32, description="按操作类型过滤"),
tag_key: str | None = Query(None, max_length=64, description="按业务标签键过滤"),
tag_value: str | None = Query(None, max_length=256, description="按业务标签值过滤"),
trace_id: str | None = Query(None, max_length=64),
correlation_id: str | None = Query(None, max_length=64),
start_time: datetime | None = Query(None, description="ISO8601 开始时间"),
end_time: datetime | None = Query(None, description="ISO8601 结束时间"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""分页查询执行记录。"""
input_dto = PaginateExecutionsInput(
limit=limit,
offset=offset,
system_id=system_id,
env_key=env_key,
tool_slug=tool_slug,
status=status,
caller=caller,
caller_id=caller_id,
operation=operation,
tag_key=tag_key,
tag_value=tag_value,
trace_id=trace_id,
correlation_id=correlation_id,
start_at=start_time,
end_at=end_time,
)
output = await use_cases.execution_service.paginate_executions(input_dto)
return {"success": True, "data": output.model_dump()}
# =============================================================================
# === 扩展:统计聚合 / 慢执行 / 按 trace 查询(静态路径,必须在 /{execution_id} 之前) ===
# =============================================================================
@execution_router.get("/stats", response_model=dict)
async def get_execution_stats(
system_id: int | None = Query(None, ge=1, description="系统过滤"),
env_key: str | None = Query(None, max_length=32, description="环境过滤"),
adapter_type: str | None = Query(None, max_length=32, description="适配器类型过滤"),
start_time: datetime | None = Query(None, description="起始时间"),
end_time: datetime | None = Query(None, description="结束时间"),
tool_limit: int = Query(20, ge=1, le=100, description="Top 工具数量"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""执行统计聚合。含按状态/工具分组。"""
input_dto = GetExecutionStatsInput(
system_id=system_id,
env_key=env_key,
adapter_type=adapter_type,
start=start_time,
end=end_time,
limit=tool_limit,
)
output = await use_cases.execution_service.get_execution_stats(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/stats/grouped", response_model=dict)
async def get_grouped_stats(
group_by: Literal["tool"] = Query(..., description="分组维度,仅支持 tool"),
system_id: int | None = Query(None, ge=1, description="系统过滤"),
env_key: str | None = Query(None, max_length=32, description="环境过滤"),
adapter_type: str | None = Query(None, max_length=32, description="适配器类型过滤"),
start_time: datetime | None = Query(None, description="起始时间"),
end_time: datetime | None = Query(None, description="结束时间"),
limit: int = Query(20, ge=1, le=100, description="返回数量"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""分组统计。仅支持 group_by=tool仓储仅 aggregate_by_tool 支持分组)。"""
input_dto = GetExecutionStatsInput(
system_id=system_id,
env_key=env_key,
adapter_type=adapter_type,
start=start_time,
end=end_time,
limit=limit,
)
output = await use_cases.execution_service.get_grouped_stats(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/stats/error-analysis", response_model=dict)
async def get_error_analysis(
system_id: int = Query(..., ge=1, description="系统 ID必填"),
start_time: datetime | None = Query(None, description="起始时间"),
end_time: datetime | None = Query(None, description="结束时间"),
limit: int = Query(5, ge=1, le=50, description="返回数量"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""错误分析。按错误码分组统计失败执行数量。"""
input_dto = GetExecutionStatsInput(
system_id=system_id,
start=start_time,
end=end_time,
limit=limit,
)
output = await use_cases.execution_service.get_error_analysis(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/stats/trend", response_model=dict)
async def get_execution_trend(
tool_slug: str = Query(..., min_length=1, max_length=128, description="工具标识(必填)"),
system_id: int | None = Query(None, ge=1, description="系统过滤"),
env_key: str | None = Query(None, max_length=32, description="环境过滤"),
start_time: datetime | None = Query(None, description="起始时间"),
end_time: datetime | None = Query(None, description="结束时间"),
interval: Literal["hour", "day"] = Query("day", description="时间桶粒度hour / day"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""执行趋势。按时间桶 + 工具聚合。"""
input_dto = GetExecutionTrendInput(
tool_slug=tool_slug,
system_id=system_id,
env_key=env_key,
start=start_time,
end=end_time,
interval=interval,
)
output = await use_cases.execution_service.get_execution_trend(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/slow", response_model=dict)
async def list_slow_executions(
system_id: int = Query(..., ge=1, description="系统 ID必填"),
min_duration_ms: int = Query(..., ge=0, description="耗时阈值(毫秒,必填)"),
limit: int = Query(50, ge=1, le=200, description="返回数量"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""慢执行查询。按耗时降序。"""
input_dto = GetSlowExecutionsInput(
system_id=system_id,
min_duration_ms=min_duration_ms,
limit=limit,
)
output = await use_cases.execution_service.list_slow_executions(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/by-trace/{trace_id}", response_model=dict)
async def list_by_trace(
trace_id: str = Path(..., min_length=1, max_length=64, description="调用链路 ID"),
limit: int = Query(50, ge=1, le=200, description="返回数量"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""按 trace 查询。返回 execution + trace 组合。"""
input_dto = GetTraceExecutionsInput(
trace_id=trace_id,
limit=limit,
)
output = await use_cases.execution_service.list_by_trace(input_dto)
return {"success": True, "data": output.model_dump()}
# =============================================================================
# === 动态路径端点 /{execution_id} ===
# =============================================================================
@execution_router.get("/{execution_id}", response_model=dict)
async def get_execution_detail(
execution_id: str = Path(..., min_length=1, max_length=128, description="执行记录 ID"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""获取执行详情。"""
input_dto = GetExecutionDetailInput(execution_id=execution_id)
output = await use_cases.execution_service.get_execution_detail(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/{execution_id}/trace", response_model=dict)
async def get_execution_trace(
execution_id: str = Path(..., min_length=1, max_length=128, description="执行记录 ID"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""获取执行 trace。"""
input_dto = GetExecutionTraceInput(execution_id=execution_id)
output = await use_cases.execution_service.get_execution_trace(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/{execution_id}/retries", response_model=dict)
async def list_retries(
execution_id: str = Path(..., min_length=1, max_length=128, description="执行记录 ID"),
limit: int = Query(50, ge=1, le=200, description="返回数量"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""重试链查询。基于 retry_of 字段,按 started_at 升序。"""
input_dto = ListRetriesInput(
execution_id=execution_id,
limit=limit,
)
output = await use_cases.execution_service.list_retries(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.get("/{execution_id}/context", response_model=dict)
async def get_execution_context(
execution_id: str = Path(..., min_length=1, max_length=128, description="执行记录 ID"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_required_user),
) -> dict[str, Any]:
"""执行上下文。含 trace + 关联告警。"""
input_dto = GetExecutionContextInput(execution_id=execution_id)
output = await use_cases.execution_service.get_execution_context(input_dto)
return {"success": True, "data": output.model_dump()}
@execution_router.post("/{execution_id}/retry", response_model=dict)
async def retry_execution(
execution_id: str = Path(..., min_length=1, max_length=128, description="执行记录 ID"),
use_cases: UseCases = Depends(get_use_cases),
current_user: User = Depends(get_admin_user),
) -> dict[str, Any]:
"""重试执行(管理员)。
``caller`` / ``caller_id`` 均使用当前管理员 uid用于审计日志。
执行记录的 ``caller`` 由 service 层固定为 ``retry``,确保重试语义准确。
"""
input_dto = RetryExecutionInput(
execution_id=execution_id,
caller=current_user.uid,
caller_id=current_user.uid,
)
output = await use_cases.execution_service.retry_execution(input_dto)
return {"success": True, "data": output.model_dump()}