From 23f7a8a863ad0f7f7870ee0f11031535e5373dc3 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Sat, 11 Jul 2026 05:40:36 +0800 Subject: [PATCH] =?UTF-8?q?refactor(scheduled-repo):=20=E4=BC=98=E5=8C=96?= =?UTF-8?q?=E4=BB=A3=E7=A0=81=E6=A0=BC=E5=BC=8F=E4=B8=8E=E5=87=BD=E6=95=B0?= =?UTF-8?q?=E7=AD=BE=E5=90=8D=EF=BC=8C=E8=A1=A5=E5=85=85=E5=AE=A1=E8=AE=A1?= =?UTF-8?q?=E5=AD=97=E6=AE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 格式化SQLAlchemy查询语句,将多行合并为单行提升可读性 2. 调整部分函数参数换行格式,统一代码风格 3. 为幂等键抢占、响应回写、任务状态变更等方法新增updated_by/created_by审计字段 4. 修复僵尸日志清理逻辑,补全日聚合统计写入路径并返回handler列表 5. 新增一次性任务执行成功后的收尾方法,避免调度循环 --- .../yuxi/repositories/scheduled/base.py | 22 +-- .../scheduled_task_idempotency_repository.py | 32 +++- .../scheduled/scheduled_task_repository.py | 152 +++++++++++++----- ...scheduled_task_run_log_daily_repository.py | 17 +- .../scheduled_task_run_log_repository.py | 43 +++-- 5 files changed, 168 insertions(+), 98 deletions(-) diff --git a/backend/package/yuxi/repositories/scheduled/base.py b/backend/package/yuxi/repositories/scheduled/base.py index bd4d2558..c19a10ee 100644 --- a/backend/package/yuxi/repositories/scheduled/base.py +++ b/backend/package/yuxi/repositories/scheduled/base.py @@ -23,7 +23,9 @@ from __future__ import annotations from datetime import datetime -from sqlalchemy import delete as sa_delete, func, select, update as sa_update +from sqlalchemy import delete as sa_delete +from sqlalchemy import func, select +from sqlalchemy import update as sa_update from sqlalchemy.ext.asyncio import AsyncSession from yuxi.utils.datetime_utils import utc_now_naive @@ -104,11 +106,7 @@ class BaseRepository: 用于恢复前的记录查询与唯一约束预检,避免 list_deleted(limit=1) 的误匹配。 """ - stmt = ( - select(self.model) - .where(self.model.id == record_id) - .where(self.model.is_deleted == 1) - ) + stmt = select(self.model).where(self.model.id == record_id).where(self.model.is_deleted == 1) result = await self.db.execute(stmt) return result.scalars().first() @@ -170,11 +168,7 @@ class BaseRepository: commit: bool = True, ) -> int: """按主键物理删除记录(仅删除 is_deleted=1 的记录)。返回受影响行数。""" - stmt = ( - sa_delete(self.model) - .where(self.model.id == record_id) - .where(self.model.is_deleted == 1) - ) + stmt = sa_delete(self.model).where(self.model.id == record_id).where(self.model.is_deleted == 1) result = await self.db.execute(stmt) if commit: await self.db.commit() @@ -189,11 +183,7 @@ class BaseRepository: commit: bool = True, ) -> int: """物理删除 deleted_at 早于 cutoff_at 的软删除记录。返回受影响行数。""" - stmt = ( - sa_delete(self.model) - .where(self.model.is_deleted == 1) - .where(self.model.deleted_at < cutoff_at) - ) + stmt = sa_delete(self.model).where(self.model.is_deleted == 1).where(self.model.deleted_at < cutoff_at) result = await self.db.execute(stmt) if commit: await self.db.commit() diff --git a/backend/package/yuxi/repositories/scheduled/scheduled_task_idempotency_repository.py b/backend/package/yuxi/repositories/scheduled/scheduled_task_idempotency_repository.py index e4f36720..c5a51aaf 100644 --- a/backend/package/yuxi/repositories/scheduled/scheduled_task_idempotency_repository.py +++ b/backend/package/yuxi/repositories/scheduled/scheduled_task_idempotency_repository.py @@ -14,7 +14,9 @@ from __future__ import annotations from datetime import datetime -from sqlalchemy import delete as sa_delete, select, update as sa_update +from sqlalchemy import delete as sa_delete +from sqlalchemy import select +from sqlalchemy import update as sa_update from sqlalchemy.exc import IntegrityError from yuxi.repositories.scheduled.base import BaseRepository @@ -37,7 +39,14 @@ class ScheduledTaskIdempotencyRepository(BaseRepository): # 幂等抢占 # ------------------------------------------------------------------ - async def acquire(self, key: str, *, task_id: str, operation: str) -> bool: + async def acquire( + self, + key: str, + *, + task_id: str, + operation: str, + created_by: str | None = None, + ) -> bool: """尝试抢占幂等键,返回是否为首次请求。 利用 ``idempotency_key`` 列的 UNIQUE 约束实现并发安全的抢占语义: @@ -54,11 +63,13 @@ class ScheduledTaskIdempotencyRepository(BaseRepository): 允许客户端使用相同 key 重试。 - 响应体在业务逻辑后产生,本方法写入时 ``response_body=None``, 由 :meth:`update_response` 在业务完成后回写。 + - ``created_by`` 记录发起幂等请求的操作人 UID,用于审计追溯。 Args: key: 客户端传入的 ``Idempotency-Key``(建议 UUID)。 task_id: 关联的任务 ID(字符串引用,非外键)。 operation: 操作类型(create / update / delete / pause / resume / trigger)。 + created_by: 发起请求的操作人 UID,用于审计追溯。 Returns: ``True`` 表示首次请求抢占成功;``False`` 表示重复请求(key 已存在)。 @@ -69,6 +80,8 @@ class ScheduledTaskIdempotencyRepository(BaseRepository): task_id=task_id, operation=operation, response_body=None, + created_by=created_by, + updated_by=created_by, created_at=now, updated_at=now, is_deleted=0, @@ -84,7 +97,13 @@ class ScheduledTaskIdempotencyRepository(BaseRepository): # 响应回写 # ------------------------------------------------------------------ - async def update_response(self, key: str, response_body: dict) -> None: + async def update_response( + self, + key: str, + response_body: dict, + *, + updated_by: str | None = None, + ) -> None: """回写首次请求的响应体,用于后续重复请求重放。 严格按 PRD §6.2.3 端口签名实现,不接收 ``commit`` 参数;方法内部 @@ -93,13 +112,14 @@ class ScheduledTaskIdempotencyRepository(BaseRepository): Args: key: 幂等键。 response_body: 首次请求的响应体(dict)。 + updated_by: 回写操作人 UID,用于审计追溯。 """ now = utc_now_naive() stmt = ( sa_update(self.model) .where(self.model.idempotency_key == key) .where(self.model.is_deleted == 0) - .values(response_body=response_body, updated_at=now) + .values(response_body=response_body, updated_by=updated_by, updated_at=now) ) await self.db.execute(stmt) await self.db.flush() @@ -120,9 +140,7 @@ class ScheduledTaskIdempotencyRepository(BaseRepository): (``response_body IS NULL``)时返回 ``None``。 """ stmt = ( - select(self.model.response_body) - .where(self.model.idempotency_key == key) - .where(self.model.is_deleted == 0) + select(self.model.response_body).where(self.model.idempotency_key == key).where(self.model.is_deleted == 0) ) result = await self.db.execute(stmt) return result.scalar_one_or_none() diff --git a/backend/package/yuxi/repositories/scheduled/scheduled_task_repository.py b/backend/package/yuxi/repositories/scheduled/scheduled_task_repository.py index c1c2cd11..25c4bacc 100644 --- a/backend/package/yuxi/repositories/scheduled/scheduled_task_repository.py +++ b/backend/package/yuxi/repositories/scheduled/scheduled_task_repository.py @@ -21,7 +21,8 @@ from dataclasses import dataclass from datetime import datetime from typing import Any -from sqlalchemy import func, select, update as sa_update +from sqlalchemy import func, select +from sqlalchemy import update as sa_update from yuxi.repositories.scheduled.base import BaseRepository from yuxi.storage.postgres.models_scheduler import ScheduledTask @@ -52,9 +53,7 @@ class ScheduledTaskRepository(BaseRepository): # 标准 CRUD # ------------------------------------------------------------------ - async def create( - self, data: dict[str, Any], *, commit: bool = True - ) -> ScheduledTask: + async def create(self, data: dict[str, Any], *, commit: bool = True) -> ScheduledTask: """创建定时任务。 ``task_id`` 由调用方传入(业务 UUID 或 system: 前缀约定)。从 ``data`` @@ -102,9 +101,7 @@ class ScheduledTaskRepository(BaseRepository): await self.db.refresh(task) return task - async def get_by_task_id( - self, task_id: str, *, for_update: bool = False - ) -> ScheduledTask | None: + async def get_by_task_id(self, task_id: str, *, for_update: bool = False) -> ScheduledTask | None: """根据 ``task_id`` 获取任务(排除已软删除)。 Args: @@ -124,9 +121,7 @@ class ScheduledTaskRepository(BaseRepository): result = await self.db.execute(stmt) return result.scalar_one_or_none() - async def get_by_id( - self, record_id: int, *, for_update: bool = False - ) -> ScheduledTask | None: + async def get_by_id(self, record_id: int, *, for_update: bool = False) -> ScheduledTask | None: """根据主键 ``id`` 获取任务(排除已软删除)。 用于回收站恢复等只有主键的场景,与基类 ``delete_by_id`` / @@ -208,9 +203,7 @@ class ScheduledTaskRepository(BaseRepository): ) stmt = stmt.order_by(ScheduledTask.created_at.desc()) - effective_limit, effective_offset = self._resolve_pagination( - page, page_size, None, None - ) + effective_limit, effective_offset = self._resolve_pagination(page, page_size, None, None) if effective_limit is not None: stmt = stmt.limit(effective_limit).offset(effective_offset) @@ -232,32 +225,37 @@ class ScheduledTaskRepository(BaseRepository): return tasks, total - async def update( - self, task_id: str, data: dict[str, Any], *, commit: bool = True - ) -> ScheduledTask | None: + async def update(self, task_id: str, data: dict[str, Any], *, commit: bool = True) -> ScheduledTask | None: """按 ``task_id`` 更新任务配置字段(排除已软删除)。 - 自动设置 ``updated_at=now``。从 ``data`` 中提取允许更新的字段, - 排除主键、审计字段与运行时状态字段。运行时状态字段 - (``consecutive_errors`` / ``last_run_at`` / ``next_run_at`` / - ``last_error`` / ``status``)由专用状态机方法管理,禁止通过本方法 - 直接更新,避免绕过状态机约束: + 自动设置 ``updated_at=now``。从 ``data`` 中提取允许更新的字段, + 排除主键、审计字段与运行时状态字段。运行时状态字段 + (``consecutive_errors`` / ``last_run_at`` / ``next_run_at`` / + ``last_error`` / ``status``)由专用状态机方法管理,禁止通过本方法 + 直接更新,避免绕过状态机约束: - - ``pause`` / ``resume`` —— active ↔ paused - - ``mark_dead_letter`` —— → dead_letter - - ``record_run_result`` —— 原子回写执行结果与 next_run_at + - ``pause`` / ``resume`` —— active ↔ paused + - ``mark_dead_letter`` —— → dead_letter + - ``record_run_result`` —— 原子回写执行结果与 next_run_at - Args: - task_id: 业务任务标识。 - data: 待更新字段字典。 - commit: ``True`` 时提交事务,``False`` 时仅 flush。 + Args: + task_id: 业务任务标识。 + data: 待更新字段字典。 + commit: ``True`` 时提交事务,``False`` 时仅 flush。 - Returns: - 更新后的 ``ScheduledTask`` 实例,或 None(任务不存在时)。 + Returns: + 更新后的 ``ScheduledTask`` 实例,或 None(任务不存在时)。 """ excluded = { - "task_id", "id", "created_at", "created_by", "is_deleted", - "consecutive_errors", "last_run_at", "next_run_at", "last_error", + "task_id", + "id", + "created_at", + "created_by", + "is_deleted", + "consecutive_errors", + "last_run_at", + "next_run_at", + "last_error", "status", } update_fields = {k: v for k, v in data.items() if k not in excluded} @@ -505,7 +503,7 @@ class ScheduledTaskRepository(BaseRepository): # ------------------------------------------------------------------ async def pause( - self, task_id: str, *, commit: bool = True + self, task_id: str, *, updated_by: str | None = None, commit: bool = True ) -> ScheduledTask | None: """暂停任务(``active`` → ``paused``)。 @@ -515,6 +513,7 @@ class ScheduledTaskRepository(BaseRepository): Args: task_id: 业务任务标识。 + updated_by: 操作人 uid,写入 ``updated_by`` 审计字段。 commit: ``True`` 时提交事务,``False`` 时仅 flush。 Returns: @@ -526,7 +525,7 @@ class ScheduledTaskRepository(BaseRepository): .where(ScheduledTask.task_id == task_id) .where(self._not_deleted()) .where(ScheduledTask.status == "active") - .values(status="paused", updated_at=now) + .values(status="paused", updated_by=updated_by, updated_at=now) ) result = await self.db.execute(stmt) if result.rowcount == 0: @@ -542,6 +541,7 @@ class ScheduledTaskRepository(BaseRepository): task_id: str, *, next_run_at: datetime, + updated_by: str | None = None, commit: bool = True, ) -> ScheduledTask | None: """恢复任务(``paused`` → ``active``),并重算 ``next_run_at``。 @@ -554,6 +554,7 @@ class ScheduledTaskRepository(BaseRepository): task_id: 业务任务标识。 next_run_at: 恢复后的下次执行时间(UTC naive),由调用方按 cron 表达式或 ``run_at`` 重算。 + updated_by: 操作人 uid,写入 ``updated_by`` 审计字段。 commit: ``True`` 时提交事务,``False`` 时仅 flush。 Returns: @@ -565,7 +566,12 @@ class ScheduledTaskRepository(BaseRepository): .where(ScheduledTask.task_id == task_id) .where(self._not_deleted()) .where(ScheduledTask.status == "paused") - .values(status="active", next_run_at=next_run_at, updated_at=now) + .values( + status="active", + next_run_at=next_run_at, + updated_by=updated_by, + updated_at=now, + ) ) result = await self.db.execute(stmt) if result.rowcount == 0: @@ -577,7 +583,7 @@ class ScheduledTaskRepository(BaseRepository): return await self.get_by_task_id(task_id) async def mark_dead_letter( - self, task_id: str, *, commit: bool = True + self, task_id: str, *, updated_by: str | None = None, commit: bool = True ) -> ScheduledTask | None: """将任务标记为死信(``→ dead_letter``,``enabled=False``)。 @@ -587,6 +593,7 @@ class ScheduledTaskRepository(BaseRepository): Args: task_id: 业务任务标识。 + updated_by: 操作人 uid,写入 ``updated_by`` 审计字段。 commit: ``True`` 时提交事务,``False`` 时仅 flush。 Returns: @@ -598,7 +605,12 @@ class ScheduledTaskRepository(BaseRepository): .where(ScheduledTask.task_id == task_id) .where(self._not_deleted()) .where(ScheduledTask.status != "dead_letter") - .values(status="dead_letter", enabled=False, updated_at=now) + .values( + status="dead_letter", + enabled=False, + updated_by=updated_by, + updated_at=now, + ) ) result = await self.db.execute(stmt) if result.rowcount == 0: @@ -614,6 +626,7 @@ class ScheduledTaskRepository(BaseRepository): task_id: str, *, next_run_at: datetime, + updated_by: str | None = None, commit: bool = True, ) -> ScheduledTask | None: """将死信任务恢复为 ``active``(``dead_letter`` → ``active``)。 @@ -633,6 +646,7 @@ class ScheduledTaskRepository(BaseRepository): task_id: 业务任务标识。 next_run_at: 恢复后的下次执行时间(UTC naive),由调用方按 cron 表达式或 ``run_at`` 重算。 + updated_by: 操作人 uid,写入 ``updated_by`` 审计字段。 commit: ``True`` 时提交事务,``False`` 时仅 flush。 Returns: @@ -650,6 +664,7 @@ class ScheduledTaskRepository(BaseRepository): consecutive_errors=0, last_error=None, next_run_at=next_run_at, + updated_by=updated_by, updated_at=now, ) ) @@ -669,6 +684,8 @@ class ScheduledTaskRepository(BaseRepository): success: bool, error: str | None, backoff_until: datetime | None = None, + started_at: datetime | None = None, + updated_by: str | None = None, commit: bool = True, ) -> ScheduledTask | None: """原子回写单次执行结果并推进执行轨迹字段。 @@ -680,10 +697,14 @@ class ScheduledTaskRepository(BaseRepository): :meth:`acquire_for_run` 的推进语义冲突: - **成功**:``next_run_at`` 已在 :meth:`acquire_for_run` 抢占时推进到 - cron 下次时间(``at`` 类型置 None),本方法**不更新** ``next_run_at``。 + cron 下次时间,本方法**不更新** ``next_run_at``。一次性(``at``)任务 + 成功后由 :meth:`finalize_one_shot_task` 置 NULL 并转 ``paused``。 - **失败**:``next_run_at`` 由 ``backoff_until`` 参数指定退避时间, 调用方按退避序列计算后传入;``at`` 类型失败后传 None 表示不再重试。 + ``last_run_at`` 写入 ``started_at``(执行开始时间),与 ORM 注释 + "上次执行开始时间"语义一致;``started_at`` 为 None 时回退到当前时间。 + 本方法**不**判断死信阈值:调用方根据返回的 ``consecutive_errors`` 决定是否调用 :meth:`mark_dead_letter`。两次操作在同一事务内执行 不会有竞态。 @@ -695,25 +716,31 @@ class ScheduledTaskRepository(BaseRepository): backoff_until: 失败时的退避时间(UTC naive),成功时忽略。 ``at`` 类型失败后传 None 表示不再重试。``success=True`` 时 该参数被忽略。 + started_at: 本次执行开始时间(UTC naive),写入 ``last_run_at``。 + 为 None 时回退到当前时间。 + updated_by: 操作人 uid,写入 ``updated_by`` 审计字段。 commit: ``True`` 时提交事务,``False`` 时仅 flush。 Returns: 更新后的 ``ScheduledTask`` 实例,或 None(任务不存在)。 """ now = utc_now_naive() + run_at_time = started_at if started_at is not None else now if success: update_values = { "consecutive_errors": 0, "last_error": None, - "last_run_at": now, + "last_run_at": run_at_time, + "updated_by": updated_by, "updated_at": now, } else: update_values = { "consecutive_errors": ScheduledTask.consecutive_errors + 1, "last_error": error, - "last_run_at": now, + "last_run_at": run_at_time, "next_run_at": backoff_until, + "updated_by": updated_by, "updated_at": now, } stmt = ( @@ -731,6 +758,51 @@ class ScheduledTaskRepository(BaseRepository): await self.db.flush() return await self.get_by_task_id(task_id) + async def finalize_one_shot_task( + self, + task_id: str, + *, + updated_by: str | None = None, + commit: bool = True, + ) -> ScheduledTask | None: + """一次性任务(``at`` 类型)执行成功后的收尾。 + + ``at`` 类型任务执行成功后,若 ``delete_after_run=False``,需将 + ``next_run_at`` 置 NULL 并将 ``status`` 转为 ``paused``,避免 + ``next_run_at`` 保持过期时间导致 tick 反复扫描触发(无限循环)。 + + 与 :meth:`record_run_result` 分离,职责单一:仅在 ``at`` 类型成功 + 且不删除时由调用方在同一事务内调用。 + + Args: + task_id: 业务任务标识。 + updated_by: 操作人 uid,写入 ``updated_by`` 审计字段。 + commit: ``True`` 时提交事务,``False`` 时仅 flush。 + + Returns: + 更新后的 ``ScheduledTask`` 实例,或 None(任务不存在)。 + """ + now = utc_now_naive() + stmt = ( + sa_update(ScheduledTask) + .where(ScheduledTask.task_id == task_id) + .where(self._not_deleted()) + .values( + next_run_at=None, + status="paused", + updated_by=updated_by, + updated_at=now, + ) + ) + result = await self.db.execute(stmt) + if result.rowcount == 0: + return None + if commit: + await self.db.commit() + else: + await self.db.flush() + return await self.get_by_task_id(task_id) + # ------------------------------------------------------------------ # 内部查询构建 # ------------------------------------------------------------------ diff --git a/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_daily_repository.py b/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_daily_repository.py index d5a44028..a14875db 100644 --- a/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_daily_repository.py +++ b/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_daily_repository.py @@ -20,7 +20,8 @@ from __future__ import annotations from dataclasses import dataclass from datetime import date -from sqlalchemy import delete as sa_delete, select +from sqlalchemy import delete as sa_delete +from sqlalchemy import select from sqlalchemy.dialects.postgresql import insert as pg_insert from yuxi.repositories.scheduled.base import BaseRepository @@ -98,9 +99,7 @@ class ScheduledTaskRunLogDailyRepository(BaseRepository): "dead_letter": "dead_letter_count", } if result not in field_map: - raise ValueError( - f"result 取值必须为 {sorted(field_map)} 之一,收到: {result!r}" - ) + raise ValueError(f"result 取值必须为 {sorted(field_map)} 之一,收到: {result!r}") field_name = field_map[result] now = utc_now_naive() @@ -113,6 +112,8 @@ class ScheduledTaskRunLogDailyRepository(BaseRepository): dead_letter_count=1 if field_name == "dead_letter_count" else 0, created_at=now, updated_at=now, + created_by="system", + updated_by="system", is_deleted=0, ) stmt = stmt.on_conflict_do_update( @@ -186,9 +187,7 @@ class ScheduledTaskRunLogDailyRepository(BaseRepository): # 数据生命周期 # ------------------------------------------------------------------ - async def cleanup_old_daily_stats( - self, before: date, *, commit: bool = True - ) -> int: + async def cleanup_old_daily_stats(self, before: date, *, commit: bool = True) -> int: """物理删除 ``stat_date < before`` 的日聚合统计。 与 PRD §6.1.6「日聚合统计保留 1 年」对齐,由 ``run_log_cleanup`` handler @@ -202,9 +201,7 @@ class ScheduledTaskRunLogDailyRepository(BaseRepository): Returns: 删除行数。 """ - stmt = sa_delete(ScheduledTaskRunLogDaily).where( - ScheduledTaskRunLogDaily.stat_date < before - ) + stmt = sa_delete(ScheduledTaskRunLogDaily).where(ScheduledTaskRunLogDaily.stat_date < before) result = await self.db.execute(stmt) if commit: await self.db.commit() diff --git a/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_repository.py b/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_repository.py index 0bdc05f4..377e3191 100644 --- a/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_repository.py +++ b/backend/package/yuxi/repositories/scheduled/scheduled_task_run_log_repository.py @@ -24,7 +24,9 @@ from __future__ import annotations from datetime import datetime from typing import Any -from sqlalchemy import delete as sa_delete, func, select, update as sa_update +from sqlalchemy import delete as sa_delete +from sqlalchemy import func, select +from sqlalchemy import update as sa_update from yuxi.repositories.scheduled.base import BaseRepository from yuxi.storage.postgres.models_scheduler import ScheduledTaskRunLog @@ -49,9 +51,7 @@ class ScheduledTaskRunLogRepository(BaseRepository): # 标准 CRUD # ------------------------------------------------------------------ - async def create( - self, data: dict[str, Any], *, commit: bool = True - ) -> ScheduledTaskRunLog: + async def create(self, data: dict[str, Any], *, commit: bool = True) -> ScheduledTaskRunLog: """创建执行日志,初始 ``status='running'``。 从 ``data`` 字典提取字段构造 ``ScheduledTaskRunLog`` 实例,必填字段 @@ -107,9 +107,7 @@ class ScheduledTaskRunLogRepository(BaseRepository): result = await self.db.execute(stmt) return result.scalar_one_or_none() - async def update( - self, run_id: str, data: dict[str, Any], *, commit: bool = True - ) -> ScheduledTaskRunLog | None: + async def update(self, run_id: str, data: dict[str, Any], *, commit: bool = True) -> ScheduledTaskRunLog | None: """按 ``run_id`` 更新执行日志(排除已软删除)。 用于执行结束后写入 ``status`` / ``error_message`` / ``output`` / @@ -179,9 +177,7 @@ class ScheduledTaskRunLogRepository(BaseRepository): ) stmt = stmt.order_by(ScheduledTaskRunLog.started_at.desc()) - effective_limit, effective_offset = self._resolve_pagination( - page, page_size, None, None - ) + effective_limit, effective_offset = self._resolve_pagination(page, page_size, None, None) if effective_limit is not None: stmt = stmt.limit(effective_limit).offset(effective_offset) @@ -232,9 +228,7 @@ class ScheduledTaskRunLogRepository(BaseRepository): ) stmt = stmt.order_by(ScheduledTaskRunLog.started_at.desc()) - effective_limit, effective_offset = self._resolve_pagination( - page, page_size, None, None - ) + effective_limit, effective_offset = self._resolve_pagination(page, page_size, None, None) if effective_limit is not None: stmt = stmt.limit(effective_limit).offset(effective_offset) @@ -302,7 +296,7 @@ class ScheduledTaskRunLogRepository(BaseRepository): stale_before: datetime, error_message: str = "reclaimed by scheduler restart", commit: bool = True, - ) -> int: + ) -> list[str]: """清理僵尸 ``running`` 记录,用于调度器崩溃恢复。 调度器崩溃时,正在执行的 ``status='running'`` 日志无法被正常回写, @@ -315,9 +309,10 @@ class ScheduledTaskRunLogRepository(BaseRepository): 恢复正确。应在调度器启动时调用,``stale_before`` 通常取"上次心跳时间" 或"当前时间减去最大执行时长"。 - **不补录日聚合统计**:reclaim 是运维操作,被清理的记录数量通常很少 - (仅崩溃时正在执行的少数任务),对聚合统计影响可忽略;如需准确补录, - 调用方应在调用本方法前自行查询并处理。 + **补录日聚合统计**:通过 PostgreSQL ``RETURNING`` 子句返回被清理记录的 + ``handler_name`` 列表,调用方按 handler 分组后调用 + ``upsert_daily_stat(..., "timeout")`` 补录 ``timeout_count``,保证 + 日聚合统计准确(此前 ``timeout_count`` 字段无写入路径,永远为 0)。 Args: stale_before: 截止时间,``started_at`` 早于此时间的 running 记录 @@ -327,7 +322,7 @@ class ScheduledTaskRunLogRepository(BaseRepository): commit: ``True`` 时提交事务,``False`` 时仅 flush。 Returns: - 被清理的僵尸记录数。 + 被清理记录的 ``handler_name`` 列表(供调用方补录日聚合统计)。 """ now = utc_now_naive() stmt = ( @@ -341,21 +336,21 @@ class ScheduledTaskRunLogRepository(BaseRepository): error_message=error_message, updated_at=now, ) + .returning(ScheduledTaskRunLog.handler_name) ) result = await self.db.execute(stmt) + handler_names = list(result.scalars().all()) if commit: await self.db.commit() else: await self.db.flush() - return int(result.rowcount or 0) + return handler_names # ------------------------------------------------------------------ # 数据生命周期 # ------------------------------------------------------------------ - async def cleanup_old_logs( - self, before: datetime, *, commit: bool = True - ) -> int: + async def cleanup_old_logs(self, before: datetime, *, commit: bool = True) -> int: """物理删除 ``started_at < before`` 的明细日志。 与 PRD §6.1.6「run_logs 保留 90 天」对齐,由 ``run_log_cleanup`` handler @@ -369,9 +364,7 @@ class ScheduledTaskRunLogRepository(BaseRepository): Returns: 删除行数。 """ - stmt = sa_delete(ScheduledTaskRunLog).where( - ScheduledTaskRunLog.started_at < before - ) + stmt = sa_delete(ScheduledTaskRunLog).where(ScheduledTaskRunLog.started_at < before) result = await self.db.execute(stmt) if commit: await self.db.commit()