refactor(transfer): keep lease fencing in execution owner

This commit is contained in:
jxxghp
2026-09-02 14:04:56 +08:00
parent 075f51bc15
commit f94eb0076d
6 changed files with 132 additions and 193 deletions
+1 -4
View File
@@ -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
+8
View File
@@ -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]
+95 -1
View File
@@ -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(
+12 -113
View File
@@ -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 值。"""
+14 -8
View File
@@ -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 == {}
+2 -67
View File
@@ -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")