diff --git a/backend/server/routers/scheduler_router.py b/backend/server/routers/scheduler_router.py index 63759334..64ffcb46 100644 --- a/backend/server/routers/scheduler_router.py +++ b/backend/server/routers/scheduler_router.py @@ -13,7 +13,7 @@ from __future__ import annotations from typing import Any -from fastapi import APIRouter, Depends, Query +from fastapi import APIRouter, Depends, Header, Query from pydantic import BaseModel, ConfigDict, Field from sqlalchemy.ext.asyncio import AsyncSession from yuxi.scheduler.infrastructure.container import create_scheduler_service @@ -149,12 +149,22 @@ async def list_tasks( @scheduler_router.post("/tasks", response_model=dict) async def create_task( body: CreateTaskRequest, + idempotency_key: str | None = Header( + default=None, + alias="Idempotency-Key", + max_length=128, + description="幂等键,启用后重复请求回放缓存响应", + ), 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) + input_dto = CreateTaskInput( + **body.model_dump(), + created_by=current_user.uid, + idempotency_key=idempotency_key, + ) output = await service.create_task(input_dto) return {"success": True, "data": output.model_dump()} @@ -231,6 +241,12 @@ async def get_task( async def update_task( task_id: str, body: UpdateTaskRequest, + idempotency_key: str | None = Header( + default=None, + alias="Idempotency-Key", + max_length=128, + description="幂等键,启用后重复请求回放缓存响应", + ), db: AsyncSession = Depends(get_db), current_user: User = Depends(get_admin_user), ) -> dict[str, Any]: @@ -239,6 +255,7 @@ async def update_task( 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) @@ -248,12 +265,22 @@ async def update_task( @scheduler_router.delete("/tasks/{task_id}", response_model=dict) async def delete_task( task_id: str, + idempotency_key: str | None = Header( + default=None, + alias="Idempotency-Key", + max_length=128, + description="幂等键,启用后重复请求回放缓存响应", + ), 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) + input_dto = DeleteTaskInput( + task_id=task_id, + updated_by=current_user.uid, + idempotency_key=idempotency_key, + ) await service.delete_task(input_dto) return {"success": True, "data": {"task_id": task_id, "status": "deleted"}} @@ -261,12 +288,22 @@ async def delete_task( @scheduler_router.post("/tasks/{task_id}/pause", response_model=dict) async def pause_task( task_id: str, + idempotency_key: str | None = Header( + default=None, + alias="Idempotency-Key", + max_length=128, + description="幂等键,启用后重复请求回放缓存响应", + ), 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) + input_dto = PauseTaskInput( + task_id=task_id, + updated_by=current_user.uid, + idempotency_key=idempotency_key, + ) output = await service.pause_task(input_dto) return {"success": True, "data": output.model_dump()} @@ -274,6 +311,12 @@ async def pause_task( @scheduler_router.post("/tasks/{task_id}/resume", response_model=dict) async def resume_task( task_id: str, + idempotency_key: str | None = Header( + default=None, + alias="Idempotency-Key", + max_length=128, + description="幂等键,启用后重复请求回放缓存响应", + ), db: AsyncSession = Depends(get_db), current_user: User = Depends(get_admin_user), ) -> dict[str, Any]: @@ -282,7 +325,11 @@ async def resume_task( ``paused`` -> ``active`` / ``dead_letter`` -> ``active``。 """ service = create_scheduler_service(db) - input_dto = ResumeTaskInput(task_id=task_id, updated_by=current_user.uid) + input_dto = ResumeTaskInput( + task_id=task_id, + updated_by=current_user.uid, + idempotency_key=idempotency_key, + ) output = await service.resume_task(input_dto) return {"success": True, "data": output.model_dump()} @@ -291,6 +338,12 @@ async def resume_task( async def trigger_task( task_id: str, body: TriggerTaskRequest | None = None, + idempotency_key: str | None = Header( + default=None, + alias="Idempotency-Key", + max_length=128, + description="幂等键,启用后重复请求回放缓存响应", + ), db: AsyncSession = Depends(get_db), current_user: User = Depends(get_admin_user), ) -> dict[str, Any]: @@ -298,6 +351,7 @@ async def trigger_task( 生成 ``run_id`` 后入队 ARQ 执行,不影响下一次自动执行。 可通过 body 传入 payload 覆盖任务定义中的 payload,None 时使用任务原 payload。 + 触发人 ``current_user.uid`` 透传到 worker 写入 ``run_log.created_by`` 审计字段。 """ arq_pool = await get_arq_pool() service = create_scheduler_service(db, arq_pool=arq_pool) @@ -306,6 +360,8 @@ async def trigger_task( 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 {"success": True, "data": output.model_dump()}