ForcePilot/backend/package/yuxi/scheduler/adapters/persistence/mappers.py
Kris 6498af03e3 feat(scheduler): 完整实现定时任务调度限界上下文基础架构
本次提交完成了定时任务调度限界上下文的全量基础架构搭建,包括:
1.  基于六边形架构的完整分层(core/use_cases/framework/adapters/infrastructure)
2.  任务调度核心领域模型、端口契约与校验工具
3.  持久化适配器层与SQLAlchemy仓储实现
4.  调度运行时核心组件(handler注册表、执行引擎)
5.  内置维护型任务handler(日志清理、幂等记录清理)
6.  全局异常处理器与PostgreSQL表结构适配
7.  ARQ worker调度任务集成与启动装配逻辑
2026-06-22 20:07:20 +08:00

126 lines
4.5 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""ORM ↔ dataclass 映射边界。
本模块是 ``adapters/persistence`` 层负责 ORM 实例 ↔ core dataclass 转换的位置。
scheduler 限界上下文无敏感字段mapper 仅做字段映射,不调用 ``encrypt`` /
``decrypt``(与 ``external_systems/adapters/persistence/mappers.py`` 不同)。
设计原则(见设计方案 §10
- **mapper 不调用 ORM 的 ``to_dict`` 方法**:直接读 ORM 字段构造 dataclass。
原因:``to_dict`` 返回的 dict 结构与 dataclass 字段不完全一致(如 DateTime
字段在 ``to_dict`` 中被 ``format_utc_datetime`` 格式化为字符串)。
- **无敏感字段**scheduler 域所有字段均为明文存储,无需加解密。
- **处理 None 入参**ORM 实例为 None 时返回 None避免调用方需额外判空。
- **仓储层 dataclass → core dataclass**``DailyStat`` / ``HandlerSummary`` 在
仓储层与 core 层均有同名 dataclassmapper 按字段名转换。
"""
from __future__ import annotations
from typing import TYPE_CHECKING, Any
from yuxi.scheduler.core.models import (
DailyStat as DailyStatDC,
HandlerSummary as HandlerSummaryDC,
ScheduledTask as ScheduledTaskDC,
ScheduledTaskRunLog as ScheduledTaskRunLogDC,
)
if TYPE_CHECKING:
from yuxi.storage.postgres.models_scheduler import (
ScheduledTask as ScheduledTaskORM,
ScheduledTaskRunLog as ScheduledTaskRunLogORM,
)
def orm_to_task(orm: ScheduledTaskORM | None) -> ScheduledTaskDC | None:
"""ORM → dataclass``ScheduledTask`` 字段一一映射。
DateTime 字段直接赋值core 用 ``Any``JSON 字段直接赋值。
入参为 None 时返回 None避免调用方需额外判空。
"""
if orm is None:
return None
return ScheduledTaskDC(
id=orm.id,
task_id=orm.task_id,
handler_name=orm.handler_name,
owner_scope=orm.owner_scope,
owner_id=orm.owner_id,
schedule_kind=orm.schedule_kind,
cron_expression=orm.cron_expression,
run_at=orm.run_at,
tz=orm.tz,
payload=orm.payload or {},
enabled=orm.enabled,
delete_after_run=orm.delete_after_run,
block_strategy=orm.block_strategy,
stagger_seconds=orm.stagger_seconds,
consecutive_errors=orm.consecutive_errors,
status=orm.status,
last_run_at=orm.last_run_at,
next_run_at=orm.next_run_at,
last_error=orm.last_error,
created_by=orm.created_by,
updated_by=orm.updated_by,
created_at=orm.created_at,
updated_at=orm.updated_at,
is_deleted=orm.is_deleted,
deleted_at=orm.deleted_at,
)
def orm_to_run_log(orm: ScheduledTaskRunLogORM | None) -> ScheduledTaskRunLogDC | None:
"""ORM → dataclass``ScheduledTaskRunLog`` 字段一一映射。
入参为 None 时返回 None避免调用方需额外判空。
"""
if orm is None:
return None
return ScheduledTaskRunLogDC(
id=orm.id,
task_id=orm.task_id,
run_id=orm.run_id,
triggered_by=orm.triggered_by,
status=orm.status,
error_message=orm.error_message,
output=orm.output,
started_at=orm.started_at,
finished_at=orm.finished_at,
created_by=orm.created_by,
updated_by=orm.updated_by,
created_at=orm.created_at,
updated_at=orm.updated_at,
is_deleted=orm.is_deleted,
deleted_at=orm.deleted_at,
)
def daily_stat_row_to_dataclass(row: Any) -> DailyStatDC:
"""仓储层 ``DailyStat`` → core ``DailyStat`` dataclass按字段名映射。
仓储层 ``DailyStat`` 是 dataclassrow 可能是 NamedTuple 或 dataclass 实例,
统一按字段名读取并转换为 core dataclass。
"""
return DailyStatDC(
stat_date=row.stat_date,
handler_name=row.handler_name,
success_count=int(row.success_count),
failure_count=int(row.failure_count),
timeout_count=int(row.timeout_count),
dead_letter_count=int(row.dead_letter_count),
)
def handler_summary_row_to_dataclass(row: Any) -> HandlerSummaryDC:
"""仓储层 ``HandlerSummary`` → core ``HandlerSummary`` dataclass按字段名映射。
仓储层 ``HandlerSummary`` 是 dataclassrow 可能是 NamedTuple 或 dataclass 实例,
统一按字段名读取并转换为 core dataclass。
"""
return HandlerSummaryDC(
name=row.name,
task_count=int(row.task_count),
last_active_at=row.last_active_at,
)