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.
This commit is contained in:
1
contributors/emails/tartakovsky.steven@gmail.com
Normal file
1
contributors/emails/tartakovsky.steven@gmail.com
Normal file
@@ -0,0 +1 @@
|
||||
startakovsky
|
||||
@@ -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"
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user