From 590af946935fedcbb1ea9f363b1af283cc9df4b3 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 10:41:07 +0800 Subject: [PATCH] fix(outbox): distinguish repeated downloads --- app/application/chain/durable_events.py | 9 ++--- app/db/adapters/chain.py | 20 ++++++---- tests/test_chain_durable_events.py | 52 +++++++++++++++---------- 3 files changed, 46 insertions(+), 35 deletions(-) diff --git a/app/application/chain/durable_events.py b/app/application/chain/durable_events.py index d52c7f4c8..d1d61cf3e 100644 --- a/app/application/chain/durable_events.py +++ b/app/application/chain/durable_events.py @@ -55,12 +55,9 @@ class TransferHistoryRef: src_fileitem: dict[str, Any] | None -def download_added_event_key(payload: dict[str, Any]) -> str: - """由下载器与任务 hash 构造重试期间稳定的 DownloadAdded 幂等键。""" - return ( - f"download.added:{payload.get('downloader') or 'unknown'}:" - f"{payload.get('hash') or 'unknown'}:v1" - ) +def download_added_event_key(history_id: int) -> str: + """由下载历史 ID 构造本次下载事实稳定的 DownloadAdded 幂等键。""" + return f"download.added:{history_id}:v1" def transfer_result_event_key(topic: str, history_id: int) -> str: diff --git a/app/db/adapters/chain.py b/app/db/adapters/chain.py index c8b147afd..d3175f17c 100644 --- a/app/db/adapters/chain.py +++ b/app/db/adapters/chain.py @@ -80,21 +80,25 @@ class TransactionalChainDurableEventWriter(ChainDurableEventWriter): unit_of_work=SqlAlchemyUnitOfWork(session), outbox=outbox, ) - event_key = download_added_event_key(event_payload) - event_payload["idempotency_key"] = event_key - - def stage_business() -> None: + def stage_business() -> int: """在同一事务暂存下载历史和可选文件清单。""" - repository.stage_add(history_payload) + history = repository.stage_add(history_payload) if file_payloads: repository.stage_add_files(file_payloads) + return int(history.id) - command.execute( - intent=OutboxIntent( + def build_intent(history_id: int) -> OutboxIntent: + """历史 ID 确定后构造本次下载事实的稳定事件键。""" + event_key = download_added_event_key(history_id) + event_payload["idempotency_key"] = event_key + return OutboxIntent( event_key=event_key, topic=DOWNLOAD_ADDED_TOPIC, payload=snapshot_download_added(event_payload), - ), + ) + + command.execute( + intent=build_intent, stage_business=stage_business, after_commit=after_commit, publish=lambda: publish(event_payload), diff --git a/tests/test_chain_durable_events.py b/tests/test_chain_durable_events.py index 357d07bc3..8adc632a7 100644 --- a/tests/test_chain_durable_events.py +++ b/tests/test_chain_durable_events.py @@ -4,7 +4,6 @@ import json import pytest from sqlalchemy import create_engine, select -from sqlalchemy.exc import IntegrityError from sqlalchemy.orm import sessionmaker from app.application.chain.durable_events import ( @@ -145,32 +144,43 @@ def test_download_history_and_event_intent_share_one_transaction(): assert history.download_hash == "hash-2" assert download_file.fullpath == "/downloads/Demo.mkv" assert outbox.status == "completed" - assert outbox.event_key == "download.added:qb:hash-2:v1" + assert outbox.event_key == f"download.added:{history.id}:v1" assert [call if isinstance(call, str) else call[0] for call in calls] == [ "after_commit", "event", ] - with pytest.raises(IntegrityError): - writer.download_added( - history_payload={ - "path": "/downloads/duplicate.mkv", - "type": MediaType.MOVIE.value, - "title": "Duplicate", - "download_hash": "hash-2", - }, - file_payloads=[], - event_payload={ - "hash": "hash-2", - "context": context, - "downloader": "qb", - "episodes": [], - }, - after_commit=lambda: None, - publish=lambda _payload: None, - ) + writer.download_added( + history_payload={ + "path": "/downloads/duplicate.mkv", + "type": MediaType.MOVIE.value, + "title": "Duplicate", + "download_hash": "hash-2", + }, + file_payloads=[], + event_payload={ + "hash": "hash-2", + "context": context, + "downloader": "qb", + "episodes": [], + }, + after_commit=lambda: None, + publish=lambda _payload: None, + ) with factory() as session: - assert len(session.execute(select(DownloadHistory)).scalars().all()) == 1 + histories = session.execute( + select(DownloadHistory).order_by(DownloadHistory.id) + ).scalars().all() + outboxes = session.execute( + select(OutboxMessage).order_by(OutboxMessage.id) + ).scalars().all() + assert len(histories) == 2 + assert len(outboxes) == 2 + assert all(message.status == "completed" for message in outboxes) + assert [message.event_key for message in outboxes] == [ + f"download.added:{history.id}:v1" for history in histories + ] + assert outboxes[0].event_key != outboxes[1].event_key def test_transfer_event_failure_leaves_committed_intent_pending():