From be8f87d6fb89965914d10b7f2b72217370787726 Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Thu, 25 Jun 2026 00:10:06 -0700 Subject: [PATCH] =?UTF-8?q?feat(prax):=20monitor=E2=86=92Telegram=20relay?= =?UTF-8?q?=20(prax=5Fmonitor=20bot)=20+=20fix=20service=20feedback=20loop?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Mirrors the live 'drone @prax monitor run' Mission-Control feed to a dedicated Telegram bot (DPLAN-0221). New monitoring/telegram_relay.py taps _render_event, batches every 5s (4000-split, 150 flood-cap, fail-silent-once), gated by --relay/env so local monitor stays console-only. Reboot-survivable prax-monitor.service. 937 prax tests green (31 new). Deploy fixes (devpulse): ExecStart -> 'monitor run' (module __main__ rejects 'run all --relay'); service log moved out of system_logs/ to ~/.aipass/ to break a monitor<->@trigger feedback loop. --- CHANGELOG.md | 25 ++ .../handlers/monitoring/telegram_relay.py | 217 +++++++++ src/aipass/prax/apps/modules/monitor.py | 36 +- src/aipass/prax/prax-monitor.service | 34 ++ src/aipass/prax/tests/test_monitor_module.py | 1 + src/aipass/prax/tests/test_telegram_relay.py | 411 ++++++++++++++++++ 6 files changed, 722 insertions(+), 2 deletions(-) create mode 100644 src/aipass/prax/apps/handlers/monitoring/telegram_relay.py create mode 100644 src/aipass/prax/prax-monitor.service create mode 100644 src/aipass/prax/tests/test_telegram_relay.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 59ec8b6e..e19ad26f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,31 @@ PyPI version — not the changelog header. --- +## [2026-06-25] + +### Added + +- **Prax monitor → Telegram relay (`prax_monitor` bot)** — the live + `drone @prax monitor run` Mission-Control feed now mirrors to a dedicated + Telegram bot, so the whole-system monitor is watchable from a phone ("same + monitor, different window"). New `monitoring/telegram_relay.py` taps the single + render seam (`_render_event`), buffers events, and flushes every 5s (4000-char + split, 150-line flood cap, `disable_notification`); fail-silent-once when + unconfigured. Gated behind `--relay` / `AIPASS_PRAX_MONITOR_RELAY=1` so a local + `monitor run` stays console-only (no double-send). Bot config (token + chat_id) + loads from the @api secret `telegram/prax_monitor`. Ships a reboot-survivable + `prax-monitor.service` user unit. 937 prax tests green (31 new). (DPLAN-0221) + +### Fixed + +- **prax-monitor service feedback loop** — the unit wrote its own stdout into + `system_logs/`, the very directory the monitor tails *and* @trigger watches, + creating a self-reinforcing loop (monitor output → re-tailed and recorded by + @trigger into `trigger_data.json` → reported as a file change → more output). + Moved the service log to `~/.aipass/` to break the cycle. Also corrected the + ExecStart to `monitor run` (relay enabled via env) — the module `__main__` + rejects the drone-style `run all --relay` argument form. (DPLAN-0221) + ## [2026-06-24] ### Changed diff --git a/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py b/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py new file mode 100644 index 00000000..6def5e8f --- /dev/null +++ b/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py @@ -0,0 +1,217 @@ +# =================== AIPass ==================== +# Name: telegram_relay.py +# Description: Telegram relay for prax monitor feed +# Version: 1.0.0 +# Created: 2026-06-24 +# Modified: 2026-06-24 +# ============================================= + +""" +Telegram relay for prax monitor — mirrors the live console feed to a Telegram chat. + +Buffers events in a thread-safe list, flushes every 5s via a daemon thread. +Activation requires --relay flag or AIPASS_PRAX_MONITOR_RELAY=1, plus a valid +bot config passed by the module layer (monitor.py loads from @api secrets). +""" + +import json +import os +import threading +from datetime import datetime +from typing import Optional +from urllib.error import URLError +from urllib.request import Request +from urllib.request import urlopen as _http_fetch + +from aipass.prax.apps.modules.logger import get_direct_logger +from aipass.prax.apps.handlers.json import json_handler + +logger = get_direct_logger() + +BATCH_INTERVAL = 5.0 +TELEGRAM_MAX_LENGTH = 4000 +FLOOD_CAP = 150 + +_lock = threading.Lock() +_buffer: list[str] = [] +_thread: Optional[threading.Thread] = None +_stop_event = threading.Event() +_bot_token: Optional[str] = None +_chat_id: Optional[int] = None +_RELAY_ACTIVE = False + + +def init_relay(enabled: bool, config: Optional[dict] = None) -> None: + """Start the relay if enabled and config is valid. Safe no-op otherwise. + + Args: + enabled: Whether the relay flag is set. + config: Bot config dict with 'bot_token' and 'chat_id' keys. + Loaded by the module layer from @api secrets. + """ + global _RELAY_ACTIVE, _bot_token, _chat_id, _thread + + if not enabled: + return + + json_handler.log_operation("relay_init", {"enabled": enabled, "has_config": config is not None}) + + if config is None: + logger.info("[telegram_relay] No config provided — relay inactive") + return + + token = config.get("bot_token", "") + chat = config.get("chat_id") + if not token or not chat: + logger.info("[telegram_relay] Incomplete config (missing bot_token or chat_id) — relay inactive") + return + + _bot_token = token + _chat_id = int(chat) + _RELAY_ACTIVE = True + + _stop_event.clear() + _thread = threading.Thread(target=_flush_loop, name="telegram-relay", daemon=True) + _thread.start() + logger.info("[telegram_relay] Relay started (chat_id=%s)", _chat_id) + + +def relay_event(event) -> None: + """Format a MonitoringEvent and buffer for Telegram delivery. No-op when inactive.""" + if not _RELAY_ACTIVE: + return + + line = _format_event(event) + if not line: + return + + with _lock: + _buffer.append(line) + + +def stop_relay() -> None: + """Final flush and thread shutdown. Safe no-op when inactive.""" + global _RELAY_ACTIVE, _thread + + if not _RELAY_ACTIVE: + return + + _RELAY_ACTIVE = False + _stop_event.set() + + _flush_buffer() + + if _thread is not None: + _thread.join(timeout=BATCH_INTERVAL + 2) + _thread = None + + json_handler.log_operation("relay_stopped", {}) + logger.info("[telegram_relay] Relay stopped") + + +def _format_event(event) -> Optional[str]: + """Format a MonitoringEvent as a plain-text line (no Rich markup).""" + ts = datetime.now().strftime("%H:%M:%S") + + if event.event_type == "command": + parts = [f"▶ {event.message}"] + caller = getattr(event, "caller", None) + target = None + if hasattr(event, "action") and event.action and ":" in event.action: + action_parts = event.action.split(":", 1) + if len(action_parts) == 2 and action_parts[1]: + target = action_parts[1] + if caller and caller.upper() != "UNKNOWN": + attr = caller + if target: + attr = f"{caller} → {target}" + parts.insert(0, f" {attr}") + return "\n".join(parts) + + if event.event_type == "hook": + action = getattr(event, "action", "unknown") + if action == "fired": + return f"⚡ HOOK {event.message}" + if action == "skipped": + return f"· HOOK {event.message}" + return f"? HOOK {event.message}" + + branch_label = event.branch.upper() + pid = getattr(event, "pid", None) + if pid: + branch_label = f"{branch_label}:{pid}" + return f"[{ts}] [{branch_label}] {event.message}" + + +def _flush_buffer() -> None: + """Drain buffer and send to Telegram.""" + with _lock: + if not _buffer: + return + lines = list(_buffer) + _buffer.clear() + + if not _bot_token or not _chat_id: + return + + if len(lines) > FLOOD_CAP: + suppressed = len(lines) - FLOOD_CAP + lines = lines[:FLOOD_CAP] + lines.append(f"…({suppressed} more suppressed)") + + _send_batched(lines) + + +def _flush_loop() -> None: + """Background thread: flush buffer every BATCH_INTERVAL seconds.""" + while not _stop_event.is_set(): + _stop_event.wait(BATCH_INTERVAL) + if _buffer: + _flush_buffer() + + +def _send_batched(lines: list[str]) -> None: + """Split lines into ≤4000-char chunks and POST each to Telegram.""" + if not lines: + return + + batch: list[str] = [] + batch_len = 0 + + for line in lines: + line_len = len(line) + (1 if batch else 0) + if batch_len + line_len > TELEGRAM_MAX_LENGTH and batch: + _send_message("\n".join(batch)) + batch = [] + batch_len = 0 + batch.append(line) + batch_len += line_len + + if batch: + _send_message("\n".join(batch)) + + +def _send_message(text: str) -> bool: + """POST a single message to the Telegram Bot API.""" + url = f"https://api.telegram.org/bot{_bot_token}/sendMessage" + payload = json.dumps( + { + "chat_id": _chat_id, + "text": text, + "disable_notification": True, + } + ).encode("utf-8") + req = Request(url, data=payload, headers={"Content-Type": "application/json"}) + + try: + with _http_fetch(req, timeout=10) as resp: + result = json.loads(resp.read().decode("utf-8")) + return result.get("ok", False) + except (URLError, Exception) as e: + logger.warning("[telegram_relay] Send failed: %s", e) + return False + + +def is_relay_enabled_by_env() -> bool: + """Check if the relay is enabled via environment variable.""" + return os.environ.get("AIPASS_PRAX_MONITOR_RELAY", "").strip() in ("1", "true", "yes") diff --git a/src/aipass/prax/apps/modules/monitor.py b/src/aipass/prax/apps/modules/monitor.py index 4b04251d..373c7287 100755 --- a/src/aipass/prax/apps/modules/monitor.py +++ b/src/aipass/prax/apps/modules/monitor.py @@ -43,6 +43,12 @@ from aipass.prax.apps.handlers.monitoring import ( ModuleTracker, # module_tracker.py ) from aipass.prax.apps.handlers.monitoring.event_queue import MonitoringEvent +from aipass.prax.apps.handlers.monitoring.telegram_relay import ( + init_relay, + relay_event, + stop_relay, + is_relay_enabled_by_env, +) # ============================================================================= @@ -188,6 +194,10 @@ def print_help(): console.print(" Monitor specific branches (comma-separated)") console.print(" Example: drone @prax monitor run seedgo,cli,flow") console.print() + console.print(" [cyan]drone @prax monitor run --relay[/cyan]") + console.print(" Enable Telegram relay (mirrors feed to prax_monitor bot)") + console.print(" Also enabled by env AIPASS_PRAX_MONITOR_RELAY=1") + console.print() console.print(" [cyan]drone @prax monitor --help[/cyan]") console.print(" Show this help") console.print() @@ -249,6 +259,17 @@ def handle_command(command: str, args: List[str]) -> bool: return True +def _load_relay_config() -> Optional[dict]: + """Load Telegram relay config from @api secrets.""" + try: + from aipass.api.apps.modules.secrets import get_secret + + return get_secret("telegram", "prax_monitor", as_json=True) + except Exception as e: + logger.info("[monitor] Could not load relay config: %s", e) + return None + + def _run_monitor(args: List[str]) -> bool: """Launch Mission Control live monitoring.""" global _event_queue, _module_tracker @@ -262,6 +283,14 @@ def _run_monitor(args: List[str]) -> bool: _module_tracker = ModuleTracker() _stop_event.clear() + # Initialize Telegram relay (--relay flag or env var) + _relay_enabled = "--relay" in args or is_relay_enabled_by_env() + if _relay_enabled: + args = [a for a in args if a != "--relay"] + init_relay(_relay_enabled, _load_relay_config() if _relay_enabled else None) + if _relay_enabled: + console.print("[green]monitor → Telegram relay ON (prax_monitor)[/green]") + _is_tty = sys.stdin.isatty() # Display header @@ -311,10 +340,11 @@ def _start_threads(): def _stop_threads(): - """Stop all monitoring threads""" + """Stop all monitoring threads and Telegram relay""" global _event_queue _stop_event.set() + stop_relay() if _event_queue: _event_queue.stop() @@ -328,7 +358,7 @@ def _stop_threads(): def _render_event(event) -> None: - """Render a single monitoring event to the console.""" + """Render a single monitoring event to the console, and relay to Telegram.""" branch_pid = _get_pid_for_branch(event.branch) if event.event_type == "command": @@ -344,6 +374,8 @@ def _render_event(event) -> None: else: print_event(event.event_type, event.branch, event.message, event.level, pid=branch_pid) + relay_event(event) + def _display_worker(): """Display thread - pulls events from queue and displays them. No filtering.""" diff --git a/src/aipass/prax/prax-monitor.service b/src/aipass/prax/prax-monitor.service new file mode 100644 index 00000000..a1fd6b4b --- /dev/null +++ b/src/aipass/prax/prax-monitor.service @@ -0,0 +1,34 @@ +# Systemd user service for the prax monitor with Telegram relay. +# +# Install: +# cp prax-monitor.service ~/.config/systemd/user/ +# systemctl --user daemon-reload +# +# Usage: +# systemctl --user start prax-monitor +# systemctl --user enable prax-monitor # auto-start on login +# systemctl --user status prax-monitor + +[Unit] +Description=AIPass Prax Monitor — Telegram relay +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +# NOTE: the module __main__ takes a single positional (the subcommand); 'run' → +# _run_monitor([]) = all branches. The relay is enabled by the env var below +# (the __main__ argparse rejects extra tokens like 'all --relay'). +ExecStart=%h/Projects/AIPass/.venv/bin/python3 -m aipass.prax.apps.modules.monitor run +WorkingDirectory=%h/Projects/AIPass +Environment=AIPASS_PRAX_MONITOR_RELAY=1 +Restart=always +RestartSec=5 +# IMPORTANT: log OUTSIDE system_logs/ — the monitor tails system_logs/*.log and +# @trigger watches it too. Writing the monitor's own output there creates a +# feedback loop (monitor output → re-tailed/recorded → reported → more output). +StandardOutput=append:%h/.aipass/prax-monitor.service.log +StandardError=append:%h/.aipass/prax-monitor.service.log + +[Install] +WantedBy=default.target diff --git a/src/aipass/prax/tests/test_monitor_module.py b/src/aipass/prax/tests/test_monitor_module.py index 4beffded..31b7332f 100644 --- a/src/aipass/prax/tests/test_monitor_module.py +++ b/src/aipass/prax/tests/test_monitor_module.py @@ -40,6 +40,7 @@ _MONITORING_MOCKS = { "aipass.prax.apps.handlers.monitoring.interactive_filter": MagicMock(), "aipass.prax.apps.handlers.monitoring.monitoring_filters": MagicMock(), "aipass.prax.apps.handlers.monitoring.file_watcher_integration": MagicMock(), + "aipass.prax.apps.handlers.monitoring.telegram_relay": MagicMock(), } diff --git a/src/aipass/prax/tests/test_telegram_relay.py b/src/aipass/prax/tests/test_telegram_relay.py new file mode 100644 index 00000000..fdba8e93 --- /dev/null +++ b/src/aipass/prax/tests/test_telegram_relay.py @@ -0,0 +1,411 @@ +# =================== AIPass ==================== +# Name: test_telegram_relay.py +# Description: Tests for the Telegram relay handler +# Version: 1.0.0 +# Created: 2026-06-24 +# Modified: 2026-06-24 +# ============================================= + +"""Tests for apps/handlers/monitoring/telegram_relay.py + +Covers: +- Event formatting (log, command, hook types, PID labels) +- init_relay: disabled path, missing config, incomplete config, successful start +- Fail-silent-once: exactly one log line when config absent, zero sends +- stop_relay: final flush and thread join, no-op when inactive +- relay_event: buffering when active, no-op when inactive +- Batching: 4000-char split across messages +- 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 +""" + +import importlib +import sys +from dataclasses import dataclass, field +from datetime import datetime +from typing import Optional +from unittest.mock import MagicMock, patch + + +@dataclass +class FakeEvent: + """Minimal MonitoringEvent stand-in for tests.""" + + priority: int = 3 + timestamp: datetime = field(default_factory=datetime.now) + event_type: str = "" + branch: str = "" + action: str = "" + message: str = "" + level: str = "info" + caller: Optional[str] = None + pid: Optional[int] = None + + +def _import_relay(): + """Import (or reload) telegram_relay with mocked dependencies.""" + fresh_mocks = { + "aipass.prax.apps.modules.logger": MagicMock(), + "aipass.prax.apps.handlers.json": MagicMock(), + "aipass.prax.apps.handlers.json.json_handler": MagicMock(), + } + with patch.dict(sys.modules, fresh_mocks): + if "aipass.prax.apps.handlers.monitoring.telegram_relay" in sys.modules: + mod = importlib.reload(sys.modules["aipass.prax.apps.handlers.monitoring.telegram_relay"]) + else: + mod = importlib.import_module("aipass.prax.apps.handlers.monitoring.telegram_relay") + + 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) + return mod + + +# --------------------------------------------------------------------------- +# Event formatting +# --------------------------------------------------------------------------- + + +class TestFormatEvent: + """Test _format_event produces correct plain-text lines.""" + + def test_log_event_basic(self): + """Log event includes uppercased branch and message.""" + relay = _import_relay() + event = FakeEvent(event_type="log", branch="seedgo", message="Audit started") + result = relay._format_event(event) + assert "[SEEDGO]" in result + assert "Audit started" in result + + def test_log_event_with_pid(self): + """Log event with PID shows BRANCH:PID label.""" + relay = _import_relay() + event = FakeEvent(event_type="log", branch="devpulse", message="Working", pid=12345) + result = relay._format_event(event) + assert "[DEVPULSE:12345]" in result + + def test_command_event(self): + """Command event shows arrow prefix.""" + relay = _import_relay() + event = FakeEvent(event_type="command", branch="drone", message="seedgo audit", action="") + result = relay._format_event(event) + assert "▶ seedgo audit" in result + + def test_command_event_with_caller_and_target(self): + """Command event with caller and target shows attribution line.""" + relay = _import_relay() + event = FakeEvent( + event_type="command", + branch="drone", + message="seedgo audit", + action="run:prax", + caller="devpulse", + ) + result = relay._format_event(event) + assert "devpulse → prax" in result + assert "▶ seedgo audit" in result + + def test_command_event_unknown_caller_omitted(self): + """Command event with UNKNOWN caller omits caller line.""" + relay = _import_relay() + event = FakeEvent( + event_type="command", + branch="drone", + message="test", + caller="UNKNOWN", + ) + result = relay._format_event(event) + assert "UNKNOWN" not in result + + def test_hook_event_fired(self): + """Fired hook event uses lightning bolt symbol.""" + relay = _import_relay() + event = FakeEvent(event_type="hook", branch="hooks", message="cadence:fired", action="fired") + result = relay._format_event(event) + assert result == "⚡ HOOK cadence:fired" + + def test_hook_event_skipped(self): + """Skipped hook event uses dot symbol.""" + relay = _import_relay() + event = FakeEvent(event_type="hook", branch="hooks", message="cadence:skipped", action="skipped") + result = relay._format_event(event) + assert result == "· HOOK cadence:skipped" + + def test_hook_event_unknown_action(self): + """Unknown hook action uses question mark symbol.""" + relay = _import_relay() + event = FakeEvent(event_type="hook", branch="hooks", message="something", action="other") + result = relay._format_event(event) + assert result.startswith("? HOOK") + + def test_file_event(self): + """File event formats like a log event with branch and message.""" + relay = _import_relay() + event = FakeEvent(event_type="file", branch="prax", message="monitor.py modified") + result = relay._format_event(event) + assert "[PRAX]" in result + assert "monitor.py modified" in result + + +# --------------------------------------------------------------------------- +# init_relay +# --------------------------------------------------------------------------- + + +class TestInitRelay: + """Test init_relay activation and fail-silent behavior.""" + + def test_disabled_is_noop(self): + """Disabled flag skips all initialization.""" + relay = _import_relay() + relay.init_relay(enabled=False) + assert relay._RELAY_ACTIVE is False + assert relay._thread is None + + def test_no_config_stays_inactive(self): + """None config keeps relay inactive.""" + relay = _import_relay() + relay.init_relay(enabled=True, config=None) + assert relay._RELAY_ACTIVE is False + + def test_no_config_logs_exactly_once(self): + """Missing config produces exactly one info log, not per-cycle spam.""" + relay = _import_relay() + relay.init_relay(enabled=True, config=None) + calls = [c for c in relay.logger.info.call_args_list if "inactive" in str(c)] + assert len(calls) == 1 + + def test_incomplete_config_missing_chat_id(self): + """Config without chat_id stays inactive.""" + relay = _import_relay() + relay.init_relay(enabled=True, config={"bot_token": "tok123"}) + assert relay._RELAY_ACTIVE is False + + def test_incomplete_config_missing_token(self): + """Config without bot_token stays inactive.""" + relay = _import_relay() + relay.init_relay(enabled=True, config={"chat_id": 123}) + assert relay._RELAY_ACTIVE is False + + def test_valid_config_activates(self): + """Valid config sets active flag and stores credentials.""" + relay = _import_relay() + relay.init_relay(enabled=True, config={"bot_token": "tok123", "chat_id": 456}) + assert relay._RELAY_ACTIVE is True + assert relay._bot_token == "tok123" + assert relay._chat_id == 456 + assert relay._thread is not None + relay.stop_relay() + + def test_valid_config_starts_daemon_thread(self): + """Valid config starts a named daemon thread.""" + relay = _import_relay() + relay.init_relay(enabled=True, config={"bot_token": "t", "chat_id": 1}) + assert relay._thread.daemon is True + assert relay._thread.name == "telegram-relay" + relay.stop_relay() + + +# --------------------------------------------------------------------------- +# relay_event +# --------------------------------------------------------------------------- + + +class TestRelayEvent: + """Test relay_event buffering.""" + + def test_noop_when_inactive(self): + """Inactive relay does not buffer events.""" + relay = _import_relay() + event = FakeEvent(event_type="log", branch="test", message="hello") + relay.relay_event(event) + assert len(relay._buffer) == 0 + + def test_buffers_when_active(self): + """Active relay appends formatted line to buffer.""" + relay = _import_relay() + setattr(relay, "_RELAY_ACTIVE", True) + event = FakeEvent(event_type="log", branch="test", message="hello") + relay.relay_event(event) + assert len(relay._buffer) == 1 + assert "hello" in relay._buffer[0] + + +# --------------------------------------------------------------------------- +# stop_relay +# --------------------------------------------------------------------------- + + +class TestStopRelay: + """Test stop_relay cleanup.""" + + def test_noop_when_inactive(self): + """Stopping an inactive relay is a safe no-op.""" + relay = _import_relay() + relay.stop_relay() + assert relay._RELAY_ACTIVE is False + + def test_flushes_and_joins(self): + """Stopping an active relay clears state and joins thread.""" + relay = _import_relay() + relay.init_relay(enabled=True, config={"bot_token": "t", "chat_id": 1}) + assert relay._RELAY_ACTIVE is True + relay.stop_relay() + assert relay._RELAY_ACTIVE is False + assert relay._thread is None + + +# --------------------------------------------------------------------------- +# Batching +# --------------------------------------------------------------------------- + + +class TestBatching: + """Test _send_batched 4000-char splitting.""" + + def test_single_message_under_limit(self): + """Short content sends as one message.""" + relay = _import_relay() + sent = [] + setattr(relay, "_send_message", lambda text: sent.append(text) or True) + relay._send_batched(["short line"]) + assert len(sent) == 1 + assert sent[0] == "short line" + + def test_splits_at_4000_chars(self): + """Lines exceeding 4000 chars are split across multiple messages.""" + relay = _import_relay() + sent = [] + setattr(relay, "_send_message", lambda text: sent.append(text) or True) + lines = [f"line-{i:04d}-" + "x" * 90 for i in range(50)] + relay._send_batched(lines) + assert len(sent) > 1 + for msg in sent: + assert len(msg) <= relay.TELEGRAM_MAX_LENGTH + 200 + + def test_empty_lines_sends_nothing(self): + """Empty list sends no messages.""" + relay = _import_relay() + sent = [] + setattr(relay, "_send_message", lambda text: sent.append(text) or True) + relay._send_batched([]) + assert len(sent) == 0 + + +# --------------------------------------------------------------------------- +# Flood cap +# --------------------------------------------------------------------------- + + +class TestFloodCap: + """Test flood cap truncation.""" + + def test_under_cap_sends_all(self): + """Lines under FLOOD_CAP are all delivered.""" + relay = _import_relay() + setattr(relay, "_bot_token", "t") + setattr(relay, "_chat_id", 1) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._buffer.extend([f"line {i}" for i in range(100)]) + setattr(relay, "_RELAY_ACTIVE", True) + relay._flush_buffer() + assert len(sent) == 100 + + def test_over_cap_truncates_with_notice(self): + """Lines over FLOOD_CAP are truncated with a suppression notice.""" + relay = _import_relay() + setattr(relay, "_bot_token", "t") + setattr(relay, "_chat_id", 1) + sent_lines = [] + setattr(relay, "_send_batched", lambda lines: sent_lines.extend(lines)) + relay._buffer.extend([f"line {i}" for i in range(200)]) + setattr(relay, "_RELAY_ACTIVE", True) + relay._flush_buffer() + assert len(sent_lines) == relay.FLOOD_CAP + 1 + assert "50 more suppressed" in sent_lines[-1] + + +# --------------------------------------------------------------------------- +# _render_event calls relay_event +# --------------------------------------------------------------------------- + + +class TestRenderEventCallsRelay: + """Test that monitor._render_event calls relay_event.""" + + def test_render_event_calls_relay(self): + """_render_event in monitor.py calls relay_event after console render.""" + fresh_mocks = { + "aipass.prax.apps.handlers.monitoring": MagicMock(), + "aipass.prax.apps.handlers.monitoring.event_queue": MagicMock(), + "aipass.prax.apps.handlers.monitoring.filesystem_handler": MagicMock(), + "aipass.prax.apps.handlers.monitoring.log_watcher": MagicMock(), + "aipass.prax.apps.handlers.monitoring.unified_stream": MagicMock(), + "aipass.prax.apps.handlers.monitoring.module_tracker": MagicMock(), + "aipass.prax.apps.handlers.monitoring.branch_detector": MagicMock(), + "aipass.prax.apps.handlers.monitoring.interactive_filter": MagicMock(), + "aipass.prax.apps.handlers.monitoring.monitoring_filters": MagicMock(), + "aipass.prax.apps.handlers.monitoring.file_watcher_integration": MagicMock(), + "aipass.prax.apps.handlers.monitoring.telegram_relay": MagicMock(), + } + with patch.dict(sys.modules, fresh_mocks): + if "aipass.prax.apps.modules.monitor" in sys.modules: + mod = importlib.reload(sys.modules["aipass.prax.apps.modules.monitor"]) + else: + mod = importlib.import_module("aipass.prax.apps.modules.monitor") + + event = MagicMock() + event.event_type = "log" + event.branch = "TEST" + event.message = "hello" + event.level = "info" + event.action = "" + + with patch.object(mod, "_get_pid_for_branch", return_value=None): + mod._render_event(event) + + mod.relay_event.assert_called_once_with(event) + + +# --------------------------------------------------------------------------- +# is_relay_enabled_by_env +# --------------------------------------------------------------------------- + + +class TestRelayEnabledByEnv: + """Test environment variable detection.""" + + def test_not_set(self): + """Unset env var returns False.""" + relay = _import_relay() + with patch.dict("os.environ", {}, clear=True): + assert relay.is_relay_enabled_by_env() is False + + def test_set_to_1(self): + """AIPASS_PRAX_MONITOR_RELAY=1 returns True.""" + relay = _import_relay() + with patch.dict("os.environ", {"AIPASS_PRAX_MONITOR_RELAY": "1"}): + assert relay.is_relay_enabled_by_env() is True + + def test_set_to_true(self): + """AIPASS_PRAX_MONITOR_RELAY=true returns True.""" + relay = _import_relay() + with patch.dict("os.environ", {"AIPASS_PRAX_MONITOR_RELAY": "true"}): + assert relay.is_relay_enabled_by_env() is True + + def test_set_to_yes(self): + """AIPASS_PRAX_MONITOR_RELAY=yes returns True.""" + relay = _import_relay() + with patch.dict("os.environ", {"AIPASS_PRAX_MONITOR_RELAY": "yes"}): + assert relay.is_relay_enabled_by_env() is True + + def test_set_to_0(self): + """AIPASS_PRAX_MONITOR_RELAY=0 returns False.""" + relay = _import_relay() + with patch.dict("os.environ", {"AIPASS_PRAX_MONITOR_RELAY": "0"}): + assert relay.is_relay_enabled_by_env() is False