mirror of
https://github.com/HKUDS/nanobot.git
synced 2026-08-07 01:48:53 +00:00
fix(trigger): clean up deleted trigger deliveries
This commit is contained in:
parent
acb0e853ff
commit
b941233138
@ -143,6 +143,7 @@ class LocalTriggerStore:
|
|||||||
|
|
||||||
def delete(self, trigger_id: str) -> bool:
|
def delete(self, trigger_id: str) -> bool:
|
||||||
"""Delete a trigger by ID."""
|
"""Delete a trigger by ID."""
|
||||||
|
trigger_id = trigger_id.strip()
|
||||||
self._ensure_dirs()
|
self._ensure_dirs()
|
||||||
with self._lock:
|
with self._lock:
|
||||||
triggers = self._load_triggers_unlocked()
|
triggers = self._load_triggers_unlocked()
|
||||||
@ -150,6 +151,7 @@ class LocalTriggerStore:
|
|||||||
if len(remaining) == len(triggers):
|
if len(remaining) == len(triggers):
|
||||||
return False
|
return False
|
||||||
self._save_triggers_unlocked(remaining)
|
self._save_triggers_unlocked(remaining)
|
||||||
|
self._delete_delivery_files_for_trigger_unlocked(trigger_id)
|
||||||
return True
|
return True
|
||||||
|
|
||||||
def enqueue(self, trigger_id: str, content: str) -> TriggerDelivery:
|
def enqueue(self, trigger_id: str, content: str) -> TriggerDelivery:
|
||||||
@ -324,6 +326,35 @@ class LocalTriggerStore:
|
|||||||
delivery.path.unlink(missing_ok=True)
|
delivery.path.unlink(missing_ok=True)
|
||||||
return 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
|
@staticmethod
|
||||||
def _atomic_write(path: Path, content: str) -> None:
|
def _atomic_write(path: Path, content: str) -> None:
|
||||||
path.parent.mkdir(parents=True, exist_ok=True)
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
|||||||
@ -1150,6 +1150,8 @@ async def test_webui_automations_route_manages_local_triggers(
|
|||||||
chat_id="abc",
|
chat_id="abc",
|
||||||
session_key="websocket:abc",
|
session_key="websocket:abc",
|
||||||
)
|
)
|
||||||
|
delivery = trigger_store.enqueue(trigger.id, "Review queued PR")
|
||||||
|
assert delivery.path is not None
|
||||||
channel = _ch(
|
channel = _ch(
|
||||||
bus,
|
bus,
|
||||||
session_manager=_seed_session(tmp_path, key="websocket:abc"),
|
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 deleted.status_code == 200
|
||||||
assert trigger_store.get(trigger.id) is None
|
assert trigger_store.get(trigger.id) is None
|
||||||
|
assert not delivery.path.exists()
|
||||||
finally:
|
finally:
|
||||||
await channel.stop()
|
await channel.stop()
|
||||||
await server_task
|
await server_task
|
||||||
|
|||||||
@ -1,6 +1,7 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import json
|
||||||
from contextlib import suppress
|
from contextlib import suppress
|
||||||
from pathlib import Path
|
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
|
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:
|
def test_trigger_store_allows_multiple_triggers_per_session(tmp_path: Path) -> None:
|
||||||
store = LocalTriggerStore(tmp_path)
|
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")
|
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:
|
def test_recover_processing_deliveries_requeues_claimed_delivery(tmp_path: Path) -> None:
|
||||||
store = LocalTriggerStore(tmp_path)
|
store = LocalTriggerStore(tmp_path)
|
||||||
trigger = store.create(
|
trigger = store.create(
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user