fix(transfer): preserve frozen source step identity

This commit is contained in:
jxxghp
2026-09-07 12:20:44 +08:00
parent 90275e195e
commit 6ae5b9d9f0
2 changed files with 97 additions and 10 deletions
+16 -9
View File
@@ -794,20 +794,25 @@ class TransHandler:
cls,
*,
step_runner: Optional[TransferStepRunner],
fileitem: FileItem,
source_fileitem: dict[str, Any],
target_storage: str,
source_oper: StorageBase,
target_oper: StorageBase,
target_file: Path,
transfer_type: str,
) -> tuple[Optional[FileItem], str]:
"""执行稳定传输步骤,并把跨存储 move 拆为落地与源删除。"""
"""按冻结叶节点执行传输,并把跨存储 move 拆为落地与源删除。
意图直接消费原始快照,避免旧快照补默认字段或适配器更新运行期文件信息
导致步骤身份漂移;每个外部操作单独恢复文件对象。
"""
fileitem = FileItem(**source_fileitem)
cross_storage_move = (
transfer_type == "move" and fileitem.storage != target_storage
)
materialize_type = "copy" if cross_storage_move else transfer_type
intent_payload = {
"source": fileitem.model_dump(mode="json"),
"source": source_fileitem,
"target_storage": target_storage,
"target_path": target_file.as_posix(),
"transfer_type": materialize_type,
@@ -853,7 +858,7 @@ class TransHandler:
def execute_source_delete() -> TransferStepResult:
"""在目标已落地后单独删除跨存储 move 的源文件。"""
if not source_oper.delete(fileitem):
if not source_oper.delete(FileItem(**source_fileitem)):
raise RuntimeError(f"{fileitem.path} 源文件删除失败")
return TransferStepResult(payload={
"source_path": fileitem.path,
@@ -865,7 +870,7 @@ class TransHandler:
phase="transfer",
kind="delete_move_source",
payload={
"source": fileitem.model_dump(mode="json"),
"source": source_fileitem,
"target_storage": target_storage,
"target_path": target_file.as_posix(),
},
@@ -998,6 +1003,7 @@ class TransHandler:
*,
step_runner: Optional[TransferStepRunner],
fileitem: FileItem,
source_fileitem: dict[str, Any],
meta: MetaBase,
mediainfo: MediaInfo | MusicInfo,
target_oper: StorageBase,
@@ -1007,7 +1013,7 @@ class TransHandler:
overwrite_mode: Optional[str],
need_notify: bool,
) -> tuple[bool, bool, Optional[TransferInfo]]:
"""冻结覆盖策略判定,避免目标变化后重启得到不同步骤序列"""
"""以原始源快照冻结覆盖判定,避免模型补字段改变旧任务的步骤身份"""
def execute() -> TransferStepResult:
"""执行一次覆盖策略判定并冻结完整裁决。"""
over_flag, delete_versions, failure = self.__resolve_overwrite(
@@ -1032,7 +1038,7 @@ class TransHandler:
phase="decision",
kind="resolve_overwrite",
payload={
"source": fileitem.model_dump(mode="json"),
"source": source_fileitem,
"target_storage": target_storage,
"target_path": target_file.as_posix(),
"transfer_type": transfer_type,
@@ -1224,7 +1230,7 @@ class TransHandler:
source_item = FileItem(**planned_item.source_fileitem)
new_item, error = self.__execute_transfer_with_steps(
step_runner=step_runner,
fileitem=source_item,
source_fileitem=planned_item.source_fileitem,
target_storage=planned_item.target_storage,
source_oper=source_oper,
target_oper=target_oper,
@@ -1283,6 +1289,7 @@ class TransHandler:
over_flag, delete_versions, overwrite_failure = self.__resolve_overwrite_with_step(
step_runner=step_runner,
fileitem=fileitem,
source_fileitem=frozen_source_payload,
meta=meta,
mediainfo=mediainfo,
target_oper=target_oper,
@@ -1375,7 +1382,7 @@ class TransHandler:
)
new_item, error = self.__execute_transfer_with_steps(
step_runner=step_runner,
fileitem=fileitem,
source_fileitem=planned_item.source_fileitem,
target_storage=target_storage,
source_oper=source_oper,
target_oper=target_oper,
+81 -1
View File
@@ -1,5 +1,6 @@
"""验证 TransferChain 步骤 runner 与文件执行器的崩溃恢复边界。"""
from dataclasses import replace
from datetime import datetime, timezone
from pathlib import Path
from unittest.mock import Mock
@@ -30,7 +31,10 @@ from app.db.base import Base
from app.db.models.transferexecutionstep import TransferExecutionStep
from app.db.models.transferhistory import TransferHistory
from app.db.models.transferpending import TransferPending
from app.domain.context import MediaInfo
from app.domain.meta.metabase import MetaBase
from app.modules.filemanager.transhandler import TransHandler
from app.schemas.types import MediaType
from app.schemas.workflow import FileItem
@@ -384,7 +388,7 @@ def test_cross_storage_move_materializes_before_independent_source_delete(tmp_pa
result, error = TransHandler._TransHandler__execute_transfer_with_steps(
step_runner=runner,
fileitem=source_item,
source_fileitem=source_item.model_dump(mode="json"),
target_storage="remote",
source_oper=source_oper,
target_oper=target_oper,
@@ -444,3 +448,79 @@ def test_remote_to_local_transfer_creates_target_directory_before_download(tmp_p
path=target_file.parent,
)
source_oper.delete.assert_called_once_with(source_item)
@pytest.mark.parametrize("sparse", [False, True])
@pytest.mark.parametrize("directory", [False, True])
def test_frozen_disc_plan_executes_and_replays_with_real_step_ledger(
execution_repository, monkeypatch, sparse, directory,
):
"""原盘及旧快照必须通过真实意图校验,重启后不再复制或重复删除源。"""
source = FileItem(
storage="local", path="/disc" if directory else "/disc.iso",
name="disc" if directory else "disc.iso", type="dir" if directory else "file",
size=100,
).model_dump(mode="json", exclude_unset=sparse)
leaf = (
FileItem(storage="local", path="/disc/BDMV/STREAM/00001.m2ts",
name="00001.m2ts", type="file", size=200).model_dump(
mode="json", exclude_unset=sparse,
)
if directory else dict(source)
)
target = "/library/disc" if directory else "/library/disc.iso"
planning_input = TransferPlanningInput(source_fileitem=source)
checkpoint = replace(
_runner_plan_checkpoint(), planning_input=planning_input,
target_storage="remote", final_target_path=target,
resolved_transfer_type="move",
items=(TransferPlanItem(
sequence=0, source_fileitem=leaf, target_storage="remote",
target_path=f"{target}/BDMV/STREAM/00001.m2ts" if directory else target,
),),
)
with execution_repository._session_factory() as session:
pending = session.query(TransferPending).one()
pending.planning_input = planning_input.to_payload()
pending.input_fingerprint = planning_input.fingerprint
pending.checkpoint_payload = checkpoint.to_payload()
session.commit()
def intercept(**kwargs):
"""模拟插件修正运行期大小,冻结计划和步骤身份不应随之变化。"""
kwargs["fileitem"].size = 999
return True, ""
def materialize(**kwargs):
"""模拟存储适配器更新临时链接,源删除步骤仍须沿用冻结输入。"""
assert kwargs["fileitem"].size == leaf["size"]
kwargs["fileitem"].url = "https://temporary.invalid/refreshed"
return FileItem(storage="remote", path=kwargs["target_file"].as_posix()), ""
handler = TransHandler()
monkeypatch.setattr(handler, "_TransHandler__intercept_transfer", intercept)
transfer = Mock(side_effect=materialize)
monkeypatch.setattr(TransHandler, "_TransHandler__transfer_command", transfer)
source_oper = Mock()
target_oper = Mock()
target_oper.get_item_strict.return_value = None
target_oper.get_folder.return_value = FileItem(storage="remote", path="/library", type="dir")
for _ in range(2):
runner = transfer_chain_module._DurableTransferStepRunner(
task_id="task-runner", lease_token="lease",
checkpoint_fingerprint=checkpoint.fingerprint,
repository=execution_repository,
)
result = handler.execute_transfer_plan(
checkpoint, meta=MetaBase("disc.iso"),
mediainfo=MediaInfo(type=MediaType.MOVIE, title="Disc"),
source_oper=source_oper, target_oper=target_oper, step_runner=runner,
)
assert result.success
assert runner.checkpoint(result).payload["outcome"] == "succeeded"
transfer.assert_called_once()
source_oper.delete.assert_called_once()
assert source_oper.delete.call_args.args[0].url is None
assert checkpoint.planning_input.source_fileitem == source
assert checkpoint.items[0].source_fileitem == leaf