mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-05-24 05:41:40 +00:00
fix: dedupe kanban notifier delivery claims
This commit is contained in:
parent
373c4d6647
commit
861ce7c0b6
5 changed files with 411 additions and 7 deletions
138
tests/gateway/test_kanban_notifier.py
Normal file
138
tests/gateway/test_kanban_notifier.py
Normal file
|
|
@ -0,0 +1,138 @@
|
|||
import asyncio
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from gateway.config import Platform
|
||||
from gateway.run import GatewayRunner
|
||||
from hermes_cli import kanban_db as kb
|
||||
|
||||
|
||||
class RecordingAdapter:
|
||||
def __init__(self):
|
||||
self.sent = []
|
||||
|
||||
async def send(self, chat_id, text, metadata=None):
|
||||
self.sent.append({"chat_id": chat_id, "text": text, "metadata": metadata or {}})
|
||||
|
||||
|
||||
class DisconnectedAdapters(dict):
|
||||
"""Expose a platform during collection, then simulate disconnect on get()."""
|
||||
|
||||
def get(self, key, default=None):
|
||||
return None
|
||||
|
||||
|
||||
async def _run_one_notifier_tick(monkeypatch, runner):
|
||||
real_sleep = asyncio.sleep
|
||||
|
||||
async def fake_sleep(delay):
|
||||
if delay == 5:
|
||||
return None
|
||||
runner._running = False
|
||||
await real_sleep(0)
|
||||
|
||||
monkeypatch.setattr(asyncio, "sleep", fake_sleep)
|
||||
await runner._kanban_notifier_watcher(interval=1)
|
||||
|
||||
|
||||
def _make_runner(adapter):
|
||||
runner = GatewayRunner.__new__(GatewayRunner)
|
||||
runner._running = True
|
||||
runner.adapters = {Platform.TELEGRAM: adapter}
|
||||
runner._kanban_sub_fail_counts = {}
|
||||
return runner
|
||||
|
||||
|
||||
def _create_completed_subscription(summary="done once"):
|
||||
conn = kb.connect()
|
||||
try:
|
||||
tid = kb.create_task(conn, title="notify once", assignee="worker")
|
||||
kb.add_notify_sub(conn, task_id=tid, platform="telegram", chat_id="chat-1")
|
||||
kb.complete_task(conn, tid, summary=summary)
|
||||
return tid
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _unseen_terminal_events(tid):
|
||||
conn = kb.connect()
|
||||
try:
|
||||
_, events = kb.unseen_events_for_sub(
|
||||
conn,
|
||||
task_id=tid,
|
||||
platform="telegram",
|
||||
chat_id="chat-1",
|
||||
kinds=["completed", "blocked", "gave_up", "crashed", "timed_out"],
|
||||
)
|
||||
return events
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def test_kanban_notifier_dedupes_board_slugs_pointing_to_same_db(tmp_path, monkeypatch):
|
||||
db_path = tmp_path / "shared-kanban.db"
|
||||
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
|
||||
kb.init_db()
|
||||
kb.write_board_metadata("alias-a", name="Alias A")
|
||||
kb.write_board_metadata("alias-b", name="Alias B")
|
||||
|
||||
tid = _create_completed_subscription()
|
||||
|
||||
adapter = RecordingAdapter()
|
||||
runner = _make_runner(adapter)
|
||||
|
||||
asyncio.run(_run_one_notifier_tick(monkeypatch, runner))
|
||||
|
||||
assert len(adapter.sent) == 1
|
||||
assert "Kanban" in adapter.sent[0]["text"]
|
||||
assert tid in adapter.sent[0]["text"]
|
||||
|
||||
|
||||
def test_kanban_notifier_claim_prevents_second_watcher_send(tmp_path, monkeypatch):
|
||||
db_path = tmp_path / "single-owner.db"
|
||||
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
|
||||
kb.init_db()
|
||||
|
||||
tid = _create_completed_subscription()
|
||||
|
||||
adapter1 = RecordingAdapter()
|
||||
adapter2 = RecordingAdapter()
|
||||
|
||||
asyncio.run(_run_one_notifier_tick(monkeypatch, _make_runner(adapter1)))
|
||||
asyncio.run(_run_one_notifier_tick(monkeypatch, _make_runner(adapter2)))
|
||||
|
||||
assert len(adapter1.sent) == 1
|
||||
assert adapter2.sent == []
|
||||
|
||||
|
||||
def test_kanban_notifier_rewinds_claim_if_adapter_disconnects(tmp_path, monkeypatch):
|
||||
db_path = tmp_path / "adapter-disconnect.db"
|
||||
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
|
||||
kb.init_db()
|
||||
tid = _create_completed_subscription()
|
||||
|
||||
runner = GatewayRunner.__new__(GatewayRunner)
|
||||
runner._running = True
|
||||
runner.adapters = DisconnectedAdapters({Platform.TELEGRAM: RecordingAdapter()})
|
||||
runner._kanban_sub_fail_counts = {}
|
||||
|
||||
asyncio.run(_run_one_notifier_tick(monkeypatch, runner))
|
||||
|
||||
assert [ev.kind for ev in _unseen_terminal_events(tid)] == ["completed"]
|
||||
|
||||
|
||||
def test_kanban_db_path_is_test_isolated_from_real_home():
|
||||
hermes_home = Path(kb.kanban_home())
|
||||
production_db = Path.home() / ".hermes" / "kanban.db"
|
||||
assert kb.kanban_db_path().resolve() != production_db.resolve()
|
||||
|
||||
conn = kb.connect()
|
||||
try:
|
||||
tid = kb.create_task(conn, title="x", assignee="worker")
|
||||
kb.add_notify_sub(conn, task_id=tid, platform="telegram", chat_id="chat-1")
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
assert kb.kanban_db_path().resolve().is_relative_to(hermes_home.resolve())
|
||||
assert kb.kanban_db_path().resolve() != production_db.resolve()
|
||||
Loading…
Add table
Add a link
Reference in a new issue