From fcd4778e1b01c15b9268e8e5a8c5cf730bf9961b Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Tue, 15 Sep 2026 12:14:15 -0700 Subject: [PATCH] fix(gateway): trim planned-restart notice replay to a single write-once marker pass Slim the salvaged #112111 mechanism (issue #112109) while keeping its behaviour: the boot pass and the reconnect hook both run _replay_pending_planned_restart_notification, which sends to every home channel still owed an online notice, records delivered targets in .restart_pending.json and unlinks the marker only when every owed target (configured home with gateway_restart_notification=true) has been reached. Dropped from the contributor diff: - the per-target on_delivered checkpoint callback and pending_targets field: delivered targets are written once after the pass. Residual is a benign duplicate notice only if the process dies mid send-loop. - getattr-based lazy lock -> class attribute default, same idiom as run_profile_reconcile._reconcile_lock. - _clear_planned_restart_notification in gateway/run.py: no production caller remained; the roundtrip test unlinks the path directly. - tests trimmed to two invariants: offline-at-boot is replayed once on reconnect (with live-at-boot control), and partial delivery is persisted so a fresh process does not re-notify and an opted-out home never keeps the marker alive. Live probe (temp HERMES_HOME, Discord home, adapter absent at boot then reconnected): base consumed the marker with 0 sends; fixed head retains it and sends the online notice exactly once on reconnect, then clears it. --- .../emails/tartakovsky.steven@gmail.com | 1 + gateway/run.py | 4 - gateway/run_notifications.py | 55 ++--- tests/gateway/test_restart_notice_replay.py | 211 +++++++----------- tests/gateway/test_restart_notification.py | 2 +- 5 files changed, 101 insertions(+), 172 deletions(-) create mode 100644 contributors/emails/tartakovsky.steven@gmail.com diff --git a/contributors/emails/tartakovsky.steven@gmail.com b/contributors/emails/tartakovsky.steven@gmail.com new file mode 100644 index 0000000000..10554e557c --- /dev/null +++ b/contributors/emails/tartakovsky.steven@gmail.com @@ -0,0 +1 @@ +startakovsky diff --git a/gateway/run.py b/gateway/run.py index 93365348cb..f505eb66ea 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -1550,10 +1550,6 @@ def _planned_restart_notification_pending() -> bool: return _planned_restart_notification_path().exists() -def _clear_planned_restart_notification() -> None: - _planned_restart_notification_path().unlink(missing_ok=True) - - # Gateway marker so a lazily imported cli.py load_cli_config() doesn't clobber TERMINAL_CWD. os.environ["_HERMES_GATEWAY"] = "1" diff --git a/gateway/run_notifications.py b/gateway/run_notifications.py index 5dd898d779..9810eaa7cf 100644 --- a/gateway/run_notifications.py +++ b/gateway/run_notifications.py @@ -15,7 +15,7 @@ import logging import time from contextlib import suppress from pathlib import Path -from typing import Any, Callable, Dict, Optional, cast +from typing import Any, Dict, Optional, cast from gateway.config import Platform, _BUILTIN_PLATFORM_VALUES from gateway.platforms.base import BasePlatformAdapter, _mark_notify_metadata @@ -799,60 +799,53 @@ class GatewayNotificationsMixin: return None return "Inference: Nous free tier (nous/welcome). Sign in for more: /login" - async def _replay_pending_planned_restart_notification(self) -> None: - """Checkpoint each successful home notice so unavailable targets survive boot/reconnect. + _planned_restart_notice_lock: Optional[asyncio.Lock] = None - Boot sends may outlive the restore gate and overlap reconnects. Serialize the read/send/ack - sequence; the marker also carries acknowledgments across process restarts. + async def _replay_pending_planned_restart_notification(self) -> None: + """Send the planned-restart online notice to every home channel still owed one; clear + ``.restart_pending.json`` only once all of them were reached. + + Runs from the boot pass and again from ``_install_reconnected_adapter``, so a home whose + platform was down at boot gets its notice when the platform comes back (#112109). Delivered + targets are recorded in the marker so neither a later replay nor the next process (if this + one restarts first) notifies a home twice. The lock serializes a boot pass that outlived the + restore gate against a concurrent reconnect replay. """ from gateway.run import _planned_restart_notification_path from utils import atomic_json_write - lock = getattr(self, "_planned_restart_notice_lock", None) - if lock is None: - lock = self._planned_restart_notice_lock = asyncio.Lock() - async with lock: + if self._planned_restart_notice_lock is None: + self._planned_restart_notice_lock = asyncio.Lock() + async with self._planned_restart_notice_lock: path = _planned_restart_notification_path() if not path.exists(): return try: data = json.loads(path.read_text(encoding="utf-8")) delivered = {tuple(target) for target in data.get("delivered_targets", [])} - # Resolve obligations from configuration, never from the currently live transports. - # Removed homes and explicit notification opt-outs no longer owe a notice. - targets = { + # Owed targets come from config, not live transports: a removed home or an opt-out + # (gateway_restart_notification=false) must not keep the marker alive forever. + owed = { _notice_target_key(platform.value, cfg.home_channel.chat_id, cfg.home_channel.thread_id) for platform, cfg in self.config.platforms.items() if cfg.home_channel and cfg.home_channel.chat_id and cfg.gateway_restart_notification } - pending = targets - delivered - - def checkpoint(target=None): - if target is not None: - delivered.add(target) - pending.discard(target) - data["delivered_targets"] = list(delivered) - data["pending_targets"] = list(pending) - atomic_json_write(path, data) - - checkpoint() - await self._send_home_channel_startup_notifications( - skip_targets=delivered, on_delivered=checkpoint, - ) - if not pending: + delivered |= await self._send_home_channel_startup_notifications(skip_targets=delivered) + if owed <= delivered: path.unlink(missing_ok=True) + return + data["delivered_targets"] = [list(target) for target in delivered] + atomic_json_write(path, data, indent=None) except Exception: logger.warning("Planned-restart notification remains pending", exc_info=True) async def _send_home_channel_startup_notifications( - self, *, skip_targets: Optional[set[tuple[str, str, Optional[str]]]] = None, - on_delivered: Optional[Callable[[tuple[str, str, Optional[str]]], None]] = None, + self, *, skip_targets: Optional[set[tuple[str, str, Optional[str]]]] = None ) -> set[tuple[str, str, Optional[str]]]: """Notify configured home channels that the gateway is back online. Best-effort, once per connected platform home channel. ``skip_targets`` lets startup avoid duplicate messages when a more specific restart notification is queued for the same chat. - ``on_delivered`` persists each acknowledgment before attempting the next transport. """ delivered: set[tuple[str, str, Optional[str]]] = set() skipped = skip_targets or set() @@ -874,8 +867,6 @@ class GatewayNotificationsMixin: platform, home, transport, message, "Home-channel startup notification failed for %s:%s: %s", ): delivered.add(target) - if on_delivered is not None: - on_delivered(target) logger.info("Sent home-channel startup notification to %s:%s", platform.value, home.chat_id) return delivered diff --git a/tests/gateway/test_restart_notice_replay.py b/tests/gateway/test_restart_notice_replay.py index a9d0ef36a8..9db4b44be5 100644 --- a/tests/gateway/test_restart_notice_replay.py +++ b/tests/gateway/test_restart_notice_replay.py @@ -1,15 +1,14 @@ -"""Regression for #112109: retain planned-restart notices until delivery succeeds. +"""Planned-restart online notice survives an offline home-channel transport at boot. -Uses the real boot notification pass, marker helpers, home-channel sender, and -DeliveryTransport. Patterns follow tests/gateway/test_restart_notification.py -and test_restart_resume_pending.py. Run with scripts/run_tests.sh for isolation. +The ``.restart_pending.json`` marker used to be consumed in ``finally`` even when no live transport +existed for the home channel, so the "Gateway online" notice was never sent and never replayed. +See #112109. Runs the real boot pass, marker helpers, home-channel sender and DeliveryTransport. """ -from types import SimpleNamespace -from unittest.mock import AsyncMock, Mock - import asyncio import json +from types import SimpleNamespace +from unittest.mock import AsyncMock, Mock import pytest @@ -18,6 +17,15 @@ import gateway.run as gateway_run from gateway.config import GatewayConfig, HomeChannel, Platform, PlatformConfig from gateway.platforms.base import SendResult +ONLINE_NOTICE = "♻️ Gateway online — Hermes is back and ready." + + +def _adapter(): + return SimpleNamespace( + send_path_degraded=False, + send=AsyncMock(return_value=SendResult(success=True, message_id="unit-test-notice")), + ) + @pytest.fixture def boot_notice(tmp_path, monkeypatch): @@ -26,15 +34,14 @@ def boot_notice(tmp_path, monkeypatch): # Await the boot task to completion and propagate failures deterministically. monkeypatch.setattr(gateway_run, "_startup_restore_drain_timeout_secs", lambda: 0) runner = object.__new__(gateway_run.GatewayRunner) - platform_config = PlatformConfig( - enabled=True, - gateway_restart_notification=True, - home_channel=HomeChannel( - platform=Platform.DISCORD, chat_id="unit-test-home", name="Test home" - ), - ) runner.config = GatewayConfig( - platforms={Platform.DISCORD: platform_config}, + platforms={ + Platform.DISCORD: PlatformConfig( + enabled=True, + gateway_restart_notification=True, + home_channel=HomeChannel(platform=Platform.DISCORD, chat_id="unit-test-home", name="Test home"), + ), + }, sessions_dir=tmp_path / "sessions", ) runner.adapters = {} @@ -50,140 +57,74 @@ def boot_notice(tmp_path, monkeypatch): runner._claim_pending_obligations = AsyncMock(return_value=[]) runner._redeliver_claimed_obligations = AsyncMock(return_value=0) runner._free_tier_startup_line = Mock(return_value=None) - # Keep the real requester-marker check; this case has only the planned marker. - assert not (tmp_path / ".restart_notify.json").exists() marker = tmp_path / ".restart_pending.json" marker.write_text("{}", encoding="utf-8") - adapter = SimpleNamespace(send_path_degraded=False, send=AsyncMock( - return_value=SendResult(success=True, message_id="unit-test-notice") - )) - return runner, platform_config, marker, adapter + return runner, marker -@pytest.mark.asyncio -@pytest.mark.parametrize("live", [False, True], ids=[ - "no-live-transport-retained-and-replayed", - "live-transport-notice-sent-marker-consumed", -]) -async def test_planned_restart_boot_notice(boot_notice, monkeypatch, live): - runner, platform_config, marker, adapter = boot_notice - if live: - runner.adapters[Platform.DISCORD] = adapter - transport = ( - delivery.DeliveryTransport(adapter, platform_config, Platform.DISCORD) - if live else None - ) - # Confirm the stub represents the real resolver's result for this adapter map. - resolved = delivery.resolve_delivery_transport( - Platform.DISCORD, runner.config, runner.adapters - ) - assert resolved == transport - resolver = Mock(return_value=transport) - monkeypatch.setattr(delivery, "resolve_delivery_transport", resolver) - assert marker.exists() - assert gateway_run._planned_restart_notification_pending() - +async def _boot(runner): await runner._await_startup_boot_sends( planned_restart_notification_pending=gateway_run._planned_restart_notification_pending() ) - resolver.assert_called_once_with(Platform.DISCORD, runner.config, runner.adapters) - runner._claim_pending_obligations.assert_awaited_once_with() - runner._redeliver_claimed_obligations.assert_awaited_once_with([]) - if live: - adapter.send.assert_awaited_once_with( - "unit-test-home", - "♻️ Gateway online — Hermes is back and ready.", - metadata={"non_conversational": True}, - ) - else: - adapter.send.assert_not_called() - if not live: - assert marker.exists() - assert gateway_run._planned_restart_notification_pending() - resolver.return_value = delivery.DeliveryTransport(adapter, platform_config, Platform.DISCORD) - runner._failed_platforms[Platform.DISCORD] = {} - await runner._install_reconnected_adapter(Platform.DISCORD, adapter) - await asyncio.gather(*runner._background_tasks) - adapter.send.assert_awaited_once_with( - "unit-test-home", "♻️ Gateway online — Hermes is back and ready.", - metadata={"non_conversational": True}, - ) - assert not marker.exists() - assert not gateway_run._planned_restart_notification_pending() + +async def _reconnect(runner, platform, adapter): + runner._failed_platforms[platform] = {} + await runner._install_reconnected_adapter(platform, adapter) + await asyncio.gather(*runner._background_tasks) @pytest.mark.asyncio -@pytest.mark.parametrize("outage", ["unavailable", "rejected", "exception", "cancelled"]) -async def test_partial_notice_delivery_survives_restart_and_concurrent_replay(boot_notice, outage): - runner, _, marker, adapter = boot_notice - other = SimpleNamespace(send=AsyncMock(return_value=SendResult(success=True))) - # Deliver one destination first, then encounter an unavailable/failing/hung transport. - runner.config.platforms = { - Platform.TELEGRAM: PlatformConfig( - enabled=True, - home_channel=HomeChannel(platform=Platform.TELEGRAM, chat_id="other-home", thread_id="7", name="Other home"), - ), - **runner.config.platforms, - Platform.SLACK: PlatformConfig( - enabled=True, gateway_restart_notification=False, - home_channel=HomeChannel(platform=Platform.SLACK, chat_id="muted-home", name="Muted home"), - ), - } - runner.adapters[Platform.TELEGRAM] = other - started = asyncio.Event() - release = asyncio.Event() - - async def slow_send(*args, **kwargs): - started.set() - await release.wait() - return SendResult(success=True) - - if outage != "unavailable": +@pytest.mark.parametrize("live", [False, True], ids=["offline-at-boot-replayed-on-reconnect", "live-at-boot"]) +async def test_planned_restart_notice_reaches_home_channel(boot_notice, live): + runner, marker = boot_notice + adapter = _adapter() + if live: runner.adapters[Platform.DISCORD] = adapter - if outage == "rejected": - adapter.send.return_value = SendResult(success=False, error="temporarily unavailable") - elif outage == "exception": - adapter.send.side_effect = RuntimeError("transport disconnected") - elif outage == "cancelled": - adapter.send.side_effect = slow_send + transport = delivery.resolve_delivery_transport(Platform.DISCORD, runner.config, runner.adapters) + assert (transport is not None) is live - boot = asyncio.create_task(runner._await_startup_boot_sends(planned_restart_notification_pending=True)) - if outage == "cancelled": - await asyncio.wait_for(started.wait(), timeout=5) - # Acknowledgments must already be on disk while a later destination is hung. - assert json.loads(marker.read_text())["delivered_targets"] == [["telegram", "other-home", "7"]] - boot.cancel() - with pytest.raises(asyncio.CancelledError): - await boot - else: - await boot - data = json.loads(marker.read_text()) - assert data["delivered_targets"] == [["telegram", "other-home", "7"]] - assert data["pending_targets"] == [["discord", "unit-test-home", None]] - other.send.assert_awaited_once() + await _boot(runner) - # A new runner has no in-memory delivery history: dedupe must come from the marker. + runner._redeliver_claimed_obligations.assert_awaited_once_with([]) + if not live: + adapter.send.assert_not_called() + assert marker.exists(), "marker must survive a boot with no live transport" + await _reconnect(runner, Platform.DISCORD, adapter) + adapter.send.assert_awaited_once_with("unit-test-home", ONLINE_NOTICE, metadata={"non_conversational": True}) + assert not marker.exists() + + +@pytest.mark.asyncio +async def test_partial_delivery_is_persisted_and_not_repeated(boot_notice): + runner, marker = boot_notice + telegram, discord = _adapter(), _adapter() + runner.config.platforms[Platform.TELEGRAM] = PlatformConfig( + enabled=True, + home_channel=HomeChannel(platform=Platform.TELEGRAM, chat_id="other-home", thread_id="7", name="Other"), + ) + runner.config.platforms[Platform.SLACK] = PlatformConfig( + enabled=True, gateway_restart_notification=False, + home_channel=HomeChannel(platform=Platform.SLACK, chat_id="muted-home", name="Muted"), + ) + runner.adapters[Platform.TELEGRAM] = telegram + + await _boot(runner) + + telegram.send.assert_awaited_once() + discord.send.assert_not_called() + assert json.loads(marker.read_text(encoding="utf-8"))["delivered_targets"] == [["telegram", "other-home", "7"]] + + # A fresh process has no in-memory history: dedupe must come from the marker. The opted-out + # Slack home is never owed a notice, so Discord's delivery completes the set. recovered = object.__new__(gateway_run.GatewayRunner) recovered.__dict__.update(runner.__dict__) - recovered.__dict__.pop("_planned_restart_notice_lock", None) - started.clear() - adapter.send.reset_mock() - adapter.send.side_effect = slow_send - recovered._failed_platforms[Platform.DISCORD] = {} - await asyncio.wait_for(recovered._install_reconnected_adapter(Platform.DISCORD, adapter), timeout=5) - await asyncio.wait_for(started.wait(), timeout=5) - # Installation completed even though notification delivery is still blocked. - assert marker.exists() - concurrent = asyncio.create_task(recovered._replay_pending_planned_restart_notification()) - release.set() - await asyncio.gather(concurrent, *recovered._background_tasks) - assert not marker.exists() - adapter.send.assert_awaited_once() - other.send.assert_awaited_once() + await _reconnect(recovered, Platform.DISCORD, discord) - recovered._failed_platforms[Platform.DISCORD] = {} - await recovered._install_reconnected_adapter(Platform.DISCORD, adapter) - await asyncio.gather(*recovered._background_tasks) - adapter.send.assert_awaited_once() - other.send.assert_awaited_once() + discord.send.assert_awaited_once() + telegram.send.assert_awaited_once() + assert not marker.exists() + + # Nothing pending: a later reconnect stays silent. + await _reconnect(recovered, Platform.DISCORD, discord) + discord.send.assert_awaited_once() diff --git a/tests/gateway/test_restart_notification.py b/tests/gateway/test_restart_notification.py index 1b81cb4188..7af3dbf478 100644 --- a/tests/gateway/test_restart_notification.py +++ b/tests/gateway/test_restart_notification.py @@ -28,7 +28,7 @@ def test_planned_restart_notification_pending_roundtrip(tmp_path, monkeypatch): marker.write_text("{}") assert gateway_run._planned_restart_notification_pending() is True - gateway_run._clear_planned_restart_notification() + gateway_run._planned_restart_notification_path().unlink() assert gateway_run._planned_restart_notification_pending() is False