ForcePilot/backend/server/routers/scheduler_router.py

861 lines
30 KiB
Python
Raw Normal View History

"""定时任务调度限界上下文 Router。
挂载到 /scheduler 前缀下覆盖任务 CRUD / 状态机 / 手动触发 / 执行日志查询 /
可观测性 / 健康检查 / 运维恢复用例PRD §FR-ST-01 ~ §FR-ST-09
Request Schema Input DTO 不共享类Router 内显式构造 DTO操作人字段
``created_by`` / ``updated_by`` / ``triggered_by`` ``current_user.uid`` 填充
响应统一使用 ``SchedulerResponse[T]`` 泛型模型所有端点返回 ``{"success": True, "data": ...}``
结构错误响应由全局异常处理器统一输出 ``{"success": False, "error": ...}``
对齐 ``external_systems/system_router.py`` 写法
"""
from __future__ import annotations
from typing import TYPE_CHECKING, Any, Literal
from fastapi import APIRouter, Body, Depends, Header, Path, Query
from pydantic import BaseModel, ConfigDict, Field
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.scheduler.infrastructure.container import create_scheduler_service
from yuxi.scheduler.use_cases.dto.scheduler import (
BatchTaskOperationOutput,
CountByStatusInput,
CountByStatusOutput,
CreateTaskInput,
DeleteTaskInput,
GetHealthOutput,
GetRunLogInput,
GetTaskInput,
HardDeleteTaskInput,
ListAllRunLogsInput,
ListAnomaliesOutput,
ListDailyStatsInput,
ListDailyStatsOutput,
ListDeletedTasksInput,
ListHandlerSummaryOutput,
ListRunLogsInput,
ListRunLogsOutput,
ListTasksInput,
ListTasksOutput,
ListUpcomingInput,
ListUpcomingOutput,
PauseTaskInput,
ReclaimStaleRunsInput,
ReclaimStaleRunsOutput,
RestoreTaskInput,
ResumeTaskInput,
RunLogOutput,
TaskOutput,
TriggerTaskInput,
TriggerTaskOutput,
UpdateTaskInput,
)
from yuxi.scheduler.use_cases.services.scheduler_service import SchedulerService
from yuxi.services.run_queue_service import get_arq_pool
from yuxi.storage.postgres.models_business import User
from server.utils.auth_middleware import get_admin_user, get_db, get_required_user
if TYPE_CHECKING:
from arq import ArqRedis
scheduler_router = APIRouter(prefix="/scheduler", tags=["scheduler"])
# =============================================================================
# === 响应模型 ===
# =============================================================================
class SchedulerResponse[T](BaseModel):
"""调度器统一成功响应模型。"""
model_config = ConfigDict(frozen=True)
success: bool = True
data: T
class TaskOperationData(BaseModel):
"""删除 / 硬删除等操作的轻量响应数据。"""
model_config = ConfigDict(frozen=True)
task_id: str
status: str
# =============================================================================
# === 校验常量 ===
# =============================================================================
_ISO_DATE_PATTERN = r"^\d{4}-\d{2}-\d{2}$"
_ISO_DATETIME_PATTERN = r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(\.\d+)?(Z|[+-]\d{2}:\d{2})?$"
# =============================================================================
# === Request Schemas与 Input DTO 不共享类) ===
# =============================================================================
class BatchTaskRequest(BaseModel):
"""批量任务操作请求体。"""
model_config = ConfigDict(frozen=True)
task_ids: list[str] = Field(
...,
min_length=1,
max_length=100,
description="目标任务 ID 列表,最多 100 个",
)
class CreateTaskRequest(BaseModel):
"""创建定时任务请求体。字段对齐 ``CreateTaskInput``(不含 ``created_by``)。"""
model_config = ConfigDict(frozen=True)
handler_name: str = Field(
...,
min_length=1,
max_length=128,
pattern=r"^[a-zA-Z_][a-zA-Z0-9_-]*$",
description="handler 标识,与 TaskHandler.name 引用一致",
)
owner_scope: str = Field(..., min_length=1, max_length=64)
owner_id: str = Field(..., min_length=1, max_length=128)
schedule_kind: Literal["cron", "at"] = Field(..., description="调度类型cron周期/ at一次性")
cron_expression: str | None = Field(default=None, max_length=128)
run_at: str | None = None
tz: str = Field(default="Asia/Shanghai", max_length=64)
payload: dict[str, Any] = Field(default_factory=dict)
enabled: bool = True
delete_after_run: bool = False
block_strategy: Literal["discard_later"] = Field(default="discard_later", description="阻塞策略discard_later")
class UpdateTaskRequest(BaseModel):
"""更新定时任务请求体。字段对齐 ``UpdateTaskInput``(不含 ``task_id`` 与 ``updated_by``)。
仅透传客户端显式设置的字段通过 ``exclude_unset=True``未设置字段保持 ``None``
以保留部分更新语义
"""
model_config = ConfigDict(frozen=True)
cron_expression: str | None = Field(default=None, max_length=128)
run_at: str | None = None
tz: str | None = Field(default=None, max_length=64)
payload: dict[str, Any] | None = None
enabled: bool | None = None
delete_after_run: bool | None = None
owner_scope: str | None = Field(default=None, min_length=1, max_length=64)
owner_id: str | None = Field(default=None, min_length=1, max_length=128)
class ReclaimStaleRunsRequest(BaseModel):
"""回收僵尸执行请求体。字段对齐 ``ReclaimStaleRunsInput``。"""
model_config = ConfigDict(frozen=True)
timeout_seconds: int | None = Field(
default=None,
ge=1,
le=86400,
description="执行超时阈值(秒),默认从配置读取,最大 864001 天)",
)
class TriggerTaskRequest(BaseModel):
"""手动触发任务请求体。字段对齐 ``TriggerTaskInput``(不含 ``task_id`` 与 ``triggered_by``)。
``payload`` 为可选覆盖值None 时使用任务定义中的 payload
"""
model_config = ConfigDict(frozen=True)
payload: dict[str, Any] | None = Field(
default=None,
description="可选 payload 覆盖None 时使用任务定义中的 payload",
)
# =============================================================================
# === 依赖注入 ===
# =============================================================================
async def get_idempotency_key(
idempotency_key: str | None = Header(
default=None,
alias="Idempotency-Key",
max_length=128,
description="幂等键,启用后重复请求回放缓存响应",
),
) -> str | None:
"""提取幂等键 Header 依赖。"""
return idempotency_key
async def get_scheduler_service(
db: AsyncSession = Depends(get_db),
) -> SchedulerService:
"""构造查询 / CRUD / 状态机类端点使用的 SchedulerService无 handler_registry"""
return create_scheduler_service(db)
async def get_scheduler_service_with_arq_pool(
db: AsyncSession = Depends(get_db),
arq_pool: ArqRedis = Depends(get_arq_pool),
) -> SchedulerService:
"""构造手动触发端点使用的 SchedulerService注入 arq_pool"""
return create_scheduler_service(db, arq_pool=arq_pool)
# =============================================================================
# === 静态路径端点(必须在 /tasks/{task_id} 之前声明) ===
# =============================================================================
@scheduler_router.get(
"/tasks",
response_model=SchedulerResponse[ListTasksOutput],
summary="分页列出定时任务",
operation_id="list_scheduler_tasks",
)
async def list_tasks(
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
owner_scope: str | None = Query(None, max_length=64),
owner_id: str | None = Query(None, max_length=128),
handler_name: str | None = Query(None, max_length=128),
keyword: str | None = Query(None, max_length=128, description="关键词模糊搜索task_id / handler_name"),
enabled: bool | None = Query(None),
status: Literal["active", "paused", "dead_letter"] | None = Query(None, description="任务状态过滤"),
sort_by: Literal["created_at", "next_run_at", "last_run_at", "consecutive_errors", "updated_at"] | None = Query(
None, description="排序字段,默认 created_at"
),
sort_order: Literal["asc", "desc"] | None = Query(None, description="排序方向,默认 desc"),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[ListTasksOutput]:
"""分页列出定时任务FR-ST-05 / FR-ST-09
数据可见性本端点为已登录用户可读的管理面查询owner_scope/owner_id
作为业务筛选条件由前端控制API 层不做强制数据隔离如需按用户角色
限制可见范围应在业务层通过 owner 维度过滤
"""
input_dto = ListTasksInput(
page=page,
page_size=page_size,
owner_scope=owner_scope,
owner_id=owner_id,
handler_name=handler_name,
keyword=keyword,
enabled=enabled,
status=status,
sort_by=sort_by,
sort_order=sort_order,
)
output = await service.list_tasks(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.post(
"/tasks/batch-pause",
response_model=SchedulerResponse[BatchTaskOperationOutput],
summary="批量暂停任务",
operation_id="batch_pause_scheduler_tasks",
)
async def batch_pause_tasks(
body: BatchTaskRequest,
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[BatchTaskOperationOutput]:
"""批量暂停任务(仅 ``active`` 状态被转换FR-ST-05"""
output = await service.batch_pause_tasks(
body.task_ids,
operator=current_user.uid,
)
return SchedulerResponse(data=output)
@scheduler_router.post(
"/tasks/batch-resume",
response_model=SchedulerResponse[BatchTaskOperationOutput],
summary="批量恢复任务",
operation_id="batch_resume_scheduler_tasks",
)
async def batch_resume_tasks(
body: BatchTaskRequest,
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[BatchTaskOperationOutput]:
"""批量恢复任务(``paused`` / ``dead_letter`` → ``active``FR-ST-05"""
output = await service.batch_resume_tasks(
body.task_ids,
operator=current_user.uid,
)
return SchedulerResponse(data=output)
@scheduler_router.post(
"/tasks/batch-delete",
response_model=SchedulerResponse[BatchTaskOperationOutput],
summary="批量删除任务",
operation_id="batch_delete_scheduler_tasks",
)
async def batch_delete_tasks(
body: BatchTaskRequest,
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[BatchTaskOperationOutput]:
"""批量软删除任务FR-ST-05"""
output = await service.batch_delete_tasks(
body.task_ids,
operator=current_user.uid,
)
return SchedulerResponse(data=output)
@scheduler_router.post(
"/tasks",
response_model=SchedulerResponse[TaskOutput],
summary="创建定时任务",
operation_id="create_scheduler_task",
)
async def create_task(
body: CreateTaskRequest,
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TaskOutput]:
"""创建定时任务FR-ST-01 / FR-ST-05"""
input_dto = CreateTaskInput(
**body.model_dump(),
created_by=current_user.uid,
idempotency_key=idempotency_key,
)
output = await service.create_task(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.get(
"/tasks/upcoming",
response_model=SchedulerResponse[ListUpcomingOutput],
summary="列出即将执行的任务",
operation_id="list_upcoming_scheduler_tasks",
)
async def list_upcoming(
limit: int = Query(100, ge=1, le=1000),
handler_name: str | None = Query(None, max_length=128),
hours_ahead: int = Query(
24,
ge=1,
le=168,
description="预览窗口小时数(默认 24最大 168=7 天)",
),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[ListUpcomingOutput]:
"""列出即将执行的任务未来指定小时内FR-ST-09"""
input_dto = ListUpcomingInput(
limit=limit,
handler_name=handler_name,
hours_ahead=hours_ahead,
)
output = await service.list_upcoming(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.get(
"/tasks/count-by-status",
response_model=SchedulerResponse[CountByStatusOutput],
summary="按状态计数任务",
operation_id="count_scheduler_tasks_by_status",
)
async def count_by_status(
owner_scope: str | None = Query(None, max_length=64),
owner_id: str | None = Query(None, max_length=128),
handler_name: str | None = Query(None, max_length=128),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[CountByStatusOutput]:
"""按状态计数任务FR-ST-07。支持 owner/handler 下钻过滤。"""
input_dto = CountByStatusInput(
owner_scope=owner_scope,
owner_id=owner_id,
handler_name=handler_name,
)
output = await service.count_by_status(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.get(
"/tasks/anomalies",
response_model=SchedulerResponse[ListAnomaliesOutput],
summary="列出工作台异常任务聚合",
operation_id="list_scheduler_anomalies",
)
async def list_anomalies(
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[ListAnomaliesOutput]:
"""列出工作台异常任务聚合(死信 / 连续失败 / 长期未执行)。
一次调用返回三类异常任务供工作台 AnomalyTaskCard 直接消费
避免前端依赖活跃任务列表前 N 条做前端过滤导致的漏报
"""
output = await service.list_anomalies()
return SchedulerResponse(data=output)
@scheduler_router.get(
"/tasks/deleted",
response_model=SchedulerResponse[ListTasksOutput],
summary="列出已删除任务(回收站)",
operation_id="list_deleted_scheduler_tasks",
)
async def list_deleted_tasks(
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
start_date: str | None = Query(
None,
pattern=_ISO_DATETIME_PATTERN,
description="ISO 格式起始时间,按 deleted_at 过滤",
),
end_date: str | None = Query(
None,
pattern=_ISO_DATETIME_PATTERN,
description="ISO 格式截止时间,按 deleted_at 过滤",
),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[ListTasksOutput]:
"""列出已删除任务回收站FR-ST-05
返回软删除的任务列表 ``deleted_at`` 降序排序支持时间范围过滤与分页
仅管理员可访问
"""
input_dto = ListDeletedTasksInput(
page=page,
page_size=page_size,
start_date=start_date,
end_date=end_date,
)
output = await service.list_deleted_tasks(input_dto)
return SchedulerResponse(data=output)
# =============================================================================
# === 动态路径端点 /tasks/{task_id} ===
# =============================================================================
@scheduler_router.get(
"/tasks/{task_id}",
response_model=SchedulerResponse[TaskOutput],
summary="获取任务详情",
operation_id="get_scheduler_task",
)
async def get_task(
task_id: str = Path(..., min_length=1, max_length=64),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[TaskOutput]:
"""获取任务详情FR-ST-05"""
input_dto = GetTaskInput(task_id=task_id)
output = await service.get_task(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.put(
"/tasks/{task_id}",
response_model=SchedulerResponse[TaskOutput],
summary="更新定时任务",
operation_id="update_scheduler_task",
)
async def update_task(
task_id: str = Path(..., min_length=1, max_length=64),
body: UpdateTaskRequest = Body(...),
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TaskOutput]:
"""更新定时任务FR-ST-05。仅透传客户端显式设置的字段保留部分更新语义。"""
input_dto = UpdateTaskInput(
task_id=task_id,
updated_by=current_user.uid,
idempotency_key=idempotency_key,
**body.model_dump(exclude_unset=True),
)
output = await service.update_task(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.delete(
"/tasks/{task_id}",
response_model=SchedulerResponse[TaskOperationData],
summary="删除定时任务(软删除)",
operation_id="delete_scheduler_task",
)
async def delete_task(
task_id: str = Path(..., min_length=1, max_length=64),
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TaskOperationData]:
"""删除定时任务软删除FR-ST-05"""
input_dto = DeleteTaskInput(
task_id=task_id,
updated_by=current_user.uid,
idempotency_key=idempotency_key,
)
await service.delete_task(input_dto)
return SchedulerResponse(data=TaskOperationData(task_id=task_id, status="deleted"))
@scheduler_router.post(
"/tasks/{task_id}/pause",
response_model=SchedulerResponse[TaskOutput],
summary="暂停任务",
operation_id="pause_scheduler_task",
)
async def pause_task(
task_id: str = Path(..., min_length=1, max_length=64),
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TaskOutput]:
"""暂停任务FR-ST-05"""
input_dto = PauseTaskInput(
task_id=task_id,
updated_by=current_user.uid,
idempotency_key=idempotency_key,
)
output = await service.pause_task(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.post(
"/tasks/{task_id}/resume",
response_model=SchedulerResponse[TaskOutput],
summary="恢复任务",
operation_id="resume_scheduler_task",
)
async def resume_task(
task_id: str = Path(..., min_length=1, max_length=64),
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TaskOutput]:
"""恢复任务FR-ST-05 / FR-ST-08
``paused`` -> ``active`` / ``dead_letter`` -> ``active``
"""
input_dto = ResumeTaskInput(
task_id=task_id,
updated_by=current_user.uid,
idempotency_key=idempotency_key,
)
output = await service.resume_task(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.post(
"/tasks/{task_id}/trigger",
response_model=SchedulerResponse[TriggerTaskOutput],
summary="手动触发任务执行",
operation_id="trigger_scheduler_task",
)
async def trigger_task(
task_id: str = Path(..., min_length=1, max_length=64),
body: TriggerTaskRequest | None = Body(None),
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service_with_arq_pool),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TriggerTaskOutput]:
"""手动触发任务执行FR-ST-05
生成 ``run_id`` 后入队 ARQ 执行不影响下一次自动执行
可通过 body 传入 payload 覆盖任务定义中的 payloadNone 时使用任务原 payload
触发人 ``current_user.uid`` 透传到 worker 写入 ``run_log.created_by`` 审计字段
"""
payload = body.payload if body else None
input_dto = TriggerTaskInput(
task_id=task_id,
triggered_by="manual",
payload=payload,
operator=current_user.uid,
idempotency_key=idempotency_key,
)
output = await service.trigger_task(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.post(
"/tasks/{task_id}/restore",
response_model=SchedulerResponse[TaskOutput],
summary="恢复已删除任务",
operation_id="restore_scheduler_task",
)
async def restore_task(
task_id: str = Path(..., min_length=1, max_length=64),
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TaskOutput]:
"""恢复已删除任务FR-ST-05
将软删除的任务恢复为 ``is_deleted=0``若恢复前 ``status='active'``
自动置为 ``paused`` 避免立即被 tick 扫描执行需管理员确认后手动 ``resume``
仅管理员可操作
"""
input_dto = RestoreTaskInput(
task_id=task_id,
updated_by=current_user.uid,
idempotency_key=idempotency_key,
)
output = await service.restore_task(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.delete(
"/tasks/{task_id}/hard",
response_model=SchedulerResponse[TaskOperationData],
summary="硬删除任务(不可恢复)",
operation_id="hard_delete_scheduler_task",
)
async def hard_delete_task(
task_id: str = Path(..., min_length=1, max_length=64),
idempotency_key: str | None = Depends(get_idempotency_key),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[TaskOperationData]:
"""硬删除任务不可恢复FR-ST-05
物理删除已软删除的任务记录操作不可逆仅管理员可操作
"""
input_dto = HardDeleteTaskInput(
task_id=task_id,
updated_by=current_user.uid,
idempotency_key=idempotency_key,
)
await service.hard_delete_task(input_dto)
return SchedulerResponse(data=TaskOperationData(task_id=task_id, status="hard_deleted"))
@scheduler_router.get(
"/tasks/{task_id}/run-logs",
response_model=SchedulerResponse[ListRunLogsOutput],
summary="列出任务执行日志",
operation_id="list_scheduler_run_logs",
)
async def list_run_logs(
task_id: str = Path(..., min_length=1, max_length=64),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
status: Literal["running", "success", "failure", "timeout", "skipped"] | None = Query(
None, description="执行状态过滤"
),
start_date: str | None = Query(
None,
pattern=_ISO_DATETIME_PATTERN,
description="ISO 格式起始时间,按 started_at 过滤",
),
end_date: str | None = Query(
None,
pattern=_ISO_DATETIME_PATTERN,
description="ISO 格式截止时间,按 started_at 过滤",
),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[ListRunLogsOutput]:
"""列出任务执行日志FR-ST-05 / FR-ST-09"""
input_dto = ListRunLogsInput(
task_id=task_id,
page=page,
page_size=page_size,
status=status,
start_date=start_date,
end_date=end_date,
)
output = await service.list_run_logs(input_dto)
return SchedulerResponse(data=output)
# =============================================================================
# === 跨任务执行日志查询(静态路径,须在 /run-logs/{run_id} 之前声明) ===
# =============================================================================
@scheduler_router.get(
"/run-logs",
response_model=SchedulerResponse[ListRunLogsOutput],
summary="跨任务列出执行日志",
operation_id="list_all_scheduler_run_logs",
)
async def list_all_run_logs(
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
status: Literal["running", "success", "failure", "timeout", "skipped"] | None = Query(
None, description="执行状态过滤"
),
handler_name: str | None = Query(
None, min_length=1, max_length=128, description="handler 名称过滤(跨任务按 handler 维度查询)"
),
start_date: str | None = Query(
None,
pattern=_ISO_DATETIME_PATTERN,
description="ISO 格式起始时间,按 started_at 过滤",
),
end_date: str | None = Query(
None,
pattern=_ISO_DATETIME_PATTERN,
description="ISO 格式截止时间,按 started_at 过滤",
),
task_id: str | None = Query(None, min_length=1, max_length=64),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[ListRunLogsOutput]:
"""跨任务列出执行日志FR-ST-09
至少需指定 ``task_id````status`` 或时间范围之一由用例层校验
避免无过滤的全表扫描``handler_name`` 为可选过滤维度
"""
input_dto = ListAllRunLogsInput(
page=page,
page_size=page_size,
status=status,
handler_name=handler_name,
start_date=start_date,
end_date=end_date,
task_id=task_id,
)
output = await service.list_all_run_logs(input_dto)
return SchedulerResponse(data=output)
@scheduler_router.get(
"/run-logs/{run_id}",
response_model=SchedulerResponse[RunLogOutput],
summary="获取执行日志详情",
operation_id="get_scheduler_run_log",
)
async def get_run_log(
run_id: str = Path(..., min_length=1, max_length=64),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[RunLogOutput]:
"""获取执行日志详情FR-ST-05
``run_id`` 查询单条执行日志供告警跳转 / 排查失败执行使用
"""
input_dto = GetRunLogInput(run_id=run_id)
output = await service.get_run_log(input_dto)
return SchedulerResponse(data=output)
# =============================================================================
# === 可观测性 ===
# =============================================================================
@scheduler_router.get(
"/handlers/summary",
response_model=SchedulerResponse[ListHandlerSummaryOutput],
summary="列出 handler 聚合摘要",
operation_id="list_scheduler_handler_summary",
)
async def list_handler_summary(
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[ListHandlerSummaryOutput]:
"""列出 handler 聚合摘要FR-ST-05"""
output = await service.list_handler_summary()
return SchedulerResponse(data=output)
# =============================================================================
# === 日聚合统计 ===
# =============================================================================
@scheduler_router.get(
"/stats/daily",
response_model=SchedulerResponse[ListDailyStatsOutput],
summary="列出日聚合统计",
operation_id="list_scheduler_daily_stats",
)
async def list_daily_stats(
start_date: str = Query(
...,
pattern=_ISO_DATE_PATTERN,
description="ISO 日期 YYYY-MM-DD必填",
),
end_date: str = Query(
...,
pattern=_ISO_DATE_PATTERN,
description="ISO 日期 YYYY-MM-DD必填",
),
handler_name: str | None = Query(None, max_length=128),
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_required_user),
) -> SchedulerResponse[ListDailyStatsOutput]:
"""列出日聚合统计FR-ST-09
按日期范围 + handler 维度查询日聚合统计返回含 ``total_count`` /
``success_rate`` 派生字段日期范围上限 90
"""
input_dto = ListDailyStatsInput(
start_date=start_date,
end_date=end_date,
handler_name=handler_name,
)
output = await service.list_daily_stats(input_dto)
return SchedulerResponse(data=output)
# =============================================================================
# === 健康检查与运维恢复 ===
# =============================================================================
@scheduler_router.get(
"/health",
response_model=SchedulerResponse[GetHealthOutput],
summary="获取调度器健康状态",
operation_id="get_scheduler_health",
)
async def get_health(
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[GetHealthOutput]:
"""获取调度器健康状态FR-ST-07"""
output = await service.get_health()
return SchedulerResponse(data=output)
@scheduler_router.post(
"/reclaim-stale-runs",
response_model=SchedulerResponse[ReclaimStaleRunsOutput],
summary="回收僵尸执行记录",
operation_id="reclaim_scheduler_stale_runs",
)
async def reclaim_stale_runs(
body: ReclaimStaleRunsRequest | None = None,
service: SchedulerService = Depends(get_scheduler_service),
current_user: User = Depends(get_admin_user),
) -> SchedulerResponse[ReclaimStaleRunsOutput]:
"""回收僵尸执行记录(运维恢复)。
回收 ``status=running`` 但实际已超时的执行记录标记为 ``timeout``
"""
timeout_seconds = body.timeout_seconds if body else None
input_dto = ReclaimStaleRunsInput(timeout_seconds=timeout_seconds)
output = await service.reclaim_stale_runs(input_dto)
return SchedulerResponse(data=output)