diff --git a/app/application/subscription/delete.py b/app/application/subscription/delete.py index 4ba3a59fc..5a8c9334f 100644 --- a/app/application/subscription/delete.py +++ b/app/application/subscription/delete.py @@ -98,6 +98,7 @@ class DeleteSubscribeCommand: candidate.event_payload, ) event_key = event_payload["idempotency_key"] + report_key = f"{event_key}:report" try: if self._outbox: await self._outbox.stage( @@ -108,6 +109,17 @@ class DeleteSubscribeCommand: ), datetime.now(timezone.utc), ) + await self._outbox.stage( + OutboxIntent( + event_key=report_key, + topic="subscribe.deleted.report", + payload={ + "idempotency_key": report_key, + "subscribe_info": dict(candidate.event_payload), + }, + ), + datetime.now(timezone.utc), + ) await self._unit_of_work.commit() except Exception: await self._unit_of_work.rollback() @@ -121,7 +133,14 @@ class DeleteSubscribeCommand: ) # 上报适配器会自行白名单过滤公开字段;传完整删除前快照可保留音乐实体维度, # 避免 Agent 与 API 入口收敛后丢失 music_type / total_tracks。 - self._report_deleted(dict(candidate.event_payload)) + report_result = self._report_deleted(dict(candidate.event_payload)) + if report_result is False: + raise RuntimeError("订阅删除统计上报未确认") + if self._outbox: + await self._outbox.complete_by_event_key( + report_key, + datetime.now(timezone.utc), + ) return True @staticmethod diff --git a/app/startup/modules_initializer.py b/app/startup/modules_initializer.py index b24df9b36..bb51fde56 100644 --- a/app/startup/modules_initializer.py +++ b/app/startup/modules_initializer.py @@ -235,6 +235,13 @@ def configure_runtime_data_providers() -> None: def _build_outbox_dispatcher() -> OutboxDispatcher: """创建一次恢复批次独占的 Session、Repository 和事件 handler。""" + def dispatch_subscribe_deleted_report(message) -> None: + """重放订阅删除统计;未确认时抛错以进入有限重试。""" + if not MoviePilotServerHelper.sub_done( + message.payload.get("subscribe_info") or {} + ): + raise RuntimeError("订阅删除统计上报未确认") + session = SessionFactory() return OutboxDispatcher( repository=SqlAlchemyOutboxRepository(session), @@ -251,6 +258,7 @@ def _build_outbox_dispatcher() -> OutboxDispatcher: EventType.SubscribeDeleted, message.payload, ), + "subscribe.deleted.report": dispatch_subscribe_deleted_report, "download.added": lambda message: EventManager().send_event( EventType.DownloadAdded, restore_download_added(message.payload), diff --git a/tests/test_subscribe_delete_command.py b/tests/test_subscribe_delete_command.py index faf1070fd..0c49827a8 100644 --- a/tests/test_subscribe_delete_command.py +++ b/tests/test_subscribe_delete_command.py @@ -105,6 +105,7 @@ def _command( calls.append(("report", payload)) if report_error: raise report_error + return True return DeleteSubscribeCommand( repository=_Repository(candidate, calls), @@ -224,15 +225,20 @@ async def test_delete_stages_outbox_before_commit_and_completes_after_event(): "get", "delete", "outbox_stage", + "outbox_stage", "commit", "event", "outbox_complete", "report", + "outbox_complete", ] intent = calls[2][1] assert intent.topic == "subscribe.deleted" - assert intent.event_key == calls[4][2]["idempotency_key"] - assert calls[5][1] == intent.event_key + report_intent = calls[3][1] + assert report_intent.topic == "subscribe.deleted.report" + assert intent.event_key == calls[5][2]["idempotency_key"] + assert calls[6][1] == intent.event_key + assert calls[8][1] == report_intent.event_key @pytest.mark.asyncio