diff --git a/CHANGELOG.md b/CHANGELOG.md index a107c975..f0db794f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,23 @@ PyPI version — not the changelog header. ## [2026-07-12] +### Added + +- **Telegram log-stream control: `/logs` on branch bots + interactive Prax + Monitor chat.** The per-branch session LogStreamer auto-started on first + message hardwired to full firehose with no off switch; the Prax Monitor + relay chat was send-only — no command menu, and anything typed there was + silently never read (nothing polled that token). @skills added `/logs + on|errors|off|status` to all branch bots (preference persisted per chat, + honored by the auto-start; 33 tests) and a new `PraxMonitorBot` receiver + service (`telegram-bot@prax_monitor`) with `/pause /resume /errors /all + /status` and a registered command menu (34 tests). @prax made the relay + honor the shared control file (`~/.aipass/telegram_bots/ + prax_monitor_control.json`, frozen contract: paused + level) each 5s flush — + paused discards, `errors` filters to WARNING/ERROR/CRITICAL (17 tests). + Live-verified end-to-end from Telegram Web: `/errors` silenced INFO batches + within one flush, `/all` restored them. + ### Fixed - **Legacy `builder` citizen_class migration + birth-certificate template diff --git a/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py b/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py index 6def5e8f..67372456 100644 --- a/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py +++ b/src/aipass/prax/apps/handlers/monitoring/telegram_relay.py @@ -18,6 +18,7 @@ import json import os import threading from datetime import datetime +from pathlib import Path from typing import Optional from urllib.error import URLError from urllib.request import Request @@ -32,6 +33,9 @@ BATCH_INTERVAL = 5.0 TELEGRAM_MAX_LENGTH = 4000 FLOOD_CAP = 150 +CONTROL_FILE = Path.home() / ".aipass" / "telegram_bots" / "prax_monitor_control.json" +_ERROR_MARKERS = ("WARNING", "ERROR", "CRITICAL") + _lock = threading.Lock() _buffer: list[str] = [] _thread: Optional[threading.Thread] = None @@ -39,6 +43,8 @@ _stop_event = threading.Event() _bot_token: Optional[str] = None _chat_id: Optional[int] = None _RELAY_ACTIVE = False +_control_mtime: float = 0.0 +_control_cache: dict = {} def init_relay(enabled: bool, config: Optional[dict] = None) -> None: @@ -143,8 +149,41 @@ def _format_event(event) -> Optional[str]: return f"[{ts}] [{branch_label}] {event.message}" +def _read_control() -> dict: + """Read the control file, caching by mtime. Returns defaults on missing/invalid file.""" + global _control_mtime, _control_cache + + try: + stat = CONTROL_FILE.stat() + except OSError: + logger.info("[telegram_relay] Control file not found, using defaults") + return {} + + if stat.st_mtime == _control_mtime: + return _control_cache + + try: + data = json.loads(CONTROL_FILE.read_text(encoding="utf-8")) + if not isinstance(data, dict): + raise ValueError("not a dict") + _control_mtime = stat.st_mtime + _control_cache = data + return data + except (json.JSONDecodeError, ValueError, OSError) as exc: + logger.warning("[telegram_relay] Control file parse error, using defaults: %s", exc) + _control_mtime = stat.st_mtime + _control_cache = {} + return {} + + def _flush_buffer() -> None: - """Drain buffer and send to Telegram.""" + """Drain buffer, apply control-file pause/filter, and send to Telegram. + + Control semantics (written by the @skills TG bot): + - paused=true: buffer is discarded, nothing sent. + - level="errors": only lines containing WARNING/ERROR/CRITICAL are sent. + - level="all" (default): everything is sent. + """ with _lock: if not _buffer: return @@ -154,6 +193,17 @@ def _flush_buffer() -> None: if not _bot_token or not _chat_id: return + ctrl = _read_control() + + if ctrl.get("paused", False): + return + + level = ctrl.get("level", "all") + if level == "errors": + lines = [ln for ln in lines if any(m in ln for m in _ERROR_MARKERS)] + if not lines: + return + if len(lines) > FLOOD_CAP: suppressed = len(lines) - FLOOD_CAP lines = lines[:FLOOD_CAP] diff --git a/src/aipass/prax/tests/test_telegram_relay.py b/src/aipass/prax/tests/test_telegram_relay.py index 71627774..2168c1b8 100644 --- a/src/aipass/prax/tests/test_telegram_relay.py +++ b/src/aipass/prax/tests/test_telegram_relay.py @@ -21,6 +21,7 @@ Covers: """ import importlib +import json import sys from dataclasses import dataclass, field from datetime import datetime @@ -411,3 +412,169 @@ class TestRelayEnabledByEnv: relay = _import_relay() with patch.dict("os.environ", {"AIPASS_PRAX_MONITOR_RELAY": "0"}): assert relay.is_relay_enabled_by_env() is False + + +# --------------------------------------------------------------------------- +# Control file: _read_control +# --------------------------------------------------------------------------- + + +class TestReadControl: + """Test _read_control reads, caches, and handles errors.""" + + def test_missing_file_returns_empty(self, tmp_path): + """Missing control file returns empty dict (defaults).""" + relay = _import_relay() + setattr(relay, "CONTROL_FILE", tmp_path / "nonexistent.json") + assert relay._read_control() == {} + + def test_valid_file_returns_content(self, tmp_path): + """Valid JSON control file is read and returned.""" + relay = _import_relay() + ctrl = tmp_path / "control.json" + ctrl.write_text(json.dumps({"paused": True, "level": "errors"})) + setattr(relay, "CONTROL_FILE", ctrl) + result = relay._read_control() + assert result["paused"] is True + assert result["level"] == "errors" + + def test_mtime_cache_avoids_reread(self, tmp_path): + """Same mtime returns cached result without re-reading the file.""" + relay = _import_relay() + ctrl = tmp_path / "control.json" + ctrl.write_text(json.dumps({"paused": False})) + setattr(relay, "CONTROL_FILE", ctrl) + first = relay._read_control() + ctrl.write_text("INVALID JSON") + result = relay._read_control() + assert result == first + + def test_mtime_change_triggers_reread(self, tmp_path): + """Changed mtime causes re-read of the control file.""" + import os + + relay = _import_relay() + ctrl = tmp_path / "control.json" + ctrl.write_text(json.dumps({"paused": False, "level": "all"})) + setattr(relay, "CONTROL_FILE", ctrl) + relay._read_control() + ctrl.write_text(json.dumps({"paused": True, "level": "errors"})) + os.utime(ctrl, (ctrl.stat().st_mtime + 1, ctrl.stat().st_mtime + 1)) + result = relay._read_control() + assert result["paused"] is True + assert result["level"] == "errors" + + def test_invalid_json_returns_empty_and_warns(self, tmp_path): + """Malformed JSON returns empty dict and logs a warning.""" + relay = _import_relay() + ctrl = tmp_path / "control.json" + ctrl.write_text("{bad json!!!") + setattr(relay, "CONTROL_FILE", ctrl) + result = relay._read_control() + assert result == {} + relay.logger.warning.assert_called_once() + + def test_non_dict_json_returns_empty_and_warns(self, tmp_path): + """JSON that isn't a dict returns empty and logs warning.""" + relay = _import_relay() + ctrl = tmp_path / "control.json" + ctrl.write_text(json.dumps([1, 2, 3])) + setattr(relay, "CONTROL_FILE", ctrl) + result = relay._read_control() + assert result == {} + relay.logger.warning.assert_called_once() + + +# --------------------------------------------------------------------------- +# Control file: flush behavior with pause / level filter +# --------------------------------------------------------------------------- + + +class TestFlushControl: + """Test _flush_buffer honors control file pause and level filtering.""" + + def _make_relay(self, tmp_path, control_data=None): + """Helper: import relay, wire credentials, point control file at tmp_path.""" + relay = _import_relay() + setattr(relay, "_bot_token", "t") + setattr(relay, "_chat_id", 1) + setattr(relay, "_RELAY_ACTIVE", True) + ctrl = tmp_path / "control.json" + if control_data is not None: + ctrl.write_text(json.dumps(control_data)) + setattr(relay, "CONTROL_FILE", ctrl) + return relay + + def test_paused_discards_buffer(self, tmp_path): + """Paused=true discards buffered lines, nothing sent.""" + relay = self._make_relay(tmp_path, {"paused": True, "level": "all"}) + relay._buffer.extend(["line 1", "line 2"]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert sent == [] + assert len(relay._buffer) == 0 + + def test_unpaused_sends_all(self, tmp_path): + """Paused=false with level=all sends everything.""" + relay = self._make_relay(tmp_path, {"paused": False, "level": "all"}) + relay._buffer.extend(["info line", "WARNING alert"]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert len(sent) == 2 + + def test_level_errors_filters_info_lines(self, tmp_path): + """Level=errors drops lines without WARNING/ERROR/CRITICAL markers.""" + relay = self._make_relay(tmp_path, {"paused": False, "level": "errors"}) + relay._buffer.extend( + [ + "normal info line", + "[10:00:00] [PRAX] WARNING disk full", + "just a log", + "[10:00:01] [FLOW] ERROR crash", + "[10:00:02] [CLI] CRITICAL meltdown", + ] + ) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert len(sent) == 3 + assert all(any(m in ln for m in ("WARNING", "ERROR", "CRITICAL")) for ln in sent) + + def test_level_errors_all_filtered_sends_nothing(self, tmp_path): + """Level=errors with no matching lines sends nothing.""" + relay = self._make_relay(tmp_path, {"paused": False, "level": "errors"}) + relay._buffer.extend(["info line 1", "info line 2"]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert sent == [] + + def test_missing_control_file_sends_all(self, tmp_path): + """Missing control file = defaults (not paused, level=all).""" + relay = self._make_relay(tmp_path) + relay._buffer.extend(["line 1", "line 2"]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert len(sent) == 2 + + def test_missing_keys_use_defaults(self, tmp_path): + """Control file with empty dict = not paused, level=all.""" + relay = self._make_relay(tmp_path, {}) + relay._buffer.extend(["line 1"]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert len(sent) == 1 + + def test_flood_cap_still_applies_after_filter(self, tmp_path): + """FLOOD_CAP is enforced after level filtering.""" + relay = self._make_relay(tmp_path, {"paused": False, "level": "all"}) + relay._buffer.extend([f"line {i}" for i in range(200)]) + sent = [] + setattr(relay, "_send_batched", lambda lines: sent.extend(lines)) + relay._flush_buffer() + assert len(sent) == relay.FLOOD_CAP + 1 + assert "suppressed" in sent[-1] diff --git a/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py b/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py index 451dc12d..43c5b9e8 100644 --- a/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py +++ b/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py @@ -509,9 +509,12 @@ class BaseBot: self._active_chat_id = chat_id self._write_mirror_mapping() if self.branch_name is not None and self._log_streamer is None: - self._log_streamer = LogStreamer(self.bot_token, chat_id, self.branch_name) - self._log_streamer.start() - logger.info("Log streamer started for branch: %s", self.branch_name) + pref = self._load_logs_preference() + pref_mode = pref.get("mode", "all") if pref else "all" + if pref_mode != "off": + self._log_streamer = LogStreamer(self.bot_token, chat_id, self.branch_name, level_filter=pref_mode) + self._log_streamer.start() + logger.info("Log streamer started for branch: %s (mode=%s)", self.branch_name, pref_mode) # Allowlist check if not self.is_user_allowed(user_id): @@ -571,6 +574,11 @@ class BaseBot: """ cmd_name, cmd_args = parsed + # /logs command — session log stream control + if cmd_name == "logs": + self._handle_logs_command(chat_id, cmd_args) + return True + # /monitor command — system-wide log subscription if cmd_name == "monitor": self._handle_monitor_command(chat_id, cmd_args) @@ -2327,6 +2335,121 @@ class BaseBot: logger.error("Failed to clear monitor subscription: %s", e) return False + # ============================================= + # /LOGS — SESSION LOG STREAM CONTROL + # ============================================= + + def _handle_logs_command(self, chat_id: int, args: str) -> None: + """Route /logs subcommands: on, off, errors, status.""" + if self.branch_name is None: + self.send_message(chat_id, "Not available — this bot has no branch log stream.") + return + + subcmd = args.strip().lower().split()[0] if args.strip() else "" + + if subcmd == "on": + self._logs_start(chat_id, mode="all") + elif subcmd == "off": + self._logs_stop(chat_id) + elif subcmd == "errors": + self._logs_start(chat_id, mode="default") + elif subcmd == "status": + self._logs_status(chat_id) + else: + self.send_message( + chat_id, + "/logs on — full log stream\n" + "/logs errors — warnings & errors only\n" + "/logs off — stop streaming\n" + "/logs status — current state", + ) + + def _logs_start(self, chat_id: int, mode: str) -> None: + """Start or restart the session log streamer with the given mode.""" + if self._log_streamer is not None: + self._log_streamer.stop() + self._log_streamer = None + + if not self._save_logs_preference(chat_id, mode): + self.send_message(chat_id, "Failed to save logs preference.") + return + + self._log_streamer = LogStreamer( + self.bot_token, + chat_id, + self.branch_name, # type: ignore[arg-type] # guarded by _handle_logs_command + level_filter=mode, + ) + self._log_streamer.start() + + mode_label = "all levels" if mode == "all" else "errors & warnings" + self.send_message( + chat_id, + f"Log streaming: {mode_label}\n\n/logs off to stop\n/logs errors for filtered mode", + ) + logger.info("Log streaming started: chat_id=%s, mode=%s, branch=%s", chat_id, mode, self.branch_name) + + def _logs_stop(self, chat_id: int) -> None: + """Stop the session log streamer.""" + if self._log_streamer is not None: + self._log_streamer.stop() + self._log_streamer = None + + self._save_logs_preference(chat_id, "off") + self.send_message(chat_id, "Log streaming stopped.\n\n/logs on to resume.") + logger.info("Log streaming stopped: chat_id=%s, branch=%s", chat_id, self.branch_name) + + def _logs_status(self, chat_id: int) -> None: + """Show current log streaming status.""" + pref = self._load_logs_preference() + running = self._log_streamer is not None and self._log_streamer._running + + if not running: + mode_info = "" + if pref and pref.get("mode") == "off": + mode_info = " (disabled)" + self.send_message(chat_id, f"Log streaming: stopped{mode_info}\n\n/logs on to start.") + return + + mode = pref.get("mode", "all") if pref else "all" + mode_label = "all levels" if mode == "all" else "errors & warnings" + self.send_message( + chat_id, + f"Log streaming: active\nMode: {mode_label}\nBranch: {self.branch_name}", + ) + + def _logs_preference_file(self) -> Path: + """Return path to the local logs preference file.""" + return Path.home() / ".aipass" / "telegram_bots" / f".{self.bot_id}_logs.json" + + def _load_logs_preference(self) -> dict | None: + """Load logs preference from local state file.""" + pref_file = self._logs_preference_file() + if not pref_file.exists(): + return None + try: + data = json.loads(pref_file.read_text(encoding="utf-8")) + if isinstance(data, dict): + return data + return None + except (json.JSONDecodeError, OSError) as e: + logger.warning("Failed to load logs preference: %s", e) + return None + + def _save_logs_preference(self, chat_id: int, mode: str) -> bool: + """Persist logs preference to local state file.""" + pref_file = self._logs_preference_file() + try: + pref_file.parent.mkdir(parents=True, exist_ok=True) + pref_file.write_text( + json.dumps({"chat_id": chat_id, "mode": mode}, indent=2), + encoding="utf-8", + ) + return True + except OSError as e: + logger.error("Failed to save logs preference: %s", e) + return False + def get_custom_commands(self) -> dict: """ Hook: return additional bot-specific commands. @@ -2343,6 +2466,11 @@ class BaseBot: "menu_text": "Log monitor", }, } + if self.branch_name is not None: + commands["logs"] = { + "description": "Control branch log streaming — /logs on, off, errors, status", + "menu_text": "Log streaming", + } if self.branch_name is None: commands["create"] = { "description": "Create a Telegram bot for a branch — e.g. /create chat devpulse", @@ -2484,6 +2612,7 @@ class BaseBot: _BOT_CLASSES = { "scheduler": ".scheduler_bot:SchedulerBot", + "prax_monitor": ".prax_monitor_bot:PraxMonitorBot", } diff --git a/src/aipass/skills/lib/telegram/apps/handlers/prax_monitor_bot.py b/src/aipass/skills/lib/telegram/apps/handlers/prax_monitor_bot.py new file mode 100644 index 00000000..44a8a397 --- /dev/null +++ b/src/aipass/skills/lib/telegram/apps/handlers/prax_monitor_bot.py @@ -0,0 +1,186 @@ +# =================== AIPass ==================== +# Name: prax_monitor_bot.py +# Description: Telegram bot for the Prax Monitor chat — command receiver for relay control +# Version: 1.0.0 +# Created: 2026-07-12 +# Modified: 2026-07-12 +# ============================================= + +""" +PraxMonitorBot — a BaseBot subclass for the Prax Monitor TG chat. + +Receives commands from the Prax Monitor Telegram chat and writes a control file +that the prax relay reads each flush cycle (~5s). No tmux/Claude sessions. + +Control file: ~/.aipass/telegram_bots/prax_monitor_control.json +Schema: {"paused": bool, "level": "all"|"errors", "updated_at": iso8601} +""" + +import json +import subprocess +from datetime import datetime, timezone +from pathlib import Path + +from aipass.prax import logger + +from .base_bot import BaseBot + + +CONTROL_FILE = Path.home() / ".aipass" / "telegram_bots" / "prax_monitor_control.json" + + +class PraxMonitorBot(BaseBot): + """Prax Monitor command bot — controls the relay via a shared control file.""" + + def handle_message(self, chat_id: int, text: str, message: dict) -> None: + """Reject free-text — this bot only serves commands.""" + self.send_message( + chat_id, + "I only handle commands.\nTry /pause, /resume, /errors, /all, or /status", + ) + + def handle_file(self, chat_id: int, message: dict) -> None: + """Reject files.""" + self.send_message(chat_id, "I don't process files. Try /status or /help") + + def _dispatch_command(self, chat_id: int, parsed: tuple) -> bool: + cmd_name, cmd_args = parsed + if cmd_name == "pause": + self._handle_pause(chat_id) + return True + if cmd_name == "resume": + self._handle_resume(chat_id) + return True + if cmd_name == "errors": + self._handle_errors(chat_id) + return True + if cmd_name == "all": + self._handle_all(chat_id) + return True + if cmd_name == "status": + self._handle_prax_status(chat_id) + return True + return super()._dispatch_command(chat_id, parsed) + + def get_custom_commands(self) -> dict: + cmds = super().get_custom_commands() + cmds["pause"] = { + "description": "Pause the prax log relay", + "menu_text": "Pause relay", + } + cmds["resume"] = { + "description": "Resume the prax log relay", + "menu_text": "Resume relay", + } + cmds["errors"] = { + "description": "Show errors & warnings only", + "menu_text": "Errors only", + } + cmds["all"] = { + "description": "Show all log levels", + "menu_text": "All levels", + } + return cmds + + # ============================================= + # COMMAND HANDLERS + # ============================================= + + def _handle_pause(self, chat_id: int) -> None: + ctrl = self._read_control() + ctrl["paused"] = True + if self._write_control(ctrl): + self.send_message(chat_id, "Relay paused.\n\n/resume to restart.") + else: + self.send_message(chat_id, "Failed to write control file.") + logger.info("Prax monitor paused (chat_id=%s)", chat_id) + + def _handle_resume(self, chat_id: int) -> None: + ctrl = self._read_control() + ctrl["paused"] = False + if self._write_control(ctrl): + level = ctrl.get("level", "all") + self.send_message(chat_id, f"Relay resumed (level: {level}).\n\n/pause to stop.") + else: + self.send_message(chat_id, "Failed to write control file.") + logger.info("Prax monitor resumed (chat_id=%s)", chat_id) + + def _handle_errors(self, chat_id: int) -> None: + ctrl = self._read_control() + ctrl["level"] = "errors" + if self._write_control(ctrl): + self.send_message(chat_id, "Level set to errors & warnings only.\n\n/all for full firehose.") + else: + self.send_message(chat_id, "Failed to write control file.") + logger.info("Prax monitor level=errors (chat_id=%s)", chat_id) + + def _handle_all(self, chat_id: int) -> None: + ctrl = self._read_control() + ctrl["level"] = "all" + if self._write_control(ctrl): + self.send_message(chat_id, "Level set to all.\n\n/errors for filtered mode.") + else: + self.send_message(chat_id, "Failed to write control file.") + logger.info("Prax monitor level=all (chat_id=%s)", chat_id) + + def _handle_prax_status(self, chat_id: int) -> None: + ctrl = self._read_control() + paused = ctrl.get("paused", False) + level = ctrl.get("level", "all") + updated = ctrl.get("updated_at", "never") + + relay_alive = self._check_relay_alive() + relay_status = "running" if relay_alive else "not detected" + + state = "paused" if paused else "active" + level_label = "errors & warnings" if level == "errors" else "all levels" + + self.send_message( + chat_id, + f"Prax Monitor\nState: {state}\nLevel: {level_label}\nRelay: {relay_status}\nLast update: {updated}", + ) + + # ============================================= + # CONTROL FILE I/O + # ============================================= + + def _read_control(self) -> dict: + """Read the control file; return defaults if missing or corrupt.""" + if not CONTROL_FILE.exists(): + return {"paused": False, "level": "all"} + try: + data = json.loads(CONTROL_FILE.read_text(encoding="utf-8")) + if isinstance(data, dict): + return data + return {"paused": False, "level": "all"} + except (json.JSONDecodeError, OSError) as e: + logger.warning("Failed to read control file: %s", e) + return {"paused": False, "level": "all"} + + def _write_control(self, ctrl: dict) -> bool: + """Write the control file with updated_at timestamp.""" + ctrl["updated_at"] = datetime.now(timezone.utc).isoformat() + try: + CONTROL_FILE.parent.mkdir(parents=True, exist_ok=True) + CONTROL_FILE.write_text( + json.dumps(ctrl, indent=2), + encoding="utf-8", + ) + return True + except OSError as e: + logger.error("Failed to write control file: %s", e) + return False + + @staticmethod + def _check_relay_alive() -> bool: + """Check if the prax-monitor systemd service is active.""" + try: + result = subprocess.run( + ["systemctl", "--user", "is-active", "prax-monitor"], + capture_output=True, + text=True, + timeout=5, + ) + return result.stdout.strip() == "active" + except (subprocess.TimeoutExpired, OSError): + return False diff --git a/src/aipass/skills/lib/telegram/tests/test_log_streamer.py b/src/aipass/skills/lib/telegram/tests/test_log_streamer.py index 7d58a64c..43c36149 100644 --- a/src/aipass/skills/lib/telegram/tests/test_log_streamer.py +++ b/src/aipass/skills/lib/telegram/tests/test_log_streamer.py @@ -517,7 +517,7 @@ class TestBaseBotIntegration: bot.process_update(fake_update) - MockStreamer.assert_called_once_with("123:FAKETOKEN", 42, "api") + MockStreamer.assert_called_once_with("123:FAKETOKEN", 42, "api", level_filter="all") mock_instance.start.assert_called_once() def test_streamer_not_started_when_branch_name_is_none(self, tmp_path, _patch_base_bot_deps): diff --git a/src/aipass/skills/lib/telegram/tests/test_logs.py b/src/aipass/skills/lib/telegram/tests/test_logs.py new file mode 100644 index 00000000..6537c12e --- /dev/null +++ b/src/aipass/skills/lib/telegram/tests/test_logs.py @@ -0,0 +1,466 @@ +# =================== AIPass ==================== +# Name: test_logs.py +# Description: Tests for /logs command — session log stream control +# Version: 1.0.0 +# Created: 2026-07-12 +# Modified: 2026-07-12 +# ============================================= + +""" +Tests for /logs command — session log stream control. + +Tests cover: + - /logs on persists preference and starts streamer + - /logs off stops streamer and persists "off" + - /logs errors starts with level_filter="default" + - /logs status shows correct state + - /logs command routing (on, off, errors, status, bare, unknown) + - Auto-start at handle_update honors saved preference + - /logs unavailable on base bot (branch_name=None) + - Persistence roundtrip (save + reload) +""" + +import json +from pathlib import Path + +import pytest +from unittest.mock import patch, MagicMock + + +# ============================================= +# HELPERS +# ============================================= + + +@pytest.fixture +def _patch_base_bot_deps(tmp_path): + """Patch heavy BaseBot dependencies to allow lightweight instantiation.""" + pref_file = tmp_path / "logs_pref.json" + patches = [ + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.PENDING_DIR", tmp_path), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.signal.signal"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.atexit.register"), + ] + for p in patches: + p.start() + yield pref_file + for p in patches: + p.stop() + + +def _make_bot(tmp_path, _patch_base_bot_deps, branch_name: str | None = "testbranch"): + """Create a BaseBot with logs preference redirected to tmp_path.""" + from aipass.skills.lib.telegram.apps.handlers.base_bot import BaseBot + + workdir = tmp_path / "workdir" + workdir.mkdir(exist_ok=True) + bot = BaseBot( + bot_id="logs_test", + bot_token="123:FAKETOKEN", + work_dir=workdir, + bot_name="Logs Test Bot", + allowed_user_ids=[111], + branch_name=branch_name, + ) + pref_file: Path = _patch_base_bot_deps + bot._logs_preference_file = lambda: pref_file # type: ignore[assignment] + return bot + + +# ============================================= +# 1. /LOGS ON — PERSIST + START STREAMER +# ============================================= + + +class TestLogsOn: + """Verify /logs on persists preference and starts streamer.""" + + def test_on_writes_preference(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + with ( + patch.object(bot, "send_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer") as MockStreamer, + ): + MockStreamer.return_value = MagicMock() + bot._logs_start(42, "all") + + data = json.loads(pref_file.read_text()) + assert data == {"chat_id": 42, "mode": "all"} + + def test_on_starts_streamer(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "send_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer") as MockStreamer, + ): + mock_instance = MagicMock() + MockStreamer.return_value = mock_instance + + bot._logs_start(42, "all") + + MockStreamer.assert_called_once_with( + "123:FAKETOKEN", + 42, + "testbranch", + level_filter="all", + ) + mock_instance.start.assert_called_once() + assert bot._log_streamer is mock_instance + + def test_on_sends_confirmation(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "send_message") as mock_send, + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer", return_value=MagicMock()), + ): + bot._logs_start(42, "all") + mock_send.assert_called_once() + msg = mock_send.call_args[0][1] + assert "all levels" in msg + + def test_on_stops_existing_streamer(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + old_streamer = MagicMock() + bot._log_streamer = old_streamer + + with ( + patch.object(bot, "send_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer", return_value=MagicMock()), + ): + bot._logs_start(42, "all") + old_streamer.stop.assert_called_once() + + def test_on_aborts_on_save_failure(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "_save_logs_preference", return_value=False), + patch.object(bot, "send_message") as mock_send, + ): + bot._logs_start(42, "all") + msg = mock_send.call_args[0][1] + assert "Failed" in msg + assert bot._log_streamer is None + + +# ============================================= +# 2. /LOGS ERRORS — FILTERED MODE +# ============================================= + + +class TestLogsErrors: + """Verify /logs errors starts with level_filter="default".""" + + def test_errors_starts_with_default_filter(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "send_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer") as MockStreamer, + ): + MockStreamer.return_value = MagicMock() + bot._logs_start(42, "default") + + MockStreamer.assert_called_once_with( + "123:FAKETOKEN", + 42, + "testbranch", + level_filter="default", + ) + + def test_errors_confirmation_message(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "send_message") as mock_send, + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer", return_value=MagicMock()), + ): + bot._logs_start(42, "default") + msg = mock_send.call_args[0][1] + assert "errors & warnings" in msg + + def test_errors_persists_mode(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + with ( + patch.object(bot, "send_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer", return_value=MagicMock()), + ): + bot._logs_start(42, "default") + data = json.loads(pref_file.read_text()) + assert data["mode"] == "default" + + +# ============================================= +# 3. /LOGS OFF — STOP + PERSIST +# ============================================= + + +class TestLogsOff: + """Verify /logs off stops streamer and persists 'off'.""" + + def test_off_stops_streamer(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + mock_streamer = MagicMock() + bot._log_streamer = mock_streamer + + with patch.object(bot, "send_message"): + bot._logs_stop(42) + + mock_streamer.stop.assert_called_once() + assert bot._log_streamer is None + + def test_off_persists_off(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + + with patch.object(bot, "send_message"): + bot._logs_stop(42) + + data = json.loads(pref_file.read_text()) + assert data["mode"] == "off" + + def test_off_sends_confirmation(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._logs_stop(42) + + mock_send.assert_called_once() + assert "stopped" in mock_send.call_args[0][1].lower() + + def test_off_safe_when_no_streamer(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + assert bot._log_streamer is None + + with patch.object(bot, "send_message"): + bot._logs_stop(42) + + assert bot._log_streamer is None + + +# ============================================= +# 4. /LOGS STATUS +# ============================================= + + +class TestLogsStatus: + """Verify /logs status shows correct state.""" + + def test_status_when_stopped(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._logs_status(42) + msg = mock_send.call_args[0][1] + assert "stopped" in msg + + def test_status_when_disabled(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + pref_file.write_text(json.dumps({"chat_id": 42, "mode": "off"})) + + with patch.object(bot, "send_message") as mock_send: + bot._logs_status(42) + msg = mock_send.call_args[0][1] + assert "disabled" in msg + + def test_status_when_active(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + pref_file.write_text(json.dumps({"chat_id": 42, "mode": "all"})) + bot._log_streamer = MagicMock(_running=True) + + with patch.object(bot, "send_message") as mock_send: + bot._logs_status(42) + msg = mock_send.call_args[0][1] + assert "active" in msg + assert "all levels" in msg + + def test_status_shows_errors_mode(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + pref_file.write_text(json.dumps({"chat_id": 42, "mode": "default"})) + bot._log_streamer = MagicMock(_running=True) + + with patch.object(bot, "send_message") as mock_send: + bot._logs_status(42) + msg = mock_send.call_args[0][1] + assert "errors & warnings" in msg + + def test_status_shows_branch_name(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + pref_file.write_text(json.dumps({"chat_id": 42, "mode": "all"})) + bot._log_streamer = MagicMock(_running=True) + + with patch.object(bot, "send_message") as mock_send: + bot._logs_status(42) + msg = mock_send.call_args[0][1] + assert "testbranch" in msg + + +# ============================================= +# 5. COMMAND ROUTING +# ============================================= + + +class TestLogsCommandRouting: + """Verify _handle_logs_command routes subcommands correctly.""" + + def test_on_routes_to_start_all(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_logs_start") as mock_start: + bot._handle_logs_command(42, "on") + mock_start.assert_called_once_with(42, mode="all") + + def test_errors_routes_to_start_default(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_logs_start") as mock_start: + bot._handle_logs_command(42, "errors") + mock_start.assert_called_once_with(42, mode="default") + + def test_off_routes_to_stop(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_logs_stop") as mock_stop: + bot._handle_logs_command(42, "off") + mock_stop.assert_called_once_with(42) + + def test_status_routes_to_status(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_logs_status") as mock_stat: + bot._handle_logs_command(42, "status") + mock_stat.assert_called_once_with(42) + + def test_bare_logs_shows_help(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._handle_logs_command(42, "") + msg = mock_send.call_args[0][1] + assert "/logs on" in msg + assert "/logs off" in msg + assert "/logs errors" in msg + + def test_unknown_subcommand_shows_help(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._handle_logs_command(42, "banana") + msg = mock_send.call_args[0][1] + assert "/logs on" in msg + + +# ============================================= +# 6. BASE BOT GUARD (branch_name=None) +# ============================================= + + +class TestLogsBaseBotGuard: + """Verify /logs is unavailable on the base bot.""" + + def test_no_branch_shows_unavailable(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps, branch_name=None) + with patch.object(bot, "send_message") as mock_send: + bot._handle_logs_command(42, "on") + msg = mock_send.call_args[0][1] + assert "not available" in msg.lower() + + def test_no_branch_does_not_start_streamer(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps, branch_name=None) + with patch.object(bot, "send_message"): + bot._handle_logs_command(42, "on") + assert bot._log_streamer is None + + def test_logs_in_custom_commands_only_for_branch_bots(self, tmp_path, _patch_base_bot_deps): + branch_bot = _make_bot(tmp_path, _patch_base_bot_deps, branch_name="mybranch") + assert "logs" in branch_bot.get_custom_commands() + + base_bot = _make_bot(tmp_path, _patch_base_bot_deps, branch_name=None) + assert "logs" not in base_bot.get_custom_commands() + + +# ============================================= +# 7. PERSISTENCE ROUNDTRIP +# ============================================= + + +class TestLogsPersistence: + """Verify save + load roundtrip.""" + + def test_save_and_load(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + bot._save_logs_preference(42, "default") + + result = bot._load_logs_preference() + assert result == {"chat_id": 42, "mode": "default"} + + def test_load_returns_none_when_no_file(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + assert bot._load_logs_preference() is None + + def test_load_returns_none_on_corrupt_json(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + pref_file.write_text("not json{{{") + + assert bot._load_logs_preference() is None + + def test_load_returns_none_on_non_dict(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + pref_file: Path = _patch_base_bot_deps + pref_file.write_text('"just a string"') + + assert bot._load_logs_preference() is None + + +# ============================================= +# 8. AUTO-START HONORS PREFERENCE +# ============================================= + + +class TestAutoStartPreference: + """Verify handle_update auto-start respects saved preference.""" + + def test_autostart_defaults_to_all(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + bot._active_chat_id = None + + with ( + patch.object(bot, "_write_mirror_mapping"), + patch.object(bot, "is_user_allowed", return_value=True), + patch.object(bot, "check_rate_limit", return_value=True), + patch.object(bot, "handle_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer") as MockStreamer, + ): + MockStreamer.return_value = MagicMock() + bot.process_update({"message": {"text": "hi", "chat": {"id": 42}, "from": {"id": 111, "username": "u"}}}) + MockStreamer.assert_called_once_with("123:FAKETOKEN", 42, "testbranch", level_filter="all") + + def test_autostart_honors_errors_preference(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + bot._active_chat_id = None + pref_file: Path = _patch_base_bot_deps + pref_file.write_text(json.dumps({"chat_id": 42, "mode": "default"})) + + with ( + patch.object(bot, "_write_mirror_mapping"), + patch.object(bot, "is_user_allowed", return_value=True), + patch.object(bot, "check_rate_limit", return_value=True), + patch.object(bot, "handle_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer") as MockStreamer, + ): + MockStreamer.return_value = MagicMock() + bot.process_update({"message": {"text": "hi", "chat": {"id": 42}, "from": {"id": 111, "username": "u"}}}) + MockStreamer.assert_called_once_with("123:FAKETOKEN", 42, "testbranch", level_filter="default") + + def test_autostart_skips_when_off(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + bot._active_chat_id = None + pref_file: Path = _patch_base_bot_deps + pref_file.write_text(json.dumps({"chat_id": 42, "mode": "off"})) + + with ( + patch.object(bot, "_write_mirror_mapping"), + patch.object(bot, "is_user_allowed", return_value=True), + patch.object(bot, "check_rate_limit", return_value=True), + patch.object(bot, "handle_message"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.LogStreamer") as MockStreamer, + ): + bot.process_update({"message": {"text": "hi", "chat": {"id": 42}, "from": {"id": 111, "username": "u"}}}) + MockStreamer.assert_not_called() + assert bot._log_streamer is None diff --git a/src/aipass/skills/lib/telegram/tests/test_prax_monitor_bot.py b/src/aipass/skills/lib/telegram/tests/test_prax_monitor_bot.py new file mode 100644 index 00000000..421e6450 --- /dev/null +++ b/src/aipass/skills/lib/telegram/tests/test_prax_monitor_bot.py @@ -0,0 +1,400 @@ +# =================== AIPass ==================== +# Name: test_prax_monitor_bot.py +# Description: Tests for PraxMonitorBot — Prax Monitor TG chat command receiver +# Version: 1.0.0 +# Created: 2026-07-12 +# Modified: 2026-07-12 +# ============================================= + +""" +Tests for PraxMonitorBot — command receiver for the Prax Monitor TG chat. + +Tests cover: + - /pause writes paused=true to control file + - /resume writes paused=false to control file + - /errors writes level=errors + - /all writes level=all + - /status shows current state and relay liveness + - Free-text and file uploads rejected + - Command routing dispatches correctly + - Control file I/O: read defaults on missing/corrupt, write adds updated_at + - Slash-menu includes all custom commands + - Write failure sends error message +""" + +import json + +import pytest +from unittest.mock import patch, MagicMock + +from aipass.skills.lib.telegram.apps.handlers.prax_monitor_bot import PraxMonitorBot + + +# ============================================= +# HELPERS +# ============================================= + + +@pytest.fixture +def _patch_base_bot_deps(tmp_path): + """Patch heavy BaseBot dependencies to allow lightweight instantiation.""" + patches = [ + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.PENDING_DIR", tmp_path), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.signal.signal"), + patch("aipass.skills.lib.telegram.apps.handlers.base_bot.atexit.register"), + ] + for p in patches: + p.start() + yield + for p in patches: + p.stop() + + +@pytest.fixture +def ctrl_file(tmp_path): + """Provide a tmp control file path and patch the module constant.""" + f = tmp_path / "prax_monitor_control.json" + with patch("aipass.skills.lib.telegram.apps.handlers.prax_monitor_bot.CONTROL_FILE", f): + yield f + + +def _make_bot(tmp_path, _patch_base_bot_deps): + workdir = tmp_path / "workdir" + workdir.mkdir(exist_ok=True) + return PraxMonitorBot( + bot_id="prax_monitor", + bot_token="123:FAKETOKEN", + work_dir=workdir, + bot_name="Prax Monitor Bot", + allowed_user_ids=[111], + branch_name=None, + ) + + +# ============================================= +# 1. /PAUSE +# ============================================= + + +class TestPause: + def test_pause_writes_control(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message"): + bot._handle_pause(42) + + data = json.loads(ctrl_file.read_text()) + assert data["paused"] is True + assert "updated_at" in data + + def test_pause_sends_confirmation(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._handle_pause(42) + assert "paused" in mock_send.call_args[0][1].lower() + + def test_pause_preserves_level(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": False, "level": "errors"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message"): + bot._handle_pause(42) + + data = json.loads(ctrl_file.read_text()) + assert data["paused"] is True + assert data["level"] == "errors" + + def test_pause_write_failure(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "_write_control", return_value=False), + patch.object(bot, "send_message") as mock_send, + ): + bot._handle_pause(42) + assert "failed" in mock_send.call_args[0][1].lower() + + +# ============================================= +# 2. /RESUME +# ============================================= + + +class TestResume: + def test_resume_writes_control(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": True, "level": "all"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message"): + bot._handle_resume(42) + + data = json.loads(ctrl_file.read_text()) + assert data["paused"] is False + + def test_resume_sends_confirmation(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._handle_resume(42) + assert "resumed" in mock_send.call_args[0][1].lower() + + def test_resume_shows_level_in_message(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": True, "level": "errors"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._handle_resume(42) + assert "errors" in mock_send.call_args[0][1] + + +# ============================================= +# 3. /ERRORS +# ============================================= + + +class TestErrors: + def test_errors_writes_level(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message"): + bot._handle_errors(42) + + data = json.loads(ctrl_file.read_text()) + assert data["level"] == "errors" + + def test_errors_sends_confirmation(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._handle_errors(42) + assert "errors & warnings" in mock_send.call_args[0][1].lower() + + def test_errors_preserves_paused(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": True, "level": "all"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message"): + bot._handle_errors(42) + + data = json.loads(ctrl_file.read_text()) + assert data["paused"] is True + assert data["level"] == "errors" + + +# ============================================= +# 4. /ALL +# ============================================= + + +class TestAll: + def test_all_writes_level(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": False, "level": "errors"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message"): + bot._handle_all(42) + + data = json.loads(ctrl_file.read_text()) + assert data["level"] == "all" + + def test_all_sends_confirmation(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot._handle_all(42) + msg = mock_send.call_args[0][1] + assert "all" in msg.lower() + + +# ============================================= +# 5. /STATUS +# ============================================= + + +class TestStatus: + def test_status_shows_state(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": False, "level": "all", "updated_at": "2026-07-12T10:00:00Z"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "send_message") as mock_send, + patch.object(PraxMonitorBot, "_check_relay_alive", return_value=True), + ): + bot._handle_prax_status(42) + msg = mock_send.call_args[0][1] + assert "active" in msg + assert "all levels" in msg + assert "running" in msg + + def test_status_shows_paused(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": True, "level": "errors"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "send_message") as mock_send, + patch.object(PraxMonitorBot, "_check_relay_alive", return_value=False), + ): + bot._handle_prax_status(42) + msg = mock_send.call_args[0][1] + assert "paused" in msg + assert "errors & warnings" in msg + assert "not detected" in msg + + def test_status_defaults_when_no_file(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with ( + patch.object(bot, "send_message") as mock_send, + patch.object(PraxMonitorBot, "_check_relay_alive", return_value=False), + ): + bot._handle_prax_status(42) + msg = mock_send.call_args[0][1] + assert "active" in msg + assert "all levels" in msg + + +# ============================================= +# 6. FREE-TEXT & FILE REJECTION +# ============================================= + + +class TestRejection: + def test_freetext_rejected(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot.handle_message(42, "hello there", {}) + assert "commands" in mock_send.call_args[0][1].lower() + + def test_file_rejected(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message") as mock_send: + bot.handle_file(42, {"document": {}}) + assert "don't process files" in mock_send.call_args[0][1].lower() + + +# ============================================= +# 7. COMMAND ROUTING +# ============================================= + + +class TestCommandRouting: + def test_pause_dispatched(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_handle_pause") as mock: + assert bot._dispatch_command(42, ("pause", "")) is True + mock.assert_called_once_with(42) + + def test_resume_dispatched(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_handle_resume") as mock: + assert bot._dispatch_command(42, ("resume", "")) is True + mock.assert_called_once_with(42) + + def test_errors_dispatched(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_handle_errors") as mock: + assert bot._dispatch_command(42, ("errors", "")) is True + mock.assert_called_once_with(42) + + def test_all_dispatched(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_handle_all") as mock: + assert bot._dispatch_command(42, ("all", "")) is True + mock.assert_called_once_with(42) + + def test_status_dispatched(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "_handle_prax_status") as mock: + assert bot._dispatch_command(42, ("status", "")) is True + mock.assert_called_once_with(42) + + def test_unknown_falls_through_to_parent(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + with patch.object(bot, "send_message"): + result = bot._dispatch_command(42, ("help", "")) + assert result is True + + +# ============================================= +# 8. CONTROL FILE I/O +# ============================================= + + +class TestControlFileIO: + def test_read_defaults_when_missing(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + result = bot._read_control() + assert result == {"paused": False, "level": "all"} + + def test_read_defaults_on_corrupt(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text("not json{{{") + bot = _make_bot(tmp_path, _patch_base_bot_deps) + result = bot._read_control() + assert result == {"paused": False, "level": "all"} + + def test_read_defaults_on_non_dict(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text('"just a string"') + bot = _make_bot(tmp_path, _patch_base_bot_deps) + result = bot._read_control() + assert result == {"paused": False, "level": "all"} + + def test_read_returns_data(self, tmp_path, _patch_base_bot_deps, ctrl_file): + ctrl_file.write_text(json.dumps({"paused": True, "level": "errors"})) + bot = _make_bot(tmp_path, _patch_base_bot_deps) + result = bot._read_control() + assert result["paused"] is True + assert result["level"] == "errors" + + def test_write_adds_timestamp(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + bot._write_control({"paused": False, "level": "all"}) + data = json.loads(ctrl_file.read_text()) + assert "updated_at" in data + assert "T" in data["updated_at"] + + def test_write_roundtrip(self, tmp_path, _patch_base_bot_deps, ctrl_file): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + bot._write_control({"paused": True, "level": "errors"}) + result = bot._read_control() + assert result["paused"] is True + assert result["level"] == "errors" + assert "updated_at" in result + + +# ============================================= +# 9. SLASH MENU +# ============================================= + + +class TestSlashMenu: + def test_custom_commands_include_all(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + cmds = bot.get_custom_commands() + assert "pause" in cmds + assert "resume" in cmds + assert "errors" in cmds + assert "all" in cmds + assert "monitor" in cmds + + def test_custom_commands_have_descriptions(self, tmp_path, _patch_base_bot_deps): + bot = _make_bot(tmp_path, _patch_base_bot_deps) + cmds = bot.get_custom_commands() + for cmd in ("pause", "resume", "errors", "all"): + assert "description" in cmds[cmd] + assert "menu_text" in cmds[cmd] + + +# ============================================= +# 10. RELAY LIVENESS CHECK +# ============================================= + + +class TestRelayAlive: + def test_alive_when_active(self): + mock_result = MagicMock() + mock_result.stdout = "active\n" + with patch( + "aipass.skills.lib.telegram.apps.handlers.prax_monitor_bot.subprocess.run", return_value=mock_result + ): + assert PraxMonitorBot._check_relay_alive() is True + + def test_not_alive_when_inactive(self): + mock_result = MagicMock() + mock_result.stdout = "inactive\n" + with patch( + "aipass.skills.lib.telegram.apps.handlers.prax_monitor_bot.subprocess.run", return_value=mock_result + ): + assert PraxMonitorBot._check_relay_alive() is False + + def test_not_alive_on_error(self): + with patch( + "aipass.skills.lib.telegram.apps.handlers.prax_monitor_bot.subprocess.run", + side_effect=OSError("no systemctl"), + ): + assert PraxMonitorBot._check_relay_alive() is False