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