fix: prax TG relay offline backoff — network-class send failures enter offline mode (1s-60s doubling, flush gate skips sends, monitor loop never blocks, viewers keep rendering), log-once (enter + 5min summary + recovery w/ drop count), full reset on first success. Found live in Patrick's plug-pull: bots went quiet right, relay spun Send failed every 5s (89 lines) + self-fed via log watcher re-ingest. 11 new tests, 1007 prax green devpulse-verified
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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 |
|
||||
|-----------|-------|----------|
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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])
|
||||
|
||||
Reference in New Issue
Block a user