fix(outbox): distinguish repeated downloads

This commit is contained in:
jxxghp
2026-08-24 10:41:14 +08:00
parent 473c0588cc
commit 590af94693
3 changed files with 46 additions and 35 deletions
+3 -6
View File
@@ -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:
+12 -8
View File
@@ -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),
+31 -21
View File
@@ -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():