diff --git a/app/application/chain/data.py b/app/application/chain/data.py index d5aedd4bf..645ab0c1f 100644 --- a/app/application/chain/data.py +++ b/app/application/chain/data.py @@ -10,8 +10,11 @@ from collections.abc import Callable from dataclasses import dataclass from typing import Any, Optional +from app.application.transfer import TransferAdmissionRepository + OperFactory = Callable[[], Any] +TransferAdmissionRepositoryFactory = Callable[[], TransferAdmissionRepository] @dataclass(frozen=True, slots=True) @@ -23,7 +26,7 @@ class ChainDataPorts: workflow: OperFactory download_history: OperFactory transfer_history: OperFactory - transfer_pending: OperFactory + transfer_pending: TransferAdmissionRepositoryFactory media_server: OperFactory download_failure: OperFactory user: OperFactory @@ -77,12 +80,6 @@ class TransferHistoryPortProxy(_ChainDataPortProxy): port_name = "transfer_history" -class TransferPendingPortProxy(_ChainDataPortProxy): - """待整理数据端口代理。""" - - port_name = "transfer_pending" - - class MediaServerPortProxy(_ChainDataPortProxy): """媒体服务器数据端口代理。""" @@ -156,8 +153,8 @@ def get_chain_transfer_history_port() -> Any: return get_chain_data_ports().transfer_history() -def get_chain_transfer_pending_port() -> Any: - """创建待整理数据端口实例。""" +def get_chain_transfer_pending_port() -> TransferAdmissionRepository: + """创建类型化的整理任务 durable admission 仓储。""" return get_chain_data_ports().transfer_pending() diff --git a/app/application/transfer.py b/app/application/transfer.py index c57308e49..ab93a5ba3 100644 --- a/app/application/transfer.py +++ b/app/application/transfer.py @@ -19,14 +19,10 @@ from copy import deepcopy from dataclasses import dataclass from pathlib import Path from time import monotonic -from typing import Callable, Dict, List, Optional, Tuple, Union +from typing import Callable, Dict, List, Optional, Protocol, Tuple, Union -from pydantic import BaseModel, ConfigDict +from pydantic import BaseModel, ConfigDict, PrivateAttr -from app.schemas.transfer import MetaInfo as _SchemaMetaInfo -from app.schemas.transfer import MusicInfo as _SchemaMusicInfo -from app.schemas.transfer import MusicMeta as _SchemaMusicMeta -from app.schemas.workflow import MediaInfo as _SchemaMediaInfo from app.adapters.system.host import SystemUtils from app.application.agent import get_prompt_manager, get_running_agent_manager from app.domain.context import MediaInfo, MusicInfo @@ -40,6 +36,9 @@ from app.schemas.history import DownloadHistory from app.schemas.media import OptionalMediaIdentityMixin, resolve_media_identity from app.schemas.system import TransferDirectoryConf from app.schemas.tmdb import TmdbEpisode +from app.schemas.transfer import MetaInfo as _SchemaMetaInfo +from app.schemas.transfer import MusicInfo as _SchemaMusicInfo +from app.schemas.transfer import MusicMeta as _SchemaMusicMeta from app.schemas.transfer import TransferInfo, TransferJob, TransferJobTask from app.schemas.types import ( MUSIC_ENTITY_ALBUM, @@ -48,7 +47,7 @@ from app.schemas.types import ( MediaType, ReplyMode, ) - +from app.schemas.workflow import MediaInfo as _SchemaMediaInfo class TransferTask(OptionalMediaIdentityMixin, BaseModel): @@ -81,6 +80,16 @@ class TransferTask(OptionalMediaIdentityMixin, BaseModel): manual: Optional[bool] = False background: Optional[bool] = True preview: Optional[bool] = False + _admission_task_id: Optional[str] = PrivateAttr(default=None) + + @property + def admission_task_id(self) -> Optional[str]: + """返回仅供宿主持久准入和终态结算使用的内部任务标识。""" + return self._admission_task_id + + def bind_admission_task_id(self, task_id: str) -> None: + """绑定持久准入生成的稳定身份,不改变插件可见序列化字段。""" + self._admission_task_id = task_id def to_dict(self): """ @@ -113,6 +122,42 @@ class TransferQueue(BaseModel): result: Optional[TransferInfo] = None +TRANSFER_ADMISSION_ACCEPTED = "accepted" + + +@dataclass(frozen=True, slots=True) +class TransferAdmission: + """描述已经持久化、可在进程退出后恢复的整理任务准入事实。""" + + task_id: str + storage: str + src_path: str + state: str + created_at: str + updated_at: str + last_error: Optional[str] = None + + +class TransferAdmissionRepository(Protocol): + """整理任务 durable admission 所需的类型化持久化端口。""" + + def admit(self, *, storage: str, src_path: str) -> TransferAdmission: + """幂等登记源文件并返回稳定任务身份。""" + ... + + def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]: + """按登记顺序返回等待恢复或执行的任务。""" + ... + + def record_enqueue_failure(self, *, task_id: str, error: str) -> None: + """记录内存队列接收失败,保留任务供后续恢复。""" + ... + + def discard_task(self, *, task_id: str) -> int: + """按稳定任务身份删除已经到达终态的登记。""" + ... + + class TransferQueueService: """协调整理任务登记、入队、移除和队列视图查询。""" @@ -120,29 +165,43 @@ class TransferQueueService: self, *, register_task: Callable[[TransferTask], bool], + admit_task: Callable[[TransferTask], TransferAdmission], enqueue: Callable[[TransferQueue], None], before_enqueue: Callable[[TransferTask], None], - after_enqueue: Callable[[TransferTask], None], + enqueue_failed: Callable[[TransferTask, Exception], None], remove_task: Callable[[FileItem], None], list_tasks: Callable[[], List[TransferJob]], expire_tasks: Callable[[], None], ) -> None: """保存队列用例依赖,避免 Application 服务绑定具体线程队列实现。""" self._register_task = register_task + self._admit_task = admit_task self._enqueue = enqueue self._before_enqueue = before_enqueue - self._after_enqueue = after_enqueue + self._enqueue_failed = enqueue_failed self._remove_task = remove_task self._list_tasks = list_tasks self._expire_tasks = expire_tasks def put(self, task: TransferTask, callback: Callable) -> bool: - """登记并入队一个整理任务,保持原有副作用顺序。""" + """先持久化准入事实再入队;任何前置失败都撤销内存作业视图。""" if not task or not self._register_task(task): return False - self._before_enqueue(task) - self._enqueue(TransferQueue(task=task, callback=callback)) - self._after_enqueue(task) + try: + admission = self._admit_task(task) + task.bind_admission_task_id(admission.task_id) + except Exception: + self._remove_task(task.fileitem) + raise + try: + self._before_enqueue(task) + self._enqueue(TransferQueue(task=task, callback=callback)) + except Exception as err: + try: + self._enqueue_failed(task, err) + finally: + self._remove_task(task.fileitem) + raise return True def remove(self, fileitem: FileItem) -> None: diff --git a/app/chain/transfer.py b/app/chain/transfer.py index 9a84e3071..1d1c0da4a 100755 --- a/app/chain/transfer.py +++ b/app/chain/transfer.py @@ -52,6 +52,7 @@ from app.application.outbox import ( from app.application.transfer import ( FailedRetryScheduler, JobManager, + TransferAdmission, TransferFailureNotification, TransferFailureNotificationAggregator, TransferQueue, @@ -225,8 +226,8 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo self.retry_scheduler = FailedRetryScheduler() # 整理失败通知聚合器 self.failure_notification_aggregator = TransferFailureNotificationAggregator() - # 待整理文件落盘登记,用于进程重启后回放内存队列里未完成的任务 - self._pendingoper = get_chain_transfer_pending_port() + # durable admission 仓储先于内存入队保存任务,进程退出后仍可恢复。 + self._transfer_admissions = get_chain_transfer_pending_port() # 转移成功的文件清单 self._success_target_files: Dict[Tuple, List[str]] = {} # 批次级刮削缓冲,避免同一批多文件入库重复触发目录刮削 @@ -867,7 +868,8 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo """ 添加到待整理队列 :param task: 任务信息 - :return: True表示任务已添加到队列,False表示任务无效或已存在(重复) + :return: True表示任务已添加,False表示链已关闭或任务无效/重复 + :raises Exception: 持久准入、批次登记或内存入队失败 """ with self._worker_state_lock: if self._closing: @@ -879,9 +881,10 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo """构建保持旧队列对象和私有兼容接缝的应用服务。""" return TransferQueueService( register_task=self.__put_to_jobview, + admit_task=self.__admit_transfer, enqueue=self._queue.put, before_enqueue=self._register_scrape_batch_task, - after_enqueue=self.__register_pending, + enqueue_failed=self.__record_enqueue_failure, remove_task=self.jobview.remove_task, list_tasks=self.jobview.list_jobs, expire_tasks=self.__expire_stale_transfer_tasks, @@ -936,7 +939,7 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo if stop_event.is_set(): return try: - pendings = self._pendingoper.list_all() + pendings = self._transfer_admissions.list_accepted() except Exception as err: logger.error(f"读取待整理文件登记失败:{err}") return @@ -944,9 +947,11 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo return logger.info(f"发现 {len(pendings)} 个上次未整理完的文件,正在重新送入整理链 ...") replayed = 0 - for storage, src_path in pendings: + for admission in pendings: if stop_event.is_set(): break + storage = admission.storage + src_path = admission.src_path try: fileitem, should_discard = self.__build_replay_fileitem(storage, src_path) # stat 等同步 I/O 返回后重新检查,关闭期间不得注销尚未完成的登记。 @@ -955,7 +960,9 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo if not fileitem: if should_discard: # 源文件确认已消失,注销登记避免每次启动重复回放 - self._pendingoper.discard(storage=storage, src_path=src_path) + self._transfer_admissions.discard_task( + task_id=admission.task_id + ) continue self.do_transfer(fileitem=fileitem) replayed += 1 @@ -1009,19 +1016,35 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo modify_time=modify_time, ), False - def __register_pending(self, task: TransferTask): - """ - 落盘登记一个待整理文件,登记失败不影响正常入队。 - :param task: 任务信息 - """ + def __admit_transfer(self, task: TransferTask) -> TransferAdmission: + """在内存入队前持久化源文件并返回稳定任务身份。""" fileitem = task.fileitem if task else None - if not fileitem or not fileitem.path: - return + if not fileitem or not fileitem.storage or not fileitem.path: + raise ValueError("整理任务缺少源文件身份") + return self._transfer_admissions.admit( + storage=fileitem.storage, + src_path=fileitem.path, + ) + + def __record_enqueue_failure( + self, + task: TransferTask, + error: Exception, + ) -> None: + """记录内存入队失败并撤销该任务的批次占位。""" try: - self._pendingoper.register(storage=fileitem.storage, src_path=fileitem.path) - except Exception as err: - # 登记只是重启后的补救手段,失败不能阻断正常整理 - logger.debug(f"登记待整理文件失败: {fileitem.path} - {err}") + if task.admission_task_id: + self._transfer_admissions.record_enqueue_failure( + task_id=task.admission_task_id, + error=str(error), + ) + except Exception as record_error: + logger.error( + "记录整理任务入队失败原因异常:" + f"{task.admission_task_id} - {record_error}" + ) + finally: + self._finish_scrape_batch_task(task) def __discard_pending(self, task: TransferTask): """ @@ -1031,13 +1054,17 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo 重试预算重新送入整理链;留在本表反而会每次重启都重复回放。 :param task: 任务信息 """ - fileitem = task.fileitem if task else None - if not fileitem or not fileitem.path: + if not task or not task.admission_task_id: return try: - self._pendingoper.discard(storage=fileitem.storage, src_path=fileitem.path) + self._transfer_admissions.discard_task( + task_id=task.admission_task_id + ) except Exception as err: - logger.debug(f"注销待整理文件登记失败: {fileitem.path} - {err}") + logger.error( + "注销整理任务 durable admission 失败: " + f"{task.admission_task_id} - {err}" + ) def __put_to_jobview(self, task: TransferTask) -> bool: """ @@ -2751,7 +2778,15 @@ class TransferChain(FileFilterMixin, ScrapeBatchMixin, EpisodeFormatMixin, Histo preview=preview, ) if background: - if self.put_to_queue(task=transfer_task): + try: + queued = self.put_to_queue(task=transfer_task) + except Exception as err: + all_success = False + message = f"{file_path.name} 加入整理队列失败:{err}" + err_msgs.append(message) + logger.error(message) + continue + if queued: logger.info(f"{file_path.name} 已添加到整理队列") else: logger.debug(f"{file_path.name} 已在整理队列中,跳过") diff --git a/app/db/adapters/transfer.py b/app/db/adapters/transfer.py new file mode 100644 index 000000000..b89163d39 --- /dev/null +++ b/app/db/adapters/transfer.py @@ -0,0 +1,115 @@ +"""整理任务持久准入端口的 SQLAlchemy 适配器。""" + +from collections.abc import Callable +from datetime import datetime +from uuid import uuid4 + +from sqlalchemy.exc import IntegrityError +from sqlalchemy.orm import Session + +from app.application.transfer import ( + TRANSFER_ADMISSION_ACCEPTED, + TransferAdmission, +) +from app.db.models.transferpending import TransferPending +from app.db.oper.transferpending import TransferPendingOper +from app.db.uow import SqlAlchemyUnitOfWork + + +class TransactionalTransferAdmissionRepository: + """以短生命周期 Session 实现整理任务持久准入端口。""" + + def __init__(self, session_factory: Callable[[], Session]) -> None: + """保存由组合根提供的同步会话工厂。""" + self._session_factory = session_factory + + @staticmethod + def _now() -> str: + """生成与历史登记时间可按字典序比较的当前时间。""" + return datetime.now().strftime("%Y-%m-%d %H:%M:%S") + + @staticmethod + def _project(pending: TransferPending) -> TransferAdmission: + """在 Session 有效期内把 ORM 行冻结为应用层 DTO。""" + created_at = pending.created_at or pending.updated_at + return TransferAdmission( + task_id=pending.task_id, + storage=pending.storage, + src_path=pending.src_path, + state=pending.state, + created_at=created_at, + updated_at=pending.updated_at, + last_error=pending.last_error, + ) + + def admit(self, *, storage: str, src_path: str) -> TransferAdmission: + """幂等持久化准入事实,并返回跨重启稳定的任务标识。""" + now_time = self._now() + try: + with self._session_factory() as session: + transaction = SqlAlchemyUnitOfWork(session) + try: + pending = TransferPendingOper(db=session).stage_admit( + task_id=uuid4().hex, + storage=storage, + src_path=src_path, + state=TRANSFER_ADMISSION_ACCEPTED, + now_time=now_time, + ) + if pending is None: + raise ValueError("整理任务的存储与源路径不能为空") + session.flush() + admission = self._project(pending) + transaction.commit() + return admission + except Exception: + transaction.rollback() + raise + except IntegrityError as error: + # 并发准入可能同时通过查询;唯一约束决定赢家,输家回读稳定身份。 + with self._session_factory() as session: + pending = TransferPendingOper(db=session).get_by_identity( + storage=storage, + src_path=src_path, + ) + if pending is None: + raise RuntimeError("并发准入冲突后未找到已提交记录") from error + return self._project(pending) + + def list_accepted(self, limit: int = 5000) -> list[TransferAdmission]: + """在独立只读会话中投影等待恢复或执行的准入记录。""" + with self._session_factory() as session: + pending_items = TransferPendingOper(db=session).list_by_state( + state=TRANSFER_ADMISSION_ACCEPTED, + 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: + transaction = SqlAlchemyUnitOfWork(session) + try: + TransferPendingOper(db=session).stage_record_enqueue_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: + transaction = SqlAlchemyUnitOfWork(session) + try: + deleted = TransferPendingOper(db=session).stage_discard_task( + task_id=task_id, + ) + transaction.commit() + return deleted + except Exception: + transaction.rollback() + raise diff --git a/app/db/models/transferpending.py b/app/db/models/transferpending.py index 9cfd9211d..bfa870623 100644 --- a/app/db/models/transferpending.py +++ b/app/db/models/transferpending.py @@ -1,6 +1,8 @@ -from typing import List, Optional +from datetime import datetime +from typing import List, Optional, cast +from uuid import uuid4 -from sqlalchemy import Index, String, delete, select +from sqlalchemy import Index, 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 @@ -15,22 +17,43 @@ class TransferPending(Base): 蒸发。而已经稳定落地的文件不会再产生任何监控事件,也不会有新的补偿扫描起点 ——结果就是永久漏件,只能靠人工比对补整理。 - 这里只落盘最小事实:存储与源文件路径。重启后重新走一遍整理入口,由整理历史 - 查重挡掉已经完成的,因此不需要序列化 meta/mediainfo 这些重对象,也不存在 - 识别结果陈旧的问题。 + 这里只落盘恢复所需的最小事实:稳定任务身份、存储、源文件路径、准入状态和 + 最近入队错误。重启后重新走一遍整理入口,由整理历史查重挡掉已经完成的, + 因此不需要序列化 meta/mediainfo 这些重对象,也不存在识别结果陈旧的问题。 """ id = get_id_column() + # 稳定任务标识 + task_id: Mapped[str] = mapped_column( + String(64), nullable=False, default=lambda: uuid4().hex + ) # 存储 storage: Mapped[str] = mapped_column(String, nullable=False) # 源文件路径 src_path: Mapped[str] = mapped_column(String, nullable=False) # 登记时间 created_at: Mapped[Optional[str]] = mapped_column(String) + # 持久状态 + state: Mapped[str] = mapped_column(String(32), nullable=False, default="accepted") + # 最后更新时间 + updated_at: Mapped[str] = mapped_column( + String(40), nullable=False, + default=lambda: datetime.now().strftime("%Y-%m-%d %H:%M:%S"), + ) + # 最近一次入队失败原因 + last_error: Mapped[Optional[str]] = mapped_column(Text) __table_args__ = ( # 同一个文件重复入队只保留一条,回放时不会重复送入整理链 Index("ux_transferpending_storage_path", "storage", "src_path", unique=True), + # 恢复主查询按状态过滤、登记时间与主键稳定排序 + Index( + "ix_transferpending_state_created", + "state", + "created_at", + "id", + ), + UniqueConstraint("task_id", name="uq_transferpending_task_id"), ) @classmethod @@ -50,11 +73,128 @@ class TransferPending(Base): select(cls).where(cls.storage == storage, cls.src_path == src_path) ).scalars().first() if pending: - return pending - pending = cls(storage=storage, src_path=src_path, created_at=now_time) + return cast("TransferPending", pending) + pending = cls( + storage=storage, + src_path=src_path, + state="accepted", + created_at=now_time, + updated_at=now_time, + ) db.add(pending) return pending + @classmethod + def stage_admit(cls, db: Session, *, task_id: str, storage: str, + src_path: str, state: str, + now_time: str) -> Optional["TransferPending"]: + """ + 在调用方会话中暂存一条持久接纳记录。 + + 相同存储与路径已存在时返回原记录,确保重复监控事件复用同一个任务标识。 + :param db: 数据库会话 + :param task_id: 任务标识 + :param storage: 存储 + :param src_path: 源文件路径 + :param state: 持久状态 + :param now_time: 当前时间 + :return: 接纳记录 + """ + if not task_id or not storage or not src_path or not state: + return None + pending = db.execute( + select(cls).where(cls.storage == storage, cls.src_path == src_path) + ).scalars().first() + if pending: + return cast("TransferPending", pending) + pending = cls( + task_id=task_id, + storage=storage, + src_path=src_path, + state=state, + created_at=now_time, + updated_at=now_time, + ) + db.add(pending) + return pending + + @classmethod + def list_by_state(cls, db: Session, *, state: str, + limit: Optional[int] = 5000) -> List["TransferPending"]: + """ + 按登记顺序列出指定持久状态的接纳记录。 + :param db: 数据库会话 + :param state: 持久状态 + :param limit: 单次读取上限 + :return: 接纳记录列表 + """ + if not state: + return [] + return list(db.execute( + select(cls) + .where(cls.state == state) + .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"]: + """ + 按存储与源路径查询一条持久接纳记录。 + :param db: 数据库会话 + :param storage: 存储 + :param src_path: 源文件路径 + :return: 接纳记录 + """ + if not storage or not src_path: + return None + return cast( + Optional["TransferPending"], + db.execute( + select(cls).where( + cls.storage == storage, + cls.src_path == src_path, + ) + ).scalars().first(), + ) + + @classmethod + def record_enqueue_failure(cls, db: Session, *, task_id: str, + error: str, now_time: str) -> int: + """ + 在调用方会话中记录任务最近一次入队失败。 + :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) + .values(last_error=error, updated_at=now_time), + execution_options={"synchronize_session": False}, + ) + + @classmethod + def discard_task(cls, db: Session, *, task_id: str) -> int: + """ + 在调用方会话中按任务标识删除接纳记录。 + :param db: 数据库会话 + :param task_id: 任务标识 + :return: 删除的记录数 + """ + if not task_id: + return 0 + return execute_dml( + db, delete(cls).where(cls.task_id == task_id), + execution_options={"synchronize_session": False}, + ) + @classmethod def discard(cls, db: Session, storage: str, src_path: str) -> int: """ diff --git a/app/db/oper/transferpending.py b/app/db/oper/transferpending.py index 38a9df933..204cffdcb 100644 --- a/app/db/oper/transferpending.py +++ b/app/db/oper/transferpending.py @@ -9,8 +9,9 @@ class TransferPendingOper(DbOper): """ 待整理文件登记管理。 - 只保存「存储 + 源文件路径」这一最小事实,用于在进程重启后把没走完整理链的 - 文件重新送回去,避免挂载故障重启后永久漏件。 + 保存稳定任务身份、存储、源文件路径和准入状态,用于在进程重启后把没走完 + 整理链的文件重新送回去,避免挂载故障重启后永久漏件。旧版路径登记接口继续 + 保留,供插件和兼容调用方使用。 """ def register(self, storage: str, src_path: str) -> Optional[TransferPending]: @@ -30,6 +31,93 @@ class TransferPendingOper(DbOper): ) ) + def stage_admit(self, *, task_id: str, storage: str, src_path: str, + state: str, now_time: str) -> Optional[TransferPending]: + """ + 在当前会话中暂存一条持久接纳记录。 + + 适配器应传入显式 Session,使提交与回滚仍由应用用例对应的 UoW 管理。 + :param task_id: 任务标识 + :param storage: 存储 + :param src_path: 源文件路径 + :param state: 持久状态 + :param now_time: 当前时间 + :return: 接纳记录 + """ + return self._execute_sync_write( + lambda session: TransferPending.stage_admit( + session, + task_id=task_id, + storage=storage, + src_path=src_path, + state=state, + now_time=now_time, + ) + ) + + def list_by_state(self, *, state: str, + limit: Optional[int] = 5000) -> List[TransferPending]: + """ + 使用当前会话列出指定状态记录。 + :param state: 持久状态 + :param limit: 单次读取上限 + :return: ORM 接纳记录列表 + """ + return self._execute_sync_query( + lambda session: TransferPending.list_by_state( + session, + state=state, + limit=limit, + ) + ) or [] + + def get_by_identity(self, *, storage: str, + src_path: str) -> Optional[TransferPending]: + """ + 使用当前会话按存储与源路径查询接纳记录。 + :param storage: 存储 + :param src_path: 源文件路径 + :return: 接纳记录 + """ + return self._execute_sync_query( + lambda session: TransferPending.get_by_identity( + session, + storage=storage, + src_path=src_path, + ) + ) + + def stage_record_enqueue_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_enqueue_failure( + session, + task_id=task_id, + error=error, + now_time=now_time, + ) + ) + + def stage_discard_task(self, *, task_id: str) -> int: + """ + 在当前会话中暂存按任务标识删除接纳记录。 + :param task_id: 任务标识 + :return: 删除的记录数 + """ + return self._execute_sync_write( + lambda session: TransferPending.discard_task( + session, + task_id=task_id, + ) + ) + def discard(self, storage: str, src_path: str) -> int: """ 注销一个待整理文件登记。 diff --git a/app/startup/initializers/modules.py b/app/startup/initializers/modules.py index b4cb5f87b..9a43cf897 100644 --- a/app/startup/initializers/modules.py +++ b/app/startup/initializers/modules.py @@ -110,6 +110,7 @@ from app.db.adapters.outbox import SqlAlchemyAsyncOutboxStager, SqlAlchemyOutbox from app.db.adapters.site import TransactionalSiteRepository from app.db.adapters.subscription import TransactionalSubscribeWriter from app.db.adapters.transaction import TransactionalWriteRunner +from app.db.adapters.transfer import TransactionalTransferAdmissionRepository from app.db.adapters.workflow import TransactionalWorkflowExecutionService from app.db.oper.agentchat import AgentChatOper from app.db.oper.agenttask import AgentTaskOper @@ -123,7 +124,6 @@ from app.db.oper.subscribe import SubscribeOper from app.db.oper.subscribehistory import SubscribeHistoryOper from app.db.oper.systemconfig import SystemConfigOper from app.db.oper.transferhistory import TransferHistoryOper -from app.db.oper.transferpending import TransferPendingOper from app.db.oper.user import UserOper from app.db.oper.userconfig import UserConfigOper from app.db.oper.workflow import WorkflowOper, configure_workflow_legacy_writer @@ -854,7 +854,9 @@ async def init_modules() -> HostRuntime: workflow=lambda: WorkflowOper(), download_history=lambda: DownloadHistoryOper(), transfer_history=lambda: TransferHistoryOper(), - transfer_pending=lambda: TransferPendingOper(), + transfer_pending=lambda: TransactionalTransferAdmissionRepository( + SessionFactory + ), media_server=lambda: MediaServerOper(), download_failure=lambda: TransactionalDownloadFailureRepository( SessionFactory diff --git a/database/versions/b1e7d3f5a9c2_3_0_13.py b/database/versions/b1e7d3f5a9c2_3_0_13.py new file mode 100644 index 000000000..4a3f66b88 --- /dev/null +++ b/database/versions/b1e7d3f5a9c2_3_0_13.py @@ -0,0 +1,165 @@ +"""3.0.13 为整理准入记录增加稳定身份与持久状态。 + +Revision ID: b1e7d3f5a9c2 +Revises: 5f2a9c1e7b4d +Create Date: 2026-08-27 +""" + +from datetime import datetime +from uuid import NAMESPACE_URL, uuid5 + +import sqlalchemy as sa +from alembic import op + +revision = "b1e7d3f5a9c2" +down_revision = "5f2a9c1e7b4d" +branch_labels = None +depends_on = None + +_TABLE_NAME = "transferpending" +_TASK_ID_CONSTRAINT = "uq_transferpending_task_id" +_STATE_CREATED_INDEX = "ix_transferpending_state_created" +_NEW_COLUMNS = {"task_id", "state", "updated_at", "last_error"} + + +def _column_names() -> set[str]: + """返回当前待整理登记表的字段集合。""" + inspector = sa.inspect(op.get_bind()) + if _TABLE_NAME not in inspector.get_table_names(): + return set() + return { + column["name"] + for column in inspector.get_columns(_TABLE_NAME) + } + + +def _has_task_id_constraint() -> bool: + """判断稳定任务标识唯一约束是否已经存在。""" + inspector = sa.inspect(op.get_bind()) + return any( + constraint.get("name") == _TASK_ID_CONSTRAINT + for constraint in inspector.get_unique_constraints(_TABLE_NAME) + ) + + +def _has_state_created_index() -> bool: + """判断恢复主查询的复合索引是否已经存在。""" + inspector = sa.inspect(op.get_bind()) + return any( + index.get("name") == _STATE_CREATED_INDEX + for index in inspector.get_indexes(_TABLE_NAME) + ) + + +def _backfill_admission_state() -> None: + """为旧登记保守生成稳定身份、接纳状态与更新时间。""" + pending = sa.table( + _TABLE_NAME, + sa.column("id", sa.Integer()), + sa.column("storage", sa.String()), + sa.column("src_path", sa.String()), + sa.column("created_at", sa.String()), + sa.column("task_id", sa.String()), + sa.column("state", sa.String()), + sa.column("updated_at", sa.String()), + ) + connection = op.get_bind() + rows = connection.execute( + sa.select( + pending.c.id, + pending.c.storage, + pending.c.src_path, + pending.c.created_at, + pending.c.task_id, + pending.c.state, + pending.c.updated_at, + ) + ).mappings().all() + fallback_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S") + for row in rows: + task_id = row["task_id"] or uuid5( + NAMESPACE_URL, + ( + "moviepilot:transferpending:" + f"{row['id']}:{row['storage']}:{row['src_path']}" + ), + ).hex + values = { + "task_id": task_id, + "state": row["state"] or "accepted", + "updated_at": ( + row["updated_at"] or row["created_at"] or fallback_time + ), + } + connection.execute( + pending.update() + .where(pending.c.id == row["id"]) + .values(**values) + ) + + +def upgrade() -> None: + """增加准入状态字段,回填旧行并建立稳定任务标识约束。""" + columns = _column_names() + if not columns: + return + if "task_id" not in columns: + op.add_column( + _TABLE_NAME, + sa.Column("task_id", sa.String(length=64), nullable=True), + ) + if "state" not in columns: + op.add_column( + _TABLE_NAME, + sa.Column("state", sa.String(length=32), nullable=True), + ) + if "updated_at" not in columns: + op.add_column( + _TABLE_NAME, + sa.Column("updated_at", sa.String(length=40), nullable=True), + ) + if "last_error" not in columns: + op.add_column( + _TABLE_NAME, + sa.Column("last_error", sa.Text(), nullable=True), + ) + + _backfill_admission_state() + with op.batch_alter_table(_TABLE_NAME) as batch_op: + batch_op.alter_column( + "task_id", existing_type=sa.String(length=64), nullable=False + ) + batch_op.alter_column( + "state", existing_type=sa.String(length=32), nullable=False + ) + batch_op.alter_column( + "updated_at", existing_type=sa.String(length=40), nullable=False + ) + if not _has_task_id_constraint(): + with op.batch_alter_table(_TABLE_NAME) as batch_op: + batch_op.create_unique_constraint( + _TASK_ID_CONSTRAINT, + ["task_id"], + ) + if not _has_state_created_index(): + op.create_index( + _STATE_CREATED_INDEX, + _TABLE_NAME, + ["state", "created_at", "id"], + unique=False, + ) + + +def downgrade() -> None: + """移除准入状态字段并保留旧版可识别的登记事实。""" + columns = _column_names() + if not columns or not (_NEW_COLUMNS & columns): + return + if _has_state_created_index(): + op.drop_index(_STATE_CREATED_INDEX, table_name=_TABLE_NAME) + with op.batch_alter_table(_TABLE_NAME) as batch_op: + if "task_id" in columns and _has_task_id_constraint(): + batch_op.drop_constraint(_TASK_ID_CONSTRAINT, type_="unique") + for column_name in ("last_error", "updated_at", "state", "task_id"): + if column_name in columns: + batch_op.drop_column(column_name) diff --git a/docs/architecture-optimization-checklist.md b/docs/architecture-optimization-checklist.md index bbcc0359b..2b0179a0b 100644 --- a/docs/architecture-optimization-checklist.md +++ b/docs/architecture-optimization-checklist.md @@ -69,7 +69,7 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` | 指标 | 当前值 | 解释 | |---|---:|---| -| 宿主 Python 模块 / 内部依赖边 | 835 / 6,817 | `dependency-baseline.json` 当前快照 | +| 宿主 Python 模块 / 内部依赖边 | 836 / 6,827 | `dependency-baseline.json` 当前快照 | | 非平凡 SCC | 2 | 新增 Chain 包根环;另一个是隔离的 29 模块 TMDB 移植包环 | | 跨层 DB 边界债务 | 0 | Application、Chain、API、Agent、Runtime、Workflow 到 DB 的受控债务均为零 | | Model/Oper 事务债务 | 0 | 自建 Session、自动事务装饰器、直接 commit/rollback 等基线均为零 | @@ -78,7 +78,7 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` | Python 源码量 | 约 271,400 行 | 60 个文件超过 1,000 行,14 个超过 2,000 行 | | 长方法 | 281 个超过 80 行 | 67 个超过 150 行,23 个超过 250 行;大量是私有方法 | | 全量 mypy 历史债务 | 11,983 / 601 文件 | strict frontier 当前只覆盖 41 个文件,且 ratchet 已新增 2 个错误 | -| Ruff 历史诊断 | 972 | 低水位门禁通过,但规则集只覆盖 `E4/E7/E9/F/I` | +| Ruff 历史诊断 | 967 | 低水位门禁通过,但规则集只覆盖 `E4/E7/E9/F/I` | | 覆盖率低水位 | Application 77.82%,Domain 79.24% | Chain、Runtime、Agent、Adapter、Startup 未进入包级覆盖率门禁 | ### 3.3 热点文件 @@ -106,8 +106,8 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` | ID | 优先级 | 状态 | 事项 | 目标结果 | |---|---|---|---|---| | ARCH-001 | P0 | 已交付 | 恢复 mypy ratchet | `5df388719` 已推送,主线既有 CI gate 通过 | -| ARCH-101 | P1 | 执行中 | 统一规则、总览、基线和语义门禁 | 文档声明与机器拒绝条件一一对应 | -| ARCH-102 | P1 | 待执行 | 将 Transfer pending 升级为真实 E3 状态机 | 崩溃窗口可判定恢复,结果未知时进入人工确认 | +| ARCH-101 | P1 | 已交付 | 统一规则、总览、基线和语义门禁 | `113355784` 已推送,Unit Tests `33031697902`、Pylint `33031697785` 全绿,远端 `0/0` | +| ARCH-102 | P1 | 执行中 | 将 Transfer pending 升级为真实 E3 状态机 | `S1-L1.1` 至 `S1-L1.5` 全部交付后,崩溃窗口可判定恢复,结果未知时进入人工确认 | | ARCH-103 | P1 | 待执行 | 类型化 Chain/Agent 数据 Port 与 DTO | 宿主主路径不再注入无 Session Oper,不向入口泄漏 ORM | | ARCH-104 | P1 | 待执行 | 收口跨多次写入的业务事务 | 站点/规则引用清理可整体回滚或幂等恢复 | | ARCH-105 | P1 | 待执行 | 明确 post-commit 与 Outbox 完成语义 | “业务已提交、后置效果 pending”可被调用方正确识别 | @@ -184,11 +184,23 @@ MoviePilot V3 已经形成较清晰的模块化单体:`foundation`、`domain` ### ARCH-102 将 Transfer pending 升级为真实 E3 状态机 +**分叶状态** + +- `S1-L1.1 Durable admission`:`VERIFIED`。已交付 persist-before-enqueue、Application-owned typed Port、 + DB adapter 与可逆 migration,宿主退出 raw/`Any` `TransferPendingOper` admission 路径。 +- `S1-L1.2 Planning checkpoint`:`PLANNED`。持久化稳定任务身份、整理模式、规划状态和目标 checkpoint。 +- `S1-L1.3 Lease 与恢复调度`:`PLANNED`。交付 claim/lease/heartbeat/attempt、过期接管与唯一恢复入口。 +- `S1-L1.4 幂等执行与终态结算`:`PLANNED`。交付文件/历史幂等、唯一 retry owner 和 + `manual_review` 语义。 +- `S1-L1.5 E3 全链收口`:`PLANNED`。完成崩溃矩阵、兼容验收与旧路径删除。此叶交付前, + ARCH-102 父项保持“执行中”,不得以局部绿色宣称 E3 完成。 + **问题与证据** -- `app/application/transfer.py:139-145` 当前先把任务接受到内存结构,再执行持久登记回调。 -- `app/chain/transfer.py:1012-1024` 吞掉登记失败并继续执行; - `tests/test_transfer_pending_replay.py:62-71` 固化了这一 fail-open 行为。 +- `S1-L1.1` 已将任务准入改为独立事务先 commit、后写内存队列;admission、batch 或 enqueue + 异常均不会伪装成重复任务成功,失败记录可供恢复。 +- 宿主 canonical Chain 只取得类型化 `TransferAdmissionRepository`;旧 Oper API 仅保留给统一兼容层, + 插件公开 `TransferTask.to_dict()` 字段未增加内部任务标识。 - worker 未知异常最终也会在 `app/chain/transfer.py:1107-1113,1262-1272` 删除 pending。 - `app/db/models/transferpending.py:9-34` 只有 `storage/src_path/created_at`,没有目标、模式、 step、lease、attempt、last_error,无法判定“文件已移动、历史未提交”等中间态。 diff --git a/docs/architecture-overview.md b/docs/architecture-overview.md index a646f8142..f175fc9b7 100644 --- a/docs/architecture-overview.md +++ b/docs/architecture-overview.md @@ -705,7 +705,7 @@ flowchart LR | 指标 | 当前值 | |---|---:| | Python 模块 | 835 | -| 内部导入边 | 6,817 | +| 内部导入边 | 6,827 | | 非平凡 SCC | 2(`ARCH-107` 临时 Chain 包根环;精确 containment 的 TMDB 移植包环) | | Direct egress | 66(12 条待迁移债务,54 条精确 containment) | | Module Contract V2 spec | 215(其中 214 个进入 `run_module` 观察面) | diff --git a/docs/architecture-refactor-roadmap.md b/docs/architecture-refactor-roadmap.md index eb117df62..8e9af3375 100644 --- a/docs/architecture-refactor-roadmap.md +++ b/docs/architecture-refactor-roadmap.md @@ -79,16 +79,24 @@ G-ARCH 只有在以下条件全部满足后才可完成: | S0-L2.4 Adapter zero-growth | `DELIVERED` | S0-L2.3 | `2553226f3`:冻结 28 条直连及 owner,收缩/新增/stale policy 门禁生效,远端 `0/0` | | S0-L2.4b HTTP/Egress 事实与政策 | `DELIVERED` | S0-L2.4 | `47f0de745`、`43d52a35b`、`8d602149f`:冻结 66 条出口事实,消除 CI 类型/覆盖率漂移;远端全绿且 `0/0` | | S0-L2.5 Event consumer 识别 | `DELIVERED` | S0-L2.1 | `86157be2a`:consumer 只识别可静态证明的 EventManager 注册;全量 6,459 passed / 6 skipped,CI `33029645165`/`33029645254` 全绿,远端 `0/0` | -| S0-L2.6 事实源与 CI 投影 | `VERIFIED` | S0-L2.2,S0-L2.4b,S0-L2.5 | 99/17 条逐调用事实、17 条 consumer policy 与 CI 分层本地通过;全量 6,481 passed / 6 skipped,待推送 CI | +| S0-L2.6 事实源与 CI 投影 | `DELIVERED` | S0-L2.2,S0-L2.4b,S0-L2.5 | `113355784`:99/17 条逐调用事实、17 条 consumer policy 与 CI 分层交付;Unit Tests `33031697902`、Pylint `33031697785` 全绿,远端 `0/0` | ### S1:可靠性、事务与数据合同 退出条件:Transfer 达到 E3;正式数据 Port 类型化;跨表业务操作有单一 UoW;post-commit/Outbox 竞争、失败呈现、at-least-once 和幂等语义全部闭环。 +`S1-L1` 是 ARCH-102 的 Transfer E3 父项,当前状态为**执行中**。只有 `S1-L1.1` 至 +`S1-L1.5` 全部 `DELIVERED`,真实调用链完成迁移且旧 fail-open 路径退出 canonical 主程序后, +父项和 ARCH-102 才能标记已交付。 + | Leaf | 状态 | 依赖 | 完成定义 | |---|---|---|---| -| S1-L1 Transfer E3 状态机 | `PLANNED` | S0 | `ARCH-102` 全部完成:migration、持久状态、checkpoint、lease、唯一 retry owner、重放和 manual review 一次交付 | +| S1-L1.1 Durable admission | `VERIFIED` | S0 | Application-owned typed Port + DB adapter + migration 落地;先持久 commit 再入队,入队失败保留可恢复记录;宿主不再通过 raw/`Any` `TransferPendingOper` 处理 admission | +| S1-L1.2 Planning checkpoint | `PLANNED` | S1-L1.1 | 稳定任务身份、整理模式和 planning 状态持久化;目标路径只在规划完成后写入 checkpoint,任何文件副作用前已有可判定状态 | +| S1-L1.3 Lease 与恢复调度 | `PLANNED` | S1-L1.2 | claim/lease/heartbeat/attempt 与过期接管规则落地;启动回放和同进程恢复共用唯一调度入口,同一任务同时只有一个 worker owner | +| S1-L1.4 幂等执行与终态结算 | `PLANNED` | S1-L1.3 | 文件操作、历史提交和 checkpoint 可重放;唯一 retry owner 生效,未知外部结果进入 `manual_review`,仅完整终态删除 pending | +| S1-L1.5 E3 全链收口 | `PLANNED` | S1-L1.4 | 崩溃矩阵、升级/降级、重复回放和插件 ABI 验收完整;旧 fail-open、重复状态与兼容层外旧入口删除,ARCH-102 债务归零 | | S1-L2 Workflow typed query | `PLANNED` | S0 | Workflow Application Port 不返回 `Any`/ORM,Session 内投影 DTO,正式调用方全部切换 | | S1-L3 Chain/Agent typed data ports | `PLANNED` | S1-L2 | `ChainDataPorts`/`AgentDataPorts` 的 raw Oper/`Any` factory 全部清零,兼容调用进入 Legacy 层 | | S1-L4 Subscription mutation UoW | `PLANNED` | S1-L3 | Subscription mutation 不跨 Session 传 ORM,正式写路径一个 UoW,旧自动事务入口退出 canonical 路径 | @@ -103,10 +111,10 @@ G-ARCH 只有在以下条件全部满足后才可完成: | Leaf | 状态 | 依赖 | 完成定义 | |---|---|---|---| | S2-L1 日志/消息资源显式生命周期 | `PLANNED` | S0 | import 和非消息 Chain 构造零新增线程;bootstrap 显式创建,失败和正常关闭均收口 | -| S2-L2 ChainBase 与 SCC 清零 | `PLANNED` | S0-L3 | canonical `app.chain.base` 落地,包根无 eager/重复导出,宿主包根导入清零,Chain SCC 消失 | +| S2-L2 ChainBase 与 SCC 清零 | `PLANNED` | S0-L2.2 | canonical `app.chain.base` 落地,包根无 eager/重复导出,宿主包根导入清零,Chain SCC 消失 | | S2-L3 GlobalVar/provider 注册收口 | `PLANNED` | S2-L1 | `global_vars` canonical 消费清零,provider 注册进入显式装配阶段并可 reset;Legacy 入口精确保留 | -| S2-L4 Passkey 缓存边界 | `PLANNED` | S0-L4 | Application 不识别 Redis;原子 consume 由 runtime cache contract + backend 实现 | -| S2-L5 Backup artifact Port | `PLANNED` | S0-L4 | Application 不构造 `BackupFiles`,文件 I/O 由注入 Adapter 拥有 | +| S2-L4 Passkey 缓存边界 | `PLANNED` | S0-L2.4 | Application 不识别 Redis;原子 consume 由 runtime cache contract + backend 实现 | +| S2-L5 Backup artifact Port | `PLANNED` | S0-L2.4 | Application 不构造 `BackupFiles`,文件 I/O 由注入 Adapter 拥有 | | S2-L6 Application Adapter/DNS 债务清零 | `PLANNED` | S2-L4,S2-L5 | Application 到具体 Adapter 的未批准边归零,SSRF DNS I/O 进入注入 Port,批准通用机制有精确规则和门禁 | | S2-L7 Chain Adapter/宿主 HTTP 债务清零 | `PLANNED` | S2-L6 | Chain 具体 Adapter 与 11 条普通 direct HTTP/Session bridge 归零;SDK/stream/vendor 例外保持精确 containment | @@ -133,10 +141,10 @@ G-ARCH 只有在以下条件全部满足后才可完成: | Leaf | 状态 | 依赖 | 完成定义 | |---|---|---|---| | S4-L1 Module strict contract | `PLANNED` | S2-L7 | 宿主 provider `ANY` 结果归零,宿主 admission strict,官方/第三方插件分级兼容 | -| S4-L2 Event strict contract | `PLANNED` | S0-L5,S1-L6 | 宿主事件输入/输出按风险 strict,诊断例外只属于第三方插件兼容 | +| S4-L2 Event strict contract | `PLANNED` | S0-L2.6,S1-L6 | 宿主事件输入/输出按风险 strict,诊断例外只属于第三方插件兼容 | | S4-L3 Complexity v2 | `PLANNED` | S3 | 私有方法、class/file、圈复杂度进入门禁;所有超限通过职责拆分归零 | | S4-L4 全量 mypy 清零 | `PLANNED` | S3,S4-L1,S4-L2 | `mypy-baseline.json` 归零并删除债务接受路径,全宿主 strict 类型通过 | -| S4-L5 Ruff 治理债务清零 | `PLANNED` | S3 | 当前受控 972 条诊断归零,规则集扩展经过独立审查且新增诊断为零 | +| S4-L5 Ruff 治理债务清零 | `PLANNED` | S3 | 当前受控 967 条诊断归零,规则集扩展经过独立审查且新增诊断为零 | | S4-L6 Coverage/并发/质量证据 | `PLANNED` | S3,S4-L1,S4-L2 | 高风险包纳入 coverage;raw concurrency 分类清零;Module Quality 有真实 evidence test | ### S5:Plugin、Agent、Domain、Startup 与最终收口 @@ -154,57 +162,68 @@ G-ARCH 只有在以下条件全部满足后才可完成: ## 4. 当前活动叶子 -### S0-L2.6 事实源与 CI 投影 +### S1-L1.1 Durable admission -**Status:** `VERIFIED`(本地验收完成,等待提交、推送和远端 CI 确认) +**Status:** `VERIFIED` **Outcome** -统一 Event producer/consumer 的 AST 事实源,完整解析 positional/keyword 参数、别名、重绑定和 -有限条件表达式。生成快照保存逐调用 line-free 事实及 multiplicity;consumer 由独立人工 policy -按 exact fingerprint set 准入,任何刷新快照的操作都不能自动接受新消费注册。 +把“接受整理任务”变成真正的 durable admission:Application 先通过类型化 Port 在独立事务中 +commit pending,再尝试写入进程内队列。队列写入失败或进程在 commit 后退出时,持久记录仍能成为 +后续恢复起点;持久化失败则不允许任务进入队列。该叶只结清 admission 所有权和顺序,不预先宣称 +planning、lease、幂等执行或终态恢复已经完成。 **Ownership** -- `scripts/architecture/event_facts.py` 的统一 producer/consumer provenance collector。 -- `scripts/architecture/event_policy.py` 与 `runtime-contract-policy.json` 的只读人工 consumer policy。 -- `scripts/architecture/baseline.py` 的 runtime schema v3、迁移链、事实索引与 diagnostics。 -- producer/consumer、policy、baseline/CLI 和 CI 分层测试。 -- 架构规范、总览、优化清单与本路线图的单一事实说明。 +- `app/application/transfer.py` 拥有 admission DTO、Protocol、结果语义与 persist-before-enqueue 编排。 +- `app/db/adapters/` 提供短 Session/UoW 的 Transfer pending 持久化实现;`app/db/oper/` 只接收 + adapter 拥有的 Session 并 stage/flush。 +- `app/startup/` 负责构造并注入 adapter,宿主 Chain 不再取得 raw/`Any` `TransferPendingOper`。 +- `app/db/models/transferpending.py` 与配套 Alembic migration 只承载该叶需要的 durable admission + schema,并验证既有记录升级及 downgrade。 +- Transfer queue、pending repository、migration、架构边界及兼容回归测试。 **Excluded** -- 不修改生产 EventManager、事件 ABI、handler 执行顺序或插件消费者。 -- 不把 `app/plugins/**` 副本纳入宿主扫描。 -- 不使用源码行号、通配符或自动写入 policy 接受新 consumer。 +- 不在本叶实现 planning checkpoint、lease/heartbeat、文件操作幂等、历史结算或 `manual_review`; + 这些分别由 `S1-L1.2` 至 `S1-L1.4` 完整交付。 +- 不修改插件公开 Transfer ABI、旧插件导入路径或第三方插件行为;必要兼容只通过统一 Compat/Legacy + 层委托新的 canonical admission 实现。 +- 不修改或扫描 `app/plugins/**` 插件副本。 +- 不保留宿主 canonical raw/`Any` pending port 与新 typed Port 两套正式入口。 **Acceptance** ```bash .venv/bin/python -m pytest \ - tests/test_architecture_event_facts.py \ - tests/test_architecture_event_policy.py \ + tests/test_transfer_queue_service.py \ + tests/test_transfer_pending_replay.py \ + tests/test_transfer_admission_migration.py \ + tests/test_db_transferpending_queries.py \ + tests/test_database_migration_startup.py \ tests/test_architecture_dependencies.py \ - tests/test_architecture_contract_baseline.py \ - tests/test_architecture_baseline_cli.py \ - tests/test_architecture_ci.py -q -.venv/bin/python scripts/architecture/event_policy.py + tests/test_legacy_import_compat.py -q +.venv/bin/mypy --config-file mypy.ini .venv/bin/python scripts/architecture/baseline.py --check-host --diagnostics .venv/bin/python scripts/architecture/ruff_ratchet.py .venv/bin/python scripts/architecture/mypy_ratchet.py -.venv/bin/pylint scripts/architecture/baseline.py \ - scripts/architecture/event_facts.py \ - scripts/architecture/event_policy.py \ - tests/test_architecture_event_facts.py \ - tests/test_architecture_event_policy.py \ - tests/test_architecture_dependencies.py \ - tests/test_architecture_contract_baseline.py \ - tests/test_architecture_baseline_cli.py \ - tests/test_architecture_ci.py +uv run --locked --no-sync python tests/run.py git diff --check ``` **Delivery** -- 单一提交主题:统一 Event facts、锁定 consumer policy 并拆分 CI 语义/快照投影。 +- 单一提交主题:交付 Transfer durable admission、类型化持久化边界及可逆 migration。 +- 提交前证明 pending commit 发生在 enqueue 之前;持久化失败不入队,enqueue 失败或 commit 后崩溃 + 均保留可恢复记录;宿主 canonical 路径不再导入或取得 raw/`Any` `TransferPendingOper`。 +- 运行锁定全量测试与 scoped Pylint;插件公开 Transfer ABI 和统一兼容导入测试必须保持通过。 - 推送 `origin/v3` 后确认提交祖先关系、远端 SHA 和 ahead/behind `0/0`。 + +**Local verification (2026-08-27)** + +- 锁定全量:`6,496 passed, 7 skipped`;跳过项包含本机未配置隔离库的 PostgreSQL migration + 用例,SQLite upgrade/downgrade/re-upgrade 已真实执行。 +- scoped Pylint:`10.00/10`;host dependency baseline、Ruff ratchet、mypy ratchet 与 + `git diff --check` 全部通过。 +- failure injection 已覆盖 admission 失败不入队、batch/enqueue 失败保留记录、批次返回失败、 + queue -> worker -> terminal discard 稳定身份,以及 Legacy TransferTask 序列化字段不变。 diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index 3203d48ac..9ce8f1142 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -465,6 +465,13 @@ Durable post-commit side effects have a separate boundary: cancellation and bounded shutdown waiting, but it is not a durable queue and must not replace an Outbox or persistent task table. +Transfer durable admission follows the same ownership direction without using +the Outbox as an execution queue: `app/application/transfer.py` owns the typed +admission contract and persist-before-enqueue orchestration, while +`app/db/adapters/transfer.py` commits it in a short Session/UoW. Canonical host +chains never obtain `TransferPendingOper`; its no-Session API remains only for +the exact legacy plugin import contract. + ## Composition and Compatibility Boundaries - Startup registers concrete cache factories before decorated business modules @@ -597,6 +604,8 @@ driven workflow registration. | `app/application/subscription/write.py` | Subscription media translation and sync/async write-port orchestration | | `app/application/outbox.py` | Durable intent, topic handler and Outbox repository contracts | | `app/db/adapters/outbox.py` | SQLAlchemy Outbox persistence, claim/lease and retry state adapter | +| `app/application/transfer.py` | Transfer task, durable admission contract and persist-before-enqueue use case | +| `app/db/adapters/transfer.py` | SQLAlchemy durable admission persistence and detached snapshot adapter | | `app/application/scheduling.py` | Runtime scheduler facade for Agent tools and endpoints; `Scheduler` class registered by `app/startup/initializers/scheduler.py` | | `app/application/commands.py` | Command registry facade for Agent tools and endpoints; `Command` class registered by `app/startup/initializers/command.py` | | `app/application/workflow.py` | Workflow use cases plus the runtime port consumed by API and Chain; `WorkFlowManager` is registered by `app/startup/initializers/workflow.py` | diff --git a/tests/conftest.py b/tests/conftest.py index b12474988..02dc5cecd 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -230,7 +230,7 @@ def configure_plugin_system_services(): from app.db.oper.subscribe import SubscribeOper from app.db.oper.subscribehistory import SubscribeHistoryOper from app.db.oper.transferhistory import TransferHistoryOper - from app.db.oper.transferpending import TransferPendingOper + from app.db.adapters.transfer import TransactionalTransferAdmissionRepository from app.db.oper.user import UserOper from app.db.oper.workflow import WorkflowOper, configure_workflow_legacy_writer from app.db.oper.message import MessageOper @@ -303,7 +303,9 @@ def configure_plugin_system_services(): workflow=lambda: WorkflowOper(), download_history=lambda: DownloadHistoryOper(), transfer_history=lambda: TransferHistoryOper(), - transfer_pending=lambda: TransferPendingOper(), + transfer_pending=lambda: TransactionalTransferAdmissionRepository( + SessionFactory + ), media_server=lambda: MediaServerOper(), download_failure=lambda: TransactionalDownloadFailureRepository( SessionFactory diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index bdce75911..e0128f404 100644 --- a/tests/fixtures/architecture/dependency-baseline.json +++ b/tests/fixtures/architecture/dependency-baseline.json @@ -1441,8 +1441,8 @@ "runtime_only": true } }, - "edge_count": 6817, - "edge_sha256": "e3d43fec9f7bc936ef5a2ffe7ba11ea054d1ba7d1c42e7101978480cd99a63fe", + "edge_count": 6827, + "edge_sha256": "34e2be621e40f0f07ff04065655446c4ad0062883701f9e1e2235c29359900f7", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -3973,6 +3973,8 @@ "app.application.chain.context -> app.application.configuration", "app.application.chain.context -> app.runtime", "app.application.chain.context -> app.runtime.stop", + "app.application.chain.data -> app.application", + "app.application.chain.data -> app.application.transfer", "app.application.chain.durable_events -> app.application", "app.application.chain.durable_events -> app.application.history", "app.application.chain.durable_events -> app.domain", @@ -5141,6 +5143,14 @@ "app.db.adapters.subscription -> app.db.uow", "app.db.adapters.transaction -> app.db", "app.db.adapters.transaction -> app.db.uow", + "app.db.adapters.transfer -> app.application", + "app.db.adapters.transfer -> app.application.transfer", + "app.db.adapters.transfer -> app.db", + "app.db.adapters.transfer -> app.db.models", + "app.db.adapters.transfer -> app.db.models.transferpending", + "app.db.adapters.transfer -> app.db.oper", + "app.db.adapters.transfer -> app.db.oper.transferpending", + "app.db.adapters.transfer -> app.db.uow", "app.db.adapters.workflow -> app.application", "app.db.adapters.workflow -> app.application.workflow", "app.db.adapters.workflow -> app.db", @@ -7905,6 +7915,7 @@ "app.startup.initializers.modules -> app.db.adapters.site", "app.startup.initializers.modules -> app.db.adapters.subscription", "app.startup.initializers.modules -> app.db.adapters.transaction", + "app.startup.initializers.modules -> app.db.adapters.transfer", "app.startup.initializers.modules -> app.db.adapters.workflow", "app.startup.initializers.modules -> app.db.oper", "app.startup.initializers.modules -> app.db.oper.agentchat", @@ -7919,7 +7930,6 @@ "app.startup.initializers.modules -> app.db.oper.subscribehistory", "app.startup.initializers.modules -> app.db.oper.systemconfig", "app.startup.initializers.modules -> app.db.oper.transferhistory", - "app.startup.initializers.modules -> app.db.oper.transferpending", "app.startup.initializers.modules -> app.db.oper.user", "app.startup.initializers.modules -> app.db.oper.userconfig", "app.startup.initializers.modules -> app.db.oper.workflow", @@ -8262,7 +8272,7 @@ "app.workflow.actions.transfer_file -> app.workflow", "app.workflow.actions.transfer_file -> app.workflow.actions" ], - "module_count": 835, + "module_count": 836, "modules": [ "app", "app.adapters", @@ -8660,6 +8670,7 @@ "app.db.adapters.site", "app.db.adapters.subscription", "app.db.adapters.transaction", + "app.db.adapters.transfer", "app.db.adapters.workflow", "app.db.base", "app.db.decorators", diff --git a/tests/fixtures/architecture/mypy-baseline.json b/tests/fixtures/architecture/mypy-baseline.json index 76c498ce8..42b173ffb 100644 --- a/tests/fixtures/architecture/mypy-baseline.json +++ b/tests/fixtures/architecture/mypy-baseline.json @@ -1595,7 +1595,7 @@ "no-any-return": 3, "no-redef": 2, "no-untyped-call": 8, - "no-untyped-def": 16, + "no-untyped-def": 15, "operator": 3, "return-value": 2, "truthy-function": 5, @@ -1704,9 +1704,6 @@ "no-untyped-def": 20, "type-arg": 1 }, - "app/db/models/transferpending.py": { - "no-any-return": 1 - }, "app/db/models/user.py": { "no-untyped-def": 10 }, diff --git a/tests/fixtures/architecture/ruff-baseline.json b/tests/fixtures/architecture/ruff-baseline.json index 25e5f41fb..49c77abf7 100644 --- a/tests/fixtures/architecture/ruff-baseline.json +++ b/tests/fixtures/architecture/ruff-baseline.json @@ -401,9 +401,6 @@ "app/application/torrent_cache.py": { "I001": 1 }, - "app/application/transfer.py": { - "I001": 1 - }, "app/application/workflow.py": { "I001": 1 }, @@ -1450,9 +1447,6 @@ "tests/test_interaction_router.py": { "E402": 5 }, - "tests/test_legacy_import_compat.py": { - "I001": 1 - }, "tests/test_lifecycle_shutdown.py": { "F841": 1 }, @@ -1873,18 +1867,12 @@ "tests/test_transfer_movie_collection.py": { "I001": 1 }, - "tests/test_transfer_pending_replay.py": { - "I001": 1 - }, "tests/test_transfer_preview.py": { "I001": 1 }, "tests/test_transfer_queue_count.py": { "I001": 1 }, - "tests/test_transfer_queue_service.py": { - "I001": 1 - }, "tests/test_transfer_rename_build_event.py": { "I001": 1 }, @@ -1894,9 +1882,6 @@ "tests/test_transfer_tmdb_category.py": { "I001": 1 }, - "tests/test_transfer_worker_lifecycle.py": { - "I001": 1 - }, "tests/test_transferhistory_media_source_migration.py": { "I001": 1 }, diff --git a/tests/test_architecture_dependencies.py b/tests/test_architecture_dependencies.py index ca4528d04..211953dff 100644 --- a/tests/test_architecture_dependencies.py +++ b/tests/test_architecture_dependencies.py @@ -478,6 +478,95 @@ def test_transfer_chains_use_explicit_data_port_getters(): assert violations == [] +def test_transfer_pending_oper_import_is_confined_to_database_boundary(): + """宿主仅允许事务适配器和兼容导出直接导入整理待处理 Oper。""" + allowed_paths = { + "app/db/adapters/transfer.py", + "app/db/oper/__init__.py", + } + violations: list[str] = [] + for path in APP_ROOT.rglob("*.py"): + relative = path.relative_to(PROJECT_ROOT).as_posix() + if relative.startswith("app/plugins/") or relative in allowed_paths: + continue + tree = ast.parse(path.read_text(encoding="utf-8-sig"), filename=str(path)) + for node in ast.walk(tree): + if isinstance(node, ast.Import): + if any( + alias.name == "app.db.oper.transferpending" + for alias in node.names + ): + violations.append(f"{relative}:{node.lineno}") + elif isinstance(node, ast.ImportFrom) and ( + node.module == "app.db.oper.transferpending" + or ( + node.module == "app.db.oper" + and any( + alias.name in {"transferpending", "TransferPendingOper"} + for alias in node.names + ) + ) + ): + violations.append(f"{relative}:{node.lineno}") + + assert violations == [] + + +def test_startup_injects_transactional_transfer_admission_repository(): + """启动组合根必须向 Chain 注入事务型整理准入仓储。""" + path = APP_ROOT / "startup" / "initializers" / "modules.py" + tree = ast.parse(path.read_text(encoding="utf-8-sig"), filename=str(path)) + imports_repository = any( + isinstance(node, ast.ImportFrom) + and node.module == "app.db.adapters.transfer" + and any( + alias.name == "TransactionalTransferAdmissionRepository" + for alias in node.names + ) + for node in ast.walk(tree) + ) + transfer_pending_factories = [ + keyword.value + for node in ast.walk(tree) + if isinstance(node, ast.Call) + and isinstance(node.func, ast.Name) + and node.func.id == "configure_chain_data_ports" + for keyword in node.keywords + if keyword.arg == "transfer_pending" + ] + + assert imports_repository is True + assert len(transfer_pending_factories) == 1 + assert any( + isinstance(node, ast.Call) + and isinstance(node.func, ast.Name) + and node.func.id == "TransactionalTransferAdmissionRepository" + for node in ast.walk(transfer_pending_factories[0]) + ) + + +def test_transfer_pending_chain_port_is_typed_without_legacy_proxy(): + """整理准入端口必须返回明确 Protocol,且旧 Proxy 不得重新出现。""" + path = APP_ROOT / "application" / "chain" / "data.py" + tree = ast.parse(path.read_text(encoding="utf-8-sig"), filename=str(path)) + class_names = { + node.name + for node in tree.body + if isinstance(node, ast.ClassDef) + } + getters = [ + node + for node in tree.body + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) + and node.name == "get_chain_transfer_pending_port" + ] + + assert "TransferPendingPortProxy" not in class_names + assert len(getters) == 1 + assert getters[0].returns is not None + assert ast.unparse(getters[0].returns) == "TransferAdmissionRepository" + + def test_agent_consumers_use_explicit_data_port_getters(): """Agent 生产模块不得把兼容数据端口代理重新伪装成数据库 Oper。""" forbidden = { diff --git a/tests/test_db_transferpending_queries.py b/tests/test_db_transferpending_queries.py index d7f3afe09..a902937cf 100644 --- a/tests/test_db_transferpending_queries.py +++ b/tests/test_db_transferpending_queries.py @@ -6,8 +6,11 @@ 因此这里对着真实数据库断言查回的内容,而不是断言调用了什么。 """ import pytest +from sqlalchemy import create_engine +from sqlalchemy.orm import sessionmaker from app.db import base as db_base +from app.db.adapters.transfer import TransactionalTransferAdmissionRepository from app.db.models.transferpending import TransferPending from app.db.oper.transferpending import TransferPendingOper @@ -195,3 +198,174 @@ def test_oper_discard_and_clear_report_counts(db): assert oper.discard(storage="local", src_path="/mnt/a.mkv") == 1 assert oper.clear() >= 1 assert oper.list_all() == [] + + +def test_stage_admit_is_idempotent_and_keeps_stable_task_id(db): + """显式准入重复执行时必须复用首个稳定任务标识。""" + first = TransferPending.stage_admit( + db.session, + task_id="task-first", + storage="local", + src_path="/mnt/durable.mkv", + state="accepted", + now_time="2026-08-27 10:00:00", + ) + second = TransferPending.stage_admit( + db.session, + task_id="task-second", + storage="local", + src_path="/mnt/durable.mkv", + state="accepted", + now_time="2026-08-27 11:00:00", + ) + + assert first is second + assert second.task_id == "task-first" + assert second.updated_at == "2026-08-27 10:00:00" + + +def test_state_queries_failure_record_and_task_discard(db): + """状态查询、失败留痕和按任务删除应共享同一稳定身份。""" + TransferPending.stage_admit( + db.session, + task_id="task-accepted", + storage="local", + src_path="/mnt/accepted.mkv", + state="accepted", + now_time="2026-08-27 10:00:00", + ) + db.add(TransferPending( + task_id="task-other", + storage="local", + src_path="/mnt/other.mkv", + state="other", + created_at="2026-08-27 10:00:01", + updated_at="2026-08-27 10:00:01", + )) + db.session.flush() + + accepted = TransferPending.list_by_state( + db.session, + state="accepted", + ) + assert [item.task_id for item in accepted] == ["task-accepted"] + assert TransferPending.record_enqueue_failure( + db.session, + task_id="task-accepted", + error="queue full", + now_time="2026-08-27 10:01:00", + ) == 1 + db.session.expire_all() + failed = TransferPending.get_by_identity( + db.session, + storage="local", + src_path="/mnt/accepted.mkv", + ) + assert failed.last_error == "queue full" + assert failed.updated_at == "2026-08-27 10:01:00" + assert TransferPending.discard_task( + db.session, + task_id="task-accepted", + ) == 1 + + +def test_oper_staging_reuses_explicit_write_session(db, monkeypatch): + """Oper 的新暂存入口必须服从调用方 Session,不得隐式提交。""" + monkeypatch.setattr( + db_base, + "run_sync_transaction", + lambda _operation: (_ for _ in ()).throw( + AssertionError("不应创建额外同步事务") + ), + ) + oper = TransferPendingOper(db.session) + + pending = oper.stage_admit( + task_id="task-explicit", + storage="local", + src_path="/mnt/explicit-stage.mkv", + state="accepted", + now_time="2026-08-27 10:00:00", + ) + assert pending.task_id == "task-explicit" + assert [item.task_id for item in oper.list_by_state(state="accepted")] == [ + "task-explicit" + ] + assert oper.stage_record_enqueue_failure( + task_id="task-explicit", + error="queue full", + now_time="2026-08-27 10:01:00", + ) == 1 + assert oper.stage_discard_task(task_id="task-explicit") == 1 + + +def test_transactional_repository_commits_frozen_projections(tmp_path): + """适配器应独立提交 UoW,并在会话关闭前冻结应用 DTO。""" + engine = create_engine(f"sqlite:///{tmp_path / 'transfer.db'}") + TransferPending.__table__.create(engine) + factory = sessionmaker(bind=engine) + repository = TransactionalTransferAdmissionRepository(factory) + + admitted = repository.admit( + storage="local", + src_path="/mnt/repository.mkv", + ) + repeated = repository.admit( + storage="local", + src_path="/mnt/repository.mkv", + ) + assert repeated == admitted + assert admitted.task_id + assert admitted.state == "accepted" + assert repository.list_accepted() == [admitted] + + repository.record_enqueue_failure( + task_id=admitted.task_id, + error="queue full", + ) + failed = repository.list_accepted()[0] + assert failed.last_error == "queue full" + assert repository.discard_task(task_id=admitted.task_id) == 1 + assert repository.list_accepted() == [] + engine.dispose() + + +def test_transactional_repository_rolls_back_failed_write(monkeypatch): + """适配器写入异常时必须回滚自身 UoW 并传播原异常。""" + class SessionContext: + """为回滚断言提供最小 Session 上下文。""" + + def __init__(self): + """初始化提交与回滚计数。""" + self.commits = 0 + self.rollbacks = 0 + + def __enter__(self): + """返回当前伪会话。""" + return self + + def __exit__(self, _exc_type, _exc_value, _traceback): + """不吞掉被测异常。""" + return False + + def commit(self): + """记录提交调用。""" + self.commits += 1 + + def rollback(self): + """记录回滚调用。""" + self.rollbacks += 1 + + session = SessionContext() + repository = TransactionalTransferAdmissionRepository(lambda: session) + monkeypatch.setattr( + TransferPendingOper, + "stage_record_enqueue_failure", + lambda self, **_kwargs: (_ for _ in ()).throw(ValueError("write failed")), + ) + + with pytest.raises(ValueError, match="write failed"): + repository.record_enqueue_failure(task_id="task", error="failure") + + assert session.rollbacks == 1 + assert session.commits == 0 diff --git a/tests/test_legacy_import_compat.py b/tests/test_legacy_import_compat.py index 64341354a..ca138ba6f 100644 --- a/tests/test_legacy_import_compat.py +++ b/tests/test_legacy_import_compat.py @@ -14,16 +14,15 @@ from app.runtime.compat.diagnostics import ( ) from app.runtime.compat.imports import install_legacy_import_hook from app.runtime.compat.manifest import ( + _MESSAGE_NOTIFICATION_SYMBOL_ALIASES, MODULE_ALIASES, PACKAGE_ALIASES, PACKAGE_EXPORTS, SYMBOL_ALIASES, VIRTUAL_PACKAGES, ModuleAlias, - _MESSAGE_NOTIFICATION_SYMBOL_ALIASES, ) - LEGACY_PACKAGE = "legacy_compat_test" LEGACY_MODULE = f"{LEGACY_PACKAGE}.target" CANONICAL_PACKAGE = "canonical_compat_test" @@ -295,6 +294,23 @@ def test_physical_modules_resolve_moved_symbols_without_reverse_imports(): assert schemas_package.TransferQueue is legacy_transfer.TransferQueue +def test_legacy_transfer_task_hides_internal_admission_identity(): + """持久准入身份不得改变插件旧任务字典的公开字段集合。""" + legacy_transfer = importlib.import_module("app.sdk._legacy.transfer") + task = legacy_transfer.TransferTask(fileitem={ + "storage": "local", + "path": "/downloads/movie.mkv", + "type": "file", + }) + public_fields = set(task.to_dict()) + + task.bind_admission_task_id("internal-task-id") + + assert set(task.to_dict()) == public_fields + assert "task_id" not in task.to_dict() + assert "admission_task_id" not in task.to_dict() + + def test_chain_media_legacy_scraping_symbols_resolve_to_scraping_chain(): """刮削拆分后,旧 app.chain.media 路径应能继续取用刮削公开符号。""" legacy_media = importlib.import_module("app.chain.media") diff --git a/tests/test_transfer_admission_migration.py b/tests/test_transfer_admission_migration.py new file mode 100644 index 000000000..060ba02ec --- /dev/null +++ b/tests/test_transfer_admission_migration.py @@ -0,0 +1,241 @@ +"""整理任务持久准入字段的 Alembic 迁移测试。""" + +import importlib +import os +import uuid + +import pytest +import sqlalchemy as sa +from alembic.migration import MigrationContext +from alembic.operations import Operations + +from app.db.models.transferpending import TransferPending + +try: + import psycopg2 as postgres_driver + from psycopg2 import sql + + POSTGRESQL_DIALECT = "postgresql+psycopg2" +except ModuleNotFoundError: + import psycopg as postgres_driver + from psycopg import sql + + POSTGRESQL_DIALECT = "postgresql+psycopg" + +MIGRATION = "database.versions.b1e7d3f5a9c2_3_0_13" + + +def _bind_migration(monkeypatch, connection): + """把迁移绑定到隔离数据库连接。""" + migration = importlib.import_module(MIGRATION) + monkeypatch.setattr( + migration, + "op", + Operations(MigrationContext.configure(connection)), + ) + return migration + + +def _create_legacy_table(connection) -> None: + """创建 3.0.12 时代的待整理登记表。""" + metadata = sa.MetaData() + table = sa.Table( + "transferpending", + metadata, + sa.Column("id", sa.Integer(), primary_key=True), + sa.Column("storage", sa.String(), nullable=False), + sa.Column("src_path", sa.String(), nullable=False), + sa.Column("created_at", sa.String(), nullable=True), + ) + sa.Index( + "ux_transferpending_storage_path", + table.c.storage, + table.c.src_path, + unique=True, + ) + metadata.create_all(connection) + connection.execute(table.insert(), [ + { + "id": 1, + "storage": "local", + "src_path": "/mnt/dated.mkv", + "created_at": "2026-08-26 10:00:00", + }, + { + "id": 2, + "storage": "alist", + "src_path": "/mnt/undated.mkv", + "created_at": None, + }, + ]) + + +def _rows(connection) -> list[dict[str, object]]: + """读取迁移后的准入字段快照。""" + pending = sa.table( + "transferpending", + sa.column("id", sa.Integer()), + sa.column("task_id", sa.String()), + sa.column("state", sa.String()), + sa.column("created_at", sa.String()), + sa.column("updated_at", sa.String()), + sa.column("last_error", sa.Text()), + ) + return [ + dict(row) + for row in connection.execute( + sa.select(pending).order_by(pending.c.id) + ).mappings().all() + ] + + +def test_transfer_admission_upgrade_downgrade_reupgrade( + monkeypatch, +) -> None: + """旧行应保守回填,且 SQLite 支持重复升级、降级和再次升级。""" + engine = sa.create_engine("sqlite://") + with engine.begin() as connection: + _create_legacy_table(connection) + migration = _bind_migration(monkeypatch, connection) + + migration.upgrade() + migration.upgrade() + + inspector = sa.inspect(connection) + assert { + column["name"] + for column in inspector.get_columns("transferpending") + } == {column.name for column in TransferPending.__table__.columns} + constraints = { + constraint["name"] + for constraint in inspector.get_unique_constraints("transferpending") + } + assert "uq_transferpending_task_id" in constraints + assert { + index["name"] + for index in inspector.get_indexes("transferpending") + } == { + "ix_transferpending_state_created", + "ux_transferpending_storage_path", + } + + upgraded = _rows(connection) + first_task_ids = [row["task_id"] for row in upgraded] + assert all(first_task_ids) + assert len(set(first_task_ids)) == 2 + assert {row["state"] for row in upgraded} == {"accepted"} + assert upgraded[0]["updated_at"] == upgraded[0]["created_at"] + assert upgraded[1]["updated_at"] + assert {row["last_error"] for row in upgraded} == {None} + + migration.downgrade() + downgraded_inspector = sa.inspect(connection) + assert { + column["name"] + for column in downgraded_inspector.get_columns("transferpending") + } == {"id", "storage", "src_path", "created_at"} + assert { + index["name"] + for index in downgraded_inspector.get_indexes("transferpending") + } == {"ux_transferpending_storage_path"} + legacy_rows = connection.execute( + sa.text( + "SELECT id, storage, src_path, created_at " + "FROM transferpending ORDER BY id" + ) + ).mappings().all() + assert [row["src_path"] for row in legacy_rows] == [ + "/mnt/dated.mkv", + "/mnt/undated.mkv", + ] + + migration.upgrade() + reupgraded = _rows(connection) + assert [row["task_id"] for row in reupgraded] == first_task_ids + assert {row["state"] for row in reupgraded} == {"accepted"} + assert { + index["name"] + for index in sa.inspect(connection).get_indexes("transferpending") + } == { + "ix_transferpending_state_created", + "ux_transferpending_storage_path", + } + + +def test_transfer_admission_migration_runs_on_postgresql(monkeypatch) -> None: + """隔离 PostgreSQL 应真实执行准入字段、约束、索引和可逆回滚。""" + prefix = "MOVIEPILOT_TEST_POSTGRESQL_" + host = os.getenv(f"{prefix}HOST") + database = os.getenv(f"{prefix}DATABASE") + username = os.getenv(f"{prefix}USERNAME") + if not host or not database or not username: + pytest.skip("未配置隔离 PostgreSQL migration 测试库") + + port = os.getenv(f"{prefix}PORT", "5432") + password = os.getenv(f"{prefix}PASSWORD", "") + schema = f"transfer_admission_{uuid.uuid4().hex}" + with postgres_driver.connect( + host=host, + port=port, + dbname=database, + user=username, + password=password, + ) as connection: + connection.autocommit = True + with connection.cursor() as cursor: + cursor.execute(sql.SQL("CREATE SCHEMA {}").format(sql.Identifier(schema))) + + engine = None + try: + engine = sa.create_engine( + sa.URL.create( + POSTGRESQL_DIALECT, + username=username, + password=password, + host=host, + port=int(port), + database=database, + ), + connect_args={"options": f"-csearch_path={schema}"}, + ) + with engine.begin() as connection: + _create_legacy_table(connection) + migration = _bind_migration(monkeypatch, connection) + migration.upgrade() + migration.upgrade() + + inspector = sa.inspect(connection) + assert { + constraint["name"] + for constraint in inspector.get_unique_constraints( + "transferpending" + ) + } >= {"uq_transferpending_task_id"} + assert { + index["name"] + for index in inspector.get_indexes("transferpending") + } >= {"ix_transferpending_state_created"} + assert all(row["task_id"] for row in _rows(connection)) + + migration.downgrade() + assert { + column["name"] + for column in sa.inspect(connection).get_columns("transferpending") + } == {"id", "storage", "src_path", "created_at"} + finally: + if engine is not None: + engine.dispose() + with postgres_driver.connect( + host=host, + port=port, + dbname=database, + user=username, + password=password, + ) as connection: + connection.autocommit = True + with connection.cursor() as cursor: + cursor.execute( + sql.SQL("DROP SCHEMA IF EXISTS {} CASCADE").format( + sql.Identifier(schema) + ) + ) diff --git a/tests/test_transfer_pending_replay.py b/tests/test_transfer_pending_replay.py index a0cde2173..0db82e83d 100644 --- a/tests/test_transfer_pending_replay.py +++ b/tests/test_transfer_pending_replay.py @@ -7,26 +7,38 @@ 这些测试固定三项不变量:入队即落盘登记、终态即注销、重启能回放。 """ -from pathlib import Path import threading +from pathlib import Path from unittest.mock import MagicMock +from app.application.transfer import TransferAdmission, TransferTask from app.chain.transfer import TransferChain -from app.application.transfer import TransferTask from app.schemas.file import FileItem -def _build_chain(pendingoper) -> TransferChain: +def _build_chain(admissions) -> TransferChain: """ 构造绕过单例初始化的 TransferChain 骨架。 - :param pendingoper: 待整理登记管理替身 + :param admissions: durable admission 仓储替身 :return: TransferChain 骨架 """ chain = object.__new__(TransferChain) - chain._pendingoper = pendingoper + chain._transfer_admissions = admissions return chain +def _admission(path: str, task_id: str = "task-1") -> TransferAdmission: + """构造一条可脱离数据库会话使用的准入快照。""" + return TransferAdmission( + task_id=task_id, + storage="local", + src_path=path, + state="accepted", + created_at="2026-08-27 10:00:00", + updated_at="2026-08-27 10:00:00", + ) + + def _task(path: str, storage: str = "local") -> TransferTask: """ 构造测试用整理任务。 @@ -45,44 +57,38 @@ def _task(path: str, storage: str = "local") -> TransferTask: )) -def test_register_pending_records_storage_and_path(): +def test_admit_transfer_records_storage_and_path(): """ 入队时必须落盘登记「存储 + 源路径」这一最小事实。 """ - pendingoper = MagicMock() - chain = _build_chain(pendingoper) + admissions = MagicMock() + admissions.admit.return_value = _admission( + "/mnt/cd2/downloads/Movie.2024.mkv" + ) + chain = _build_chain(admissions) - chain._TransferChain__register_pending(_task("/mnt/cd2/downloads/Movie.2024.mkv")) - - pendingoper.register.assert_called_once_with( - storage="local", src_path="/mnt/cd2/downloads/Movie.2024.mkv" + result = chain._TransferChain__admit_transfer( + _task("/mnt/cd2/downloads/Movie.2024.mkv") ) - -def test_register_pending_failure_does_not_break_enqueue(): - """ - 落盘登记只是重启后的补救手段,登记失败绝不能阻断正常整理。 - """ - pendingoper = MagicMock() - pendingoper.register.side_effect = RuntimeError("db locked") - chain = _build_chain(pendingoper) - - # 不抛异常即为通过 - chain._TransferChain__register_pending(_task("/mnt/cd2/downloads/Movie.2024.mkv")) + admissions.admit.assert_called_once_with( + storage="local", src_path="/mnt/cd2/downloads/Movie.2024.mkv" + ) + assert result.task_id == "task-1" def test_discard_pending_on_terminal_state(): """ 整理到达终态后必须注销登记,否则每次重启都会重复回放。 """ - pendingoper = MagicMock() - chain = _build_chain(pendingoper) + admissions = MagicMock() + chain = _build_chain(admissions) + task = _task("/mnt/cd2/downloads/Movie.2024.mkv") + task.bind_admission_task_id("task-1") - chain._TransferChain__discard_pending(_task("/mnt/cd2/downloads/Movie.2024.mkv")) + chain._TransferChain__discard_pending(task) - pendingoper.discard.assert_called_once_with( - storage="local", src_path="/mnt/cd2/downloads/Movie.2024.mkv" - ) + admissions.discard_task.assert_called_once_with(task_id="task-1") def test_replay_resends_pending_files_to_transfer(tmp_path, monkeypatch): @@ -92,9 +98,9 @@ def test_replay_resends_pending_files_to_transfer(tmp_path, monkeypatch): media = tmp_path / "Movie.2024.mkv" media.write_bytes(b"x" * 10) - pendingoper = MagicMock() - pendingoper.list_all.return_value = [("local", str(media))] - chain = _build_chain(pendingoper) + admissions = MagicMock() + admissions.list_accepted.return_value = [_admission(str(media))] + chain = _build_chain(admissions) transferred = [] monkeypatch.setattr(chain, "do_transfer", lambda **kw: transferred.append(kw["fileitem"])) @@ -114,16 +120,16 @@ def test_replay_discards_vanished_files(tmp_path): """ 源文件已消失的登记要注销,否则每次启动都会重复回放一个不存在的文件。 """ - pendingoper = MagicMock() + admissions = MagicMock() missing = tmp_path / "gone.mkv" - pendingoper.list_all.return_value = [("local", str(missing))] - chain = _build_chain(pendingoper) + admissions.list_accepted.return_value = [_admission(str(missing))] + chain = _build_chain(admissions) chain.do_transfer = MagicMock() chain._TransferChain__replay_pending() chain.do_transfer.assert_not_called() - pendingoper.discard.assert_called_once_with(storage="local", src_path=str(missing)) + admissions.discard_task.assert_called_once_with(task_id="task-1") def test_replay_keeps_registration_when_mount_unreadable(tmp_path, monkeypatch): @@ -135,9 +141,9 @@ def test_replay_keeps_registration_when_mount_unreadable(tmp_path, monkeypatch): media = tmp_path / "Movie.2024.mkv" media.write_bytes(b"x") - pendingoper = MagicMock() - pendingoper.list_all.return_value = [("local", str(media))] - chain = _build_chain(pendingoper) + admissions = MagicMock() + admissions.list_accepted.return_value = [_admission(str(media))] + chain = _build_chain(admissions) chain.do_transfer = MagicMock() def unreadable(self, *_args, **_kwargs): @@ -151,7 +157,7 @@ def test_replay_keeps_registration_when_mount_unreadable(tmp_path, monkeypatch): chain._TransferChain__replay_pending() chain.do_transfer.assert_not_called() - pendingoper.discard.assert_not_called() + admissions.discard_task.assert_not_called() def test_replay_restores_bluray_directory_type(tmp_path, monkeypatch): @@ -162,9 +168,9 @@ def test_replay_restores_bluray_directory_type(tmp_path, monkeypatch): bluray.mkdir() src_path = f"{bluray.as_posix()}/" - pendingoper = MagicMock() - pendingoper.list_all.return_value = [("local", src_path)] - chain = _build_chain(pendingoper) + admissions = MagicMock() + admissions.list_accepted.return_value = [_admission(src_path)] + chain = _build_chain(admissions) transferred = [] monkeypatch.setattr(chain, "do_transfer", lambda **kw: transferred.append(kw["fileitem"])) @@ -180,9 +186,9 @@ def test_replay_is_noop_without_registrations(): """ 没有登记时回放不应触碰整理链。 """ - pendingoper = MagicMock() - pendingoper.list_all.return_value = [] - chain = _build_chain(pendingoper) + admissions = MagicMock() + admissions.list_accepted.return_value = [] + chain = _build_chain(admissions) chain.do_transfer = MagicMock() chain._TransferChain__replay_pending() @@ -194,9 +200,9 @@ def test_replay_survives_db_failure(): """ 读取登记失败不能让启动流程报错。 """ - pendingoper = MagicMock() - pendingoper.list_all.side_effect = RuntimeError("db gone") - chain = _build_chain(pendingoper) + admissions = MagicMock() + admissions.list_accepted.side_effect = RuntimeError("db gone") + chain = _build_chain(admissions) chain.do_transfer = MagicMock() chain._TransferChain__replay_pending() @@ -213,9 +219,12 @@ def test_replay_continues_after_single_file_failure(tmp_path, monkeypatch): for item in (first, second): item.write_bytes(b"x") - pendingoper = MagicMock() - pendingoper.list_all.return_value = [("local", str(first)), ("local", str(second))] - chain = _build_chain(pendingoper) + admissions = MagicMock() + admissions.list_accepted.return_value = [ + _admission(str(first), "task-1"), + _admission(str(second), "task-2"), + ] + chain = _build_chain(admissions) handled = [] @@ -241,12 +250,12 @@ def test_replay_stop_keeps_unprocessed_registrations(tmp_path, monkeypatch): first = tmp_path / "A.mkv" first.write_bytes(b"x") missing_second = tmp_path / "gone.mkv" - pendingoper = MagicMock() - pendingoper.list_all.return_value = [ - ("local", str(first)), - ("local", str(missing_second)), + admissions = MagicMock() + admissions.list_accepted.return_value = [ + _admission(str(first), "task-1"), + _admission(str(missing_second), "task-2"), ] - chain = _build_chain(pendingoper) + chain = _build_chain(admissions) stop_event = threading.Event() transferred = [] @@ -260,4 +269,4 @@ def test_replay_stop_keeps_unprocessed_registrations(tmp_path, monkeypatch): chain._TransferChain__replay_pending(stop_event) assert transferred == [first.as_posix()] - pendingoper.discard.assert_not_called() + admissions.discard_task.assert_not_called() diff --git a/tests/test_transfer_queue_service.py b/tests/test_transfer_queue_service.py index d340d5b7e..d80a0c4f0 100644 --- a/tests/test_transfer_queue_service.py +++ b/tests/test_transfer_queue_service.py @@ -1,18 +1,32 @@ -from unittest.mock import Mock +from types import SimpleNamespace +from unittest.mock import Mock, patch -from app.application.transfer import TransferQueueService +import pytest +from sqlalchemy import create_engine +from sqlalchemy.orm import sessionmaker + +from app.application.transfer import TransferAdmission, TransferQueueService +from app.db.adapters.transfer import TransactionalTransferAdmissionRepository +from app.db.models.transferpending import TransferPending from app.schemas.file import FileItem - -from tests.test_transfer_job_manager import make_task +from tests.test_transfer_job_manager import make_task, make_transfer_chain def _service(**overrides): """构造可观测整理队列服务及其默认依赖。""" dependencies = { "register_task": Mock(return_value=True), + "admit_task": Mock(return_value=TransferAdmission( + task_id="task-1", + storage="local", + src_path="/tmp/demo.mkv", + state="accepted", + created_at="2026-08-27 10:00:00", + updated_at="2026-08-27 10:00:00", + )), "enqueue": Mock(), "before_enqueue": Mock(), - "after_enqueue": Mock(), + "enqueue_failed": Mock(), "remove_task": Mock(), "list_tasks": Mock(return_value=["job"]), "expire_tasks": Mock(), @@ -22,17 +36,26 @@ def _service(**overrides): def test_transfer_queue_service_put_preserves_registration_order(): - """入队必须先登记视图,再登记批次、写队列并落盘。""" + """入队必须先登记视图和 durable admission,再登记批次并写队列。""" calls = [] service, _ = _service( register_task=lambda _task: calls.append("register") or True, + admit_task=lambda _task: calls.append("admit") or TransferAdmission( + task_id="task-1", + storage="local", + src_path="/tmp/demo.mkv", + state="accepted", + created_at="2026-08-27 10:00:00", + updated_at="2026-08-27 10:00:00", + ), before_enqueue=lambda _task: calls.append("batch"), enqueue=lambda _item: calls.append("queue"), - after_enqueue=lambda _task: calls.append("pending"), ) - assert service.put(make_task(1), Mock()) is True - assert calls == ["register", "batch", "queue", "pending"] + task = make_task(1) + assert service.put(task, Mock()) is True + assert calls == ["register", "admit", "batch", "queue"] + assert task.admission_task_id == "task-1" def test_transfer_queue_service_rejects_duplicate_without_side_effects(): @@ -42,7 +65,81 @@ def test_transfer_queue_service_rejects_duplicate_without_side_effects(): assert service.put(make_task(1), Mock()) is False dependencies["before_enqueue"].assert_not_called() dependencies["enqueue"].assert_not_called() - dependencies["after_enqueue"].assert_not_called() + dependencies["admit_task"].assert_not_called() + + +def test_transfer_queue_service_blocks_enqueue_when_admission_fails(): + """持久化失败必须撤销作业视图,不能继续加入内存队列。""" + service, dependencies = _service( + admit_task=Mock(side_effect=RuntimeError("db locked")), + ) + task = make_task(1) + + with pytest.raises(RuntimeError, match="db locked"): + service.put(task, Mock()) + + dependencies["remove_task"].assert_called_once_with(task.fileitem) + dependencies["before_enqueue"].assert_not_called() + dependencies["enqueue"].assert_not_called() + + +def test_transfer_queue_service_keeps_admission_when_enqueue_fails(): + """内存入队失败必须记录原因并清理视图,durable admission 由仓储保留。""" + error = RuntimeError("queue closed") + service, dependencies = _service( + enqueue=Mock(side_effect=error), + ) + task = make_task(1) + + with pytest.raises(RuntimeError, match="queue closed"): + service.put(task, Mock()) + + dependencies["enqueue_failed"].assert_called_once_with(task, error) + dependencies["remove_task"].assert_called_once_with(task.fileitem) + + +def test_transfer_queue_service_cleans_up_when_batch_registration_fails(): + """准入后的批次登记异常也必须留痕并撤销作业视图。""" + error = RuntimeError("batch registration failed") + service, dependencies = _service( + before_enqueue=Mock(side_effect=error), + ) + task = make_task(1) + + with pytest.raises(RuntimeError, match="batch registration failed"): + service.put(task, Mock()) + + dependencies["enqueue_failed"].assert_called_once_with(task, error) + dependencies["remove_task"].assert_called_once_with(task.fileitem) + dependencies["enqueue"].assert_not_called() + + +def test_transfer_queue_service_commits_admission_before_failed_enqueue(tmp_path): + """真实仓储已提交后即使内存入队失败,任务也必须带原因留待恢复。""" + engine = create_engine(f"sqlite:///{tmp_path / 'durable-admission.db'}") + TransferPending.__table__.create(engine) + repository = TransactionalTransferAdmissionRepository(sessionmaker(bind=engine)) + task = make_task(1) + service, _ = _service( + admit_task=lambda item: repository.admit( + storage=item.fileitem.storage, + src_path=item.fileitem.path, + ), + enqueue=Mock(side_effect=RuntimeError("queue closed")), + enqueue_failed=lambda item, error: repository.record_enqueue_failure( + task_id=item.admission_task_id, + error=str(error), + ), + ) + + with pytest.raises(RuntimeError, match="queue closed"): + service.put(task, Mock()) + + admissions = repository.list_accepted() + assert len(admissions) == 1 + assert admissions[0].task_id == task.admission_task_id + assert admissions[0].last_error == "queue closed" + engine.dispose() def test_transfer_queue_service_lists_and_removes_through_ports(): @@ -56,3 +153,38 @@ def test_transfer_queue_service_lists_and_removes_through_ports(): dependencies["expire_tasks"].assert_called_once_with() dependencies["list_tasks"].assert_called_once_with() dependencies["remove_task"].assert_called_once_with(fileitem) + + +def test_do_transfer_reports_durable_admission_failure(): + """背景整理准入失败必须返回批次失败,不能伪装成重复任务成功。""" + chain = make_transfer_chain() + fileitem = make_task(1).fileitem + chain._TransferChain__get_trans_fileitems = lambda _item, **_kwargs: [ + (fileitem, False) + ] + chain.put_to_queue = Mock(side_effect=RuntimeError("db locked")) + no_history = SimpleNamespace( + get_by_src=lambda _src, storage=None: None, + get_success_by_src=lambda _src, storage=None: None, + ) + no_download = SimpleNamespace( + get_by_hash=lambda _hash: None, + get_file_by_fullpath=lambda _path: None, + get_files_by_savepath=lambda _path: [], + get_by_path=lambda _path: None, + ) + + with patch( + "app.chain.transfer.get_chain_transfer_history_port", + return_value=no_history, + ), patch( + "app.chain.transfer.get_chain_download_history_port", + return_value=no_download, + ), patch( + "app.chain.transfer.get_configured_system_config", + return_value=SimpleNamespace(get=lambda _key: None), + ): + state, message = chain.do_transfer(fileitem=fileitem, background=True) + + assert state is False + assert "加入整理队列失败:db locked" in message diff --git a/tests/test_transfer_worker_lifecycle.py b/tests/test_transfer_worker_lifecycle.py index 54eb46ed4..ccbe5f654 100644 --- a/tests/test_transfer_worker_lifecycle.py +++ b/tests/test_transfer_worker_lifecycle.py @@ -10,10 +10,10 @@ from unittest.mock import AsyncMock, MagicMock, patch import pytest +from app.application.transfer import TransferAdmission, TransferQueue, TransferTask from app.chain.transfer import TransferChain from app.foundation.singleton import Singleton from app.runtime.config import global_vars -from app.application.transfer import TransferQueue, TransferTask from app.schemas.file import FileItem from app.startup.initializers import transfer as transfer_initializer @@ -411,6 +411,62 @@ def test_worker_settles_progress_when_only_stop_sentinel_remains(monkeypatch) -> assert list(chain._queue.queue) == [chain._QUEUE_STOP_SENTINEL] +def test_durable_task_identity_flows_from_queue_to_terminal_discard(monkeypatch) -> None: + """准入生成的稳定身份必须随队列任务到 worker 终态并准确注销。""" + chain = _build_chain() + chain.runtime_config.transfer_task_timeout = 0 + task = TransferTask(fileitem=FileItem( + storage="local", + path="/downloads/durable.mkv", + type="file", + name="durable.mkv", + basename="durable", + extension="mkv", + )) + discarded = threading.Event() + admissions = MagicMock() + admissions.admit.return_value = TransferAdmission( + task_id="durable-task-id", + storage="local", + src_path=task.fileitem.path, + state="accepted", + created_at="2026-08-27 10:00:00", + updated_at="2026-08-27 10:00:00", + ) + admissions.discard_task.side_effect = ( + lambda **_kwargs: discarded.set() or 1 + ) + chain._transfer_admissions = admissions + chain.jobview = MagicMock() + chain.jobview.add_task.return_value = True + chain.jobview.pending_total.return_value = 1 + chain._register_scrape_batch_task = MagicMock() + chain._finish_scrape_batch_task = MagicMock() + chain._progress = MagicMock() + chain._active_tasks = 0 + chain._processed_num = 0 + chain._fail_num = 0 + chain._total_num = 0 + chain._TransferChain__handle_transfer = MagicMock(return_value=(True, "")) + monkeypatch.setattr(global_vars, "STOP_EVENT", threading.Event()) + + assert chain.put_to_queue(task) is True + stop_event = threading.Event() + worker = threading.Thread( + target=chain._TransferChain__start_transfer, + args=(stop_event,), + daemon=True, + ) + worker.start() + assert discarded.wait(timeout=1) + stop_event.set() + worker.join(timeout=1) + + assert worker.is_alive() is False + assert task.admission_task_id == "durable-task-id" + admissions.discard_task.assert_called_once_with(task_id="durable-task-id") + + def test_claimed_task_prevents_progress_settlement_before_active_registration() -> None: """其他 worker 已取走真实任务但尚未登记 active 时,当前批次不得提前结算。""" chain = _build_chain()