From de109846bf3c2f46b2e20ffc88c7d2404a1c42d3 Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Tue, 14 Jul 2026 17:01:23 -0700 Subject: [PATCH] =?UTF-8?q?fix:=20prax=20TG=20relay=20offline=20backoff=20?= =?UTF-8?q?=E2=80=94=20network-class=20send=20failures=20enter=20offline?= =?UTF-8?q?=20mode=20(1s-60s=20doubling,=20flush=20gate=20skips=20sends,?= =?UTF-8?q?=20monitor=20loop=20never=20blocks,=20viewers=20keep=20renderin?= =?UTF-8?q?g),=20log-once=20(enter=20+=205min=20summary=20+=20recovery=20w?= =?UTF-8?q?/=20drop=20count),=20full=20reset=20on=20first=20success.=20Fou?= =?UTF-8?q?nd=20live=20in=20Patrick's=20plug-pull:=20bots=20went=20quiet?= =?UTF-8?q?=20right,=20relay=20spun=20Send=20failed=20every=205s=20(89=20l?= =?UTF-8?q?ines)=20+=20self-fed=20via=20log=20watcher=20re-ingest.=2011=20?= =?UTF-8?q?new=20tests,=201007=20prax=20green=20devpulse-verified?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 11 ++ src/aipass/prax/CLOSED_PLANS.local.json | 7 + src/aipass/prax/README.md | 4 +- .../handlers/monitoring/telegram_relay.py | 52 +++++- src/aipass/prax/tests/test_telegram_relay.py | 155 ++++++++++++++++++ 5 files changed, 225 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 12ef16a6..7efd8a39 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,17 @@ PyPI version — not the changelog header. ### Fixed +- **Prax TG relay gets the same offline backoff as the bots.** Found in + Patrick's live plug-pull test: the bots went quiet correctly, but the + monitor→Telegram relay kept logging `Send failed` every ~5 seconds (89 lines, + no backoff) — and each failed-send error was re-ingested by the log watcher, + feeding the relay more events to fail on. Now network-class send failures put + the relay in offline mode: doubling backoff (1s→60s cap), flush gate skips + sends while offline so the monitor loop never blocks and viewers keep + rendering, log-once semantics (one enter line, one 5-minute summary, one + recovery line with drop count), full reset on first successful send. 11 new + tests; 1007 prax green. + - **TG bots no longer hot-spin when the internet drops.** Live find from Patrick's on-location tether outage: DNS failure makes `urlopen` fail instantly (no 30s long-poll wait), so the shared poll loop retried as fast as diff --git a/src/aipass/prax/CLOSED_PLANS.local.json b/src/aipass/prax/CLOSED_PLANS.local.json index 304a6cad..1c4aefee 100644 --- a/src/aipass/prax/CLOSED_PLANS.local.json +++ b/src/aipass/prax/CLOSED_PLANS.local.json @@ -83,6 +83,13 @@ "subject": "Remove single-instance lock from monitor display path — concurrent viewers", "date_closed": "2026-07-14", "location": "prax" + }, + { + "plan_id": "FPLAN-0325", + "type": "FPLAN", + "subject": "TG relay send backoff: offline mode + log-once on network failures", + "date_closed": "2026-07-14", + "location": "prax" } ], "document_metadata": { diff --git a/src/aipass/prax/README.md b/src/aipass/prax/README.md index ed70936d..c0538621 100644 --- a/src/aipass/prax/README.md +++ b/src/aipass/prax/README.md @@ -143,7 +143,7 @@ prax/ │ └── watcher/ # Background system watchers ├── prax_json/ # Auto-created per-module config/data/log files ├── templates/ # Dashboard template schema (DASHBOARD.template.json) -└── tests/ # 901 tests across 19 files +└── tests/ # 1007 tests across 19 files ``` ### Design Pattern @@ -171,7 +171,7 @@ drone @prax monitor run ## Tests -901 tests across 19 files, covering all major components: +1007 tests across 19 files, covering all major components: | Test File | Tests | Coverage | |-----------|-------|----------| diff --git a/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py b/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py index e09bd091..1e86cefe 100644 --- a/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py +++ b/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py @@ -17,6 +17,7 @@ bot config passed by the module layer (monitor.py loads from @api secrets). import json import os import threading +import time from datetime import datetime from pathlib import Path from typing import Optional @@ -47,6 +48,17 @@ _RELAY_ACTIVE = False _control_mtime: float = 0.0 _control_cache: dict = {} +_BACKOFF_INITIAL = 1.0 +_BACKOFF_CAP = 60.0 +_SUMMARY_INTERVAL = 300.0 + +_OFFLINE = False +_CURRENT_BACKOFF: float = _BACKOFF_INITIAL +_NEXT_RETRY: float = 0.0 +_OFFLINE_SINCE: float = 0.0 +_SUPPRESSED_COUNT = 0 +_LAST_SUMMARY: float = 0.0 + def init_relay(enabled: bool, config: Optional[dict] = None) -> None: """Start the relay if enabled and config is valid. Safe no-op otherwise. @@ -215,6 +227,10 @@ def _flush_buffer() -> None: lines = lines[:FLOOD_CAP] lines.append(f"…({suppressed} more suppressed)") + if _OFFLINE and time.monotonic() < _NEXT_RETRY: + _count_suppressed(len(lines)) + return + _send_batched(lines) @@ -247,8 +263,20 @@ def _send_batched(lines: list[str]) -> None: _send_message("\n".join(batch)) +def _count_suppressed(count: int) -> None: + """Track suppressed events during offline and log a summary at most every 5 minutes.""" + global _SUPPRESSED_COUNT, _LAST_SUMMARY + _SUPPRESSED_COUNT += count + now = time.monotonic() + if now - _LAST_SUMMARY >= _SUMMARY_INTERVAL: + logger.info("[telegram_relay] Still offline — %d events suppressed so far", _SUPPRESSED_COUNT) + _LAST_SUMMARY = now + + def _send_message(text: str) -> bool: """POST a single message to the Telegram Bot API.""" + global _OFFLINE, _CURRENT_BACKOFF, _NEXT_RETRY, _OFFLINE_SINCE, _SUPPRESSED_COUNT, _LAST_SUMMARY + url = f"https://api.telegram.org/bot{_bot_token}/sendMessage" payload = json.dumps( { @@ -262,9 +290,29 @@ def _send_message(text: str) -> bool: try: with _http_fetch(req, timeout=10) as resp: result = json.loads(resp.read().decode("utf-8")) + if _OFFLINE: + duration = time.monotonic() - _OFFLINE_SINCE + logger.info( + "[telegram_relay] Recovered (was offline %.0fs, %d events dropped from TG feed)", + duration, + _SUPPRESSED_COUNT, + ) + _OFFLINE = False + _CURRENT_BACKOFF = _BACKOFF_INITIAL + _SUPPRESSED_COUNT = 0 return result.get("ok", False) - except (URLError, Exception) as e: - logger.warning("[telegram_relay] Send failed: %s", e) + except (URLError, OSError) as e: + now = time.monotonic() + if not _OFFLINE: + logger.warning("[telegram_relay] TG relay offline: %s", e) + _OFFLINE = True + _OFFLINE_SINCE = now + _CURRENT_BACKOFF = _BACKOFF_INITIAL + _SUPPRESSED_COUNT = 0 + _LAST_SUMMARY = now + else: + _CURRENT_BACKOFF = min(_CURRENT_BACKOFF * 2, _BACKOFF_CAP) + _NEXT_RETRY = now + _CURRENT_BACKOFF return False diff --git a/src/aipass/prax/tests/test_telegram_relay.py b/src/aipass/prax/tests/test_telegram_relay.py index 2168c1b8..db35dc40 100644 --- a/src/aipass/prax/tests/test_telegram_relay.py +++ b/src/aipass/prax/tests/test_telegram_relay.py @@ -18,6 +18,7 @@ Covers: - Flood cap: truncation at 150 lines with suppression notice - _render_event calls relay_event in monitor.py - is_relay_enabled_by_env for env var detection +- Offline backoff: doubles+caps, resets on success, log-once, never blocks """ import importlib @@ -57,12 +58,21 @@ def _import_relay(): else: mod = importlib.import_module("aipass.prax.apps.handlers.monitoring.telegram_relay") + lock_mock = MagicMock() + lock_mock.try_acquire = MagicMock(return_value=True) + lock_mock.release = MagicMock() + setattr(mod, "instance_lock", lock_mock) + setattr(mod, "_RELAY_ACTIVE", False) setattr(mod, "_bot_token", None) setattr(mod, "_chat_id", None) mod._buffer.clear() mod._stop_event.clear() setattr(mod, "_thread", None) + setattr(mod, "_OFFLINE", False) + setattr(mod, "_CURRENT_BACKOFF", mod._BACKOFF_INITIAL) + setattr(mod, "_NEXT_RETRY", 0.0) + setattr(mod, "_SUPPRESSED_COUNT", 0) return mod @@ -578,3 +588,148 @@ class TestFlushControl: relay._flush_buffer() assert len(sent) == relay.FLOOD_CAP + 1 assert "suppressed" in sent[-1] + + +# --------------------------------------------------------------------------- +# Offline backoff +# --------------------------------------------------------------------------- + + +class TestOfflineBackoff: + """Network-failure backoff: doubles+caps, resets on success, log-once.""" + + def _make_relay(self): + relay = _import_relay() + setattr(relay, "_bot_token", "t") + setattr(relay, "_chat_id", 1) + setattr(relay, "_RELAY_ACTIVE", True) + return relay + + def test_network_error_enters_offline(self): + """First URLError sets _offline=True.""" + from urllib.error import URLError + + relay = self._make_relay() + setattr(relay, "_send_message", relay._send_message) + with patch.object(relay, "_http_fetch", side_effect=URLError("DNS failed")): + result = relay._send_message("hello") + assert result is False + assert relay._OFFLINE is True + + def test_backoff_doubles_on_repeated_failure(self): + """Backoff doubles: 1 → 2 → 4.""" + from urllib.error import URLError + + relay = self._make_relay() + with patch.object(relay, "_http_fetch", side_effect=URLError("offline")): + relay._send_message("a") + assert relay._CURRENT_BACKOFF == relay._BACKOFF_INITIAL + relay._send_message("b") + assert relay._CURRENT_BACKOFF == 2.0 + relay._send_message("c") + assert relay._CURRENT_BACKOFF == 4.0 + + def test_backoff_caps_at_60s(self): + """Backoff never exceeds _BACKOFF_CAP (60s).""" + from urllib.error import URLError + + relay = self._make_relay() + with patch.object(relay, "_http_fetch", side_effect=URLError("offline")): + for _ in range(20): + relay._send_message("x") + assert relay._CURRENT_BACKOFF == relay._BACKOFF_CAP + + def test_success_resets_offline(self): + """Successful send after offline resets state.""" + from urllib.error import URLError + + relay = self._make_relay() + with patch.object(relay, "_http_fetch", side_effect=URLError("offline")): + relay._send_message("a") + assert relay._OFFLINE is True + + ok_response = MagicMock() + ok_response.read.return_value = b'{"ok": true}' + ok_response.__enter__ = MagicMock(return_value=ok_response) + ok_response.__exit__ = MagicMock(return_value=False) + with patch.object(relay, "_http_fetch", return_value=ok_response): + result = relay._send_message("b") + assert result is True + assert relay._OFFLINE is False + assert relay._CURRENT_BACKOFF == relay._BACKOFF_INITIAL + assert relay._SUPPRESSED_COUNT == 0 + + def test_flush_suppresses_during_backoff(self): + """_flush_buffer skips _send_batched while offline and before next retry.""" + import time + + relay = self._make_relay() + setattr(relay, "_OFFLINE", True) + setattr(relay, "_NEXT_RETRY", time.monotonic() + 9999) + relay._buffer.extend(["line 1", "line 2", "line 3"]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert sent == [] + assert relay._SUPPRESSED_COUNT == 3 + + def test_flush_retries_after_backoff_expires(self): + """_flush_buffer attempts send when backoff period has elapsed.""" + import time + + relay = self._make_relay() + setattr(relay, "_OFFLINE", True) + setattr(relay, "_NEXT_RETRY", time.monotonic() - 1) + relay._buffer.extend(["retry line"]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert len(sent) == 1 + + def test_log_once_on_entering_offline(self): + """Only one warning logged on first network failure.""" + from urllib.error import URLError + + relay = self._make_relay() + with patch.object(relay, "_http_fetch", side_effect=URLError("offline")): + relay._send_message("a") + relay._send_message("b") + relay._send_message("c") + warning_calls = relay.logger.warning.call_args_list + offline_warnings = [c for c in warning_calls if "offline" in str(c).lower()] + assert len(offline_warnings) == 1 + + def test_summary_logged_after_interval(self): + """Suppression summary logged after _SUMMARY_INTERVAL elapses.""" + import time + + relay = self._make_relay() + setattr(relay, "_OFFLINE", True) + setattr(relay, "_NEXT_RETRY", time.monotonic() + 9999) + setattr(relay, "_LAST_SUMMARY", time.monotonic() - relay._SUMMARY_INTERVAL - 1) + relay._buffer.extend(["line"]) + setattr(relay, "_send_batched", lambda lines: None) + relay._flush_buffer() + info_calls = relay.logger.info.call_args_list + summary_calls = [c for c in info_calls if "suppressed" in str(c).lower()] + assert len(summary_calls) >= 1 + + def test_recovery_log_includes_drop_count(self): + """Recovery log line includes the number of dropped events.""" + from urllib.error import URLError + + relay = self._make_relay() + with patch.object(relay, "_http_fetch", side_effect=URLError("offline")): + relay._send_message("a") + setattr(relay, "_SUPPRESSED_COUNT", 42) + + ok_response = MagicMock() + ok_response.read.return_value = b'{"ok": true}' + ok_response.__enter__ = MagicMock(return_value=ok_response) + ok_response.__exit__ = MagicMock(return_value=False) + with patch.object(relay, "_http_fetch", return_value=ok_response): + relay._send_message("b") + info_calls = relay.logger.info.call_args_list + recovery_calls = [c for c in info_calls if "recovered" in str(c).lower()] + assert len(recovery_calls) == 1 + assert "42" in str(recovery_calls[0])