refactor(transfer): keep planning rejection within complexity budget

This commit is contained in:
jxxghp
2026-09-01 07:42:53 +08:00
parent 3161ee5576
commit c4d8185d8d
6 changed files with 43 additions and 46 deletions
+4
View File
@@ -80,6 +80,10 @@ class TransferExecutionLeaseLostError(TransferExecutionError):
"""表示持久化写入时任务租约已失效或已被其他 worker 接管。"""
class TransferPlanningRejectedError(ValueError):
"""规划输入缺少业务必需信息且继续自动重试不会改变结果。"""
def _canonical_json(payload: Mapping[str, Any]) -> str:
"""生成稳定操作身份使用的规范 JSON。"""
return json.dumps(
-4
View File
@@ -705,10 +705,6 @@ class TransferAdmissionConflictError(ValueError):
"""同一源文件以不同规划输入重复准入时抛出的冲突错误。"""
class TransferPlanningRejectedError(ValueError):
"""规划输入缺少业务必需信息且继续自动重试不会改变结果。"""
class TransferPlanningStateError(RuntimeError):
"""计划检查点无法从当前持久状态推进时抛出的状态错误。"""
+35 -38
View File
@@ -11,13 +11,13 @@ from app.application.transfer.execution import (
TransferExecutionState,
TransferOperationObservation,
TransferOperationObservationState,
TransferPlanningRejectedError,
TransferStepResult,
)
from app.application.transfer.workflow import (
TransferLeaseLostError,
TransferPlanCheckpoint,
TransferPlanningInput,
TransferPlanningRejectedError,
TransferPlanningStateError,
TransferProviderInvocationSnapshot,
TransferProviderReference,
@@ -44,6 +44,38 @@ from app.schemas.workflow import FileItem
from .execution import _DurableTransferStepRunner, _TransferRetryExhausted
def _build_planning_rejection_checkpoint(
task: TransferTask,
*,
error: str,
planning_input: TransferPlanningInput,
) -> TransferPlanCheckpoint:
"""把确定性的宿主规划错误冻结为可重放的零步骤失败计划。"""
meta_kind = planning_input.options.get("_meta_kind")
mediainfo_kind = planning_input.options.get("_mediainfo_kind")
source_path = task.fileitem.path
return TransferPlanCheckpoint(
planning_input=planning_input,
target_storage=task.fileitem.storage,
root_target_path=source_path,
final_target_path=source_path,
resolved_transfer_type=(
task.transfer_type or planning_input.requested_transfer_type or "copy"
),
items=(),
resolved_meta=planning_input.meta,
resolved_meta_kind=meta_kind if isinstance(meta_kind, str) else None,
resolved_mediainfo=planning_input.mediainfo,
resolved_mediainfo_kind=(
mediainfo_kind if isinstance(mediainfo_kind, str) else None
),
resolved_episodes_info=planning_input.episodes_info,
need_notify=planning_input.need_notify,
overwrite_mode=planning_input.overwrite_mode,
rejection_error=error,
)
class TransferPlanningOwner(_TransferOwnerBase):
"""唯一持有整理准入后的冻结计划与 provider 选择。"""
@@ -275,7 +307,7 @@ class TransferPlanningOwner(_TransferOwnerBase):
)
planning_input = task.planning_input or self._TransferChain__build_planning_input(task)
task.bind_planning_input(planning_input)
checkpoint = self._TransferChain__build_planning_rejection_checkpoint(
checkpoint = _build_planning_rejection_checkpoint(
task,
error=error,
planning_input=planning_input,
@@ -288,41 +320,6 @@ class TransferPlanningOwner(_TransferOwnerBase):
task.bind_plan_checkpoint(persisted)
return self._plan_checkpoint_and_execute(task)
def _TransferChain__build_planning_rejection_checkpoint(
self,
task: TransferTask,
*,
error: str,
planning_input: TransferPlanningInput,
) -> TransferPlanCheckpoint:
"""把确定性的宿主规划错误冻结为可重放的零步骤失败计划。"""
source_path = task.fileitem.path
resolved_transfer_type = (
task.transfer_type
or planning_input.requested_transfer_type
or "copy"
)
return TransferPlanCheckpoint(
planning_input=planning_input,
target_storage=task.fileitem.storage,
root_target_path=source_path,
final_target_path=source_path,
resolved_transfer_type=resolved_transfer_type,
items=(),
resolved_meta=self._TransferChain__json_snapshot(task.meta),
resolved_meta_kind=(type(task.meta).__name__ if task.meta else None),
resolved_mediainfo=self._TransferChain__json_snapshot(task.mediainfo),
resolved_mediainfo_kind=(
type(task.mediainfo).__name__ if task.mediainfo else None
),
resolved_episodes_info=tuple(
self._TransferChain__json_snapshot(item) for item in (task.episodes_info or [])
),
need_notify=planning_input.need_notify,
overwrite_mode=planning_input.overwrite_mode,
rejection_error=error,
)
def _TransferChain__execute_planning_rejection(
self,
task: TransferTask,
@@ -711,7 +708,7 @@ class TransferPlanningOwner(_TransferOwnerBase):
),
)
except TransferPlanningRejectedError as error:
return self._TransferChain__build_planning_rejection_checkpoint(
return _build_planning_rejection_checkpoint(
task,
error=str(error),
planning_input=planning_input,
+1 -2
View File
@@ -712,8 +712,7 @@ class TransferWorkflowOwner(_TransferOwnerBase):
if should_reorganize:
if not reorganize:
durable_retry = self._request_durable_transfer_retry(
transferd,
requested_by="manual_reorganize",
transferd, requested_by="manual_reorganize",
)
if durable_retry is not None:
accepted, message = durable_retry
+1 -1
View File
@@ -13,6 +13,7 @@ from app.application.messaging.message import TemplateHelper
from app.application.transfer.execution import (
TransferOperationObservation,
TransferOperationObservationState,
TransferPlanningRejectedError,
TransferStepResult,
TransferStepRunner,
)
@@ -20,7 +21,6 @@ from app.application.transfer.workflow import (
TransferPlanCheckpoint,
TransferPlanItem,
TransferPlanningInput,
TransferPlanningRejectedError,
)
from app.domain.context import MediaInfo, MusicInfo
from app.domain.meta.metabase import MetaBase
+2 -1
View File
@@ -19,6 +19,7 @@ from app.application.transfer.execution import (
TransferExecutionSnapshot,
TransferExecutionState,
TransferExecutionStep,
TransferPlanningRejectedError,
TransferSettlementResult,
TransferStepState,
)
@@ -1356,7 +1357,7 @@ def test_host_planning_value_error_commits_rejection_instead_of_retrying():
kwargs["checkpoint"],
)
chain = _chain(repository=repository)
chain.plan_transfer.side_effect = transfer_application.TransferPlanningRejectedError(
chain.plan_transfer.side_effect = TransferPlanningRejectedError(
"未识别到文件集数"
)