feat(skills): Scheduler Bot Phase 2 — dedicated bot /queue + hourly digest (TDPLAN-0008)

- SchedulerBot(BaseBot): rejects free-text (NO tmux/Claude spawn), /queue shells 'drone @daemon queue --json', hourly digest thread, chunked output
- base_bot __main__: _BOT_CLASSES maps bot_id->SchedulerBot for systemd launch; passes config chat_id for digest target
- bot registered (telegram-bot@scheduler), systemd unit installed (NOT started)
- 26 new tests (519 telegram / 252 skills green), seedgo 99%

fix(daemon): queue --json emits valid JSON at module level — flatten prompt_preview (collapse newlines) + soft_wrap/markup off. Verified valid via 'python -m'. NOTE: the drone router still re-wraps sub-command stdout at width 80, corrupting machine --json for ALL consumers routed through drone (separate @drone core bug; skills mitigates with strict=False JSON parse).
This commit is contained in:
AIOSAI
2026-06-25 16:56:56 -07:00
parent b51ac87eb7
commit d252403d37
4 changed files with 646 additions and 3 deletions
+6 -2
View File
@@ -102,7 +102,8 @@ def _build_queue(jobs: list, runstate: dict) -> list:
owner = f"@{owner.split('@')[0]}"
prompt = job.get("prompt", "")
preview = prompt[:80] + "..." if len(prompt) > 80 else prompt
flat = " ".join(prompt.split()) # collapse newlines/whitespace → single-line preview
preview = flat[:80] + "..." if len(flat) > 80 else flat
entries.append(
{
@@ -183,7 +184,10 @@ def handle_command(command: str, args: List[str]) -> bool:
if "--json" in args:
output = _build_json_output(entries)
console.print(json.dumps(output, indent=2))
# soft_wrap=True + markup=False — default console.print forces width-80
# wrapping on non-TTY and parses [] markup, injecting newlines mid-string
# that corrupt machine output. These flags emit clean, parseable JSON.
console.print(json.dumps(output, indent=2), soft_wrap=True, markup=False)
else:
_print_rich_table(entries)
@@ -1902,6 +1902,11 @@ class BaseBot:
# CLI ENTRY POINT
# =============================================
_BOT_CLASSES = {
"scheduler": ".scheduler_bot:SchedulerBot",
}
if __name__ == "__main__":
parser = argparse.ArgumentParser(description="AIPass Telegram Bot")
parser.add_argument("--bot-id", required=True, help="Bot identifier")
@@ -1914,7 +1919,15 @@ if __name__ == "__main__":
print(f"No config found for bot_id={args.bot_id}")
sys.exit(1)
bot = BaseBot(
bot_cls = BaseBot
if args.bot_id in _BOT_CLASSES:
mod_name, cls_name = _BOT_CLASSES[args.bot_id].rsplit(":", 1)
import importlib
mod = importlib.import_module(mod_name, package=__package__)
bot_cls = getattr(mod, cls_name)
bot = bot_cls(
bot_id=args.bot_id,
bot_token=config["bot_token"],
work_dir=Path(config.get("work_dir", str(Path.home()))),
@@ -1923,4 +1936,8 @@ if __name__ == "__main__":
branch_name=config.get("branch_name"),
shared_session=config.get("shared_session"),
)
if hasattr(bot, "_config_chat_id") and config.get("chat_id"):
bot._config_chat_id = config["chat_id"]
sys.exit(bot.run())
@@ -0,0 +1,227 @@
# =================== AIPass ====================
# Name: scheduler_bot.py
# Description: Dedicated Telegram bot for the daemon's job queue — announce/query only
# Version: 1.0.0
# Created: 2026-06-25
# Modified: 2026-06-25
# =============================================
"""
SchedulerBot — a BaseBot subclass for the AIPass Scheduler.
Pure announce/query bot: serves /queue, posts an hourly digest, receives
lifecycle pings from @daemon. Free-text messages do NOT spin up tmux/Claude.
"""
import json
import subprocess
import threading
from typing import Optional
from aipass.prax import logger
from .base_bot import BaseBot
DIGEST_INTERVAL = 3600.0
QUEUE_CMD = ["drone", "@daemon", "queue", "--json"]
class SchedulerBot(BaseBot):
"""Dedicated scheduler bot — no tmux/Claude, just /queue and hourly digest."""
def __init__(self, **kwargs) -> None:
super().__init__(**kwargs)
self._digest_thread: Optional[threading.Thread] = None
self._digest_stop = threading.Event()
self._scheduler_chat_id: Optional[int] = None
def handle_message(self, chat_id: int, text: str, message: dict) -> None:
self.send_message(
chat_id,
"I'm the scheduler bot — I only handle commands.\nTry /queue or /help",
)
def handle_file(self, chat_id: int, message: dict) -> None:
self.send_message(chat_id, "I don't process files. Try /queue or /help")
def _dispatch_command(self, chat_id: int, parsed: tuple) -> bool:
cmd_name, cmd_args = parsed
if cmd_name == "queue":
self._handle_queue_command(chat_id)
return True
return super()._dispatch_command(chat_id, parsed)
def get_custom_commands(self) -> dict:
cmds = super().get_custom_commands()
cmds["queue"] = {
"description": "Show all scheduled daemon jobs — next fire, status, type",
"menu_text": "Job queue",
}
return cmds
def _handle_queue_command(self, chat_id: int) -> None:
"""Fetch queue from daemon and send formatted message."""
queue_data = self._fetch_queue()
if queue_data is None:
self.send_message(chat_id, "Failed to fetch queue from daemon.")
return
text = self._format_queue(queue_data)
for chunk in self.chunk_text(text):
self.send_message(chat_id, chunk)
def _fetch_queue(self) -> Optional[dict]:
"""Run drone @daemon queue --json and parse the result."""
try:
result = subprocess.run(
QUEUE_CMD,
capture_output=True,
text=True,
timeout=15,
)
if result.returncode != 0:
logger.warning("queue --json failed (rc=%d): %s", result.returncode, result.stderr)
return None
decoder = json.JSONDecoder(strict=False)
return decoder.decode(result.stdout)
except (subprocess.TimeoutExpired, json.JSONDecodeError, OSError) as e:
logger.warning("Failed to fetch queue: %s", e)
return None
@staticmethod
def _format_queue(data: dict) -> str:
"""Format queue JSON into a readable Telegram message."""
jobs = data.get("jobs", [])
count = data.get("count", len(jobs))
generated = data.get("generated_at", "unknown")
if not jobs:
return f"No scheduled jobs.\n\nGenerated: {generated}"
lines = [f"Scheduled Jobs ({count})\n"]
for job in jobs:
owner = job.get("owner", "?")
job_id = job.get("id", "?")
enabled = job.get("enabled", False)
jtype = job.get("type", "?")
schedule = job.get("schedule_human", "?")
next_run = job.get("next_run") or "—"
last_status = job.get("last_status") or "never run"
last_error = job.get("last_error")
prompt = job.get("prompt_preview", "")
status_icon = {
"success": "✅",
"failed": "❌",
"dispatched": "\U0001f535",
}.get(last_status, "⚪")
enabled_tag = "" if enabled else " [disabled]"
lines.append(f"{status_icon} {owner}/{job_id}{enabled_tag}")
lines.append(f" Type: {jtype} ({schedule})")
lines.append(f" Next: {next_run}")
lines.append(f" Last: {last_status}")
if last_error:
lines.append(f" Error: {last_error}")
if prompt:
lines.append(f" Prompt: {prompt[:80]}...")
lines.append("")
lines.append(f"Generated: {generated}")
return "\n".join(lines)
# =============================================
# HOURLY DIGEST
# =============================================
def start_digest(self, chat_id: int) -> None:
"""Start the hourly digest thread."""
if self._digest_thread is not None and self._digest_thread.is_alive():
return
self._scheduler_chat_id = chat_id
self._digest_stop.clear()
self._digest_thread = threading.Thread(
target=self._digest_loop,
name="scheduler-digest",
daemon=True,
)
self._digest_thread.start()
logger.info("Digest thread started (chat_id=%s)", chat_id)
def stop_digest(self) -> None:
"""Stop the hourly digest thread."""
self._digest_stop.set()
if self._digest_thread is not None:
self._digest_thread.join(timeout=5)
self._digest_thread = None
logger.info("Digest thread stopped")
def _digest_loop(self) -> None:
"""Post a queue digest every hour."""
while not self._digest_stop.is_set():
self._digest_stop.wait(DIGEST_INTERVAL)
if self._digest_stop.is_set():
break
self._post_digest()
def _post_digest(self) -> None:
"""Fetch queue and post digest to the scheduler chat."""
if self._scheduler_chat_id is None:
return
queue_data = self._fetch_queue()
if queue_data is None:
logger.warning("Digest skipped: failed to fetch queue")
return
text = f"Hourly Queue Digest\n{'=' * 20}\n\n{self._format_queue(queue_data)}"
for chunk in self.chunk_text(text):
self.send_message(self._scheduler_chat_id, chunk)
logger.info("Hourly digest posted")
# =============================================
# LIFECYCLE
# =============================================
def run(self) -> int:
"""Start the bot with digest thread."""
logger.info("=" * 60)
logger.info("%s starting (bot_id=%s)", self.bot_name, self.bot_id)
if not self.verify_connection():
logger.error("Startup health check FAILED — cannot reach Telegram API")
return 1
logger.info("Connected to Telegram API")
from datetime import datetime
self._health["started_at"] = datetime.now().isoformat()
from .json import json_handler
json_handler.log_operation("bot_started", {"bot_id": self.bot_id})
self._set_command_menu()
self._boot_monitor()
if self._check_lock():
logger.error("Another instance of bot-%s is already running", self.bot_id)
return 1
chat_id = getattr(self, "_config_chat_id", None)
if chat_id:
self.start_digest(int(chat_id))
self._create_lock()
try:
self._poll_loop()
except KeyboardInterrupt:
logger.info("Interrupted")
finally:
self.stop_digest()
self._cleanup()
return 0
@@ -0,0 +1,395 @@
"""
Tests for SchedulerBot — dedicated Telegram bot for daemon job queue (TDPLAN-0008 P2).
Tests cover:
- /queue parses drone @daemon queue --json and formats readable output
- Free-text (non-command) does NOT launch tmux/Claude
- Hourly digest thread posts a digest (mockable clock/sender)
- >4096-char digest is chunked correctly via chunk_text
- Bot loads telegram/scheduler config; missing secret fails loud
- Slash-menu includes /queue
- Command routing (/queue dispatched correctly)
"""
import json
import pytest
from unittest.mock import patch, MagicMock
from apps.handlers.scheduler_bot import SchedulerBot, QUEUE_CMD # type: ignore[import-not-found]
# =============================================
# HELPERS
# =============================================
SAMPLE_QUEUE = {
"generated_at": "2026-06-25T15:00:00Z",
"count": 2,
"jobs": [
{
"owner": "@api",
"id": "data-check",
"enabled": True,
"type": "once",
"schedule_human": "2026-07-02",
"next_run": "2026-07-02T09:00:00Z",
"last_run": None,
"last_status": "success",
"last_error": None,
"prompt_preview": "Check the live data and report...",
"wake": {"fresh": True, "model": "haiku"},
},
{
"owner": "@backup",
"id": "weekly-snap",
"enabled": False,
"type": "interval",
"schedule_human": "every 7d",
"next_run": "2026-07-01T00:00:00Z",
"last_run": "2026-06-24T00:00:00Z",
"last_status": "failed",
"last_error": "disk full",
"prompt_preview": "Run full backup...",
"wake": {"fresh": True},
},
],
}
EMPTY_QUEUE = {"generated_at": "2026-06-25T15:00:00Z", "count": 0, "jobs": []}
@pytest.fixture
def _patch_base_bot_deps(tmp_path):
"""Patch heavy BaseBot dependencies to allow lightweight instantiation."""
patches = [
patch("apps.handlers.base_bot.PENDING_DIR", tmp_path),
patch("apps.handlers.base_bot.signal.signal"),
patch("apps.handlers.base_bot.atexit.register"),
]
for p in patches:
p.start()
yield
for p in patches:
p.stop()
def _make_scheduler_bot(tmp_path, _patch_base_bot_deps):
workdir = tmp_path / "workdir"
workdir.mkdir()
bot = SchedulerBot(
bot_id="scheduler",
bot_token="123:FAKETOKEN",
work_dir=workdir,
bot_name="AIPass Scheduler Bot",
allowed_user_ids=[111],
branch_name=None,
)
return bot
# =============================================
# 1. /queue PARSES AND FORMATS
# =============================================
class TestQueueCommand:
"""/queue fetches queue --json, formats, and sends."""
def test_queue_formats_jobs(self):
text = SchedulerBot._format_queue(SAMPLE_QUEUE)
assert "@api/data-check" in text
assert "@backup/weekly-snap" in text
assert "[disabled]" in text
assert "once" in text
assert "success" in text
def test_queue_shows_error_when_failed(self):
text = SchedulerBot._format_queue(SAMPLE_QUEUE)
assert "disk full" in text
def test_queue_empty_shows_no_jobs(self):
text = SchedulerBot._format_queue(EMPTY_QUEUE)
assert "No scheduled jobs" in text
def test_queue_shows_count(self):
text = SchedulerBot._format_queue(SAMPLE_QUEUE)
assert "2" in text
def test_queue_shows_status_icons(self):
text = SchedulerBot._format_queue(SAMPLE_QUEUE)
assert "✅" in text # success
assert "❌" in text # failed
def test_queue_command_sends_formatted(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with (
patch.object(bot, "_fetch_queue", return_value=SAMPLE_QUEUE),
patch.object(bot, "send_message") as mock_send,
):
bot._handle_queue_command(42)
mock_send.assert_called_once()
msg = mock_send.call_args[0][1]
assert "@api/data-check" in msg
def test_queue_command_handles_fetch_failure(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with (
patch.object(bot, "_fetch_queue", return_value=None),
patch.object(bot, "send_message") as mock_send,
):
bot._handle_queue_command(42)
msg = mock_send.call_args[0][1]
assert "Failed" in msg
def test_fetch_queue_parses_subprocess(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
mock_result = MagicMock()
mock_result.returncode = 0
mock_result.stdout = json.dumps(SAMPLE_QUEUE)
with patch("apps.handlers.scheduler_bot.subprocess.run", return_value=mock_result):
data = bot._fetch_queue()
assert data is not None
assert data["count"] == 2
assert len(data["jobs"]) == 2
def test_fetch_queue_returns_none_on_failure(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
mock_result = MagicMock()
mock_result.returncode = 1
mock_result.stderr = "error"
with patch("apps.handlers.scheduler_bot.subprocess.run", return_value=mock_result):
data = bot._fetch_queue()
assert data is None
def test_fetch_queue_returns_none_on_timeout(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
import subprocess as sp
with patch("apps.handlers.scheduler_bot.subprocess.run", side_effect=sp.TimeoutExpired(QUEUE_CMD, 15)):
data = bot._fetch_queue()
assert data is None
# =============================================
# 2. FREE-TEXT DOES NOT LAUNCH TMUX
# =============================================
class TestNoTmux:
"""Free-text and file messages do NOT spin up tmux/Claude."""
def test_free_text_rejected(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with (
patch.object(bot, "send_message") as mock_send,
patch.object(bot, "ensure_tmux_session") as mock_tmux,
):
bot.handle_message(42, "Hello there", {"message_id": 1})
mock_tmux.assert_not_called()
msg = mock_send.call_args[0][1]
assert "commands" in msg.lower() or "/queue" in msg
def test_file_rejected(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with (
patch.object(bot, "send_message") as mock_send,
patch("apps.handlers.base_bot.BaseBot.handle_file") as mock_parent_file,
):
bot.handle_file(42, {"message_id": 1, "document": {"file_id": "abc"}})
mock_parent_file.assert_not_called()
msg = mock_send.call_args[0][1]
assert "file" in msg.lower() or "/queue" in msg
# =============================================
# 3. HOURLY DIGEST
# =============================================
class TestDigest:
"""Hourly digest thread posts queue digest."""
def test_digest_posts_message(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
bot._scheduler_chat_id = 42
with (
patch.object(bot, "_fetch_queue", return_value=SAMPLE_QUEUE),
patch.object(bot, "send_message") as mock_send,
):
bot._post_digest()
mock_send.assert_called_once()
msg = mock_send.call_args[0][1]
assert "Hourly Queue Digest" in msg
assert "@api/data-check" in msg
def test_digest_skips_on_fetch_failure(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
bot._scheduler_chat_id = 42
with (
patch.object(bot, "_fetch_queue", return_value=None),
patch.object(bot, "send_message") as mock_send,
):
bot._post_digest()
mock_send.assert_not_called()
def test_digest_skips_when_no_chat_id(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
assert bot._scheduler_chat_id is None
with patch.object(bot, "send_message") as mock_send:
bot._post_digest()
mock_send.assert_not_called()
def test_start_digest_creates_thread(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
bot.start_digest(42)
try:
assert bot._digest_thread is not None
assert bot._digest_thread.daemon is True
assert bot._digest_thread.name == "scheduler-digest"
assert bot._scheduler_chat_id == 42
finally:
bot.stop_digest()
def test_stop_digest_cleans_up(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
bot.start_digest(42)
bot.stop_digest()
assert bot._digest_thread is None
assert bot._digest_stop.is_set()
# =============================================
# 4. CHUNKING LONG MESSAGES
# =============================================
class TestChunking:
""">4096-char messages are split via chunk_text."""
def test_long_queue_is_chunked(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
many_jobs = {
"generated_at": "now",
"count": 100,
"jobs": [
{
"owner": f"@branch{i}",
"id": f"job-{i}",
"enabled": True,
"type": "daily",
"schedule_human": "every day",
"next_run": "2026-07-01T00:00:00Z",
"last_run": None,
"last_status": None,
"last_error": None,
"prompt_preview": "A" * 80,
"wake": {},
}
for i in range(100)
],
}
with (
patch.object(bot, "_fetch_queue", return_value=many_jobs),
patch.object(bot, "send_message") as mock_send,
):
bot._handle_queue_command(42)
assert mock_send.call_count >= 2
def test_short_queue_not_chunked(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with (
patch.object(bot, "_fetch_queue", return_value=SAMPLE_QUEUE),
patch.object(bot, "send_message") as mock_send,
):
bot._handle_queue_command(42)
assert mock_send.call_count == 1
# =============================================
# 5. SECRET LOADING
# =============================================
class TestSecretLoading:
"""Bot loads telegram/scheduler config; missing fails loud."""
def test_missing_secret_returns_none(self):
"""When get_secret returns None, load_bot_config returns None — bot won't start."""
from apps.handlers.config import load_bot_config # type: ignore[import-not-found]
with patch("apps.handlers.config._get_secret", return_value=None):
config = load_bot_config("scheduler")
assert config is None
def test_secret_provides_chat_id(self):
"""The scheduler config includes chat_id for digest delivery."""
config = {
"bot_id": "scheduler",
"bot_token": "123:FAKE",
"bot_name": "AIPass Scheduler Bot",
"branch_name": "daemon",
"work_dir": "/tmp",
"allowed_user_ids": [111],
"chat_id": "7235222625",
}
assert "chat_id" in config
assert config["chat_id"] == "7235222625"
# =============================================
# 6. SLASH-MENU INCLUDES /queue
# =============================================
class TestSlashMenu:
"""/queue appears in get_custom_commands and the command menu."""
def test_queue_in_custom_commands(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
cmds = bot.get_custom_commands()
assert "queue" in cmds
assert "description" in cmds["queue"]
def test_inherited_commands_preserved(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
cmds = bot.get_custom_commands()
assert "monitor" in cmds
assert "create" in cmds
# =============================================
# 7. COMMAND ROUTING
# =============================================
class TestCommandRouting:
"""/queue is dispatched correctly through _dispatch_command."""
def test_queue_routed(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with patch.object(bot, "_handle_queue_command") as mock_q:
result = bot._dispatch_command(42, ("queue", ""))
assert result is True
mock_q.assert_called_once_with(42)
def test_other_commands_fall_through_to_parent(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with patch.object(bot, "send_message"):
result = bot._dispatch_command(42, ("status", ""))
assert result is True
def test_unknown_command_falls_through(self, tmp_path, _patch_base_bot_deps):
bot = _make_scheduler_bot(tmp_path, _patch_base_bot_deps)
with patch.object(bot, "send_message"):
result = bot._dispatch_command(42, ("nonexistent", ""))
assert result is False