fix(transfer): repair migration and failed recovery

This commit is contained in:
jxxghp
2026-09-01 13:31:52 +08:00
parent f5c73dcabc
commit 8e8ebcdff9
6 changed files with 221 additions and 16 deletions
+23 -12
View File
@@ -362,22 +362,24 @@ class FailedRetryMixin(_TransferOwnerBase):
if not history:
logger.error(f"整理记录不存在,ID{logid}")
return False, "整理记录不存在"
durable_retry = _request_durable_transfer_retry(
history,
requested_by="history_redo",
repository=self.transfer_execution_repository,
)
if durable_retry is not None:
return durable_retry
# 显式媒体身份代表用户要求重新规划;普通 /redo 才复用原 durable 计划。
explicit_identity = media_source is not None or media_id is not None
if explicit_identity and (not media_source or not media_id):
return False, "媒体重新识别需要同时提供 media_source 和 media_id"
if not explicit_identity:
durable_retry = _request_durable_transfer_retry(
history,
requested_by="history_redo",
repository=self.transfer_execution_repository,
)
if durable_retry is not None:
return durable_retry
# 按源目录路径重新整理
src_path = Path(history.src)
if not src_path.exists():
return False, f"源目录不存在:{src_path}"
# 查询媒体信息
explicit_identity = media_source is not None or media_id is not None
if explicit_identity and (not media_source or not media_id):
return False, "媒体重新识别需要同时提供 media_source 和 media_id"
if mtype and media_source and media_id:
if explicit_identity:
mediainfo = MediaChain().recognize_media(
mtype=mtype,
media_source=media_source,
@@ -433,7 +435,14 @@ class FailedRetryMixin(_TransferOwnerBase):
logger.info(f"{src_path.name} 识别为:{mediainfo.title_year}")
# 删除旧的已整理文件
if history.dest_fileitem:
if getattr(history, "transfer_task_id", None):
state, errmsg = self._delete_manual_transfer_history(
history=history,
transfer_history_oper=self.transfer_history_repository,
)
if not state:
return False, errmsg
elif history.dest_fileitem:
if not isinstance(history.dest_fileitem, dict):
return False, "目标文件历史数据无效"
# 解析目标文件对象
@@ -448,6 +457,8 @@ class FailedRetryMixin(_TransferOwnerBase):
fileitem=FileItem(**history.src_fileitem),
mediainfo=mediainfo,
mtype=mtype,
media_source=media_source,
media_id=media_id,
download_hash=history.download_hash,
force=True,
background=False,
+3 -3
View File
@@ -697,8 +697,8 @@ class TransferWorkflowOwner(_TransferOwnerBase):
# 自动整理按 app/application/history.py 的统一判定去重(失败记录放行重试、
# 成功但源文件已变化放行交 overwrite_mode 决断);手动整理可清理失败记录,
# 或按用户确认清理成功记录。
if (not force or reorganize) and not preview:
# 或按用户确认清理成功记录;手动显式指定媒体身份时,先解除旧失败任务再重新规划
if (not force or reorganize or (manual and (media_source is not None or media_id is not None))) and not preview:
transfer_history_oper = self.transfer_history_repository
transferd = self._get_manual_transfer_history(
fileitem=file_item,
@@ -710,7 +710,7 @@ class TransferWorkflowOwner(_TransferOwnerBase):
reorganize or not transferd.status
)
if should_reorganize:
if not reorganize:
if not reorganize and not (media_source is not None or media_id is not None):
durable_retry = self._request_durable_transfer_retry(
transferd, requested_by="manual_reorganize",
)
+19
View File
@@ -47,8 +47,27 @@ def _create_table() -> None:
)
def _is_current_schema_precreated() -> bool:
"""识别首次启动时由当前 ORM 预先创建的待整理结构。"""
inspector = sa.inspect(op.get_bind())
if "transferexecutionstep" not in inspector.get_table_names():
return False
# 未标记的全新数据库会先由 metadata.create_all 创建当前 transferpending
# 后续 transferexecutionstep 已带有指向 task_id 的外键;此时不能按 3.0.4
# 的旧字段集合删除父表,否则 PostgreSQL 会拒绝删除被依赖的父表。
return any(
foreign_key.get("referred_table") == _TABLE_NAME
and foreign_key.get("referred_columns") == ["task_id"]
and foreign_key.get("constrained_columns") == ["task_id"]
for foreign_key in inspector.get_foreign_keys("transferexecutionstep")
)
def _validate_or_recreate_table() -> None:
"""校验中断升级留下的表结构,仅允许空残表自动重建。"""
if _is_current_schema_precreated():
return
inspector = sa.inspect(op.get_bind())
columns = inspector.get_columns(_TABLE_NAME)
column_names = {column["name"] for column in columns}
+60
View File
@@ -18,6 +18,8 @@ from app.db.adapters.history.transfer import TransactionalTransferHistoryReposit
from app.db.session import SessionFactory, async_session_scope
from app.runtime.config import settings
from app.schemas.transfer import ManualTransferItem
from app.schemas.types import MediaSource, MediaType
from tests.test_transfer_job_manager import FakeMedia
from tests.test_transfer_sync_extra_files import (
FakeMeta,
make_fileitem,
@@ -677,6 +679,64 @@ def test_explicit_durable_reorganize_discards_old_task_and_replans(monkeypatch):
assert planned == [fileitem.path]
def test_manual_explicit_identity_replans_failed_durable_history(monkeypatch):
"""手动指定媒体身份时,应清理失败 durable 任务而不是复用原计划重试。"""
chain = make_transfer_chain()
fileitem = make_fileitem("/downloads/Test.Show.S01E01.mkv")
history = SimpleNamespace(
id=13,
transfer_task_id="transfer-task-13",
transfer_settlement_revision=2,
status=False,
mode="copy",
dest_fileitem=None,
download_hash=None,
downloader=None,
src=fileitem.path,
src_storage=fileitem.storage,
)
planned = []
deleted = []
_patch_transfer_planning(
monkeypatch,
chain,
fileitem,
history,
planned,
deleted,
)
monkeypatch.setattr(
chain,
"_request_durable_transfer_retry",
lambda *_args, **_kwargs: (_ for _ in ()).throw(
AssertionError("新的媒体身份不得复用旧 durable 计划")
),
)
monkeypatch.setattr(
chain,
"_delete_manual_transfer_history",
lambda history, transfer_history_oper: deleted.append(
("history", history.id)
) or (True, ""),
)
state, message = TransferChain.do_transfer(
chain,
fileitem=fileitem,
mediainfo=FakeMedia(286322),
mtype=MediaType.TV,
media_source=MediaSource.TMDB,
media_id="286322",
background=False,
manual=True,
)
assert state is True
assert message == ""
assert deleted == [("history", 13)]
assert planned == [fileitem.path]
def test_manual_reorganize_keeps_successful_move_target_as_source(monkeypatch):
"""成功移动后的目标是当前重整源,只能删历史记录,不能先删除文件。"""
chain = make_transfer_chain()
+78 -1
View File
@@ -8,7 +8,7 @@ from app.application.transfer.execution import (
TransferRetryRequestResult,
)
from app.chain.transfer.facade import TransferChain
from app.schemas.types import NotificationChannel
from app.schemas.types import MediaSource, MediaType, NotificationChannel
class _RetryCommand:
@@ -128,6 +128,83 @@ def test_durable_history_redo_only_requests_persistent_retry(monkeypatch):
]
def test_explicit_history_redo_discards_failed_task_before_replanning(monkeypatch):
"""显式 /redo 身份应先解除旧失败任务,再执行新的识别和整理。"""
repository = _install_discard_port(monkeypatch)
history = SimpleNamespace(
id=86,
transfer_task_id="transfer-task-86",
transfer_settlement_revision=5,
src="/downloads/source.mkv",
src_storage="local",
src_fileitem={
"storage": "local",
"path": "/downloads/source.mkv",
"type": "file",
"name": "source.mkv",
},
dest_fileitem=None,
download_hash=None,
media_source=None,
media_id=None,
episode_group=None,
)
deleted = []
planned = []
history_port = SimpleNamespace(
get=lambda history_id: history,
delete=lambda history_id: deleted.append(("history", history_id)),
)
chain = object.__new__(TransferChain)
chain.transfer_execution_repository = repository
chain.transfer_history_repository = history_port
chain.obtain_images = lambda **_kwargs: None
monkeypatch.setattr(
"app.chain.transfer.retry.Path.exists",
lambda _path: True,
)
monkeypatch.setattr(
"app.chain.transfer.retry.MediaChain",
lambda: SimpleNamespace(
recognize_media=lambda **_kwargs: SimpleNamespace(title_year="Test Show")
),
)
monkeypatch.setattr(
"app.chain.transfer.records.clear_transfer_failures",
lambda *_args: None,
)
def fake_do_transfer(**kwargs):
"""记录显式身份重整是否重新进入整理准入。"""
planned.append(kwargs)
return True, ""
monkeypatch.setattr(chain, "do_transfer", fake_do_transfer)
state, message = chain._re_transfer(
logid=86,
mtype=MediaType.TV,
media_source=MediaSource.TMDB,
media_id="286322",
)
assert state is True
assert message == ""
assert _DiscardCommand.calls == [
(
repository,
{
"task_id": "transfer-task-86",
"history_id": 86,
"settlement_revision": 5,
},
)
]
assert deleted == [("history", 86)]
assert planned[0]["media_source"] == MediaSource.TMDB
assert planned[0]["media_id"] == "286322"
def test_durable_manual_cleanup_discards_task_and_removes_old_state(monkeypatch):
"""显式重整命中 FAILED durable 历史时应放弃任务并清理旧状态。"""
repository = _install_discard_port(monkeypatch)
@@ -65,6 +65,44 @@ def test_upgrade_recreates_empty_partial_table(monkeypatch) -> None:
} == {"id", "storage", "src_path", "created_at"}
def test_upgrade_keeps_current_schema_precreated_with_execution_fk(monkeypatch) -> None:
"""首次初始化已由当前 ORM 建表时,旧迁移不得删除被步骤表依赖的父表。"""
engine = sa.create_engine("sqlite://")
with engine.connect() as connection:
connection.execute(sa.text("PRAGMA foreign_keys=ON"))
connection.commit()
with connection.begin():
connection.execute(sa.text(
"CREATE TABLE transferpending ("
"id INTEGER PRIMARY KEY, storage VARCHAR NOT NULL, "
"src_path VARCHAR NOT NULL, created_at VARCHAR, "
"task_id VARCHAR NOT NULL UNIQUE, execution_state VARCHAR NOT NULL)"
))
connection.execute(sa.text(
"CREATE TABLE transferexecutionstep ("
"id INTEGER PRIMARY KEY, task_id VARCHAR NOT NULL, "
"FOREIGN KEY (task_id) REFERENCES transferpending(task_id) "
"ON DELETE CASCADE)"
))
migration = _bind_migration(monkeypatch, connection)
migration.upgrade()
inspector = sa.inspect(connection)
assert "transferexecutionstep" in inspector.get_table_names()
assert "task_id" in {
column["name"]
for column in inspector.get_columns("transferpending")
}
index = next(
index
for index in inspector.get_indexes("transferpending")
if index["name"] == "ux_transferpending_storage_path"
)
assert index["column_names"] == ["storage", "src_path"]
assert index["unique"] == 1
def test_upgrade_rejects_nonempty_partial_table(monkeypatch) -> None:
"""含数据残表无法可靠推断源身份时必须显式拒绝迁移。"""
engine = sa.create_engine("sqlite://")