feat: TG log-stream control — /logs on branch bots (persisted pref, honored by auto-start) + prax_monitor receiver bot (/pause /resume /errors /all /status, menu registered) + relay honors shared control file each flush. 84 new tests, all suites green; live-verified end-to-end from Telegram Web (errors filter kicked in within one flush).
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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",
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user