From b941233138e591ca15505f961b8ec24f61c553b6 Mon Sep 17 00:00:00 2001 From: chengyongru Date: Thu, 2 Jul 2026 10:56:15 +0800 Subject: [PATCH] fix(trigger): clean up deleted trigger deliveries --- nanobot/triggers/local_store.py | 31 ++++++++++++ tests/channels/test_websocket_http_routes.py | 3 ++ tests/triggers/test_local_triggers.py | 53 ++++++++++++++++++++ 3 files changed, 87 insertions(+) diff --git a/nanobot/triggers/local_store.py b/nanobot/triggers/local_store.py index 41b5e5d0..35a5629b 100644 --- a/nanobot/triggers/local_store.py +++ b/nanobot/triggers/local_store.py @@ -143,6 +143,7 @@ class LocalTriggerStore: def delete(self, trigger_id: str) -> bool: """Delete a trigger by ID.""" + trigger_id = trigger_id.strip() self._ensure_dirs() with self._lock: triggers = self._load_triggers_unlocked() @@ -150,6 +151,7 @@ class LocalTriggerStore: if len(remaining) == len(triggers): return False self._save_triggers_unlocked(remaining) + self._delete_delivery_files_for_trigger_unlocked(trigger_id) return True def enqueue(self, trigger_id: str, content: str) -> TriggerDelivery: @@ -324,6 +326,35 @@ class LocalTriggerStore: delivery.path.unlink(missing_ok=True) return True + def _delete_delivery_files_for_trigger_unlocked(self, trigger_id: str) -> None: + for directory in (self.inbox_dir, self.processing_dir, self.failed_dir): + for path in directory.iterdir(): + if not path.is_file(): + continue + if self._delivery_file_trigger_id(path) != trigger_id: + continue + try: + path.unlink(missing_ok=True) + except OSError as exc: + logger.warning( + "Trigger: failed to delete delivery file {} for deleted trigger {}: {}", + path, + trigger_id, + exc, + ) + + @staticmethod + def _delivery_file_trigger_id(path: Path) -> str | None: + try: + data = json.loads(path.read_text(encoding="utf-8")) + except Exception: + return None + raw = data.get("delivery", data) if isinstance(data, dict) else None + if not isinstance(raw, dict): + return None + trigger_id = raw.get("triggerId", raw.get("trigger_id", "")) + return str(trigger_id) if trigger_id else None + @staticmethod def _atomic_write(path: Path, content: str) -> None: path.parent.mkdir(parents=True, exist_ok=True) diff --git a/tests/channels/test_websocket_http_routes.py b/tests/channels/test_websocket_http_routes.py index 2df8d577..c967a5d6 100644 --- a/tests/channels/test_websocket_http_routes.py +++ b/tests/channels/test_websocket_http_routes.py @@ -1150,6 +1150,8 @@ async def test_webui_automations_route_manages_local_triggers( chat_id="abc", session_key="websocket:abc", ) + delivery = trigger_store.enqueue(trigger.id, "Review queued PR") + assert delivery.path is not None channel = _ch( bus, session_manager=_seed_session(tmp_path, key="websocket:abc"), @@ -1216,6 +1218,7 @@ async def test_webui_automations_route_manages_local_triggers( ) assert deleted.status_code == 200 assert trigger_store.get(trigger.id) is None + assert not delivery.path.exists() finally: await channel.stop() await server_task diff --git a/tests/triggers/test_local_triggers.py b/tests/triggers/test_local_triggers.py index c5f6c921..5541ecfe 100644 --- a/tests/triggers/test_local_triggers.py +++ b/tests/triggers/test_local_triggers.py @@ -1,6 +1,7 @@ from __future__ import annotations import asyncio +import json from contextlib import suppress from pathlib import Path @@ -13,6 +14,25 @@ from nanobot.triggers.local_store import LocalTriggerStore, TriggerDisabledError from nanobot.webui.metadata import WEBUI_MESSAGE_SOURCE_METADATA_KEY, WEBUI_TURN_METADATA_KEY +def _write_delivery_file(path: Path, *, trigger_id: str, delivery_id: str) -> None: + path.write_text( + json.dumps( + { + "version": 1, + "delivery": { + "id": delivery_id, + "triggerId": trigger_id, + "content": "queued", + "createdAtMs": 1, + "attempts": 0, + "lastError": None, + }, + } + ), + encoding="utf-8", + ) + + def test_trigger_store_allows_multiple_triggers_per_session(tmp_path: Path) -> None: store = LocalTriggerStore(tmp_path) @@ -50,6 +70,39 @@ def test_enqueue_rejects_disabled_trigger(tmp_path: Path) -> None: store.enqueue(trigger.id, "Review PR #4502") +def test_delete_removes_delivery_files_for_trigger(tmp_path: Path) -> None: + store = LocalTriggerStore(tmp_path) + trigger = store.create( + name="PR review", + channel="websocket", + chat_id="chat-1", + session_key="websocket:chat-1", + ) + other = store.create( + name="CI summary", + channel="websocket", + chat_id="chat-2", + session_key="websocket:chat-2", + ) + inbox = store.inbox_dir / "1-tdl_inbox.json" + processing = store.processing_dir / "2-tdl_processing.json" + failed = store.failed_dir / "3-tdl_failed.json" + other_inbox = store.inbox_dir / "4-tdl_other.json" + _write_delivery_file(inbox, trigger_id=trigger.id, delivery_id="tdl_inbox") + _write_delivery_file(processing, trigger_id=trigger.id, delivery_id="tdl_processing") + _write_delivery_file(failed, trigger_id=trigger.id, delivery_id="tdl_failed") + _write_delivery_file(other_inbox, trigger_id=other.id, delivery_id="tdl_other") + + assert store.delete(trigger.id) is True + + assert store.get(trigger.id) is None + assert not inbox.exists() + assert not processing.exists() + assert not failed.exists() + assert other_inbox.exists() + assert store.get(other.id) is not None + + def test_recover_processing_deliveries_requeues_claimed_delivery(tmp_path: Path) -> None: store = LocalTriggerStore(tmp_path) trigger = store.create(