refactor: add durable transfer planning checkpoints

This commit is contained in:
jxxghp
2026-08-27 13:16:33 +08:00
parent 49d30e5fdf
commit 23bff6b2bb
41 changed files with 7401 additions and 1098 deletions
+163 -2
View File
@@ -2,6 +2,7 @@
from collections.abc import Callable
from datetime import datetime
from typing import Optional
from uuid import uuid4
from sqlalchemy.exc import IntegrityError
@@ -9,7 +10,13 @@ from sqlalchemy.orm import Session
from app.application.transfer import (
TRANSFER_ADMISSION_ACCEPTED,
TRANSFER_ADMISSION_PLANNED,
TRANSFER_ADMISSION_PROVIDER_PENDING,
TransferAdmission,
TransferAdmissionConflictError,
TransferPlanCheckpoint,
TransferPlanningInput,
TransferPlanningStateError,
)
from app.db.models.transferpending import TransferPending
from app.db.oper.transferpending import TransferPendingOper
@@ -32,6 +39,40 @@ class TransactionalTransferAdmissionRepository:
def _project(pending: TransferPending) -> TransferAdmission:
"""在 Session 有效期内把 ORM 行冻结为应用层 DTO。"""
created_at = pending.created_at or pending.updated_at
planning_input = TransferPlanningInput.from_payload(pending.planning_input)
if pending.input_version != planning_input.schema_version:
raise TransferPlanningStateError("整理规划输入列版本与 JSON 版本不一致")
if pending.input_fingerprint != planning_input.fingerprint:
raise TransferAdmissionConflictError("整理规划输入 JSON 与持久指纹不一致")
checkpoint = (
TransferPlanCheckpoint.from_payload(pending.checkpoint_payload)
if pending.checkpoint_payload is not None
else None
)
if checkpoint is not None:
if pending.checkpoint_version != checkpoint.schema_version:
raise TransferPlanningStateError("整理检查点列版本与 JSON 版本不一致")
if checkpoint.planning_input.fingerprint != pending.input_fingerprint:
raise TransferAdmissionConflictError("整理检查点内嵌输入与准入指纹不一致")
if pending.state in {
TRANSFER_ADMISSION_PROVIDER_PENDING,
TRANSFER_ADMISSION_PLANNED,
} and checkpoint is None:
raise TransferPlanningStateError("待执行任务缺少完整检查点")
if pending.state == TRANSFER_ADMISSION_ACCEPTED and checkpoint is not None:
raise TransferPlanningStateError("接纳态任务不能携带计划检查点")
if (
pending.state == TRANSFER_ADMISSION_PROVIDER_PENDING
and checkpoint is not None
and not checkpoint.is_provider_pending
):
raise TransferPlanningStateError("provider_pending 状态缺少 provider 调用快照")
if (
pending.state == TRANSFER_ADMISSION_PLANNED
and checkpoint is not None
and checkpoint.is_provider_pending
):
raise TransferPlanningStateError("planned 状态不能携带 provider-only 检查点")
return TransferAdmission(
task_id=pending.task_id,
storage=pending.storage,
@@ -40,10 +81,39 @@ class TransactionalTransferAdmissionRepository:
created_at=created_at,
updated_at=pending.updated_at,
last_error=pending.last_error,
input_fingerprint=pending.input_fingerprint,
planning_input=planning_input,
checkpoint=checkpoint,
)
def admit(self, *, storage: str, src_path: str) -> TransferAdmission:
"""幂等持久化准入事实,并返回跨重启稳定的任务标识。"""
@staticmethod
def _assert_input_match(
pending: TransferPending,
planning_input: TransferPlanningInput,
) -> None:
"""拒绝同一源文件以不同规划输入复用既有任务身份。"""
if pending.input_fingerprint != planning_input.fingerprint:
raise TransferAdmissionConflictError(
f"整理源文件已按不同输入准入: {pending.storage}:{pending.src_path}"
)
def admit(
self,
*,
storage: str,
src_path: str,
planning_input: Optional[TransferPlanningInput] = None,
) -> TransferAdmission:
"""按输入指纹幂等持久化准入事实,并返回跨重启稳定身份。"""
effective_input = planning_input or TransferPlanningInput.legacy(
storage=storage,
src_path=src_path,
)
if (
effective_input.source_fileitem.get("storage") != storage
or effective_input.source_fileitem.get("path") != src_path
):
raise ValueError("整理规划输入的源文件身份与准入参数不一致")
now_time = self._now()
try:
with self._session_factory() as session:
@@ -55,10 +125,14 @@ class TransactionalTransferAdmissionRepository:
src_path=src_path,
state=TRANSFER_ADMISSION_ACCEPTED,
now_time=now_time,
input_version=effective_input.schema_version,
planning_input=effective_input.to_payload(),
input_fingerprint=effective_input.fingerprint,
)
if pending is None:
raise ValueError("整理任务的存储与源路径不能为空")
session.flush()
self._assert_input_match(pending, effective_input)
admission = self._project(pending)
transaction.commit()
return admission
@@ -74,6 +148,7 @@ class TransactionalTransferAdmissionRepository:
)
if pending is None:
raise RuntimeError("并发准入冲突后未找到已提交记录") from error
self._assert_input_match(pending, effective_input)
return self._project(pending)
def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]:
@@ -85,6 +160,19 @@ class TransactionalTransferAdmissionRepository:
)
return [self._project(pending) for pending in pending_items]
def list_recoverable(self, limit: int = 5000) -> list[TransferAdmission]:
"""投影接纳、provider 待执行或已规划的全部可恢复任务。"""
with self._session_factory() as session:
pending_items = TransferPendingOper(db=session).list_by_states(
states=(
TRANSFER_ADMISSION_ACCEPTED,
TRANSFER_ADMISSION_PROVIDER_PENDING,
TRANSFER_ADMISSION_PLANNED,
),
limit=limit,
)
return [self._project(pending) for pending in pending_items]
def record_enqueue_failure(self, *, task_id: str, error: str) -> None:
"""独立提交最近一次入队失败,保留准入记录供后续恢复。"""
with self._session_factory() as session:
@@ -100,6 +188,79 @@ class TransactionalTransferAdmissionRepository:
transaction.rollback()
raise
def checkpoint_plan(
self,
*,
task_id: str,
input_fingerprint: str,
checkpoint: TransferPlanCheckpoint,
) -> TransferAdmission:
"""以输入指纹 CAS 保存 provider 调用快照或升级宿主计划。"""
if checkpoint.planning_input.fingerprint != input_fingerprint:
raise TransferAdmissionConflictError("检查点输入与准入输入指纹不一致")
checkpoint_payload = checkpoint.to_payload()
target_state = (
TRANSFER_ADMISSION_PROVIDER_PENDING
if checkpoint.is_provider_pending
else TRANSFER_ADMISSION_PLANNED
)
source_states = (
(TRANSFER_ADMISSION_ACCEPTED,)
if checkpoint.is_provider_pending
else (
TRANSFER_ADMISSION_ACCEPTED,
TRANSFER_ADMISSION_PROVIDER_PENDING,
)
)
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
oper = TransferPendingOper(db=session)
updated = oper.stage_checkpoint_plan(
task_id=task_id,
input_fingerprint=input_fingerprint,
checkpoint_version=checkpoint.schema_version,
checkpoint_payload=checkpoint_payload,
source_states=source_states,
target_state=target_state,
now_time=self._now(),
)
session.flush()
session.expire_all()
pending = oper.get_by_task_id(task_id=task_id)
if pending is None:
raise TransferPlanningStateError(f"未找到整理任务: {task_id}")
if pending.input_fingerprint != input_fingerprint:
raise TransferAdmissionConflictError("整理任务输入指纹已经改变")
if not updated and not (
pending.state == target_state
and pending.checkpoint_payload == checkpoint_payload
):
raise TransferPlanningStateError(
f"整理任务不能从状态 {pending.state} 保存 {target_state} 检查点"
)
admission = self._project(pending)
transaction.commit()
return admission
except Exception:
transaction.rollback()
raise
def record_planning_failure(self, *, task_id: str, error: str) -> None:
"""独立提交规划错误并保持任务处于接纳态供恢复重试。"""
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
TransferPendingOper(db=session).stage_record_planning_failure(
task_id=task_id,
error=error,
now_time=self._now(),
)
transaction.commit()
except Exception:
transaction.rollback()
raise
def discard_task(self, *, task_id: str) -> int:
"""在独立事务中按稳定任务标识删除已到终态的准入记录。"""
with self._session_factory() as session:
+192 -6
View File
@@ -1,13 +1,66 @@
import hashlib
import json
from datetime import datetime
from typing import List, Optional, cast
from typing import Any, List, Optional, cast
from uuid import uuid4
from sqlalchemy import Index, String, Text, UniqueConstraint, delete, select, update
from sqlalchemy import JSON, Index, Integer, String, Text, UniqueConstraint, delete, select, update
from sqlalchemy.orm import Mapped, Session, mapped_column
from app.db.base import Base, execute_dml, get_id_column
def _legacy_planning_payload(storage: str, src_path: str) -> dict[str, Any]:
return {
"schema_version": 1,
"source_fileitem": {"storage": storage, "path": src_path},
"meta": None,
"mediainfo": None,
"target_directory": None,
"target_storage": None,
"target_path": None,
"requested_transfer_type": None,
"media_source": None,
"media_id": None,
"media_type": None,
"need_scrape": False,
"need_rename": True,
"need_notify": True,
"overwrite_mode": None,
"episodes_info": [],
"preview": False,
"options": {"legacy_replan": True},
}
def _planning_fingerprint(payload: dict[str, Any]) -> str:
canonical = json.dumps(
payload,
ensure_ascii=True,
sort_keys=True,
separators=(",", ":"),
allow_nan=False,
)
return hashlib.sha256(canonical.encode("utf-8")).hexdigest()
def _default_planning_payload(context: Any) -> dict[str, Any]:
params = context.get_current_parameters()
return _legacy_planning_payload(
params.get("storage", ""),
params.get("src_path", ""),
)
def _default_planning_fingerprint(context: Any) -> str:
params = context.get_current_parameters()
payload = params.get("planning_input") or _legacy_planning_payload(
params.get("storage", ""),
params.get("src_path", ""),
)
return _planning_fingerprint(payload)
class TransferPending(Base):
"""
待整理文件登记。
@@ -17,9 +70,9 @@ class TransferPending(Base):
蒸发。而已经稳定落地的文件不会再产生任何监控事件,也不会有新的补偿扫描起点
——结果就是永久漏件,只能靠人工比对补整理。
这里只落盘恢复所需的最小事实:稳定任务身份、存储、源文件路径、准入状态和
最近入队错误。重启后重新走一遍整理入口,由整理历史查重挡掉已经完成的,
因此不需要序列化 meta/mediainfo 这些重对象,也不存在识别结果陈旧的问题
准入时保存版本化规划输入和指纹;纯规划完成后以同一行原子保存完整有序计划并
推进到 planned。重启恢复可直接消费已规划路径,避免再次触发 rename 等插件事件。
旧路径登记接口仍生成最小 legacy_replan 输入,供插件兼容调用方继续使用
"""
id = get_id_column()
@@ -42,6 +95,22 @@ class TransferPending(Base):
)
# 最近一次入队失败原因
last_error: Mapped[Optional[str]] = mapped_column(Text)
# 规划输入格式版本
input_version: Mapped[int] = mapped_column(Integer, nullable=False, default=1)
# 版本化规划输入 JSON
planning_input: Mapped[dict[str, Any]] = mapped_column(
JSON, nullable=False, default=_default_planning_payload
)
# 规划输入规范 JSON 的 SHA-256 指纹
input_fingerprint: Mapped[str] = mapped_column(
String(64), nullable=False, default=_default_planning_fingerprint
)
# 完整计划格式版本,尚未规划时为空
checkpoint_version: Mapped[Optional[int]] = mapped_column(Integer)
# 完整有序计划 JSON,尚未规划时为空
checkpoint_payload: Mapped[Optional[dict[str, Any]]] = mapped_column(JSON)
# 规划完成时间
planned_at: Mapped[Optional[str]] = mapped_column(String(40))
__table_args__ = (
# 同一个文件重复入队只保留一条,回放时不会重复送入整理链
@@ -74,12 +143,16 @@ class TransferPending(Base):
).scalars().first()
if pending:
return cast("TransferPending", pending)
planning_input = _legacy_planning_payload(storage, src_path)
pending = cls(
storage=storage,
src_path=src_path,
state="accepted",
created_at=now_time,
updated_at=now_time,
input_version=1,
planning_input=planning_input,
input_fingerprint=_planning_fingerprint(planning_input),
)
db.add(pending)
return pending
@@ -87,7 +160,9 @@ class TransferPending(Base):
@classmethod
def stage_admit(cls, db: Session, *, task_id: str, storage: str,
src_path: str, state: str,
now_time: str) -> Optional["TransferPending"]:
now_time: str, input_version: int = 1,
planning_input: Optional[dict[str, Any]] = None,
input_fingerprint: Optional[str] = None) -> Optional["TransferPending"]:
"""
在调用方会话中暂存一条持久接纳记录。
@@ -98,6 +173,9 @@ class TransferPending(Base):
:param src_path: 源文件路径
:param state: 持久状态
:param now_time: 当前时间
:param input_version: 规划输入格式版本
:param planning_input: 版本化规划输入 JSON
:param input_fingerprint: 规划输入规范 JSON 指纹
:return: 接纳记录
"""
if not task_id or not storage or not src_path or not state:
@@ -107,6 +185,8 @@ class TransferPending(Base):
).scalars().first()
if pending:
return cast("TransferPending", pending)
effective_input = planning_input or _legacy_planning_payload(storage, src_path)
effective_fingerprint = input_fingerprint or _planning_fingerprint(effective_input)
pending = cls(
task_id=task_id,
storage=storage,
@@ -114,6 +194,9 @@ class TransferPending(Base):
state=state,
created_at=now_time,
updated_at=now_time,
input_version=input_version,
planning_input=effective_input,
input_fingerprint=effective_fingerprint,
)
db.add(pending)
return pending
@@ -137,6 +220,25 @@ class TransferPending(Base):
.limit(limit)
).scalars().all())
@classmethod
def list_by_states(cls, db: Session, *, states: tuple[str, ...],
limit: Optional[int] = 5000) -> List["TransferPending"]:
"""
按登记顺序列出多个可恢复持久状态的记录。
:param db: 数据库会话
:param states: 允许恢复的状态集合
:param limit: 单次读取上限
:return: 接纳记录列表
"""
if not states:
return []
return list(db.execute(
select(cls)
.where(cls.state.in_(states))
.order_by(cls.created_at.asc(), cls.id.asc())
.limit(limit)
).scalars().all())
@classmethod
def get_by_identity(cls, db: Session, *, storage: str,
src_path: str) -> Optional["TransferPending"]:
@@ -159,6 +261,90 @@ class TransferPending(Base):
).scalars().first(),
)
@classmethod
def get_by_task_id(cls, db: Session, *, task_id: str) -> Optional["TransferPending"]:
"""
按稳定任务标识查询一条持久登记。
:param db: 数据库会话
:param task_id: 稳定任务标识
:return: 接纳记录
"""
if not task_id:
return None
return cast(
Optional["TransferPending"],
db.execute(select(cls).where(cls.task_id == task_id)).scalars().first(),
)
@classmethod
def checkpoint_plan(cls, db: Session, *, task_id: str,
input_fingerprint: str, checkpoint_version: int,
checkpoint_payload: dict[str, Any],
source_states: tuple[str, ...], target_state: str,
now_time: str) -> int:
"""
以输入指纹为 CAS 条件原子保存计划并推进到已规划。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param input_fingerprint: 规划输入规范 JSON 指纹
:param checkpoint_version: 检查点格式版本
:param checkpoint_payload: 完整有序计划 JSON
:param source_states: 允许推进检查点的起始状态
:param target_state: 检查点提交后的目标状态
:param now_time: 当前时间
:return: 更新的记录数
"""
if (
not task_id
or not input_fingerprint
or not checkpoint_payload
or not source_states
or not target_state
):
return 0
return execute_dml(
db,
update(cls)
.where(
cls.task_id == task_id,
cls.state.in_(source_states),
cls.input_fingerprint == input_fingerprint,
)
.values(
state=target_state,
checkpoint_version=checkpoint_version,
checkpoint_payload=checkpoint_payload,
planned_at=now_time,
last_error=None,
updated_at=now_time,
),
execution_options={"synchronize_session": False},
)
@classmethod
def record_planning_failure(cls, db: Session, *, task_id: str,
error: str, now_time: str) -> int:
"""
为接纳态或 provider 待执行任务记录规划失败,不改变其恢复状态。
:param db: 数据库会话
:param task_id: 稳定任务标识
:param error: 失败原因
:param now_time: 当前时间
:return: 更新的记录数
"""
if not task_id:
return 0
return execute_dml(
db,
update(cls)
.where(
cls.task_id == task_id,
cls.state.in_(("accepted", "provider_pending")),
)
.values(last_error=error, updated_at=now_time),
execution_options={"synchronize_session": False},
)
@classmethod
def record_enqueue_failure(cls, db: Session, *, task_id: str,
error: str, now_time: str) -> int:
+92 -2
View File
@@ -1,5 +1,5 @@
from datetime import datetime
from typing import List, Optional, Tuple
from typing import Any, List, Optional, Tuple
from app.db.base import DbOper
from app.db.models.transferpending import TransferPending
@@ -32,7 +32,9 @@ class TransferPendingOper(DbOper):
)
def stage_admit(self, *, task_id: str, storage: str, src_path: str,
state: str, now_time: str) -> Optional[TransferPending]:
state: str, now_time: str, input_version: int = 1,
planning_input: Optional[dict[str, Any]] = None,
input_fingerprint: Optional[str] = None) -> Optional[TransferPending]:
"""
在当前会话中暂存一条持久接纳记录。
@@ -42,6 +44,9 @@ class TransferPendingOper(DbOper):
:param src_path: 源文件路径
:param state: 持久状态
:param now_time: 当前时间
:param input_version: 规划输入格式版本
:param planning_input: 版本化规划输入 JSON
:param input_fingerprint: 规划输入规范 JSON 指纹
:return: 接纳记录
"""
return self._execute_sync_write(
@@ -52,6 +57,9 @@ class TransferPendingOper(DbOper):
src_path=src_path,
state=state,
now_time=now_time,
input_version=input_version,
planning_input=planning_input,
input_fingerprint=input_fingerprint,
)
)
@@ -71,6 +79,22 @@ class TransferPendingOper(DbOper):
)
) or []
def list_by_states(self, *, states: tuple[str, ...],
limit: Optional[int] = 5000) -> List[TransferPending]:
"""
使用当前会话列出多个可恢复状态的记录。
:param states: 允许恢复的状态集合
:param limit: 单次读取上限
:return: ORM 接纳记录列表
"""
return self._execute_sync_query(
lambda session: TransferPending.list_by_states(
session,
states=states,
limit=limit,
)
) or []
def get_by_identity(self, *, storage: str,
src_path: str) -> Optional[TransferPending]:
"""
@@ -87,6 +111,72 @@ class TransferPendingOper(DbOper):
)
)
def get_by_task_id(self, *, task_id: str) -> Optional[TransferPending]:
"""
使用当前会话按稳定任务标识查询接纳记录。
:param task_id: 稳定任务标识
:return: 接纳记录
"""
return self._execute_sync_query(
lambda session: TransferPending.get_by_task_id(
session,
task_id=task_id,
)
)
def stage_checkpoint_plan(
self,
*,
task_id: str,
input_fingerprint: str,
checkpoint_version: int,
checkpoint_payload: dict[str, Any],
source_states: tuple[str, ...],
target_state: str,
now_time: str,
) -> int:
"""
在当前会话中以输入指纹为条件暂存完整计划检查点。
:param task_id: 稳定任务标识
:param input_fingerprint: 规划输入规范 JSON 指纹
:param checkpoint_version: 检查点格式版本
:param checkpoint_payload: 完整有序计划 JSON
:param source_states: 允许执行 CAS 的起始状态
:param target_state: 检查点提交后的目标状态
:param now_time: 当前时间
:return: 更新的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.checkpoint_plan(
session,
task_id=task_id,
input_fingerprint=input_fingerprint,
checkpoint_version=checkpoint_version,
checkpoint_payload=checkpoint_payload,
source_states=source_states,
target_state=target_state,
now_time=now_time,
)
)
def stage_record_planning_failure(self, *, task_id: str, error: str,
now_time: str) -> int:
"""
在当前会话中记录规划失败并保持任务处于接纳态。
:param task_id: 稳定任务标识
:param error: 失败原因
:param now_time: 当前时间
:return: 更新的记录数
"""
return self._execute_sync_write(
lambda session: TransferPending.record_planning_failure(
session,
task_id=task_id,
error=error,
now_time=now_time,
)
)
def stage_record_enqueue_failure(self, *, task_id: str, error: str,
now_time: str) -> int:
"""