diff --git a/backend/server/routers/__init__.py b/backend/server/routers/__init__.py index a76c648b..87738daa 100644 --- a/backend/server/routers/__init__.py +++ b/backend/server/routers/__init__.py @@ -17,6 +17,7 @@ from server.routers.user_router import user_router from server.routers.filesystem_router import filesystem_router from server.routers.workspace_router import workspace from server.routers.mention_router import mention_router +from server.routers.scheduler_router import scheduler_router from server.routers.external_systems import external_systems_router _LITE_MODE = os.environ.get("LITE_MODE", "").lower() in ("true", "1") @@ -42,6 +43,7 @@ router.include_router(user_router) # /api/user/* 用户级配置与凭据 router.include_router(filesystem_router) # /api/viewer/filesystem/* 工作台文件系统视图 router.include_router(workspace) # /api/workspace/* 用户个人工作区 router.include_router(mention_router) # /api/mention/* 提及文件搜索接口 +router.include_router(scheduler_router) # /api/scheduler/* 定时任务调度管理 router.include_router(external_systems_router) # /api/system/external-systems/* 外部系统管理(骨架,子 router 待实现) if not _LITE_MODE: diff --git a/backend/server/routers/scheduler_router.py b/backend/server/routers/scheduler_router.py new file mode 100644 index 00000000..63759334 --- /dev/null +++ b/backend/server/routers/scheduler_router.py @@ -0,0 +1,500 @@ +"""定时任务调度限界上下文 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`` 填充。 + +对齐 ``external_systems/system_router.py`` 写法。 +""" + +from __future__ import annotations + +from typing import Any + +from fastapi import APIRouter, Depends, 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 ( + CreateTaskInput, + DeleteTaskInput, + GetRunLogInput, + GetTaskInput, + HardDeleteTaskInput, + ListAllRunLogsInput, + ListDailyStatsInput, + ListDeletedTasksInput, + ListRunLogsInput, + ListTasksInput, + ListUpcomingInput, + PauseTaskInput, + ReclaimStaleRunsInput, + RestoreTaskInput, + ResumeTaskInput, + TriggerTaskInput, + UpdateTaskInput, +) +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 + +scheduler_router = APIRouter(prefix="/scheduler", tags=["scheduler"]) + + +# ============================================================================= +# === Request Schemas(与 Input DTO 不共享类) === +# ============================================================================= + + +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: str = Field(..., pattern=r"^(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: str = Field(default="discard_later", pattern=r"^(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) + + +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", + ) + + +# ============================================================================= +# === 静态路径端点(必须在 /tasks/{task_id} 之前声明) === +# ============================================================================= + + +@scheduler_router.get("/tasks", response_model=dict) +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), + enabled: bool | None = Query(None), + status: str | None = Query(None, pattern=r"^(active|paused|dead_letter)$"), + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """分页列出定时任务(FR-ST-05 / FR-ST-09)。""" + service = create_scheduler_service(db) + input_dto = ListTasksInput( + page=page, + page_size=page_size, + owner_scope=owner_scope, + owner_id=owner_id, + handler_name=handler_name, + enabled=enabled, + status=status, + ) + output = await service.list_tasks(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.post("/tasks", response_model=dict) +async def create_task( + body: CreateTaskRequest, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """创建定时任务(FR-ST-01 / FR-ST-05)。""" + service = create_scheduler_service(db) + input_dto = CreateTaskInput(**body.model_dump(), created_by=current_user.uid) + output = await service.create_task(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.get("/tasks/upcoming", response_model=dict) +async def list_upcoming( + limit: int = Query(100, ge=1, le=1000), + handler_name: str | None = Query(None, max_length=128), + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """列出即将执行的任务(未来 24 小时内,FR-ST-09)。""" + service = create_scheduler_service(db) + input_dto = ListUpcomingInput(limit=limit, handler_name=handler_name) + output = await service.list_upcoming(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.get("/tasks/count-by-status", response_model=dict) +async def count_by_status( + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """按状态计数任务(FR-ST-07)。""" + service = create_scheduler_service(db) + output = await service.count_by_status() + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.get("/tasks/deleted", response_model=dict) +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, description="ISO 格式起始时间,按 deleted_at 过滤"), + end_date: str | None = Query(None, description="ISO 格式截止时间,按 deleted_at 过滤"), + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """列出已删除任务(回收站,FR-ST-05)。 + + 返回软删除的任务列表,按 ``deleted_at`` 降序排序,支持时间范围过滤与分页。 + 仅管理员可访问。 + """ + service = create_scheduler_service(db) + 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 {"success": True, "data": output.model_dump()} + + +# ============================================================================= +# === 动态路径端点 /tasks/{task_id} === +# ============================================================================= + + +@scheduler_router.get("/tasks/{task_id}", response_model=dict) +async def get_task( + task_id: str, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """获取任务详情(FR-ST-05)。""" + service = create_scheduler_service(db) + input_dto = GetTaskInput(task_id=task_id) + output = await service.get_task(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.put("/tasks/{task_id}", response_model=dict) +async def update_task( + task_id: str, + body: UpdateTaskRequest, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """更新定时任务(FR-ST-05)。仅透传客户端显式设置的字段,保留部分更新语义。""" + service = create_scheduler_service(db) + input_dto = UpdateTaskInput( + task_id=task_id, + updated_by=current_user.uid, + **body.model_dump(exclude_unset=True), + ) + output = await service.update_task(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.delete("/tasks/{task_id}", response_model=dict) +async def delete_task( + task_id: str, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """删除定时任务(软删除,FR-ST-05)。""" + service = create_scheduler_service(db) + input_dto = DeleteTaskInput(task_id=task_id, updated_by=current_user.uid) + await service.delete_task(input_dto) + return {"success": True, "data": {"task_id": task_id, "status": "deleted"}} + + +@scheduler_router.post("/tasks/{task_id}/pause", response_model=dict) +async def pause_task( + task_id: str, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """暂停任务(FR-ST-05)。""" + service = create_scheduler_service(db) + input_dto = PauseTaskInput(task_id=task_id, updated_by=current_user.uid) + output = await service.pause_task(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.post("/tasks/{task_id}/resume", response_model=dict) +async def resume_task( + task_id: str, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """恢复任务(FR-ST-05 / FR-ST-08)。 + + ``paused`` -> ``active`` / ``dead_letter`` -> ``active``。 + """ + service = create_scheduler_service(db) + input_dto = ResumeTaskInput(task_id=task_id, updated_by=current_user.uid) + output = await service.resume_task(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.post("/tasks/{task_id}/trigger", response_model=dict) +async def trigger_task( + task_id: str, + body: TriggerTaskRequest | None = None, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """手动触发任务执行(FR-ST-05)。 + + 生成 ``run_id`` 后入队 ARQ 执行,不影响下一次自动执行。 + 可通过 body 传入 payload 覆盖任务定义中的 payload,None 时使用任务原 payload。 + """ + arq_pool = await get_arq_pool() + service = create_scheduler_service(db, arq_pool=arq_pool) + payload = body.payload if body else None + input_dto = TriggerTaskInput( + task_id=task_id, + triggered_by="manual", + payload=payload, + ) + output = await service.trigger_task(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.post("/tasks/{task_id}/restore", response_model=dict) +async def restore_task( + task_id: str, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """恢复已删除任务(FR-ST-05)。 + + 将软删除的任务恢复为 ``is_deleted=0``。若恢复前 ``status='active'``, + 自动置为 ``paused`` 避免立即被 tick 扫描执行,需管理员确认后手动 ``resume``。 + 仅管理员可操作。 + """ + service = create_scheduler_service(db) + input_dto = RestoreTaskInput(task_id=task_id, updated_by=current_user.uid) + output = await service.restore_task(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.delete("/tasks/{task_id}/hard", response_model=dict) +async def hard_delete_task( + task_id: str, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """硬删除任务(不可恢复,FR-ST-05)。 + + 物理删除已软删除的任务记录,操作不可逆。仅管理员可操作。 + """ + service = create_scheduler_service(db) + input_dto = HardDeleteTaskInput(task_id=task_id, updated_by=current_user.uid) + await service.hard_delete_task(input_dto) + return {"success": True, "data": {"task_id": task_id, "status": "hard_deleted"}} + + +@scheduler_router.get("/tasks/{task_id}/run-logs", response_model=dict) +async def list_run_logs( + task_id: str, + page: int = Query(1, ge=1), + page_size: int = Query(20, ge=1, le=100), + status: str | None = Query( + None, pattern=r"^(running|success|failure|timeout|skipped)$" + ), + start_date: str | None = Query(None, description="ISO 格式起始时间"), + end_date: str | None = Query(None, description="ISO 格式截止时间"), + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """列出任务执行日志(FR-ST-05 / FR-ST-09)。""" + service = create_scheduler_service(db) + 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 {"success": True, "data": output.model_dump()} + + +# ============================================================================= +# === 跨任务执行日志查询(静态路径,须在 /run-logs/{run_id} 之前声明) === +# ============================================================================= + + +@scheduler_router.get("/run-logs", response_model=dict) +async def list_all_run_logs( + page: int = Query(1, ge=1), + page_size: int = Query(20, ge=1, le=100), + status: str | None = Query( + None, pattern=r"^(running|success|failure|timeout|skipped)$" + ), + start_date: str | None = Query(None, description="ISO 格式起始时间"), + end_date: str | None = Query(None, description="ISO 格式截止时间"), + task_id: str | None = Query(None, min_length=1, max_length=64), + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """跨任务列出执行日志(FR-ST-09)。 + + ``task_id`` 与 ``status`` 不可同时为空(由用例层校验),避免无过滤的全表扫描。 + """ + service = create_scheduler_service(db) + input_dto = ListAllRunLogsInput( + page=page, + page_size=page_size, + status=status, + start_date=start_date, + end_date=end_date, + task_id=task_id, + ) + output = await service.list_all_run_logs(input_dto) + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.get("/run-logs/{run_id}", response_model=dict) +async def get_run_log( + run_id: str, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """获取执行日志详情(FR-ST-05)。 + + 按 ``run_id`` 查询单条执行日志,供告警跳转 / 排查失败执行使用。 + """ + service = create_scheduler_service(db) + input_dto = GetRunLogInput(run_id=run_id) + output = await service.get_run_log(input_dto) + return {"success": True, "data": output.model_dump()} + + +# ============================================================================= +# === 可观测性 === +# ============================================================================= + + +@scheduler_router.get("/handlers/summary", response_model=dict) +async def list_handler_summary( + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """列出 handler 聚合摘要(FR-ST-05)。""" + service = create_scheduler_service(db) + output = await service.list_handler_summary() + return {"success": True, "data": output.model_dump()} + + +# ============================================================================= +# === 日聚合统计 === +# ============================================================================= + + +@scheduler_router.get("/stats/daily", response_model=dict) +async def list_daily_stats( + start_date: str = Query(..., description="ISO 日期 YYYY-MM-DD(必填)"), + end_date: str = Query(..., description="ISO 日期 YYYY-MM-DD(必填)"), + handler_name: str | None = Query(None, max_length=128), + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """列出日聚合统计(FR-ST-09)。 + + 按日期范围 + handler 维度查询日聚合统计,返回含 ``total_count`` / + ``success_rate`` 派生字段。日期范围上限 90 天。 + """ + service = create_scheduler_service(db) + input_dto = ListDailyStatsInput( + start_date=start_date, + end_date=end_date, + handler_name=handler_name, + ) + output = await service.list_daily_stats(input_dto) + return {"success": True, "data": output.model_dump()} + + +# ============================================================================= +# === 健康检查与运维恢复 === +# ============================================================================= + + +@scheduler_router.get("/health", response_model=dict) +async def get_health( + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_required_user), +) -> dict[str, Any]: + """获取调度器健康状态(FR-ST-07)。""" + service = create_scheduler_service(db) + output = await service.get_health() + return {"success": True, "data": output.model_dump()} + + +@scheduler_router.post("/reclaim-stale-runs", response_model=dict) +async def reclaim_stale_runs( + body: ReclaimStaleRunsRequest | None = None, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_admin_user), +) -> dict[str, Any]: + """回收僵尸执行记录(运维恢复)。 + + 回收 ``status=running`` 但实际已超时的执行记录,标记为 ``timeout``。 + """ + service = create_scheduler_service(db) + 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 {"success": True, "data": output.model_dump()} diff --git a/backend/server/utils/lifespan.py b/backend/server/utils/lifespan.py index 17d2044d..8816eeb9 100644 --- a/backend/server/utils/lifespan.py +++ b/backend/server/utils/lifespan.py @@ -26,6 +26,7 @@ async def lifespan(app: FastAPI): logger.error(f"Failed to register external system adapters during startup: {e}") + """FastAPI lifespan事件管理器""" # 初始化数据库连接 try: @@ -34,6 +35,7 @@ async def lifespan(app: FastAPI): await pg_manager.ensure_business_schema() await pg_manager.ensure_knowledge_schema() await pg_manager.ensure_external_schema() + await pg_manager.ensure_scheduler_schema() except Exception as e: logger.error(f"Failed to initialize database during startup: {e}")