refactor(scheduled-repo): 优化代码格式与函数签名,补充审计字段
1. 格式化SQLAlchemy查询语句,将多行合并为单行提升可读性 2. 调整部分函数参数换行格式,统一代码风格 3. 为幂等键抢占、响应回写、任务状态变更等方法新增updated_by/created_by审计字段 4. 修复僵尸日志清理逻辑,补全日聚合统计写入路径并返回handler列表 5. 新增一次性任务执行成功后的收尾方法,避免调度循环
This commit is contained in:
parent
f04418fb0c
commit
23f7a8a863
@ -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()
|
||||
|
||||
@ -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()
|
||||
|
||||
@ -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)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# 内部查询构建
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
@ -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()
|
||||
|
||||
@ -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()
|
||||
|
||||
Loading…
Reference in New Issue
Block a user