diff --git a/app/application/transfer/execution.py b/app/application/transfer/execution.py index 2b55c8d6f..b88560a02 100644 --- a/app/application/transfer/execution.py +++ b/app/application/transfer/execution.py @@ -80,6 +80,10 @@ class TransferExecutionLeaseLostError(TransferExecutionError): """表示持久化写入时任务租约已失效或已被其他 worker 接管。""" +class TransferPlanningRejectedError(ValueError): + """规划输入缺少业务必需信息且继续自动重试不会改变结果。""" + + def _canonical_json(payload: Mapping[str, Any]) -> str: """生成稳定操作身份使用的规范 JSON。""" return json.dumps( diff --git a/app/application/transfer/workflow.py b/app/application/transfer/workflow.py index 193ddf13c..9f0c14010 100644 --- a/app/application/transfer/workflow.py +++ b/app/application/transfer/workflow.py @@ -705,10 +705,6 @@ class TransferAdmissionConflictError(ValueError): """同一源文件以不同规划输入重复准入时抛出的冲突错误。""" -class TransferPlanningRejectedError(ValueError): - """规划输入缺少业务必需信息且继续自动重试不会改变结果。""" - - class TransferPlanningStateError(RuntimeError): """计划检查点无法从当前持久状态推进时抛出的状态错误。""" diff --git a/app/chain/transfer/plan.py b/app/chain/transfer/plan.py index 46380477c..da6415c69 100644 --- a/app/chain/transfer/plan.py +++ b/app/chain/transfer/plan.py @@ -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, diff --git a/app/chain/transfer/workflow.py b/app/chain/transfer/workflow.py index e4d1296a8..942063f1f 100644 --- a/app/chain/transfer/workflow.py +++ b/app/chain/transfer/workflow.py @@ -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 diff --git a/app/modules/filemanager/transhandler.py b/app/modules/filemanager/transhandler.py index 478271b80..be97cf817 100644 --- a/app/modules/filemanager/transhandler.py +++ b/app/modules/filemanager/transhandler.py @@ -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 diff --git a/tests/test_transfer_planning_checkpoint.py b/tests/test_transfer_planning_checkpoint.py index 9a093b7d9..b905c0f2d 100644 --- a/tests/test_transfer_planning_checkpoint.py +++ b/tests/test_transfer_planning_checkpoint.py @@ -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( "未识别到文件集数" )