fix(kanban): deliver notifications from non-dispatch gateways

This commit is contained in:
Gille
2026-07-31 13:07:28 -06:00
committed by Teknium
parent e08870f62c
commit 991f5f1e9e
9 changed files with 271 additions and 83 deletions
+11 -10
View File
@@ -10,11 +10,11 @@ Only one gateway owns the kanban dispatcher. The owning gateway keeps
`kanban.dispatch_in_gateway: true` (the default); every other gateway sets it
to `false`.
**Why this matters:** a gateway with `dispatch_in_gateway: true` opens per-board
SQLite connections for both the dispatcher and the notifier watcher. Multiple
gateways doing this concurrently multiplies the open file descriptors on each
`kanban.db` and amplifies WAL `-shm` reader contention. Gating both paths on the
same flag means exactly one process touches the kanban DBs.
**Why this matters:** dispatching is single-owner so multiple gateways do not
race to spawn the same work. Notification delivery is profile-owned instead:
each gateway polls only subscriptions for profiles whose platform adapters it
hosts. The atomic event claim prevents duplicate delivery across watcher
processes.
## Configuration
@@ -30,10 +30,11 @@ Or set the env var: `HERMES_KANBAN_DISPATCH_IN_GATEWAY=false`
## What each gateway does
| Gateway role | dispatch_in_gateway | Opens per-board DBs? | Runs dispatcher + notifier? |
|---|---|---|---|
| default (dispatch owner) | true (default) | yes | yes |
| writer, admin, coder, etc. | false | no | no |
| Gateway role | dispatch_in_gateway | Opens subscribed board DBs? | Dispatcher | Notifier |
|---|---|---|---|---|
| default (confirmed dispatch-lock owner) | true (default) | yes | yes | owned profiles + legacy unstamped subscriptions |
| writer, admin, coder, etc. | false | yes, when the profile has subscriptions | no | that gateway's owned profiles |
Non-dispatch gateways still deliver messages for their own platform adapters
(Telegram, Discord, etc.) — they just don't poll kanban boards.
(Telegram, Discord, etc.). They do not dispatch tasks, and they skip boards
that have no subscriptions owned by their profiles.
+39 -35
View File
@@ -112,6 +112,16 @@ def _release_singleton_lock(handle) -> None:
class GatewayKanbanWatchersMixin:
"""Kanban watcher / notifier / dispatcher loops for GatewayRunner."""
def _owns_kanban_dispatcher_lock(self) -> bool:
"""Return whether this gateway currently owns the singleton lock."""
return getattr(self, "_kanban_dispatcher_lock_handle", None) is not None
def _release_kanban_dispatcher_lock(self) -> None:
"""Clear notifier-visible ownership before releasing the OS lock."""
handle = getattr(self, "_kanban_dispatcher_lock_handle", None)
self._kanban_dispatcher_lock_handle = None
_release_singleton_lock(handle)
async def _kanban_notifier_watcher(self, interval: float = 5.0) -> None:
"""Poll ``kanban_notify_subs`` and deliver terminal events to users.
@@ -127,34 +137,15 @@ class GatewayKanbanWatchersMixin:
WAL lock. Failures in one tick don't stop subsequent ticks.
**Multi-board:** iterates every board discovered on disk per
tick. Subscriptions live inside each board's own DB and cannot
cross boards, so delivery semantics are unchanged — this is
purely a fan-out of the single-DB poll.
tick. Each gateway polls only subscriptions owned by profiles whose
adapters it hosts. The dispatch-owning gateway also handles legacy
subscriptions without a profile stamp.
"""
# Gate: only the dispatch-owning gateway opens kanban DBs for notifier polling.
# Non-dispatch gateways have no subscriptions to deliver — all kanban state lives
# in the dispatch owner's per-board DBs. This prevents N-gateway -shm contention.
# TODO: gate per-board when per-board dispatcher_owner tracking lands.
try:
from hermes_cli.config import load_config as _load_config
except Exception:
logger.warning("kanban notifier: config loader unavailable; disabled")
return
env_override = os.environ.get("HERMES_KANBAN_DISPATCH_IN_GATEWAY", "").strip().lower()
if env_override in {"0", "false", "no", "off"}:
logger.info("kanban notifier: disabled via HERMES_KANBAN_DISPATCH_IN_GATEWAY env")
return
try:
cfg = _load_config()
except Exception as exc:
logger.warning("kanban notifier: cannot load config (%s); disabled", exc)
return
kanban_cfg = cfg.get("kanban", {}) if isinstance(cfg, dict) else {}
if not kanban_cfg.get("dispatch_in_gateway", True):
logger.info(
"kanban notifier: disabled via config kanban.dispatch_in_gateway=false"
)
return
# Dispatch and delivery have separate ownership. A deployment may run
# one dispatcher while each profile has its own gateway credentials;
# those adapter-owning gateways must still poll and deliver their own
# subscriptions. Legacy rows without a notifier_profile are visible
# only while this process holds the actual singleton dispatcher lock.
from gateway.config import Platform as _Platform
try:
from hermes_cli import kanban_db as _kb
@@ -204,6 +195,13 @@ class GatewayKanbanWatchersMixin:
try:
def _collect():
deliveries: list[dict] = []
include_unowned = self._owns_kanban_dispatcher_lock()
notifier_profiles = {notifier_profile}
notifier_profiles.update(
str(profile).strip()
for profile in getattr(self, "_profile_adapters", {})
if str(profile).strip()
)
active_platforms = {
getattr(platform, "value", str(platform)).lower()
for platform in self.adapters.keys()
@@ -262,10 +260,14 @@ class GatewayKanbanWatchersMixin:
# checkpoint traffic) is exactly the per-tick cost
# this skip avoids.
try:
if _kb.count_notify_subs(board=slug) == 0:
if _kb.count_notify_subs(
board=slug,
notifier_profiles=notifier_profiles,
include_unowned=include_unowned,
) == 0:
logger.debug(
"kanban notifier: board %s has no subscriptions; skipping open",
slug,
"kanban notifier: board %s has no subscriptions owned by %s; skipping open",
slug, sorted(notifier_profiles),
)
continue
except Exception as exc:
@@ -292,7 +294,11 @@ class GatewayKanbanWatchersMixin:
# a legacy DB. `_add_column_if_missing` now
# tolerates that race, but we still skip the
# redundant call to avoid the wasted work.
subs = _kb.list_notify_subs(conn)
subs = _kb.list_notify_subs(
conn,
notifier_profiles=notifier_profiles,
include_unowned=include_unowned,
)
if not subs:
logger.debug("kanban notifier: board %s has no subscriptions", slug)
for sub in subs:
@@ -1472,8 +1478,7 @@ class GatewayKanbanWatchersMixin:
last_warn_at = now
except asyncio.CancelledError:
logger.debug("kanban dispatcher: cancelled")
_release_singleton_lock(self._kanban_dispatcher_lock_handle)
self._kanban_dispatcher_lock_handle = None
self._release_kanban_dispatcher_lock()
raise
except Exception:
logger.exception("kanban dispatcher: unexpected watcher error")
@@ -1485,5 +1490,4 @@ class GatewayKanbanWatchersMixin:
await asyncio.sleep(min(1.0, interval - slept))
slept += 1.0
_release_singleton_lock(self._kanban_dispatcher_lock_handle)
self._kanban_dispatcher_lock_handle = None
self._release_kanban_dispatcher_lock()
+3 -3
View File
@@ -11002,9 +11002,9 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
# Start background session expiry watcher to finalize expired sessions
self._spawn_supervised(self._session_expiry_watcher, "session_expiry_watcher")
# Start background kanban notifier — delivers `completed`, `blocked`,
# `spawn_auto_blocked`, and `crashed` events to gateway subscribers
# so human-in-the-loop workflows hear back without polling.
# Start background kanban notifier — each gateway delivers events for
# subscriptions owned by the profiles whose adapters it hosts, even
# when another gateway owns the single dispatcher.
self._spawn_supervised(self._kanban_notifier_watcher, "kanban_notifier_watcher")
# Start background kanban dispatcher — spawns workers for ready
+69 -11
View File
@@ -9598,15 +9598,65 @@ def add_notify_sub(
)
def _notify_profile_filter(
notifier_profiles: Optional[Iterable[str]],
*,
include_unowned: bool,
) -> tuple[str, list[str]]:
"""Build an optional SQL predicate for notification profile ownership."""
if notifier_profiles is None:
return "", []
profiles = sorted(
{
str(profile).strip()
for profile in notifier_profiles
if str(profile).strip()
}
)
clauses: list[str] = []
params: list[str] = []
if profiles:
clauses.append(
"notifier_profile IN (" + ",".join("?" for _ in profiles) + ")"
)
params.extend(profiles)
if include_unowned:
clauses.append("notifier_profile IS NULL OR notifier_profile = ''")
if not clauses:
return "0", []
return "(" + ") OR (".join(clauses) + ")", params
def list_notify_subs(
conn: sqlite3.Connection, task_id: Optional[str] = None,
conn: sqlite3.Connection,
task_id: Optional[str] = None,
*,
notifier_profiles: Optional[Iterable[str]] = None,
include_unowned: bool = False,
) -> list[dict]:
"""List subscriptions, optionally restricted to notifier profile owners.
Passing no ``notifier_profiles`` preserves the historical all-subscriptions
result. Gateway notifier processes pass the profiles whose adapters they
own so they cannot claim another gateway's events. ``include_unowned`` is
used by the dispatch owner for legacy rows created before profile stamping.
"""
owner_where, owner_params = _notify_profile_filter(
notifier_profiles, include_unowned=include_unowned,
)
where: list[str] = []
params: list[Any] = []
if task_id is not None:
rows = conn.execute(
"SELECT * FROM kanban_notify_subs WHERE task_id = ?", (task_id,),
).fetchall()
else:
rows = conn.execute("SELECT * FROM kanban_notify_subs").fetchall()
where.append("task_id = ?")
params.append(task_id)
if owner_where:
where.append(owner_where)
params.extend(owner_params)
sql = "SELECT * FROM kanban_notify_subs"
if where:
sql += " WHERE " + " AND ".join(f"({clause})" for clause in where)
rows = conn.execute(sql, params).fetchall()
out: list[dict] = []
for row in rows:
item = dict(row)
@@ -9622,6 +9672,8 @@ def count_notify_subs(
db_path: Optional[Path] = None,
*,
board: Optional[str] = None,
notifier_profiles: Optional[Iterable[str]] = None,
include_unowned: bool = False,
) -> int:
"""Count ``kanban_notify_subs`` rows via a read-only connection.
@@ -9633,8 +9685,10 @@ def count_notify_subs(
write table content). Rows in a not-yet-checkpointed WAL are
visible, so a freshly added subscription is never missed. A missing
DB, or a legacy DB that predates the subscriptions table, counts as
zero. Path resolution matches :func:`connect` (explicit ``db_path``,
else ``board`` via :func:`kanban_db_path`). Raises
zero. When ``notifier_profiles`` is supplied, only subscriptions owned
by those profiles are counted; ``include_unowned`` also includes legacy
rows without an owner stamp. Path resolution matches :func:`connect`
(explicit ``db_path``, else ``board`` via :func:`kanban_db_path`). Raises
:class:`sqlite3.Error` when the DB exists but cannot be read
(locked, corrupt); callers choose their own fallback.
"""
@@ -9644,9 +9698,13 @@ def count_notify_subs(
conn = sqlite3.connect(path.resolve().as_uri() + "?mode=ro", uri=True)
try:
try:
row = conn.execute(
"SELECT COUNT(*) FROM kanban_notify_subs"
).fetchone()
owner_where, owner_params = _notify_profile_filter(
notifier_profiles, include_unowned=include_unowned,
)
sql = "SELECT COUNT(*) FROM kanban_notify_subs"
if owner_where:
sql += " WHERE " + owner_where
row = conn.execute(sql, owner_params).fetchone()
except sqlite3.OperationalError as exc:
if "no such table" in str(exc).lower():
return 0
+108 -2
View File
@@ -4,6 +4,10 @@ from pathlib import Path
from gateway.config import Platform
from gateway.kanban_watchers import (
_acquire_singleton_lock,
_release_singleton_lock,
)
from gateway.run import GatewayRunner
from hermes_cli import kanban_db as kb
@@ -45,6 +49,9 @@ def _make_runner(adapter):
runner._running = True
runner.adapters = {Platform.TELEGRAM: adapter}
runner._kanban_sub_fail_counts = {}
# Most tests model the default gateway after its dispatcher acquired the
# singleton lock. Tests for startup or non-owner gateways clear this.
runner._kanban_dispatcher_lock_handle = object()
return runner
@@ -160,6 +167,105 @@ def test_active_named_profile_subscription_is_delivered(tmp_path, monkeypatch):
assert "blocked" in message
def test_non_dispatch_gateway_claims_only_its_profile_subscriptions(
tmp_path, monkeypatch,
):
"""A profile gateway delivers its events while another gateway dispatches."""
db_path = tmp_path / "cross-profile-notifier.db"
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
kb.init_db()
conn = kb.connect()
try:
foreign_tid = kb.create_task(
conn, title="default-owned", assignee="worker",
)
kb.add_notify_sub(
conn,
task_id=foreign_tid,
platform="telegram",
chat_id="default-chat",
notifier_profile="default",
)
kb.complete_task(conn, foreign_tid, summary="default done")
owned_tid = kb.create_task(
conn, title="writer-owned", assignee="worker",
)
kb.add_notify_sub(
conn,
task_id=owned_tid,
platform="telegram",
chat_id="writer-chat",
notifier_profile="writer",
)
kb.complete_task(conn, owned_tid, summary="writer done")
finally:
conn.close()
adapter = RecordingAdapter()
runner = _make_runner(adapter)
runner._active_profile_name = lambda: "writer"
runner._kanban_dispatcher_lock_handle = None
asyncio.run(_run_one_notifier_tick(monkeypatch, runner))
assert [delivery["chat_id"] for delivery in adapter.sent] == ["writer-chat"]
assert owned_tid in adapter.sent[0]["text"]
assert len(_unseen_terminal_events_for(foreign_tid, "default-chat")) == 1
def test_legacy_subscription_requires_confirmed_dispatcher_lock_owner(
tmp_path, monkeypatch,
):
"""Startup and lock-losing gateways cannot claim legacy notifications."""
db_path = tmp_path / "legacy-lock-owner.db"
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
kb.init_db()
conn = kb.connect()
try:
task_id = kb.create_task(conn, title="legacy", assignee="worker")
kb.add_notify_sub(
conn,
task_id=task_id,
platform="telegram",
chat_id="legacy-chat",
)
kb.complete_task(conn, task_id, summary="legacy done")
finally:
conn.close()
startup_adapter = RecordingAdapter()
startup_runner = _make_runner(startup_adapter)
startup_runner._kanban_dispatcher_lock_handle = None
asyncio.run(_run_one_notifier_tick(monkeypatch, startup_runner))
assert startup_adapter.sent == []
assert len(_unseen_terminal_events_for(task_id, "legacy-chat")) == 1
lock_path = tmp_path / ".dispatcher.lock"
winner_handle, winner_state = _acquire_singleton_lock(lock_path)
loser_handle, loser_state = _acquire_singleton_lock(lock_path)
try:
assert winner_state == "held"
assert loser_state == "contended"
loser_adapter = RecordingAdapter()
loser_runner = _make_runner(loser_adapter)
loser_runner._kanban_dispatcher_lock_handle = loser_handle
asyncio.run(_run_one_notifier_tick(monkeypatch, loser_runner))
assert loser_adapter.sent == []
assert len(_unseen_terminal_events_for(task_id, "legacy-chat")) == 1
winner_adapter = RecordingAdapter()
winner_runner = _make_runner(winner_adapter)
winner_runner._kanban_dispatcher_lock_handle = winner_handle
asyncio.run(_run_one_notifier_tick(monkeypatch, winner_runner))
assert [item["chat_id"] for item in winner_adapter.sent] == ["legacy-chat"]
assert task_id in winner_adapter.sent[0]["text"]
finally:
_release_singleton_lock(loser_handle)
_release_singleton_lock(winner_handle)
class FailingAdapter:
"""Adapter whose send() always raises, simulating a transient send error."""
@@ -347,8 +453,8 @@ def test_kanban_notifier_isolates_per_subscription_failure(tmp_path, monkeypatch
# unordered SELECT's scan order.
original_list = kb.list_notify_subs
def bad_first(conn, task_id=None):
subs = original_list(conn, task_id)
def bad_first(conn, task_id=None, **kwargs):
subs = original_list(conn, task_id, **kwargs)
return sorted(subs, key=lambda s: 0 if s["task_id"] == tid_bad else 1)
monkeypatch.setattr(kb, "list_notify_subs", bad_first)
@@ -69,6 +69,7 @@ def _make_runner(adapters):
runner._running = True
runner.adapters = adapters
runner._kanban_sub_fail_counts = {}
runner._kanban_dispatcher_lock_handle = object()
return runner
@@ -135,4 +136,3 @@ def test_apiserver_sub_wakes_real_session_via_self_post(tmp_path, monkeypatch):
# once the wake succeeds.
assert _unseen_terminal_events(tid, "api_server", "raw-sid-123") == []
@@ -1,9 +1,4 @@
"""Tests for the dispatch_in_gateway gate on _kanban_notifier_watcher.
- Non-dispatch gateways (dispatch_in_gateway=false) exit before opening any DB.
- HERMES_KANBAN_DISPATCH_IN_GATEWAY env var disables without loading config.
- Dispatch-owning gateways (dispatch_in_gateway=true) proceed past the gate.
"""
"""Notifier polling stays active when another gateway owns dispatching."""
import asyncio
from unittest.mock import MagicMock, patch
@@ -20,12 +15,8 @@ def _make_runner(with_adapter=False):
return runner
def _fake_config(dispatch_in_gateway):
return {"kanban": {"dispatch_in_gateway": dispatch_in_gateway}}
def test_notifier_watcher_runs_when_dispatch_enabled():
"""dispatch_in_gateway=true proceeds past the gate to the board fan-out."""
def test_notifier_watcher_polls_without_dispatch_ownership():
"""A profile gateway still polls its profile-owned subscriptions."""
runner = _make_runner(with_adapter=True)
past_gate = []
sleep_calls = []
@@ -42,13 +33,14 @@ def test_notifier_watcher_runs_when_dispatch_enabled():
import hermes_cli.kanban_db as _kb
with patch("hermes_cli.config.load_config", return_value=_fake_config(True)):
with patch.object(
_kb, "list_boards",
side_effect=lambda *a, **kw: past_gate.append(True) or [],
):
with patch("asyncio.sleep", side_effect=fake_sleep):
with patch("asyncio.to_thread", side_effect=fake_to_thread):
asyncio.run(runner._kanban_notifier_watcher())
with patch.object(
_kb, "list_boards",
side_effect=lambda *a, **kw: past_gate.append(True) or [],
):
with patch("asyncio.sleep", side_effect=fake_sleep):
with patch("asyncio.to_thread", side_effect=fake_to_thread):
asyncio.run(runner._kanban_notifier_watcher())
assert past_gate, "list_boards should be called when dispatch_in_gateway=true"
assert past_gate, (
"gateways without the dispatch lock must still poll owned subscriptions"
)
@@ -63,3 +63,29 @@ def test_legacy_db_without_subs_table_counts_zero_and_stays_unmigrated(tmp_path)
)
def test_count_notify_subs_filters_profile_owners(tmp_path):
db_path = tmp_path / "owners.db"
kb.init_db(db_path)
conn = kb.connect(db_path)
try:
for profile in ("default", "writer", None):
task_id = kb.create_task(conn, title=f"owned by {profile}")
kb.add_notify_sub(
conn,
task_id=task_id,
platform="telegram",
chat_id=f"chat-{profile}",
notifier_profile=profile,
)
finally:
conn.close()
assert kb.count_notify_subs(
db_path=db_path, notifier_profiles={"writer"},
) == 1
assert kb.count_notify_subs(
db_path=db_path,
notifier_profiles={"default"},
include_unowned=True,
) == 2
+1
View File
@@ -160,6 +160,7 @@ async def test_notifier_artifact_delivery_skips_missing_files(kanban_home, tmp_p
runner = object.__new__(GatewayRunner)
runner._running = True
runner._kanban_sub_fail_counts = {}
runner._kanban_dispatcher_lock_handle = object()
fake_adapter = MagicMock()
fake_adapter.name = "telegram"