fix(transfer): close durable error recovery loop

This commit is contained in:
jxxghp
2026-09-08 22:38:59 +08:00
parent 20facc8941
commit d0c894f5eb
17 changed files with 178 additions and 14 deletions
+25
View File
@@ -457,6 +457,31 @@ def delete_transfer_history(
)
@router.post(
"/transfer/{history_id}/discard-corrupt",
summary="放弃损坏的整理任务",
response_model=_SchemaResponse[dict],
)
async def discard_corrupt_transfer_history(
history_id: int,
query: HistoryQueryService = Depends(get_history_query_service),
execution_repository: TransferExecutionRepository = Depends(get_transfer_execution_repository),
_: object = Depends(get_current_active_manage_user),
) -> Any:
"""清理无活动租约的损坏 durable 任务,保留历史供重新生成计划。"""
history = await query.get_transfer(history_id)
if not history:
return _SchemaResponse(success=False, message="整理记录不存在")
if not history.transfer_task_id:
return _SchemaResponse(success=True, message="整理任务已清理", data={"history_id": history_id})
result = await asyncio.to_thread(
TransferExecutionCommand(execution_repository).discard_corrupt_by_history,
task_id=history.transfer_task_id,
history_id=history_id,
)
return _SchemaResponse(success=result.discarded, message=result.message, data={"history_id": history_id})
@router.post(
"/transfer/{history_id}/ai-redo",
summary="智能助手重新整理",
+1 -1
View File
@@ -116,7 +116,7 @@ def _is_valid_source_media_id(
try:
UUID(normalized_media_id)
return True
except TypeError, ValueError:
except (TypeError, ValueError):
return False
if normalized_source == MediaSource.DoubanMusic and ":" in normalized_media_id:
album_id, track_number = normalized_media_id.split(":", 1)
+2 -1
View File
@@ -1,3 +1,4 @@
from __future__ import annotations
"""插件持久化数据查询、投影与写用例。"""
import json
@@ -169,7 +170,7 @@ def plugin_data_serialized_chars(value: JsonData) -> Optional[int]:
"""计算合法 JSON 值的紧凑字符数,异常对象不执行自定义字符串化。"""
try:
return len(json.dumps(value, ensure_ascii=False, separators=(",", ":")))
except TypeError, ValueError:
except (TypeError, ValueError):
return None
@@ -1,3 +1,4 @@
from __future__ import annotations
"""订阅执行准入、搜索上下文、批次任务与持久队列端口。"""
import threading
+44
View File
@@ -608,6 +608,23 @@ class TransferExecutionRepository(Protocol):
) -> TransferRetryRequestResult:
"""仅把 FAILED 终态转入到期可 claim 的 retry_wait。"""
def discard_corrupt_task(
self,
*,
task_id: str,
lease_token: str,
error: str,
) -> bool:
"""在当前租约下原子清除无法继续执行的损坏任务及其恢复证据。"""
def discard_corrupt_by_history(
self,
*,
task_id: str,
history_id: int,
) -> TransferFailureDiscardResult:
"""清理无活动租约的损坏任务,并解除对应历史绑定。"""
def discard_failed(
self,
*,
@@ -832,6 +849,33 @@ class TransferExecutionCommand:
requested_by=requested_by,
)
def discard_corrupt_task(
self,
*,
task_id: str,
lease_token: str,
error: str,
) -> bool:
"""以当前租约清理损坏任务,避免恢复线程再次回放旧步骤。"""
if not task_id or not lease_token or not error:
raise ValueError("损坏任务收口缺少任务、租约或错误原因")
return self._repository.discard_corrupt_task(
task_id=task_id, lease_token=lease_token, error=error
)
def discard_corrupt_by_history(
self,
*,
task_id: str,
history_id: int,
) -> TransferFailureDiscardResult:
"""放弃无法重试的损坏任务,保留历史记录供重新生成计划。"""
if not task_id or history_id <= 0:
raise ValueError("放弃损坏任务缺少任务或历史")
return self._repository.discard_corrupt_by_history(
task_id=task_id, history_id=history_id
)
def discard_failed(
self,
*,
+1
View File
@@ -1030,6 +1030,7 @@ class TransferFailureNotification:
image: Optional[str]
username: Optional[str]
manual_identity: bool = False
task_id: Optional[str] = None
def build_transfer_failure_group_key(task: TransferTask) -> str:
+1
View File
@@ -1,3 +1,4 @@
from __future__ import annotations
"""字幕获取、解压和存储 owner。"""
import re
+1 -1
View File
@@ -1407,7 +1407,7 @@ class TransferQueueOwner(_TransferOwnerBase):
logger.error(
f"{fileitem.name} 整理任务处理出现错误:{e} - {traceback.format_exc()}"
)
self._TransferChain__fail_transfer_task(task)
self._TransferChain__fail_transfer_task(task, e)
with task_lock:
self._processed_num += 1
self._fail_num += 1
+6 -1
View File
@@ -78,6 +78,10 @@ class FailedRetryMixin(_TransferOwnerBase):
return [
[
{"text": "重试", "callback_data": f"transfer_retry_{history_id}"},
{
"text": "重新生成计划",
"callback_data": f"transfer_regenerate_{history_id}",
},
{
"text": "智能助手接管",
"callback_data": f"transfer_ai_retry_{history_id}",
@@ -100,6 +104,7 @@ class FailedRetryMixin(_TransferOwnerBase):
"""
for prefix, action in (
("transfer_retry_", "retry"),
("transfer_regenerate_", "regenerate"),
("transfer_ai_retry_", "ai_retry"),
):
if callback_data.startswith(prefix):
@@ -125,7 +130,7 @@ class FailedRetryMixin(_TransferOwnerBase):
return False
action, history_id = callback
if action == "retry":
if action in {"retry", "regenerate"}:
self._retry_transfer_history(
history_id=history_id,
channel=channel,
+26 -4
View File
@@ -539,6 +539,7 @@ class TransferSettlementOwner(_TransferOwnerBase):
),
username=task.username,
manual_identity=manual_identity,
task_id=task.admission_task_id,
)
if not self.runtime_config.transfer_failure_notification_aggregation:
self._send_transfer_failure_notifications([notification])
@@ -692,10 +693,31 @@ class TransferSettlementOwner(_TransferOwnerBase):
self._TransferChain__release_task_claim(task)
return True
def _TransferChain__fail_transfer_task(self, task: TransferTask):
"""
标记异常整理任务失败并清理作业视图
"""
def _TransferChain__fail_transfer_task(self, task: TransferTask, error: object = "整理任务处理失败"):
"""清理作业视图,并在执行冲突时原子删除 durable 恢复证据。"""
error_text = str(error)
corrupt_plan = any(
marker in error_text
for marker in ("记录已失效", "记录不完整", "版本不一致", "检查点", "恢复状态不完整")
)
if (
isinstance(error, TransferExecutionConflictError)
and corrupt_plan
and not task.preview
and task.admission_task_id
and task.lease_token
):
try:
self._transfer_executions.discard_corrupt_task(
task_id=task.admission_task_id,
lease_token=task.lease_token,
error=str(error),
)
except Exception as cleanup_error:
logger.error(
"清理损坏整理任务 durable 证据失败:%s - %s",
task.admission_task_id, cleanup_error,
)
self.jobview.fail_unfinished_task(task)
self.jobview.try_remove_job(task)
self._finish_scrape_batch_task(task)
+1 -1
View File
@@ -1039,7 +1039,7 @@ class TransferWorkflowOwner(_TransferOwnerBase):
f"{transfer_task.fileitem.name} 整理任务处理出现错误:{e} - {traceback.format_exc()}"
)
if not preview:
self._TransferChain__fail_transfer_task(transfer_task)
self._TransferChain__fail_transfer_task(transfer_task, e)
state, err_msg = False, "整理任务处理失败,请稍后重试"
finally:
durable_settled = self._TransferChain__finish_job_execution(
+64
View File
@@ -1209,6 +1209,70 @@ class TransactionalTransferExecutionRepository:
self._rollback(transaction)
raise
def discard_corrupt_task(
self,
*,
task_id: str,
lease_token: str,
error: str,
) -> bool:
"""在当前有效租约下清理损坏 pending、步骤及历史任务映射。"""
if not task_id or not lease_token or not error:
return False
now_utc, updated_at = self._times()
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
pending = TransferPendingOper(session).get_by_task_id(task_id=task_id)
self._require_active_lease(pending, lease_token=lease_token, now_utc=now_utc)
history = TransferHistoryOper(session).get_by_transfer_task_id(task_id=task_id)
if history is not None:
history.transfer_task_id = None
history.transfer_settlement_revision = None
TransferExecutionStepOper(session).stage_delete_task(task_id=task_id)
deleted = session.query(TransferPending).filter(
TransferPending.task_id == task_id,
TransferPending.lease_token == lease_token,
).delete(synchronize_session=False)
if deleted != 1:
raise TransferExecutionConflictError("损坏整理任务清理未完成,请刷新后重试")
transaction.commit()
return True
except Exception:
self._rollback(transaction)
raise
def discard_corrupt_by_history(
self,
*,
task_id: str,
history_id: int,
) -> TransferFailureDiscardResult:
"""在无活动租约时删除损坏任务、步骤并解除历史绑定。"""
with self._session_factory() as session:
transaction = SqlAlchemyUnitOfWork(session)
try:
pending_oper = TransferPendingOper(session)
history_oper = TransferHistoryOper(session)
pending = pending_oper.get_by_task_id(task_id=task_id)
history = history_oper.get(history_id)
if pending is None:
return TransferFailureDiscardResult(True, None, "整理任务已被清理")
state = TransferExecutionState(pending.execution_state)
if pending.lease_owner is not None or pending.lease_token is not None:
return TransferFailureDiscardResult(False, state, "整理任务正在处理中,暂时无法放弃")
if history is None or history.transfer_task_id != task_id:
return TransferFailureDiscardResult(False, state, "整理历史与任务绑定已变化,请刷新后重试")
history.transfer_task_id = None
history.transfer_settlement_revision = None
TransferExecutionStepOper(session).stage_delete_task(task_id=task_id)
session.delete(pending)
transaction.commit()
return TransferFailureDiscardResult(True, state, "已放弃损坏整理任务")
except Exception:
self._rollback(transaction)
raise
def discard_failed(
self,
*,
+1 -1
View File
@@ -231,7 +231,7 @@ def _optional_int(value: object) -> int | None:
return None
try:
return int(str(value))
except TypeError, ValueError:
except (TypeError, ValueError):
return None
+1 -1
View File
@@ -831,7 +831,7 @@ class MusicBrainzModule(_ModuleBase):
"""相关度每五分成组,组内应用发行偏好并保留原始得分。"""
try:
score = int(release.get("score") or 0)
except TypeError, ValueError:
except (TypeError, ValueError):
score = 0
region_rank, script_rank = cls._release_preference_sort_key(
release,
+1 -1
View File
@@ -164,7 +164,7 @@ class TheMovieDbModule(MediaAuxiliaryProviderMixin, _ModuleBase):
try:
tmdb_id = int(raw_tmdb_id)
media_type = MediaType(request.media_type)
except TypeError, ValueError:
except (TypeError, ValueError):
return None
info = self.tmdb_info(tmdbid=tmdb_id, mtype=media_type)
if not info:
+1 -1
View File
@@ -1809,7 +1809,7 @@ def diagnose_module_callable(method: str, callback: Callable[..., Any]) -> tuple
return ()
try:
parameters = inspect.signature(callback).parameters
except TypeError, ValueError:
except (TypeError, ValueError):
return ("signature-unavailable",)
missing = tuple(
name
+1 -1
View File
@@ -120,7 +120,7 @@ class RuntimeClassificationEnrichmentCache:
return ClassificationEnrichmentCacheEntry(response=None)
try:
response = ClassificationEnrichmentResponse.model_validate(value.get("response"))
except TypeError, ValidationError:
except (TypeError, ValidationError):
return None
return ClassificationEnrichmentCacheEntry(response=response)