fix(transfer): deduplicate recovery failures and system alerts

This commit is contained in:
jxxghp
2026-09-08 14:11:07 +08:00
parent a2f1407c95
commit f9f2601144
6 changed files with 285 additions and 24 deletions

View File

@@ -114,6 +114,7 @@ class TransferQueueOwner(_TransferOwnerBase):
self._worker_owner_id = uuid.uuid4().hex
self._owned_leases: Dict[str, Tuple[str, float]] = {}
self._queued_lease_tokens: set[Tuple[str, str]] = set()
self._resident_tasks: Dict[Tuple[str, str], TransferTask] = {}
self._replay_thread: Optional[threading.Thread] = None
self._replay_stop_event = threading.Event()
self._recovery_wakeup_event = threading.Event()
@@ -358,9 +359,10 @@ class TransferQueueOwner(_TransferOwnerBase):
"""
添加到待整理队列
:param task: 任务信息
:return: True表示任务已添加False表示链已关闭或任务无效/重复
:return: 新任务入队或相同恢复租约已被接收时返回 True关闭、无效或冲突返回 False
:raises Exception: 持久准入、批次登记或内存入队失败
"""
self._TransferChain__ensure_lease_runtime_state()
with self._worker_state_lock:
if self._closing:
logger.warning("文件整理链已关闭,拒绝新的队列任务")
@@ -374,7 +376,7 @@ class TransferQueueOwner(_TransferOwnerBase):
return TransferQueueService(
register_task=self._TransferChain__put_to_jobview,
admit_task=self._TransferChain__admit_transfer,
enqueue=self._queue.put,
enqueue=self._enqueue_transfer_task,
before_enqueue=self._register_scrape_batch_task,
enqueue_failed=self._TransferChain__record_enqueue_failure,
remove_task=self.jobview.remove_task,
@@ -445,6 +447,8 @@ class TransferQueueOwner(_TransferOwnerBase):
self._owned_leases = {}
if not hasattr(self, "_queued_lease_tokens"):
self._queued_lease_tokens = set()
if not hasattr(self, "_resident_tasks"):
self._resident_tasks = {}
if not hasattr(self, "_recovery_wakeup_event"):
self._recovery_wakeup_event = threading.Event()
if not hasattr(self, "_lease_heartbeat_thread"):
@@ -541,9 +545,12 @@ class TransferQueueOwner(_TransferOwnerBase):
self._queued_lease_tokens.discard((task_id, lease_token))
def _TransferChain__is_claimed_task_enqueued(self, task_id: str, lease_token: str) -> bool:
"""返回指定 claim 是否已经成功进入普通 worker 队列"""
"""返回 claim 是否已被队列接收,包含刚被 worker 取走但尚未完成的任务"""
with self._worker_state_lock:
return (task_id, lease_token) in self._queued_lease_tokens
return (task_id, lease_token) in self._queued_lease_tokens or any(
(task.admission_task_id, task.lease_token) == (task_id, lease_token)
for task in getattr(self, "_resident_tasks", {}).values()
)
def _TransferChain__owns_lease(self, task_id: str, lease_token: Optional[str]) -> bool:
"""返回本地续期镜像是否仍持有指定 token。"""
@@ -657,26 +664,52 @@ class TransferQueueOwner(_TransferOwnerBase):
release_thread.start()
return release_thread
@staticmethod
def _resident_task_key(task: TransferTask) -> Tuple[str, str]:
"""真实队列凭证使用存储与路径,不依赖媒体识别分组或视图状态。"""
assert task.fileitem is not None
assert task.fileitem.path
path = Path(str(task.fileitem.path).replace("\\", "/")).as_posix().rstrip("/") or "/"
return task.fileitem.storage or "local", path
def _enqueue_transfer_task(self, item: TransferQueue) -> None:
"""入队成功后登记真实任务及其租约,登记持续到 worker 完成该队列项。"""
task = item.task
assert task is not None
key = self._resident_task_key(task)
with self._worker_state_lock:
self._queue.put(item)
self._resident_tasks[key] = task
if task.admission_task_id and task.lease_token:
self._queued_lease_tokens.add((task.admission_task_id, task.lease_token))
def _finish_queue_item(self, task: TransferTask) -> None:
"""仅移除当前队列项的执行登记,避免旧 worker 清掉同源新任务。"""
with self._worker_state_lock:
residents = getattr(self, "_resident_tasks", {})
key = self._resident_task_key(task)
if residents.get(key) is task:
residents.pop(key)
if task.admission_task_id and task.lease_token:
self._queued_lease_tokens.discard((task.admission_task_id, task.lease_token))
self._queue.task_done()
def _TransferChain__enqueue_claimed_task(self, task: TransferTask) -> bool:
"""把已 claim 的恢复任务送入普通队列,禁止再次准入或二次 claim"""
"""复用持有相同租约的真实队列任务;残留视图必须重新入队才能算恢复"""
self._TransferChain__assert_owned_lease(task)
if not self._TransferChain__put_to_jobview(task):
logger.warning(
"恢复任务被内存作业视图判定为重复未进入队列task_id=%s, source=%s:%s",
task.admission_task_id,
task.fileitem.storage,
task.fileitem.path,
resident = self._resident_tasks.get(self._resident_task_key(task))
if resident is not None:
return (resident.admission_task_id, resident.lease_token) == (
task.admission_task_id, task.lease_token
)
return False
if not self._TransferChain__put_to_jobview(task):
# 作业组未完成时会保留单个文件的旧终态;视图不等于真实队列。
self.jobview.remove_task(task.fileitem)
if not self._TransferChain__put_to_jobview(task):
return False
try:
self._register_scrape_batch_task(task)
assert task.admission_task_id is not None
assert task.lease_token is not None
with self._worker_state_lock:
self._queued_lease_tokens.add(
(task.admission_task_id, task.lease_token)
)
self._queue.put(
self._enqueue_transfer_task(
TransferQueue(task=task, callback=self._TransferChain__default_callback)
)
except Exception as err:
@@ -1270,7 +1303,7 @@ class TransferQueueOwner(_TransferOwnerBase):
self._TransferChain__release_task_claim(task, error=str(err))
self.jobview.try_remove_job(task)
self._finish_scrape_batch_task(task)
self._queue.task_done()
self._finish_queue_item(task)
self._TransferChain__settle_transfer_progress_if_idle()
continue
except Exception as err:
@@ -1279,7 +1312,7 @@ class TransferQueueOwner(_TransferOwnerBase):
)
self.jobview.try_remove_job(task)
self._finish_scrape_batch_task(task)
self._queue.task_done()
self._finish_queue_item(task)
self._TransferChain__settle_transfer_progress_if_idle()
self._TransferChain__ensure_recovery_scheduler(immediate=False)
continue
@@ -1391,7 +1424,7 @@ class TransferQueueOwner(_TransferOwnerBase):
f"整理任务终态结算异常:{task.admission_task_id} - {err}"
)
finally:
self._queue.task_done()
self._finish_queue_item(task)
with task_lock:
# 减少运行中的任务数
self._active_tasks -= 1

View File

@@ -2,7 +2,9 @@
from __future__ import annotations
import threading
import traceback
from collections import OrderedDict
from collections.abc import Callable
from typing import Any, Optional
@@ -11,6 +13,7 @@ from app.schemas.types import EventType
EventErrorNotifier = Callable[[str, str], object]
_MAX_REPORTED_ERRORS = 4096
class EventErrorPolicy:
@@ -22,9 +25,27 @@ class EventErrorPolicy:
notifier: Callable[[], Optional[EventErrorNotifier]],
emit_system_error: Callable[[dict[str, Any]], object],
) -> None:
"""注入通知读取器和 SystemError 发送回调"""
"""注入通知与错误广播回调,并初始化本进程有界提示去重记录"""
self._notifier = notifier
self._emit_system_error = emit_system_error
self._reported_errors: OrderedDict[tuple[str, ...], None] = OrderedDict()
self._reported_errors_lock = threading.Lock()
def _is_repeated_error(self, event: Any, handler: str, error: Exception) -> bool:
"""按持久事件及处理器错误去重提示;普通事件保持原行为,缓存有界。"""
payload = event.event_data
event_key = payload.get("idempotency_key") if isinstance(payload, dict) else None
if not isinstance(event_key, str) or not event_key:
return False
key = (event.event_type.value, event_key, handler, str(error))
with self._reported_errors_lock:
if key in self._reported_errors:
self._reported_errors.move_to_end(key)
return True
self._reported_errors[key] = None
if len(self._reported_errors) > _MAX_REPORTED_ERRORS:
self._reported_errors.popitem(last=False)
return False
def handle(
self,
@@ -35,9 +56,17 @@ class EventErrorPolicy:
method_name: str,
error: Exception,
) -> None:
"""记录并通知异常SystemError 自身失败时只降级写日志。"""
"""每次失败保留日志,仅对同一持久事件的相同错误提示一次。
Outbox 仍按严格回执重试和结算;去重只作用于本进程的系统提示及错误广播,
不吞掉原处理器异常,也不把未完成的持久事件误标为完成。
"""
trace = traceback.format_exc()
logger.error("%s 事件处理出错:%s - %s", module_name, str(error), trace)
if self._is_repeated_error(
event, f"{module_name}.{class_name}.{method_name}", error
):
return
notifier = self._notifier()
if notifier:
try:

View File

@@ -829,6 +829,10 @@ Durable post-commit side effects have a separate boundary:
context carry the stable event key, and consumers that support deduplication
should use it. Legacy notification plugins retain their existing method
signature, so the host must not claim provider-level exactly-once delivery.
Event handler failures still propagate to the strict dispatcher. The runtime
error policy bounds and deduplicates identical system alerts by event key,
handler and error within one process; every failed attempt remains in logs.
This alert cache is not a durable delivery receipt.
- Terminal history is part of the shared data-maintenance policy and is cleaned
in bounded daily batches only when that policy is enabled. Completed intents
default to 30-day retention and dead letters to 90 days; both values are
@@ -838,6 +842,11 @@ Durable post-commit side effects have a separate boundary:
cancellation and bounded shutdown waiting, but it is not a durable queue and
must not replace an Outbox or persistent task table.
The transfer queue owner tracks actual queued and executing task objects separately
from `JobManager` display rows. Recovery reuses only an identical task/lease receipt;
when no actual owner remains, a stale display row is replaced and the recovered
work is queued. A display row alone must never acknowledge successful recovery.
Transfer durable admission follows the same ownership direction without using
the Outbox as an execution queue: `app/application/transfer/workflow.py` owns the typed
admission and versioned planning-checkpoint contracts, while

View File

@@ -767,3 +767,46 @@ async def test_async_broadcast_submission_is_registered_before_stop_snapshot(
await isolated_eventmanager.stop_async()
assert isolated_eventmanager._EventManager__async_handles == {}
@pytest.mark.parametrize("event_type", [EventType.TransferFailed, EventType.SubtitleTransferFailed,
EventType.AudioTransferFailed])
def test_outbox_retries_strict_failure_without_repeating_system_alert(
isolated_eventmanager, monkeypatch, event_type):
"""真实失败 handler 仍重试至死信,相同事件的系统错误提示只发送一次。"""
from unittest.mock import Mock
from app.application.outbox import ClaimedOutboxMessage, OutboxDispatcher
from app.runtime.event.errors import EventErrorPolicy
from app.startup.composition import outbox
isolated_eventmanager._EventManager__lifecycle_state = "running"
notify, emit = Mock(), Mock()
monkeypatch.setattr(isolated_eventmanager, "_EventManager__error_policy",
EventErrorPolicy(notifier=lambda: notify, emit_system_error=emit))
monkeypatch.setattr(outbox, "EventManager", lambda: isolated_eventmanager)
attempts = []
def failed_handler(event):
"""模拟稳定重现的插件处理器故障。"""
attempts.append(event.event_data["idempotency_key"])
raise RuntimeError("broken transfer consumer")
isolated_eventmanager.add_event_listener(event_type, failed_handler)
store = Mock()
event_key = f"{event_type.value}:task-1:1"
store.claim.side_effect = [
ClaimedOutboxMessage(1, event_key, event_type.value,
{"idempotency_key": event_key,
"transferinfo": {"success": False, "message": "failed"}}, 1, attempt)
for attempt in range(1, 6)
] + [None]
dispatcher = OutboxDispatcher(store, outbox.build_outbox_handlers())
for _ in range(5):
assert dispatcher.dispatch_one()
assert not dispatcher.dispatch_one()
assert attempts == [event_key] * 5
assert [call.kwargs["dead"] for call in store.retry.call_args_list] == [False] * 4 + [True]
store.complete.assert_not_called()
notify.assert_called_once()
emit.assert_called_once()

View File

@@ -288,3 +288,35 @@ def test_all_unmanaged_config_reload_classes_have_explicit_providers(
"SystemHelper",
"TransferChain",
}
def test_error_alert_dedup_preserves_distinct_events_handlers_and_errors():
"""只合并同一持久事件同一处理器的相同错误,不压制新故障或普通事件。"""
notify, emit = Mock(), Mock()
policy = EventErrorPolicy(notifier=lambda: notify, emit_system_error=emit)
for event_key, handler, error in [
("one", "handle", "broken"), ("one", "handle", "broken"),
("two", "handle", "broken"), ("one", "other", "broken"),
("one", "handle", "new error"), (None, "handle", "broken"),
(None, "handle", "broken"),
]:
policy.handle(event=Event(EventType.TransferFailed, {"idempotency_key": event_key}),
module_name="plugin", class_name="Handler", method_name=handler,
error=RuntimeError(error))
assert notify.call_count == 6
assert emit.call_count == 6
def test_durable_error_alert_cache_is_bounded(monkeypatch):
"""长时间运行只保留最近事件的提示记录,淘汰后允许重新提示。"""
from app.runtime.event import errors
monkeypatch.setattr(errors, "_MAX_REPORTED_ERRORS", 2)
notify = Mock()
policy = EventErrorPolicy(notifier=lambda: notify, emit_system_error=Mock())
for key in ("one", "two", "one", "three", "two"):
policy.handle(event=Event(EventType.TransferFailed, {"idempotency_key": key}),
module_name="plugin", class_name="Handler", method_name="handle",
error=RuntimeError("broken"))
assert notify.call_count == 4
assert len(policy._reported_errors) == 2

View File

@@ -553,3 +553,118 @@ def test_claimed_enqueue_failure_never_uses_unfenced_error_writer() -> None:
error="queue closed",
)
assert chain._owned_leases == {}
@pytest.mark.parametrize("state", ["waiting", "running", "failed", "completed"])
def test_recovery_requeues_orphaned_jobview_without_losing_group(state):
"""残留视图没有真实执行者时必须恢复入队,并保留同组其他文件。"""
import queue
from app.application.transfer.workflow import JobManager
from tests.test_transfer_job_manager import make_task
chain = _build_chain(MagicMock())
chain._queue = queue.Queue()
chain.jobview = JobManager()
chain._register_scrape_batch_task = MagicMock()
old_task = make_task(1)
sibling = make_task(2)
assert chain.jobview.add_task(old_task, state=state)
assert chain.jobview.add_task(sibling)
recovered = make_task(1)
recovered.bind_admission_task_id("task-1")
recovered.bind_execution_lease(owner_id="test-owner", lease_token="lease-task-1")
chain._owned_leases["task-1"] = ("lease-task-1", time.monotonic() + 60)
assert chain.put_to_queue(recovered)
assert chain._queue.qsize() == 1
assert chain._queue.get_nowait().task is recovered
assert chain.jobview.pending_total() == 2
chain._transfer_admissions.release_claim.assert_not_called()
assert chain._TransferChain__is_claimed_task_enqueued("task-1", "lease-task-1")
@pytest.mark.parametrize("running", [False, True])
def test_recovery_reuses_real_task_without_second_enqueue(running):
"""相同租约的真实任务已排队或被 worker 取走时,回放应幂等成功。"""
import queue
from app.application.transfer.workflow import JobManager
from tests.test_transfer_job_manager import make_task
chain = _build_chain(MagicMock())
chain._queue = queue.Queue()
chain.jobview = JobManager()
chain._register_scrape_batch_task = MagicMock()
task = make_task(1)
task.bind_admission_task_id("task-1")
task.bind_execution_lease(owner_id="test-owner", lease_token="lease-task-1")
chain._owned_leases["task-1"] = ("lease-task-1", time.monotonic() + 60)
assert chain.put_to_queue(task)
if running:
chain._queue.get_nowait()
chain._queued_lease_tokens.clear()
duplicate = task.model_copy(deep=True)
duplicate.fileitem.path += "/"
assert chain.put_to_queue(duplicate)
assert chain._queue.qsize() == (0 if running else 1)
chain._register_scrape_batch_task.assert_called_once_with(task)
chain._transfer_admissions.release_claim.assert_not_called()
if not running:
chain._queue.get_nowait()
chain._finish_queue_item(task)
assert not chain._resident_tasks
assert chain._queue.unfinished_tasks == 0
def test_recovery_does_not_reuse_a_different_execution_lease():
"""不能把旧执行者当作新租约已入队,也不能覆盖仍持有的真实任务。"""
import queue
from app.application.transfer.workflow import JobManager
from tests.test_transfer_job_manager import make_task
chain = _build_chain(MagicMock())
chain._queue = queue.Queue()
chain.jobview = JobManager()
chain._register_scrape_batch_task = MagicMock()
old = make_task(1)
old.bind_admission_task_id("task-1")
old.bind_execution_lease(owner_id="test-owner", lease_token="old")
chain._owned_leases["task-1"] = ("old", time.monotonic() + 60)
assert chain.put_to_queue(old)
recovered = make_task(1)
recovered.bind_admission_task_id("task-1")
recovered.bind_execution_lease(owner_id="test-owner", lease_token="new")
chain._owned_leases["task-1"] = ("new", time.monotonic() + 60)
assert chain.put_to_queue(recovered) is False
assert chain._queue.qsize() == 1
assert chain._queue.get_nowait().task is old
def test_normal_enqueue_is_visible_to_recovery_without_readmission():
"""普通准入任务也登记真实队列凭证,恢复同一租约不再次准入或入队。"""
import queue
from app.application.transfer.workflow import JobManager
from tests.test_transfer_job_manager import make_task
admissions = MagicMock()
chain = _build_chain(admissions)
chain._queue = queue.Queue()
chain.jobview = JobManager()
chain._register_scrape_batch_task = MagicMock()
task = make_task(1)
admission = _admission(task.fileitem.path)
admissions.admit.return_value = admission
admissions.claim_task.return_value = admission
task.bind_planning_input(admission.planning_input)
assert chain.put_to_queue(task)
assert chain.put_to_queue(task.model_copy())
assert chain._queue.qsize() == 1
admissions.admit.assert_called_once()
admissions.claim_task.assert_called_once()
chain._register_scrape_batch_task.assert_called_once()