diff --git a/backend/package/yuxi/scheduler/adapters/persistence/mappers.py b/backend/package/yuxi/scheduler/adapters/persistence/mappers.py index 37c147fb..e89dccfa 100644 --- a/backend/package/yuxi/scheduler/adapters/persistence/mappers.py +++ b/backend/package/yuxi/scheduler/adapters/persistence/mappers.py @@ -124,10 +124,12 @@ def handler_summary_row_to_dataclass(row: Any) -> HandlerSummaryDC: """仓储层 ``HandlerSummary`` → core ``HandlerSummary`` dataclass,按字段名映射。 仓储层 ``HandlerSummary`` 是 dataclass,row 可能是 NamedTuple 或 dataclass 实例, - 统一按字段名读取并转换为 core dataclass。 + 统一按字段名读取并转换为 core dataclass。``active_task_count`` 在仓储层聚合 + 时提供,缺失时默认 0。 """ return HandlerSummaryDC( name=row.name, task_count=int(row.task_count), + active_task_count=int(row.active_task_count) if hasattr(row, "active_task_count") else 0, last_active_at=row.last_active_at, ) diff --git a/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_idempotency_repo.py b/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_idempotency_repo.py index 2e6aeaa9..32c9317e 100644 --- a/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_idempotency_repo.py +++ b/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_idempotency_repo.py @@ -78,9 +78,9 @@ class SqlAlchemyIdempotencyRepository: # 幂等查询 # ------------------------------------------------------------------ - async def get(self, key: str) -> dict | None: + async def get(self, key: str, *, operation: str | None = None) -> dict | None: """委托给 ScheduledTaskIdempotencyRepository.get。""" - return await self._repo.get(key) + return await self._repo.get(key, operation=operation) # ------------------------------------------------------------------ # 数据生命周期 diff --git a/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_task_repo.py b/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_task_repo.py index dd64f820..51670ea1 100644 --- a/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_task_repo.py +++ b/backend/package/yuxi/scheduler/adapters/persistence/sqlalchemy_task_repo.py @@ -157,11 +157,13 @@ class SqlAlchemyTaskRepository: *, owner_scope: str | None = None, owner_id: str | None = None, + handler_name: str | None = None, ) -> dict[str, int]: """委托给 ScheduledTaskRepository.count_by_status。""" return await self._repo.count_by_status( owner_scope=owner_scope, owner_id=owner_id, + handler_name=handler_name, ) # ------------------------------------------------------------------ @@ -321,9 +323,14 @@ class SqlAlchemyTaskRepository: ) return [orm_to_task(orm) for orm in orms] - async def count_deleted(self) -> int: + async def count_deleted( + self, + *, + start_time: datetime | None = None, + end_time: datetime | None = None, + ) -> int: """委托给 BaseRepository.count_deleted。""" - return await self._repo.count_deleted() + return await self._repo.count_deleted(start_time=start_time, end_time=end_time) async def hard_delete_by_id( self, diff --git a/backend/package/yuxi/scheduler/core/models.py b/backend/package/yuxi/scheduler/core/models.py index f09d055c..0e658cfb 100644 --- a/backend/package/yuxi/scheduler/core/models.py +++ b/backend/package/yuxi/scheduler/core/models.py @@ -82,11 +82,12 @@ class ScheduledTaskRunLog: class HandlerSummary: """handler 聚合摘要(与仓储层 dataclass 同名同字段,core 层独立声明)。 - 用于 Admin API 展示各 handler 的任务数与最近执行时间。 + 用于 Admin API 展示各 handler 的任务数、活跃任务数与最近执行时间。 """ name: str task_count: int + active_task_count: int = 0 last_active_at: Any = None diff --git a/backend/package/yuxi/scheduler/core/ports.py b/backend/package/yuxi/scheduler/core/ports.py index 9cbb3b3e..a87780dc 100644 --- a/backend/package/yuxi/scheduler/core/ports.py +++ b/backend/package/yuxi/scheduler/core/ports.py @@ -75,7 +75,12 @@ class BaseRepositoryPort(Protocol): end_time: datetime | None = None, ) -> list[Any]: ... - async def count_deleted(self) -> int: ... + async def count_deleted( + self, + *, + start_time: datetime | None = None, + end_time: datetime | None = None, + ) -> int: ... async def hard_delete_by_id( self, @@ -169,6 +174,7 @@ class ScheduledTaskRepositoryPort(Protocol): *, owner_scope: str | None = None, owner_id: str | None = None, + handler_name: str | None = None, ) -> dict[str, int]: ... # 并发控制 @@ -348,6 +354,6 @@ class ScheduledTaskIdempotencyRepositoryPort(Protocol): updated_by: str | None = None, ) -> None: ... - async def get(self, key: str) -> dict | None: ... + async def get(self, key: str, *, operation: str | None = None) -> dict | None: ... async def cleanup_old(self, before: datetime, *, commit: bool = True) -> int: ... diff --git a/backend/package/yuxi/scheduler/use_cases/dto/scheduler.py b/backend/package/yuxi/scheduler/use_cases/dto/scheduler.py index 1b829119..55452ff2 100644 --- a/backend/package/yuxi/scheduler/use_cases/dto/scheduler.py +++ b/backend/package/yuxi/scheduler/use_cases/dto/scheduler.py @@ -275,6 +275,19 @@ class TaskOutput(BaseModel): deleted_at: str | None = None +class TriggerTaskOutput(BaseModel): + """手动触发任务输出 DTO。 + + 返回本次执行标识 ``run_id`` 与任务快照 ``task``,便于调用方立即追踪 + 执行日志与结果。 + """ + + model_config = ConfigDict(frozen=True) + + run_id: str = Field(..., description="本次执行唯一标识") + task: TaskOutput = Field(..., description="被触发任务的快照") + + class ListTasksInput(BaseModel): """列出任务输入 DTO。 @@ -297,7 +310,11 @@ class ListTasksInput(BaseModel): class ListTasksOutput(BaseModel): - """列出任务输出 DTO。""" + """列出任务输出 DTO。 + + ``total_pages`` 为派生字段,由 ``total / page_size`` 上取整得到, + 便于前端分页组件直接消费。 + """ model_config = ConfigDict(frozen=True) @@ -305,6 +322,7 @@ class ListTasksOutput(BaseModel): total: int = 0 page: int = 1 page_size: int = 20 + total_pages: int = 1 # --------------------------------------------------------------------------- @@ -378,7 +396,10 @@ class RunLogOutput(BaseModel): class ListRunLogsOutput(BaseModel): - """列出执行日志输出 DTO。""" + """列出执行日志输出 DTO。 + + ``total_pages`` 为派生字段,由 ``total / page_size`` 上取整得到。 + """ model_config = ConfigDict(frozen=True) @@ -386,6 +407,7 @@ class ListRunLogsOutput(BaseModel): total: int = 0 page: int = 1 page_size: int = 20 + total_pages: int = 1 class ListAllRunLogsInput(BaseModel): @@ -445,6 +467,19 @@ class ListHandlerSummaryOutput(BaseModel): items: list[HandlerSummaryOutput] = Field(default_factory=list) +class CountByStatusInput(BaseModel): + """任务状态计数输入 DTO。 + + 过滤字段与 ``ListTasksInput`` 对齐,便于仪表盘按 owner/handler 下钻。 + """ + + model_config = ConfigDict(frozen=True) + + owner_scope: str | None = Field(default=None, max_length=64) + owner_id: str | None = Field(default=None, max_length=128) + handler_name: str | None = Field(default=None, max_length=128) + + class CountByStatusOutput(BaseModel): """任务状态计数输出 DTO。 @@ -462,12 +497,19 @@ class ListUpcomingInput(BaseModel): """即将执行任务查询输入 DTO。 对齐 PRD §FR-ST-09 统计查询,按 ``next_run_at`` 升序返回即将执行的任务列表。 + ``hours_ahead`` 控制预览窗口(默认 24 小时,最大 7 天)。 """ model_config = ConfigDict(frozen=True) limit: int = Field(default=100, ge=1, le=1000, description="返回条数上限") handler_name: str | None = Field(default=None, max_length=128) + hours_ahead: int = Field( + default=24, + ge=1, + le=168, + description="预览窗口小时数(默认 24,最大 168=7 天)", + ) class ListUpcomingOutput(BaseModel): @@ -603,18 +645,28 @@ class ListDeletedTasksInput(BaseModel): class RestoreTaskInput(BaseModel): - """恢复已删除任务输入 DTO。""" + """恢复已删除任务输入 DTO。支持幂等键以避免恢复/硬删除重复执行。""" model_config = ConfigDict(frozen=True) task_id: str = Field(..., min_length=1, max_length=64) updated_by: str | None = Field(default=None, max_length=128) + idempotency_key: str | None = Field( + default=None, + max_length=128, + description="客户端传入的 Idempotency-Key,None 时不启用幂等保护", + ) class HardDeleteTaskInput(BaseModel): - """硬删除任务输入 DTO。""" + """硬删除任务输入 DTO。支持幂等键以避免重复物理删除。""" model_config = ConfigDict(frozen=True) task_id: str = Field(..., min_length=1, max_length=64) updated_by: str | None = Field(default=None, max_length=128) + idempotency_key: str | None = Field( + default=None, + max_length=128, + description="客户端传入的 Idempotency-Key,None 时不启用幂等保护", + ) diff --git a/backend/package/yuxi/scheduler/use_cases/mappers.py b/backend/package/yuxi/scheduler/use_cases/mappers.py index f64b207a..8e26da71 100644 --- a/backend/package/yuxi/scheduler/use_cases/mappers.py +++ b/backend/package/yuxi/scheduler/use_cases/mappers.py @@ -31,6 +31,7 @@ from yuxi.scheduler.use_cases.dto.scheduler import ( ListUpcomingOutput, RunLogOutput, TaskOutput, + TriggerTaskOutput, ) from yuxi.utils.datetime_utils import format_utc_datetime @@ -43,6 +44,7 @@ __all__ = [ "to_list_daily_stats_output", "to_run_log_output", "to_task_output", + "to_trigger_task_output", ] @@ -80,6 +82,11 @@ def to_task_output(task: ScheduledTask) -> TaskOutput: ) +def to_trigger_task_output(task: ScheduledTask, run_id: str) -> TriggerTaskOutput: + """``ScheduledTask`` dataclass + run_id → ``TriggerTaskOutput`` DTO。""" + return TriggerTaskOutput(run_id=run_id, task=to_task_output(task)) + + def to_run_log_output(log: ScheduledTaskRunLog) -> RunLogOutput: """``ScheduledTaskRunLog`` dataclass → ``RunLogOutput`` DTO。 @@ -113,12 +120,14 @@ def to_run_log_output(log: ScheduledTaskRunLog) -> RunLogOutput: def to_handler_summary_output(summary: HandlerSummary) -> HandlerSummaryOutput: """``HandlerSummary`` dataclass → ``HandlerSummaryOutput`` DTO。 - ``description`` 字段仓储层不聚合(DB 无此列),由调用方按需补充; - ``active_task_count`` 同理,默认 0,由调用方按需补充。 + ``description`` 字段仓储层不聚合(DB 无此列),由调用方(如 worker 进程 + 持有 ``HandlerRegistry`` 时)按需补充;API 进程无注册表时保持空字符串。 + ``active_task_count`` 已由仓储层聚合,此处直接透传。 """ return HandlerSummaryOutput( handler_name=summary.name, task_count=summary.task_count, + active_task_count=summary.active_task_count, last_active_at=format_utc_datetime(summary.last_active_at), ) diff --git a/backend/package/yuxi/scheduler/use_cases/ports/scheduler_service_port.py b/backend/package/yuxi/scheduler/use_cases/ports/scheduler_service_port.py index f2efed46..c1e5e2ba 100644 --- a/backend/package/yuxi/scheduler/use_cases/ports/scheduler_service_port.py +++ b/backend/package/yuxi/scheduler/use_cases/ports/scheduler_service_port.py @@ -56,6 +56,7 @@ if TYPE_CHECKING: RunLogOutput, TaskOutput, TriggerTaskInput, + TriggerTaskOutput, UpdateTaskInput, ) @@ -145,7 +146,7 @@ class SchedulerServicePort(Protocol): """ ... - async def trigger_task(self, input_dto: TriggerTaskInput) -> TaskOutput: + async def trigger_task(self, input_dto: TriggerTaskInput) -> TriggerTaskOutput: """手动触发任务(FR-ST-05)。 生成 ``run_id`` 后入队执行,不入正常 ``next_run_at`` 流,不影响下一次 diff --git a/backend/package/yuxi/scheduler/use_cases/services/scheduler_service.py b/backend/package/yuxi/scheduler/use_cases/services/scheduler_service.py index 3af16abc..8e8ecdbd 100644 --- a/backend/package/yuxi/scheduler/use_cases/services/scheduler_service.py +++ b/backend/package/yuxi/scheduler/use_cases/services/scheduler_service.py @@ -31,6 +31,7 @@ from yuxi.scheduler.core.contracts import ( TaskResult, UnitOfWork, ) +from yuxi.scheduler.core.models import ScheduledTask from yuxi.scheduler.core.validators import ( validate_cron_expression, validate_owner_scope, @@ -46,6 +47,7 @@ from yuxi.scheduler.exceptions import ( TaskStatusTransitionError, ) from yuxi.scheduler.use_cases.dto.scheduler import ( + CountByStatusInput, CountByStatusOutput, CreateTaskInput, DeleteTaskInput, @@ -72,6 +74,7 @@ from yuxi.scheduler.use_cases.dto.scheduler import ( RunLogOutput, TaskOutput, TriggerTaskInput, + TriggerTaskOutput, UpdateTaskInput, ) from yuxi.scheduler.use_cases.mappers import ( @@ -82,6 +85,7 @@ from yuxi.scheduler.use_cases.mappers import ( to_list_upcoming_output, to_run_log_output, to_task_output, + to_trigger_task_output, ) from yuxi.scheduler.use_cases.ports.scheduler_service_port import SchedulerServicePort from yuxi.scheduler.use_cases.utils.backoff import calc_backoff_seconds @@ -225,22 +229,48 @@ class SchedulerService(SchedulerServicePort): """ return now_utc_naive.replace(tzinfo=UTC).astimezone(_DAILY_STAT_TZ).date() + @staticmethod + def _total_pages(total: int, page_size: int) -> int: + """根据总数与每页条数计算总页数(至少 1 页)。""" + if total <= 0: + return 1 + return (total + page_size - 1) // page_size + + @staticmethod + def _build_effective_task(task: Any, update_data: dict[str, Any]) -> Any: + """把待更新字段覆盖到原任务对象上,得到用于重算 next_run_at 的生效视图。 + + 在 ``update_task`` 的 ``enabled=True & status=paused`` 分支中,同一请求可能 + 同时修改 ``cron_expression`` / ``tz`` / ``run_at``,若直接用原 ``task`` 重算 + ``next_run_at`` 会使用时区/表达式旧值。本方法生成一个只读生效视图,确保 + 重算结果与本次更新意图一致。 + """ + effective = dict(task.__dict__) + effective.update(update_data) + return ScheduledTask(**effective) + # ================================================================== # 幂等辅助 # ================================================================== - async def _check_idempotency_replay(self, key: str | None) -> dict | None: + async def _check_idempotency_replay( + self, + key: str | None, + *, + operation: str, + ) -> dict | None: """查询幂等缓存响应。 - ``key`` 为 None 时直接返回 None(未启用幂等);否则查询幂等仓储, - 命中返回已缓存的 response_body dict,未命中(首次请求或 in-flight)返回 None。 + ``key`` 为 None 时直接返回 None(未启用幂等);否则按 ``(key, operation)`` + 查询幂等仓储,命中且 operation 一致时返回已缓存的 response_body dict, + 未命中(首次请求、in-flight 或 operation 不匹配)返回 None。 - 幂等仓储的 ``get`` 方法在记录不存在或 response_body 为 NULL 时均返回 - None,调用方需通过 ``acquire`` 区分"首次"与"in-flight 重复"。 + 按 operation 隔离不同 API 操作的幂等缓存,避免 create 的响应被 update + 等其它操作错误重放。 """ if key is None: return None - return await self._repos.idempotency.get(key) + return await self._repos.idempotency.get(key, operation=operation) async def _acquire_idempotency( self, @@ -323,7 +353,10 @@ class SchedulerService(SchedulerServicePort): 首次请求抢占 key 后执行业务并回写响应体,重复请求回放缓存响应。 """ # 0. 幂等回放 - cached = await self._check_idempotency_replay(input_dto.idempotency_key) + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="create", + ) if cached is not None: return TaskOutput(**cached) @@ -378,7 +411,7 @@ class SchedulerService(SchedulerServicePort): "delete_after_run": input_dto.delete_after_run, "block_strategy": input_dto.block_strategy, "stagger_seconds": stagger_seconds, - "status": "active", + "status": "active" if input_dto.enabled else "paused", "next_run_at": next_run_at, "created_by": input_dto.created_by, } @@ -434,17 +467,25 @@ class SchedulerService(SchedulerServicePort): total=total, page=input_dto.page, page_size=input_dto.page_size, + total_pages=self._total_pages(total, input_dto.page_size), ) async def update_task(self, input_dto: UpdateTaskInput) -> TaskOutput: """更新任务(FR-ST-05)。 - 修改 ``cron_expression`` 时重新计算 ``next_run_at`` 与 ``stagger_seconds``。 + 修改 ``cron_expression`` / ``tz`` / ``run_at`` 时重新计算 ``next_run_at``; + ``stagger_seconds`` 仅依赖 ``task_id`` 与 ``schedule_kind``,不在更新时重算。 ``schedule_kind`` 与 ``handler_name`` 不可修改(仓储 ``update`` 已排除)。 ``idempotency_key`` 非空时启用幂等保护。 + + 显式传 ``null`` 的字段(``tz`` / ``cron_expression`` / ``run_at`` / ``payload``) + 会被拒绝,避免非法清空业务必填字段。 """ # 0. 幂等回放 - cached = await self._check_idempotency_replay(input_dto.idempotency_key) + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="update", + ) if cached is not None: return TaskOutput(**cached) @@ -455,26 +496,33 @@ class SchedulerService(SchedulerServicePort): details={"task_id": input_dto.task_id}, ) - update_data: dict[str, Any] = input_dto.model_dump(exclude_unset=True, exclude={"task_id", "idempotency_key"}) + update_data: dict[str, Any] = input_dto.model_dump( + exclude_unset=True, exclude={"task_id", "idempotency_key"} + ) if not update_data: return to_task_output(task) + # 不允许显式清空业务必填字段 + for field in ("cron_expression", "run_at", "tz", "payload"): + if field in update_data and update_data[field] is None: + raise SchedulerValidationError(f"{field} 不允许显式置空") + # payload 校验 - if "payload" in update_data and update_data["payload"] is not None: + if "payload" in update_data: validate_payload_size( update_data["payload"], self._config.scheduler_max_payload_size_kb, ) - # tz 校验(方案 D):传入 tz 时校验 IANA 时区合法性 - if "tz" in update_data and update_data["tz"] is not None: + # tz 校验:传入 tz 时校验 IANA 时区合法性 + if "tz" in update_data: validate_tz(update_data["tz"]) - # owner_scope 校验(方案 F1):传入 owner_scope 时校验枚举值 - if "owner_scope" in update_data and update_data["owner_scope"] is not None: + # owner_scope 校验:传入 owner_scope 时校验枚举值 + if "owner_scope" in update_data: validate_owner_scope(update_data["owner_scope"]) - # cron 类型拒绝 delete_after_run=True(方案 F2):仅对 at 类型有意义 + # cron 类型拒绝 delete_after_run=True:仅对 at 类型有意义 effective_schedule_kind = task.schedule_kind if ( "delete_after_run" in update_data @@ -484,39 +532,45 @@ class SchedulerService(SchedulerServicePort): raise SchedulerValidationError("cron 任务不支持 delete_after_run=True") # schedule_kind 与字段匹配校验:cron 任务不允许更新 run_at,at 任务不允许更新 cron_expression - if task.schedule_kind == "cron" and update_data.get("run_at") is not None: + if task.schedule_kind == "cron" and "run_at" in update_data: raise SchedulerValidationError("cron 任务不支持更新 run_at") - if task.schedule_kind == "at" and update_data.get("cron_expression") is not None: + if task.schedule_kind == "at" and "cron_expression" in update_data: raise SchedulerValidationError("at 任务不支持更新 cron_expression") - # cron_expression 变更时重算 next_run_at(方案 F3:stagger_seconds 仅依赖 - # task_id 与 schedule_kind,与 cron_expression 无关,不在此重算) - if "cron_expression" in update_data and update_data["cron_expression"] is not None: + # cron_expression 变更时重算 next_run_at(stagger_seconds 仅依赖 task_id + # 与 schedule_kind,与 cron_expression 无关,不在此重算) + if "cron_expression" in update_data: validate_cron_expression( update_data["cron_expression"], min_interval_minutes=self._config.scheduler_min_cron_interval_minutes, ) effective_tz = update_data.get("tz", task.tz) - update_data["next_run_at"] = calc_next_run_at(update_data["cron_expression"], tz=effective_tz) + update_data["next_run_at"] = calc_next_run_at( + update_data["cron_expression"], tz=effective_tz + ) # run_at 变更时重算 next_run_at(at 类型),校验未来时间 - if "run_at" in update_data and update_data["run_at"] is not None and task.schedule_kind == "at": - run_at_dt = self._parse_iso_to_utc_naive(update_data["run_at"], field_name="run_at") + if "run_at" in update_data and task.schedule_kind == "at": + run_at_dt = self._parse_iso_to_utc_naive( + update_data["run_at"], field_name="run_at" + ) if run_at_dt <= utc_now_naive(): raise SchedulerValidationError("run_at 必须为未来时间") update_data["next_run_at"] = run_at_dt - # enabled 联动 status(方案 E): + # enabled 联动 status: # - enabled=False 且 status=active → 同步 status=paused(避免 tick 继续触发) - # - enabled=True 且 status=paused → 同步 status=active 并重算 next_run_at + # - enabled=True 且 status=paused → 同步 status=active 并按本次更新后的 + # cron_expression / tz / run_at 重算 next_run_at,保证时区/表达式一致 # - dead_letter 状态拒绝通过 enabled 直接启用,必须通过 resume_task 路径恢复 - if "enabled" in update_data and update_data["enabled"] is not None: + if "enabled" in update_data: new_enabled = update_data["enabled"] if new_enabled is False and task.status == "active": update_data["status"] = "paused" elif new_enabled is True and task.status == "paused": + effective_task = self._build_effective_task(task, update_data) update_data["status"] = "active" - update_data["next_run_at"] = self._calc_next_run_at(task) + update_data["next_run_at"] = self._calc_next_run_at(effective_task) elif new_enabled is True and task.status == "dead_letter": raise SchedulerValidationError( "dead_letter 任务需通过 resume_task 恢复,不能通过 enabled 直接启用", @@ -529,7 +583,9 @@ class SchedulerService(SchedulerServicePort): operation="update", operator=input_dto.updated_by, ) - updated = await self._repos.task.update(input_dto.task_id, update_data, commit=False) + updated = await self._repos.task.update( + input_dto.task_id, update_data, commit=False + ) if updated is not None: await self._store_idempotency_response( input_dto.idempotency_key, @@ -553,7 +609,10 @@ class SchedulerService(SchedulerServicePort): 不包含领域快照,避免回放时还原已删除数据造成混淆。 """ # 0. 幂等回放:命中说明删除已完成,直接返回(无返回值) - cached = await self._check_idempotency_replay(input_dto.idempotency_key) + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="delete", + ) if cached is not None: return @@ -594,7 +653,10 @@ class SchedulerService(SchedulerServicePort): (非 ``active`` 任务暂停)抛 ``TaskStatusTransitionError``,不回写响应体。 """ # 0. 幂等回放 - cached = await self._check_idempotency_replay(input_dto.idempotency_key) + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="pause", + ) if cached is not None: return TaskOutput(**cached) @@ -642,7 +704,10 @@ class SchedulerService(SchedulerServicePort): ``idempotency_key`` 非空时启用幂等保护。 """ # 0. 幂等回放 - cached = await self._check_idempotency_replay(input_dto.idempotency_key) + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="resume", + ) if cached is not None: return TaskOutput(**cached) @@ -700,22 +765,25 @@ class SchedulerService(SchedulerServicePort): ) return to_task_output(task) - async def trigger_task(self, input_dto: TriggerTaskInput) -> TaskOutput: + async def trigger_task(self, input_dto: TriggerTaskInput) -> TriggerTaskOutput: """手动触发任务(FR-ST-05)。 生成 ``run_id`` 后入队执行,不入正常 ``next_run_at`` 流,不影响下一次 自动执行。死信任务允许手动触发但不重置 ``status``。``input_dto.operator`` 为触发人 uid,透传到 worker 写入 ``run_log.created_by`` 审计字段。 ``idempotency_key`` 非空时启用幂等保护:首次请求抢占 key 后入队, - 重复请求回放缓存的 TaskOutput。因入队与 acquire 不在同一事务 + 重复请求回放缓存的 ``TriggerTaskOutput``。因入队与 acquire 不在同一事务 (ARQ 无事务语义),入队失败时 acquire 已提交的幂等记录由 cleanup_old 物理删除回收(24 小时 TTL),客户端可用相同 key 重试但会收到 409 直到记录被清理——此为可接受的弱一致性边界。 """ # 0. 幂等回放 - cached = await self._check_idempotency_replay(input_dto.idempotency_key) + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="trigger", + ) if cached is not None: - return TaskOutput(**cached) + return TriggerTaskOutput(**cached) task = await self._repos.task.get_by_task_id(input_dto.task_id) if task is None: @@ -748,7 +816,7 @@ class SchedulerService(SchedulerServicePort): ) # 入队成功后回写幂等响应体(独立事务) - response = to_task_output(task).model_dump() + response = to_trigger_task_output(task, run_id).model_dump() await self._store_idempotency_response( input_dto.idempotency_key, response, @@ -764,7 +832,7 @@ class SchedulerService(SchedulerServicePort): "operator": input_dto.operator, }, ) - return to_task_output(task) + return to_trigger_task_output(task, run_id) # ================================================================== # 查询 @@ -792,6 +860,7 @@ class SchedulerService(SchedulerServicePort): total=total, page=input_dto.page, page_size=input_dto.page_size, + total_pages=self._total_pages(total, input_dto.page_size), ) async def get_run_log(self, input_dto: GetRunLogInput) -> RunLogOutput: @@ -848,6 +917,7 @@ class SchedulerService(SchedulerServicePort): total=total, page=input_dto.page, page_size=input_dto.page_size, + total_pages=self._total_pages(total, input_dto.page_size), ) async def list_handler_summary(self) -> ListHandlerSummaryOutput: @@ -859,15 +929,19 @@ class SchedulerService(SchedulerServicePort): summaries = await self._repos.task.list_handler_summary() return to_list_handler_summary_output(summaries) - async def count_by_status(self) -> CountByStatusOutput: + async def count_by_status(self, input_dto: CountByStatusInput) -> CountByStatusOutput: """按状态计数任务(FR-ST-07)。""" - dist = await self._repos.task.count_by_status() + dist = await self._repos.task.count_by_status( + owner_scope=input_dto.owner_scope, + owner_id=input_dto.owner_id, + handler_name=input_dto.handler_name, + ) return to_count_by_status_output(dist) async def list_upcoming(self, input_dto: ListUpcomingInput) -> ListUpcomingOutput: """列出即将执行的任务(FR-ST-09)。""" now = utc_now_naive() - before = now + timedelta(hours=24) + before = now + timedelta(hours=input_dto.hours_ahead) tasks = await self._repos.task.list_upcoming( after=now, before=before, @@ -1400,8 +1474,9 @@ class SchedulerService(SchedulerServicePort): async def list_deleted_tasks(self, input_dto: ListDeletedTasksInput) -> ListTasksOutput: """列出已删除任务(回收站,FR-ST-05)。 - ``list_deleted`` 不返回总数,``total`` 取当前页条数(客户端可按 - ``len(items) < page_size`` 判断是否还有下一页)。 + 返回软删除的任务列表,按 ``deleted_at`` 降序排序,支持时间范围过滤与真实分页。 + ``total`` 由 ``count_deleted`` 按相同过滤条件统计,确保前端分页组件可正确展示 + 总页数。 """ start_time: datetime | None = None end_time: datetime | None = None @@ -1412,25 +1487,38 @@ class SchedulerService(SchedulerServicePort): limit = input_dto.page_size offset = (input_dto.page - 1) * input_dto.page_size - tasks = await self._repos.task.list_deleted( - limit=limit, - offset=offset, - start_time=start_time, - end_time=end_time, + tasks, total = await asyncio.gather( + self._repos.task.list_deleted( + limit=limit, + offset=offset, + start_time=start_time, + end_time=end_time, + ), + self._repos.task.count_deleted(start_time=start_time, end_time=end_time), ) return ListTasksOutput( items=[to_task_output(t) for t in tasks], - total=len(tasks), + total=total, page=input_dto.page, page_size=input_dto.page_size, + total_pages=self._total_pages(total, input_dto.page_size), ) async def restore_task(self, input_dto: RestoreTaskInput) -> TaskOutput: """恢复已删除任务(FR-ST-05)。 将软删除的任务恢复为 ``is_deleted=0``。若恢复前 ``status='active'``, - 自动置为 ``paused`` 避免立即被 tick 扫描执行。 + 自动置为 ``paused`` 避免立即被 tick 扫描执行。``idempotency_key`` + 非空时启用幂等保护:首次请求抢占 key 后执行业务并回写响应体, + 重复请求回放缓存响应。 """ + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="restore", + ) + if cached is not None: + return TaskOutput(**cached) + deleted_task = await self._repos.task.get_deleted_by_task_id(input_dto.task_id) if deleted_task is None: raise SchedulerEntityNotFoundError( @@ -1441,6 +1529,12 @@ class SchedulerService(SchedulerServicePort): was_active = deleted_task.status == "active" async with self._uow: + await self._acquire_idempotency( + input_dto.idempotency_key, + task_id=input_dto.task_id, + operation="restore", + operator=input_dto.updated_by, + ) await self._repos.task.restore_by_id( deleted_task.id, updated_by=input_dto.updated_by, @@ -1452,15 +1546,35 @@ class SchedulerService(SchedulerServicePort): updated_by=input_dto.updated_by, commit=False, ) + restored = await self._repos.task.get_by_task_id(input_dto.task_id) + if restored is not None: + await self._store_idempotency_response( + input_dto.idempotency_key, + to_task_output(restored).model_dump(), + operator=input_dto.updated_by, + ) - restored = await self._repos.task.get_by_task_id(input_dto.task_id) + if restored is None: + raise SchedulerEntityNotFoundError( + f"任务 {input_dto.task_id!r} 恢复后查询失败", + details={"task_id": input_dto.task_id}, + ) return to_task_output(restored) async def hard_delete_task(self, input_dto: HardDeleteTaskInput) -> None: """硬删除任务(不可恢复,FR-ST-05)。 - 物理删除已软删除的任务记录,操作不可逆。 + 物理删除已软删除的任务记录,操作不可逆。``idempotency_key`` 非空时 + 启用幂等保护:首次请求抢占 key 后执行删除并回写标记响应,重复请求 + 直接返回(无副作用)。 """ + cached = await self._check_idempotency_replay( + input_dto.idempotency_key, + operation="hard_delete", + ) + if cached is not None: + return + deleted_task = await self._repos.task.get_deleted_by_task_id(input_dto.task_id) if deleted_task is None: raise SchedulerEntityNotFoundError( @@ -1468,7 +1582,20 @@ class SchedulerService(SchedulerServicePort): details={"task_id": input_dto.task_id}, ) - await self._repos.task.hard_delete_by_id(deleted_task.id) + async with self._uow: + await self._acquire_idempotency( + input_dto.idempotency_key, + task_id=input_dto.task_id, + operation="hard_delete", + operator=input_dto.updated_by, + ) + await self._repos.task.hard_delete_by_id(deleted_task.id, commit=False) + await self._store_idempotency_response( + input_dto.idempotency_key, + {"task_id": input_dto.task_id, "status": "hard_deleted"}, + operator=input_dto.updated_by, + ) + logger.info( "scheduler_task_hard_deleted", extra={"task_id": input_dto.task_id, "operator": input_dto.updated_by},