From f94eb0076d7b7e05a54e4c1e1aec0e0848dadac9 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Wed, 2 Sep 2026 14:04:56 +0800 Subject: [PATCH] refactor(transfer): keep lease fencing in execution owner --- app/application/transfer/workflow.py | 5 +- app/chain/transfer/contract.py | 8 ++ app/chain/transfer/execution.py | 96 +++++++++++++++++++- app/chain/transfer/queue.py | 125 +++----------------------- tests/test_transfer_pending_replay.py | 22 +++-- tests/test_transfer_queue_service.py | 69 +------------- 6 files changed, 132 insertions(+), 193 deletions(-) diff --git a/app/application/transfer/workflow.py b/app/application/transfer/workflow.py index eede0013e..9f0c14010 100644 --- a/app/application/transfer/workflow.py +++ b/app/application/transfer/workflow.py @@ -971,7 +971,6 @@ class TransferQueueService: *, register_task: Callable[[TransferTask], bool], admit_task: Callable[[TransferTask], TransferAdmission], - claim_task: Callable[[TransferTask, TransferAdmission], None], enqueue: Callable[[TransferQueue], None], before_enqueue: Callable[[TransferTask], None], enqueue_failed: Callable[[TransferTask, Exception], None], @@ -982,7 +981,6 @@ class TransferQueueService: """保存队列用例依赖,避免 Application 服务绑定具体线程队列实现。""" self._register_task = register_task self._admit_task = admit_task - self._claim_task = claim_task self._enqueue = enqueue self._before_enqueue = before_enqueue self._enqueue_failed = enqueue_failed @@ -991,13 +989,12 @@ class TransferQueueService: self._expire_tasks = expire_tasks def put(self, task: TransferTask, callback: TransferCallback) -> bool: - """先持久化并 claim 准入事实再入队;前置失败撤销内存作业视图。""" + """先持久化准入事实再入队;任何前置失败都撤销内存作业视图。""" if not task or not self._register_task(task): return False try: admission = self._admit_task(task) task.bind_admission_task_id(admission.task_id) - self._claim_task(task, admission) except Exception: self._remove_task(task.fileitem) raise diff --git a/app/chain/transfer/contract.py b/app/chain/transfer/contract.py index 7ee643b45..92571667e 100644 --- a/app/chain/transfer/contract.py +++ b/app/chain/transfer/contract.py @@ -20,6 +20,9 @@ if TYPE_CHECKING: _subtitle_exts: Any _success_target_files: Any _transfer_admissions: Any + _owned_leases: Any + _worker_owner_id: str + _worker_state_lock: Any download_history_repository: Any eventmanager: Any failure_notification_aggregator: Any @@ -30,11 +33,14 @@ if TYPE_CHECKING: transfer_history_repository: Any _TransferChain__assert_owned_lease: Callable[..., Any] + _TransferChain__admit_transfer: Callable[..., Any] _TransferChain__bind_claimed_admission: Callable[..., Any] _TransferChain__build_planning_input: Callable[..., Any] _TransferChain__checkpoint_planning_rejection: Callable[..., Any] _TransferChain__claim_task_for_execution: Callable[..., Any] _TransferChain__default_callback: Callable[..., Any] + _TransferChain__ensure_lease_heartbeat_owner: Callable[..., Any] + _TransferChain__ensure_lease_runtime_state: Callable[..., Any] _TransferChain__fail_transfer_task: Callable[..., Any] _TransferChain__finish_job_execution: Callable[..., Any] _TransferChain__forget_owned_lease: Callable[..., Any] @@ -45,6 +51,8 @@ if TYPE_CHECKING: _TransferChain__mark_torrent_completed_if_done: Callable[..., Any] _TransferChain__put_to_jobview: Callable[..., Any] _TransferChain__record_uncheckpointed_failure: Callable[..., Any] + _TransferChain__register_claimed_admission: Callable[..., Any] + _TransferChain__release_admission_claim: Callable[..., Any] _TransferChain__release_task_claim: Callable[..., Any] _TransferChain__restore_mediainfo_snapshot: Callable[..., Any] _TransferChain__restore_meta_snapshot: Callable[..., Any] diff --git a/app/chain/transfer/execution.py b/app/chain/transfer/execution.py index b2f963d69..502bfe0f1 100644 --- a/app/chain/transfer/execution.py +++ b/app/chain/transfer/execution.py @@ -1,5 +1,6 @@ """持久整理步骤执行、探测与检查点推进。""" +import time from datetime import datetime, timedelta, timezone from pathlib import Path from typing import Any, Callable, Mapping, Optional, Tuple, Union, cast @@ -20,6 +21,8 @@ from app.application.transfer.execution import ( TransferStepState, ) from app.application.transfer.workflow import ( + TransferAdmission, + TransferLeaseLostError, TransferTask, ) from app.chain.media import MediaChain @@ -284,7 +287,98 @@ class _DurableTransferStepRunner: class TransferExecutionOwner(_TransferOwnerBase): - """唯一持有整理计划的外部步骤执行与结果检查点。""" + """持有整理准入租约、外部步骤执行与结果检查点。""" + + def _TransferChain__register_claimed_admission( + self, + admission: TransferAdmission, + ) -> None: + """把仓储 claim 加入续期集合,覆盖恢复构造阶段的同步 I/O。""" + if ( + admission.lease_owner != self._worker_owner_id + or not admission.lease_token + ): + raise TransferLeaseLostError( + f"整理任务 claim 投影无效:{admission.task_id}" + ) + with self._worker_state_lock: + self._owned_leases[admission.task_id] = ( + admission.lease_token, + time.monotonic() + self._WORKER_LEASE_SECONDS, + ) + self._TransferChain__ensure_lease_heartbeat_owner() + + def _TransferChain__bind_claimed_admission( + self, + task: TransferTask, + admission: TransferAdmission, + ) -> None: + """校验仓储 claim 投影并把 token 私有绑定到执行任务。""" + if admission.task_id != task.admission_task_id: + raise TransferLeaseLostError( + f"整理任务 claim 投影无效:{admission.task_id}" + ) + self._TransferChain__register_claimed_admission(admission) + assert admission.lease_owner is not None + assert admission.lease_token is not None + task.bind_execution_lease( + owner_id=admission.lease_owner, + lease_token=admission.lease_token, + ) + + def _TransferChain__claim_admitted_task( + self, + task: TransferTask, + task_id: str, + ) -> None: + """取得已准入任务的唯一租约,并绑定到内存任务。""" + claimed = cast( + Optional[TransferAdmission], + self._transfer_admissions.claim_task( + task_id=task_id, + owner_id=self._worker_owner_id, + lease_seconds=self._WORKER_LEASE_SECONDS, + ), + ) + if claimed is None: + raise TransferLeaseLostError(f"整理任务已由其他 worker claim:{task_id}") + try: + self._TransferChain__bind_claimed_admission(task, claimed) + except Exception as err: + self._TransferChain__release_admission_claim(claimed, error=str(err)) + raise + + def _TransferChain__admit_transfer(self, task: TransferTask) -> TransferAdmission: + """持久化源文件并在内存入队前取得执行租约。""" + fileitem = task.fileitem if task else None + if not fileitem or not fileitem.storage or not fileitem.path: + raise ValueError("整理任务缺少源文件身份") + planning_input = task.planning_input or self._TransferChain__build_planning_input(task) + task.bind_planning_input(planning_input) + admission = cast( + TransferAdmission, + self._transfer_admissions.admit( + storage=fileitem.storage, + src_path=fileitem.path, + planning_input=planning_input, + ), + ) + task.bind_admission_task_id(admission.task_id) + self._TransferChain__claim_admitted_task(task, admission.task_id) + return admission + + def _TransferChain__claim_task_for_execution(self, task: TransferTask) -> None: + """校验队列任务既有租约,并为兼容旧队列项补取唯一租约。""" + if task.preview: + return + self._TransferChain__ensure_lease_runtime_state() + if task.lease_token: + self._TransferChain__assert_owned_lease(task) + return + if not task.admission_task_id: + self._TransferChain__admit_transfer(task) + return + self._TransferChain__claim_admitted_task(task, task.admission_task_id) def _TransferChain__handle_transfer( diff --git a/app/chain/transfer/queue.py b/app/chain/transfer/queue.py index d65334c65..3a9d0e8ca 100644 --- a/app/chain/transfer/queue.py +++ b/app/chain/transfer/queue.py @@ -7,7 +7,7 @@ import time import traceback import uuid from pathlib import Path -from typing import Any, Callable, Dict, List, Optional, Tuple, cast +from typing import Any, Callable, Dict, List, Optional, Tuple from pydantic_core import to_jsonable_python @@ -373,7 +373,6 @@ class TransferQueueOwner(_TransferOwnerBase): return TransferQueueService( register_task=self._TransferChain__put_to_jobview, admit_task=self._TransferChain__admit_transfer, - claim_task=self._TransferChain__claim_admitted_task, enqueue=self._queue.put, before_enqueue=self._register_scrape_batch_task, enqueue_failed=self._TransferChain__record_enqueue_failure, @@ -532,45 +531,6 @@ class TransferQueueOwner(_TransferOwnerBase): time.monotonic() + self._WORKER_LEASE_SECONDS, ) - def _TransferChain__bind_claimed_admission( - self, - task: TransferTask, - admission: TransferAdmission, - ) -> None: - """校验仓储 claim 投影并把 token 私有绑定到执行任务。""" - if ( - admission.task_id != task.admission_task_id - ): - raise TransferLeaseLostError( - f"整理任务 claim 投影无效:{admission.task_id}" - ) - self._TransferChain__register_claimed_admission(admission) - assert admission.lease_owner is not None - assert admission.lease_token is not None - task.bind_execution_lease( - owner_id=admission.lease_owner, - lease_token=admission.lease_token, - ) - - def _TransferChain__register_claimed_admission( - self, - admission: TransferAdmission, - ) -> None: - """把仓储 claim 加入续期集合,覆盖恢复构造阶段可能发生的同步 I/O。""" - if ( - admission.lease_owner != self._worker_owner_id - or not admission.lease_token - ): - raise TransferLeaseLostError( - f"整理任务 claim 投影无效:{admission.task_id}" - ) - with self._worker_state_lock: - self._owned_leases[admission.task_id] = ( - admission.lease_token, - time.monotonic() + self._WORKER_LEASE_SECONDS, - ) - self._TransferChain__ensure_lease_heartbeat_owner() - def _TransferChain__forget_owned_lease(self, task_id: str, lease_token: str) -> None: """仅在 token 仍匹配时移除本进程租约镜像,避免删掉新接管记录。""" with self._worker_state_lock: @@ -615,56 +575,6 @@ class TransferQueueOwner(_TransferOwnerBase): self._owned_leases.pop(task_id, None) raise TransferLeaseLostError(f"整理任务租约已经失效:{task_id}") - def _TransferChain__claim_task_for_execution(self, task: TransferTask) -> None: - """校验队列任务既有租约,并为兼容旧队列项补取唯一租约。""" - if task.preview: - return - self._TransferChain__ensure_lease_runtime_state() - if task.lease_token: - self._TransferChain__assert_owned_lease(task) - return - if not task.admission_task_id: - admitted = self._TransferChain__admit_transfer(task) - task.bind_admission_task_id(admitted.task_id) - task_id = task.admission_task_id - if task_id is None: - raise TransferLeaseLostError("整理任务准入后仍缺少 durable 身份") - claimed = self._transfer_admissions.claim_task( - task_id=task_id, - owner_id=self._worker_owner_id, - lease_seconds=self._WORKER_LEASE_SECONDS, - ) - if claimed is None: - raise TransferLeaseLostError( - f"整理任务已由其他 worker claim:{task_id}" - ) - try: - self._TransferChain__bind_claimed_admission(task, claimed) - except Exception as err: - self._TransferChain__release_admission_claim(claimed, error=str(err)) - raise - - def _TransferChain__claim_admitted_task( - self, - task: TransferTask, - admission: TransferAdmission, - ) -> None: - """普通任务进入内存队列前取得租约,避免恢复调度抢占等待项。""" - claimed = self._transfer_admissions.claim_task( - task_id=admission.task_id, - owner_id=self._worker_owner_id, - lease_seconds=self._WORKER_LEASE_SECONDS, - ) - if claimed is None: - raise TransferLeaseLostError( - f"整理任务已由其他 worker claim:{admission.task_id}" - ) - try: - self._TransferChain__bind_claimed_admission(task, claimed) - except Exception as err: - self._TransferChain__release_admission_claim(claimed, error=str(err)) - raise - def _TransferChain__release_task_claim( self, task: TransferTask, @@ -945,14 +855,19 @@ class TransferQueueOwner(_TransferOwnerBase): "待整理文件回放收到关闭请求,已送入 %s 个文件,其余登记保持待处理", replayed, ) - elif replayed == 0: - logger.warning( - "待整理文件回放未产生可执行队列任务:本次 claim %s 个;" - "逐任务原因已写入前序日志和 transferpending.last_error", - len(pendings), - ) else: + self._TransferChain__log_replay_summary(replayed, len(pendings)) + + @staticmethod + def _TransferChain__log_replay_summary(replayed: int, claimed: int) -> None: + """按实际入队结果记录恢复汇总,避免零任务仍输出成功标记。""" + if replayed: logger.info(f"✓ 待整理文件回放完成,{replayed} 个文件已重新送入整理链") + return + logger.warning( + f"待整理文件回放未入队:claim {claimed} 个," + "原因见前序日志与 transferpending.last_error" + ) def _TransferChain__execution_replay_snapshot( self, @@ -1209,22 +1124,6 @@ class TransferQueueOwner(_TransferOwnerBase): }) return fileitem, False - def _TransferChain__admit_transfer(self, task: TransferTask) -> TransferAdmission: - """在内存入队前持久化源文件并返回稳定任务身份。""" - fileitem = task.fileitem if task else None - if not fileitem or not fileitem.storage or not fileitem.path: - raise ValueError("整理任务缺少源文件身份") - planning_input = task.planning_input or self._TransferChain__build_planning_input(task) - task.bind_planning_input(planning_input) - return cast( - TransferAdmission, - self._transfer_admissions.admit( - storage=fileitem.storage, - src_path=fileitem.path, - planning_input=planning_input, - ), - ) - @staticmethod def _TransferChain__json_snapshot(value: Any) -> Any: """把受控领域对象投影为可持久化 JSON 值。""" diff --git a/tests/test_transfer_pending_replay.py b/tests/test_transfer_pending_replay.py index c0888e580..027255e36 100644 --- a/tests/test_transfer_pending_replay.py +++ b/tests/test_transfer_pending_replay.py @@ -119,17 +119,19 @@ def _task(path: str, storage: str = "local") -> TransferTask: def test_admit_transfer_records_storage_and_path(): """ - 入队时必须落盘登记「存储 + 源路径」这一最小事实。 + 入队时必须落盘登记源身份并在返回前取得执行租约。 """ admissions = MagicMock() admissions.admit.return_value = _admission( "/mnt/cd2/downloads/Movie.2024.mkv" ) - chain = _build_chain(admissions) - - result = chain._TransferChain__admit_transfer( - _task("/mnt/cd2/downloads/Movie.2024.mkv") + admissions.claim_task.return_value = _admission( + "/mnt/cd2/downloads/Movie.2024.mkv" ) + chain = _build_chain(admissions) + task = _task("/mnt/cd2/downloads/Movie.2024.mkv") + + result = chain._TransferChain__admit_transfer(task) call = admissions.admit.call_args.kwargs assert call["storage"] == "local" @@ -137,6 +139,12 @@ def test_admit_transfer_records_storage_and_path(): assert isinstance(call["planning_input"], TransferPlanningInput) assert call["planning_input"].source_fileitem["path"] == call["src_path"] assert result.task_id == "task-1" + assert task.lease_token == "lease-task-1" + admissions.claim_task.assert_called_once_with( + task_id="task-1", + owner_id="test-owner", + lease_seconds=120, + ) def test_terminal_without_settlement_releases_claim_and_keeps_pending(): @@ -511,9 +519,7 @@ def test_replay_releases_claim_when_jobview_rejects_recovered_task( error="恢复任务未进入内存队列", ) warning.assert_called_once_with( - "待整理文件回放未产生可执行队列任务:本次 claim %s 个;" - "逐任务原因已写入前序日志和 transferpending.last_error", - 1, + "待整理文件回放未入队:claim 1 个,原因见前序日志与 transferpending.last_error", ) assert chain._owned_leases == {} diff --git a/tests/test_transfer_queue_service.py b/tests/test_transfer_queue_service.py index 662959493..247d746ee 100644 --- a/tests/test_transfer_queue_service.py +++ b/tests/test_transfer_queue_service.py @@ -1,4 +1,3 @@ -from pathlib import Path from types import SimpleNamespace from unittest.mock import Mock, patch @@ -10,7 +9,6 @@ from app.application.transfer.workflow import ( TransferAdmission, TransferPlanningInput, TransferQueueService, - TransferTask, ) from app.db.adapters.transfer.admission import TransactionalTransferAdmissionRepository from app.db.models.transferhistory import TransferHistory @@ -46,7 +44,6 @@ def _service(**overrides): updated_at="2026-08-27 10:00:00", planning_input=_planning_input(), )), - "claim_task": Mock(), "enqueue": Mock(), "before_enqueue": Mock(), "enqueue_failed": Mock(), @@ -59,7 +56,7 @@ 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, @@ -72,14 +69,13 @@ def test_transfer_queue_service_put_preserves_registration_order(): updated_at="2026-08-27 10:00:00", planning_input=_planning_input(), ), - claim_task=lambda _task, _admission: calls.append("claim"), before_enqueue=lambda _task: calls.append("batch"), enqueue=lambda _item: calls.append("queue"), ) task = make_task(1) assert service.put(task, Mock()) is True - assert calls == ["register", "admit", "claim", "batch", "queue"] + assert calls == ["register", "admit", "batch", "queue"] assert task.admission_task_id == "task-1" @@ -91,7 +87,6 @@ def test_transfer_queue_service_rejects_duplicate_without_side_effects(): dependencies["before_enqueue"].assert_not_called() dependencies["enqueue"].assert_not_called() dependencies["admit_task"].assert_not_called() - dependencies["claim_task"].assert_not_called() def test_transfer_queue_service_blocks_enqueue_when_admission_fails(): @@ -105,70 +100,10 @@ def test_transfer_queue_service_blocks_enqueue_when_admission_fails(): service.put(task, Mock()) dependencies["remove_task"].assert_called_once_with(task.fileitem) - dependencies["claim_task"].assert_not_called() dependencies["before_enqueue"].assert_not_called() dependencies["enqueue"].assert_not_called() -def test_transfer_queue_service_blocks_enqueue_when_claim_fails() -> None: - """准入后的租约竞争失败必须撤销作业视图,不能留下等待占位。""" - service, dependencies = _service( - claim_task=Mock(side_effect=RuntimeError("already claimed")), - ) - task = make_task(1) - - with pytest.raises(RuntimeError, match="already claimed"): - 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_claim_fences_recovery_owner(tmp_path: Path) -> None: - """普通任务等待内存队列期间必须持有租约,恢复 owner 不得抢占。""" - engine = create_engine(f"sqlite:///{tmp_path / 'queue-claim.db'}") - TransferHistory.__table__.create(engine) - TransferPending.__table__.create(engine) - factory = sessionmaker(bind=engine) - repository = TransactionalTransferAdmissionRepository(factory) - task = make_task(1) - task.bind_planning_input(_planning_input(task.fileitem.path)) - - def claim_task(item: TransferTask, admission: TransferAdmission) -> None: - """使用真实仓储 claim 并把租约绑定到内存任务。""" - claimed = repository.claim_task( - task_id=admission.task_id, - owner_id="queue-worker", - lease_seconds=120, - ) - assert claimed is not None - assert claimed.lease_owner is not None - assert claimed.lease_token is not None - item.bind_execution_lease( - owner_id=claimed.lease_owner, - lease_token=claimed.lease_token, - ) - - service, _ = _service( - admit_task=lambda item: repository.admit( - storage=item.fileitem.storage, - src_path=item.fileitem.path, - planning_input=item.planning_input, - ), - claim_task=claim_task, - ) - - assert service.put(task, Mock()) is True - assert task.lease_token is not None - assert repository.claim_recoverable( - owner_id="recovery-worker", - limit=100, - lease_seconds=120, - ) == [] - engine.dispose() - - def test_transfer_queue_service_keeps_admission_when_enqueue_fails(): """内存入队失败必须记录原因并清理视图,durable admission 由仓储保留。""" error = RuntimeError("queue closed")