feat(scheduler router): 为定时任务接口添加幂等性支持和审计字段

1. 新增Header参数Idempotency-Key实现幂等校验
2. 为create/update/delete/pause/resume/trigger接口补充幂等键参数
3. 为trigger任务接口添加操作人审计字段和透传逻辑
This commit is contained in:
Kris 2026-07-11 05:45:41 +08:00
parent 6d3cf3fedd
commit 1e0f4b3b14

View File

@ -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 覆盖任务定义中的 payloadNone 时使用任务原 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()}