feat(skills,hooks): TDPLAN-0009 Telegram session mirror — bidirectional, live-proven @api

@skills: dup-spawn fix (run() lock-collision exit 0), attach_only + launch_mirror_session (--dangerously-skip-permissions), _config_chat_id init bug + /proc active-transcript baseline, systemctl start. @hooks: extract_mirror_turn cursor clamp + baseline reset on delivery, mirror-file unlink guards. +4 test files. 578 skills / 112 hooks green.
This commit is contained in:
AIOSAI
2026-06-29 08:53:03 -07:00
parent 4363f9824a
commit 88e99efe4c
18 changed files with 3157 additions and 279 deletions
+3
View File
@@ -51,6 +51,9 @@ memory = [
"chromadb>=1.0",
"fastembed>=0.4",
]
telegram = [
"telethon>=1.36",
]
seedgo = []
dev = [
"pytest>=9.0.3",
@@ -12,9 +12,9 @@ Not suggestions. Violating = bug.
- **No writes outside own `.trinity/`.** Never create, edit, delete files anywhere else. Not code, not docs, not configs, not other branches' memories.
- **No git. Ever.** Not `git status`, not `drone @git anything`. Git is drone's world.
- **No `drone @ai_mail dispatch`.** Email only test-convention body (below). Never wake agent real work.
- **Dispatch focused work via `drone @ai_mail dispatch`** — to ONE owning branch, as the user's voice with detailed feedback. Reply routes to @aipass; I track the loop and report back. Not an orchestrator (no fleets, no running the floor — that's devpulse). Test-convention pings (below) still fine.
- **No registry / hooks / bypass.json / config edits.** Spot bug → report. Never patch.
- User asks build/fix/change something: tell them who. Offer dispatch through devpulse/drone — don't do it.
- User asks build/fix/change in another branch: name the owner, then dispatch focused work to them as the user's voice. Heavy orchestration, git, and fleets stay with devpulse.
## What I Do
+1 -1
View File
@@ -82,7 +82,7 @@ src/aipass/hooks/
│ └── diagnostics.py # JSONL logging for hook execution
├── logs/
│ └── engine.jsonl # JSONL diagnostics (every hook execution)
└── tests/ # 593 tests across 23 test files
└── tests/ # 666 tests across 23 test files
```
## How It Works
@@ -1,11 +1,11 @@
# =================== AIPass ====================
# Name: telegram_response.py
# Version: 1.0.0
# Version: 2.0.0
# Description: Telegram response delivery on Stop event (ported from Dev-Pass)
# Branch: hooks
# Layer: apps/handlers/notification
# Created: 2026-06-15
# Modified: 2026-06-15
# Modified: 2026-06-29
# =============================================
"""Telegram response delivery on Stop event.
@@ -18,6 +18,7 @@ Layer 2: isSidechain filter — skips sidechain entries during transcript extrac
Layer 3: Transcript position — only extracts text after the recorded injection point
"""
import hashlib
import json
import os
import re
@@ -30,12 +31,16 @@ from urllib.request import Request, urlopen
from aipass.prax.apps.modules.logger import system_logger as logger
PENDING_DIR = Path.home() / ".aipass" / "telegram_pending"
MIRROR_DIR = Path.home() / ".aipass" / "telegram_bots"
PENDING_TTL = 3600
TELEGRAM_CHAR_LIMIT = 4096
_DELIVERY_LOG = Path(__file__).resolve().parent.parent.parent.parent / "logs" / "telegram_delivery.jsonl"
def _is_expired(data: dict) -> bool:
"""Check if pending file is expired (1-hour TTL + tmux-alive check)."""
if data.get("mirror"):
return False
timestamp = data.get("timestamp", 0)
if isinstance(timestamp, str):
try:
@@ -73,38 +78,49 @@ def _try_load_pending(path: Path) -> dict | None:
return None
def _in_mirror_dir(path: Path) -> bool:
"""Check if a path is inside the persistent mirror mapping directory."""
try:
path.relative_to(MIRROR_DIR)
return True
except ValueError:
logger.info("[HOOKS] telegram: %s is not in mirror dir", path.name)
return False
def find_pending_file(session_id: str) -> Path | None: # noqa: ARG001
"""Find pending file matching current context via multi-bot matching.
Priority 1: AIPASS_BOT_ID env var -> bot-{bot_id}.json
Priority 2: CWD relative_to work_dir -> bot-*.json
Priority 1: AIPASS_BOT_ID env var -> bot-{bot_id}.json (mirror dir first, then pending)
Priority 2: CWD relative_to work_dir -> bot-*.json (both dirs)
"""
if not PENDING_DIR.exists():
return None
cwd = Path.cwd()
env_bot_id = os.environ.get("AIPASS_BOT_ID")
if env_bot_id:
v2_path = PENDING_DIR / f"bot-{env_bot_id}.json"
if _try_load_pending(v2_path) is not None:
logger.info("[HOOKS] telegram: v2 match bot_id env -> %s", v2_path.name)
return v2_path
for search_dir in [MIRROR_DIR, PENDING_DIR]:
path = search_dir / f"bot-{env_bot_id}.json"
if _try_load_pending(path) is not None:
logger.info("[HOOKS] telegram: match bot_id env -> %s", path)
return path
for pending_path in sorted(PENDING_DIR.glob("bot-*.json")):
data = _try_load_pending(pending_path)
if data is None:
continue
work_dir = data.get("work_dir")
if not work_dir:
continue
try:
cwd.relative_to(Path(work_dir))
logger.info("[HOOKS] telegram: v2 match cwd -> %s", pending_path.name)
return pending_path
except ValueError:
logger.info("[HOOKS] telegram: cwd not relative to %s, skipping", pending_path.name)
for search_dir in [PENDING_DIR, MIRROR_DIR]:
if not search_dir.exists():
continue
for pending_path in sorted(search_dir.glob("bot-*.json")):
data = _try_load_pending(pending_path)
if data is None:
continue
work_dir = data.get("work_dir")
if not work_dir:
continue
try:
cwd.relative_to(Path(work_dir))
logger.info("[HOOKS] telegram: match cwd -> %s", pending_path)
return pending_path
except ValueError:
logger.info("[HOOKS] telegram: cwd not relative to %s, skipping", pending_path.name)
continue
return None
@@ -192,6 +208,128 @@ def _collect_assistant_text(lines: list[str]) -> list[str]:
return text_parts
def _extract_user_text(content: list | str) -> str | None:
"""Extract text from user message content, returning None for tool-result-only messages."""
if isinstance(content, str):
return content.strip() or None
if not isinstance(content, list):
return None
texts: list[str] = []
has_non_tool = False
for block in content:
if not isinstance(block, dict):
continue
if block.get("type") == "text":
text = block.get("text", "").strip()
if text:
texts.append(text)
has_non_tool = True
elif block.get("type") != "tool_result":
has_non_tool = True
if not has_non_tool:
return None
return " ".join(texts) if texts else None
def _text_from_content(content: list) -> list[str]:
"""Extract text strings from a message content block list."""
if not isinstance(content, list):
return []
parts: list[str] = []
for block in content:
if isinstance(block, dict) and block.get("type") == "text":
text = block.get("text", "").strip()
if text:
parts.append(text)
return parts
def _collect_mirror_entries(lines: list[str]) -> list[tuple[str | None, list[str]]]:
"""Collect user+assistant turn pairs from JSONL lines for mirror delivery."""
turns: list[tuple[str | None, list[str]]] = []
current_user: str | None = None
current_assistant: list[str] = []
for line in lines:
try:
entry = json.loads(line)
except json.JSONDecodeError:
logger.info("[HOOKS] telegram: skipping malformed JSONL line in mirror collection")
continue
if entry.get("isSidechain", False):
continue
entry_type = entry.get("type")
message = entry.get("message", {})
content = message.get("content", [])
if entry_type == "user":
user_text = _extract_user_text(content)
if user_text is None:
continue
if current_user is not None or current_assistant:
turns.append((current_user, current_assistant))
current_user = user_text
current_assistant = []
elif entry_type == "assistant":
current_assistant.extend(_text_from_content(content))
if current_user is not None or current_assistant:
turns.append((current_user, current_assistant))
return turns
def extract_mirror_turn(transcript_path: str, start_line: int = 0) -> str | None:
"""Extract all new turns (user input + assistant response) for mirror delivery."""
path = Path(transcript_path)
if not path.exists():
logger.warning("[HOOKS] telegram: mirror transcript not found: %s", transcript_path)
return None
try:
all_lines = path.read_text(encoding="utf-8").strip().split("\n")
except OSError as e:
logger.error("[HOOKS] telegram: failed to read mirror transcript: %s", e)
return None
if not all_lines:
return None
lines = all_lines[start_line:] if start_line > 0 else all_lines
if not lines and start_line > len(all_lines):
logger.warning(
"[HOOKS] telegram: cursor ahead of transcript (%d > %d) — clamping to deliver latest",
start_line,
len(all_lines),
)
clamped = _find_last_real_user_message(all_lines)
lines = all_lines[max(0, clamped) :]
if not lines:
return None
turns = _collect_mirror_entries(lines)
if not turns:
return None
formatted = []
for user_text, assistant_parts in turns:
parts: list[str] = []
if user_text:
parts.append(f"You: {user_text}")
if assistant_parts:
parts.append("\n\n".join(assistant_parts))
if parts:
formatted.append("\n\n".join(parts))
if not formatted:
return None
result = "\n\n---\n\n".join(formatted).strip()
return result if result else None
def chunk_text(text: str, limit: int = TELEGRAM_CHAR_LIMIT) -> list[str]:
"""Split text into chunks for Telegram's message limit."""
if len(text) <= limit:
@@ -278,8 +416,14 @@ def markdown_to_telegram_html(text: str) -> str:
return text
def send_to_telegram(bot_token: str, chat_id: int, text: str, message_id: int | None = None) -> bool:
"""Send a message to Telegram via Bot API using urllib."""
def _parse_api_message(api_result: dict) -> dict:
"""Extract message_id and text from a Telegram API response."""
msg = api_result.get("result", {})
return {"ok": True, "message_id": msg.get("message_id"), "text": msg.get("text", "")}
def send_to_telegram(bot_token: str, chat_id: int, text: str, message_id: int | None = None) -> dict:
"""Send a message to Telegram via Bot API. Returns dict with ok, message_id, text."""
url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
try:
@@ -290,10 +434,10 @@ def send_to_telegram(bot_token: str, chat_id: int, text: str, message_id: int |
data = json.dumps(html_payload).encode("utf-8")
req = Request(url, data=data, headers={"Content-Type": "application/json"})
with urlopen(req, timeout=15) as resp:
result = json.loads(resp.read())
if result.get("ok"):
return True
logger.warning("[HOOKS] telegram: HTML send failed: %s", result.get("description"))
api_result = json.loads(resp.read())
if api_result.get("ok"):
return _parse_api_message(api_result)
logger.warning("[HOOKS] telegram: HTML send failed: %s", api_result.get("description"))
except Exception as e:
logger.warning("[HOOKS] telegram: HTML send error: %s, plain text fallback", e)
@@ -306,11 +450,11 @@ def send_to_telegram(bot_token: str, chat_id: int, text: str, message_id: int |
try:
with urlopen(req, timeout=15) as resp:
result = json.loads(resp.read())
if result.get("ok"):
return True
logger.error("[HOOKS] telegram: API error: %s", result.get("description"))
return False
api_result = json.loads(resp.read())
if api_result.get("ok"):
return _parse_api_message(api_result)
logger.error("[HOOKS] telegram: API error: %s", api_result.get("description"))
return {"ok": False}
except HTTPError as e:
try:
body = json.loads(e.read().decode("utf-8"))
@@ -318,17 +462,17 @@ def send_to_telegram(bot_token: str, chat_id: int, text: str, message_id: int |
logger.error("[HOOKS] telegram: HTTP %d: %s (len=%d)", e.code, description, len(text))
except Exception:
logger.error("[HOOKS] telegram: HTTP %d: %s (len=%d)", e.code, e.reason, len(text))
return False
return {"ok": False}
except URLError as e:
logger.error("[HOOKS] telegram: send failed: %s", e)
return False
return {"ok": False}
except Exception as e:
logger.error("[HOOKS] telegram: unexpected send error: %s", e)
return False
return {"ok": False}
def edit_telegram_message(bot_token: str, chat_id: int, message_id: int, text: str) -> bool:
"""Edit an existing Telegram message via Bot API."""
def edit_telegram_message(bot_token: str, chat_id: int, message_id: int, text: str) -> dict:
"""Edit an existing Telegram message via Bot API. Returns dict with ok, message_id, text."""
url = f"https://api.telegram.org/bot{bot_token}/editMessageText"
try:
@@ -337,9 +481,9 @@ def edit_telegram_message(bot_token: str, chat_id: int, message_id: int, text: s
data = json.dumps(html_payload).encode("utf-8")
req = Request(url, data=data, headers={"Content-Type": "application/json"})
with urlopen(req, timeout=15) as resp:
result = json.loads(resp.read())
if result.get("ok", False):
return True
api_result = json.loads(resp.read())
if api_result.get("ok", False):
return _parse_api_message(api_result)
logger.warning("[HOOKS] telegram: HTML edit failed, plain text fallback")
except Exception as e:
logger.warning("[HOOKS] telegram: HTML edit error: %s, plain text fallback", e)
@@ -350,23 +494,26 @@ def edit_telegram_message(bot_token: str, chat_id: int, message_id: int, text: s
try:
with urlopen(req, timeout=15) as resp:
result = json.loads(resp.read())
return result.get("ok", False)
api_result = json.loads(resp.read())
if api_result.get("ok", False):
return _parse_api_message(api_result)
return {"ok": False}
except Exception as e:
logger.warning("[HOOKS] telegram: edit failed: %s", e)
return False
return {"ok": False}
def _send_with_retry(bot_token: str, chat_id: int, text: str, retries: int = 3) -> bool:
"""Send with retry and exponential backoff."""
def _send_with_retry(bot_token: str, chat_id: int, text: str, retries: int = 3) -> dict:
"""Send with retry and exponential backoff. Returns dict with ok, message_id, text."""
for attempt in range(retries):
if send_to_telegram(bot_token, chat_id, text):
return True
result = send_to_telegram(bot_token, chat_id, text)
if result["ok"]:
return result
if attempt < retries - 1:
delay = 1.0 * (2**attempt)
logger.info("[HOOKS] telegram: send retry %d/%d after %.0fs", attempt + 2, retries, delay)
time.sleep(delay)
return False
return {"ok": False}
def handle(hook_data: dict) -> dict:
@@ -394,16 +541,19 @@ def handle(hook_data: dict) -> dict:
pending_data = json.loads(pending_file.read_text(encoding="utf-8"))
except (json.JSONDecodeError, OSError) as e:
logger.error("[HOOKS] telegram: failed to read pending: %s", e)
pending_file.unlink(missing_ok=True)
if not _in_mirror_dir(pending_file):
pending_file.unlink(missing_ok=True)
return {"stdout": "", "exit_code": 0}
is_mirror = pending_data.get("mirror", False)
chat_id = pending_data.get("chat_id")
bot_token = pending_data.get("bot_token")
processing_message_id = pending_data.get("processing_message_id")
if not chat_id or not bot_token:
logger.error("[HOOKS] telegram: missing chat_id or bot_token in pending")
pending_file.unlink(missing_ok=True)
if not is_mirror:
pending_file.unlink(missing_ok=True)
return {"stdout": "", "exit_code": 0}
response_text = _extract_response(hook_data, transcript_path, pending_data)
@@ -417,11 +567,12 @@ def handle(hook_data: dict) -> dict:
chunks = chunk_text(response_text)
logger.info("[HOOKS] telegram: sending %d chunk(s) (logs_active=%s)", len(chunks), logs_were_active)
all_sent = _deliver_chunks(chunks, bot_token, chat_id, processing_message_id, logs_were_active)
all_sent, chunk_results = _deliver_chunks(chunks, bot_token, chat_id, processing_message_id, logs_were_active)
_write_delivery_log(response_text, chunks, chunk_results, session_id)
if all_sent:
pending_file.unlink(missing_ok=True)
logger.info("[HOOKS] telegram: response delivered, pending cleaned")
_advance_pending(pending_file, pending_data, transcript_path)
else:
logger.error("[HOOKS] telegram: delivery failed — keeping pending for retry")
@@ -430,6 +581,9 @@ def handle(hook_data: dict) -> dict:
def _extract_response(hook_data: dict, transcript_path: str, pending_data: dict) -> str | None:
"""Try JSONL transcript extraction with retries, fall back to last_assistant_message."""
if pending_data.get("mirror"):
return _extract_mirror_response(transcript_path, pending_data)
response_text = None
if transcript_path:
@@ -446,7 +600,7 @@ def _extract_response(hook_data: dict, transcript_path: str, pending_data: dict)
logger.info("[HOOKS] telegram: JSONL retry %.1fs (attempt %d/3)", delay, attempt + 1)
time.sleep(delay)
if not response_text:
if not response_text and not pending_data.get("delivered"):
response_text = (hook_data.get("last_assistant_message") or "").strip()
if response_text:
logger.info("[HOOKS] telegram: using last_assistant_message fallback (%d chars)", len(response_text))
@@ -454,6 +608,23 @@ def _extract_response(hook_data: dict, transcript_path: str, pending_data: dict)
return response_text or None
def _extract_mirror_response(transcript_path: str, pending_data: dict) -> str | None:
"""Extract all new turns for mirror delivery with retries."""
if not transcript_path:
return None
start_line = pending_data.get("transcript_line_after", 0)
for attempt in range(3):
result = extract_mirror_turn(transcript_path, start_line)
if result:
logger.info("[HOOKS] telegram: mirror extraction: %d chars (attempt %d)", len(result), attempt + 1)
return result
if attempt < 2:
delay = [0.2, 0.5][attempt]
logger.info("[HOOKS] telegram: mirror retry %.1fs (attempt %d/3)", delay, attempt + 1)
time.sleep(delay)
return None
def _prepend_branch_prefix(text: str) -> str:
"""Prepend @branch identifier so user knows which branch responded."""
try:
@@ -486,22 +657,110 @@ def _deliver_chunks(
chat_id: int,
processing_message_id: int | None,
logs_were_active: bool,
) -> bool:
"""Send all response chunks to Telegram. Returns True if all succeeded."""
) -> tuple[bool, list[dict]]:
"""Send all response chunks to Telegram. Returns (all_sent, per-chunk results)."""
chunk_results: list[dict] = []
all_sent = True
single = len(chunks) == 1
for i, chunk in enumerate(chunks):
if i == 0 and processing_message_id and not logs_were_active:
text = f"[1/{len(chunks)}]\n{chunk}" if len(chunks) > 1 else chunk
if not edit_telegram_message(bot_token, chat_id, processing_message_id, text):
if not _send_with_retry(bot_token, chat_id, text):
if i == 0 and processing_message_id:
if single and not logs_were_active:
result = edit_telegram_message(bot_token, chat_id, processing_message_id, chunk)
if result["ok"]:
chunk_results.append({"idx": i, "method": "edit", **result})
continue
result = _send_with_retry(bot_token, chat_id, chunk)
chunk_results.append({"idx": i, "method": "send", **result})
if not result["ok"]:
all_sent = False
elif i == 0 and processing_message_id and logs_were_active:
continue
edit_telegram_message(bot_token, chat_id, processing_message_id, "Done.")
text = f"[1/{len(chunks)}]\n{chunk}" if len(chunks) > 1 else chunk
if not _send_with_retry(bot_token, chat_id, text):
all_sent = False
prefix = f"[{i + 1}/{len(chunks)}]\n" if not single else ""
result = _send_with_retry(bot_token, chat_id, prefix + chunk)
chunk_results.append({"idx": i, "method": "send", **result})
if not result["ok"]:
all_sent = False
return all_sent, chunk_results
def _advance_pending(pending_file: Path, pending_data: dict, transcript_path: str) -> None:
"""Advance transcript cursor so later Stops can deliver new text."""
is_mirror = pending_data.get("mirror", False)
if not transcript_path:
if not is_mirror:
pending_file.unlink(missing_ok=True)
logger.info("[HOOKS] telegram: no transcript — pending %s", "kept (mirror)" if is_mirror else "removed")
return
try:
line_count = 0
with open(transcript_path, encoding="utf-8") as f:
for _ in f:
line_count += 1
pending_data["transcript_line_after"] = line_count
pending_data["delivered"] = True
pending_file.write_text(json.dumps(pending_data, indent=2), encoding="utf-8")
logger.info("[HOOKS] telegram: cursor advanced to line %d", line_count)
except OSError as e:
logger.warning("[HOOKS] telegram: cursor advance failed: %s", e)
if not is_mirror:
pending_file.unlink(missing_ok=True)
def _write_delivery_log(intended_text: str, chunks: list[str], chunk_results: list[dict], session_id: str) -> None:
"""Write JSONL record documenting what was delivered to Telegram."""
intended_sha = hashlib.sha256(intended_text.encode("utf-8")).hexdigest()[:16]
delivered_parts = []
for cr in chunk_results:
text = cr.get("text") or ""
text = re.sub(r"^\[\d+/\d+\]\n", "", text)
delivered_parts.append(text)
delivered_concat = "\n\n".join(delivered_parts)
delivered_sha = hashlib.sha256(delivered_concat.encode("utf-8")).hexdigest()[:16]
all_ok = all(cr.get("ok") for cr in chunk_results)
match = all_ok and intended_sha == delivered_sha
culprit = None
if not match:
failed = [cr["idx"] for cr in chunk_results if not cr.get("ok")]
if failed:
culprit = f"delivery_failed: chunks {failed}"
elif abs(len(intended_text) - len(delivered_concat)) > max(len(intended_text) * 0.1, 20):
culprit = f"length_mismatch: intended={len(intended_text)} delivered={len(delivered_concat)}"
else:
prefix = f"[{i + 1}/{len(chunks)}]\n" if len(chunks) > 1 else ""
if not _send_with_retry(bot_token, chat_id, prefix + chunk):
all_sent = False
return all_sent
culprit = "formatting_conversion"
record = {
"ts": time.time(),
"session": session_id[:8] if session_id else "",
"intended_text": intended_text,
"intended_sha256": intended_sha,
"intended_len": len(intended_text),
"chunks": [
{
"idx": cr.get("idx"),
"method": cr.get("method"),
"message_id": cr.get("message_id"),
"ok": cr.get("ok"),
"returned_len": len(cr.get("text") or ""),
}
for cr in chunk_results
],
"delivered_concat": delivered_concat,
"delivered_sha256": delivered_sha,
"delivered_len": len(delivered_concat),
"match": match,
}
if culprit:
record["culprit"] = culprit
try:
_DELIVERY_LOG.parent.mkdir(parents=True, exist_ok=True)
with open(_DELIVERY_LOG, "a", encoding="utf-8") as f:
f.write(json.dumps(record) + "\n")
except OSError as e:
logger.warning("[HOOKS] telegram: delivery log write failed: %s", e)
File diff suppressed because it is too large Load Diff
+141
View File
@@ -5,6 +5,21 @@
"description": "Standards bypass configuration for this branch"
},
"bypass": [
{
"file": "lib/telegram/tests/conftest.py",
"standard": "architecture",
"reason": "Test infrastructure — conftest.py lives in tests/ by pytest convention, not the 3-layer app structure. Test support files are exempt."
},
{
"file": "lib/telegram/tests/conftest.py",
"standard": "encapsulation",
"reason": "Test infrastructure — imports aipass.prax.apps.handlers.logging.direct to redirect log output during tests. Necessary to clear cached logger state; no module entry point exists for this internal reset."
},
{
"file": "lib/telegram/tests/conftest.py",
"standard": "imports",
"reason": "Test infrastructure — sys.path manipulation is intentional: adds src/ root for aipass.* imports and skill root for local apps.handlers.* imports. Required because tests run without a full pip install of the skill."
},
{
"file": "lib/telegram/tests/test_response_router.py",
"standard": "architecture",
@@ -15,6 +30,132 @@
"standard": "encapsulation",
"reason": "Test file — imports handler module directly for unit testing. Tests need direct access to monkeypatch module-level attributes and verify handler behavior."
},
{
"file": "lib/telegram/tests/test_monitor.py",
"standard": "architecture",
"reason": "Test file — lives in tests/ by convention. Test files are exempt from layer architecture standard."
},
{
"file": "lib/telegram/tests/test_monitor.py",
"standard": "encapsulation",
"reason": "Test file — imports handler directly for unit testing. Tests need direct access to patch module-level state and verify handler internals."
},
{
"file": "lib/telegram/tests/test_monitor.py",
"standard": "documentation",
"reason": "Test file — public functions are pytest test methods; docstrings on individual test methods are optional when class docstring and test name are self-documenting."
},
{
"file": "lib/telegram/tests/test_multi_bot.py",
"standard": "architecture",
"reason": "Test file — lives in tests/ by convention. Test files are exempt from layer architecture standard."
},
{
"file": "lib/telegram/tests/test_multi_bot.py",
"standard": "encapsulation",
"reason": "Test file — imports handler directly for unit testing. Tests need direct access to patch module-level state."
},
{
"file": "lib/telegram/tests/test_multi_bot.py",
"standard": "documentation",
"reason": "Test file — pytest test methods are self-documenting via class/method names; docstrings on individual test cases are optional."
},
{
"file": "lib/telegram/tests/test_multi_bot.py",
"standard": "hardcoded_path",
"reason": "Test data — /home/aipass/* paths are mock return values passed to validate_branch() and create_bot() mocks, not real filesystem paths. They represent synthetic registry entries in test fixtures."
},
{
"file": "lib/telegram/tests/test_multi_bot.py",
"standard": "trigger",
"lines": [409],
"reason": "Test-only teardown — base_bot.pending_file.unlink() simulates file absence to verify heartbeat backward-compat behavior. No trigger event is appropriate for test fixture manipulation."
},
{
"file": "lib/telegram/tests/test_attach_only.py",
"standard": "architecture",
"reason": "Test file — lives in tests/ by convention. Test files are exempt from layer architecture standard."
},
{
"file": "lib/telegram/tests/test_attach_only.py",
"standard": "encapsulation",
"reason": "Test file — imports handler directly for unit testing."
},
{
"file": "lib/telegram/tests/test_attach_only.py",
"standard": "documentation",
"reason": "Test file — pytest test methods are self-documenting via class/method names."
},
{
"file": "lib/telegram/tests/test_heartbeat_delivered.py",
"standard": "architecture",
"reason": "Test file — lives in tests/ by convention. Test files are exempt from layer architecture standard."
},
{
"file": "lib/telegram/tests/test_heartbeat_delivered.py",
"standard": "encapsulation",
"reason": "Test file — imports handler directly for unit testing."
},
{
"file": "lib/telegram/tests/test_heartbeat_delivered.py",
"standard": "documentation",
"reason": "Test file — pytest test methods are self-documenting via class/method names."
},
{
"file": "lib/telegram/tests/test_status_reset.py",
"standard": "architecture",
"reason": "Test file — lives in tests/ by convention. Test files are exempt from layer architecture standard."
},
{
"file": "lib/telegram/tests/test_status_reset.py",
"standard": "encapsulation",
"reason": "Test file — imports handler directly for unit testing."
},
{
"file": "lib/telegram/tests/test_status_reset.py",
"standard": "documentation",
"reason": "Test file — pytest test methods are self-documenting via class/method names."
},
{
"file": "lib/telegram/tests/test_mirror_session.py",
"standard": "architecture",
"reason": "Test file — lives in tests/ by convention. Test files are exempt from layer architecture standard."
},
{
"file": "lib/telegram/tests/test_mirror_session.py",
"standard": "encapsulation",
"reason": "Test file — imports handler directly for unit testing."
},
{
"file": "lib/telegram/tests/test_mirror_session.py",
"standard": "hardcoded_path",
"reason": "Test data — /home/test/api is a mock return value for validate_branch(), not a real filesystem path."
},
{
"file": "lib/telegram/tests/test_mirror_session.py",
"standard": "permission_flags",
"reason": "Test file — asserts that launch_mirror_session() correctly passes --dangerously-skip-permissions. Testing the TDPLAN-0009 feature, not bypassing permissions."
},
{
"file": "lib/telegram/apps/handlers/bot_factory.py",
"standard": "handlers",
"reason": "DPLAN-0220 incomplete port — bot_factory imports _api_set_secret from aipass.api.apps.modules.secrets for config persistence. Same pattern as config.py."
},
{
"file": "lib/telegram/apps/handlers/bot_factory.py",
"standard": "json_structure",
"reason": "DPLAN-0220 incomplete port — factory functions predate json_handler.log_operation convention. Adding log_operation to every function is a separate cleanup task."
},
{
"file": "lib/telegram/apps/handlers/bot_factory.py",
"standard": "meta",
"reason": "DPLAN-0220 incomplete port — file uses legacy header format from Dev-Pass port, not the META block format."
},
{
"file": "lib/telegram/apps/handlers/bot_factory.py",
"standard": "permission_flags",
"reason": "TDPLAN-0009 — launch_mirror_session() intentionally uses --dangerously-skip-permissions. Detached tmux mirror sessions have no operator to approve prompts; without it the session hangs forever. Patrick's explicit ask."
},
{
"file": "apps/handlers/loader_handler.py",
"standard": "handlers",
+1 -1
View File
@@ -4,7 +4,7 @@ description: Multi-bot Telegram bridge — routes messages between Telegram and
version: 1.0.0
tags: [communication, bridge, telegram, bot]
requires:
pip: []
pip: [telethon]
bins: [tmux, claude]
config: []
aipass: [api, prax, hooks, cli]
@@ -84,15 +84,14 @@ from .file_handler import (
)
from .bot_factory import (
create_bot,
launch_mirror_session, # noqa: F401
set_bot_commands,
start_service, # noqa: F401
validate_branch,
validate_token,
)
from .telegram_standards import build_botfather_commands
# Secrets API for monitor subscription persistence
from aipass.api.apps.modules.secrets import get_secret as _api_get_secret
from aipass.api.apps.modules.secrets import set_secret as _api_set_secret
from .bot_registry import (
list_bots as registry_list_bots,
get_bot_by_branch,
@@ -133,6 +132,7 @@ POLL_TIMEOUT = 30
SEND_KEYS_DELAY = 0.5
HEARTBEAT_INTERVAL = 30 # seconds
CLAUDE_BIN = str(Path.home() / ".local" / "bin" / "claude")
MIRROR_SESSION_TYPE = "interactive-mirror"
TEMP_DIR = Path("/tmp/telegram_uploads")
MAX_FILE_SIZE = 10 * 1024 * 1024 # 10MB
@@ -161,6 +161,7 @@ class BaseBot:
custom_commands: Optional[dict] = None,
branch_name: Optional[str] = None,
shared_session: Optional[str] = None,
attach_only: bool = False,
) -> None:
"""
Initialize BaseBot.
@@ -196,11 +197,16 @@ class BaseBot:
# Shared-session mode: inject into an existing tmux session instead of creating own
self._shared_session_name = shared_session
self._using_shared_session = False
self._attach_only = attach_only
self._mirror_mapping_written = False
self._last_transcript_path: str | None = None
self._config_chat_id: int | None = None
self.state = {
"running": True,
"message_count": 0,
"start_time": time.time(),
"conversation_start": time.time(),
"last_message_time": 0.0,
}
@@ -262,12 +268,11 @@ class BaseBot:
# Boot-start monitor if a subscription was persisted
self._boot_monitor()
# Check for existing lock
# Check for existing lock (prevents duplicate pollers for the same bot_id).
# Return 0 (not 1) so systemd Restart=on-failure does not restart-loop.
if self._check_lock():
logger.error("Another instance of bot-%s is already running", self.bot_id)
return 1
# Create lock file
logger.error("Another instance of bot-%s is already running — exiting cleanly", self.bot_id)
return 0
self._create_lock()
# Signal handlers
@@ -493,6 +498,7 @@ class BaseBot:
# Start log streamer on first valid message (if branch has a name)
if self._active_chat_id is None and chat_id:
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()
@@ -575,11 +581,16 @@ class BaseBot:
self.send_message(chat_id, "Nothing to cancel.")
return True
# Compute uptime for /status
elapsed = time.time() - self.state["start_time"]
hours, remainder = divmod(int(elapsed), 3600)
minutes, seconds = divmod(remainder, 60)
uptime_str = f"{hours}h {minutes}m {seconds}s"
# Compute conversation uptime (resets on /new) and daemon uptime (since boot)
conv_elapsed = time.time() - self.state.get("conversation_start", self.state["start_time"])
conv_h, conv_rem = divmod(int(conv_elapsed), 3600)
conv_m, conv_s = divmod(conv_rem, 60)
uptime_str = f"{conv_h}h {conv_m}m {conv_s}s"
daemon_elapsed = time.time() - self.state["start_time"]
d_h, d_rem = divmod(int(daemon_elapsed), 3600)
d_m, d_s = divmod(d_rem, 60)
daemon_uptime_str = f"{d_h}h {d_m}m {d_s}s"
# Merge custom commands from constructor and hook
merged_commands = {**self.custom_commands, **self.get_custom_commands()}
@@ -592,6 +603,7 @@ class BaseBot:
uptime=uptime_str,
message_count=self.state.get("message_count"),
chat_id=chat_id,
daemon_uptime=daemon_uptime_str,
)
registry_text = self._build_registry_status()
if registry_text:
@@ -616,8 +628,10 @@ class BaseBot:
action, response_text = result
if action == "new":
self._kill_tmux_session()
self.state["message_count"] = 0
self.state["conversation_start"] = time.time()
self.send_message(chat_id, response_text)
logger.info("Handled /new command - session killed")
logger.info("Handled /new command - session killed, counters reset")
else:
self.send_message(chat_id, result)
logger.info("Handled /%s command", cmd_name)
@@ -655,7 +669,14 @@ class BaseBot:
# Ensure tmux session
if not self.ensure_tmux_session():
logger.error("Cannot process message - tmux session unavailable")
self.send_message(chat_id, "Failed to start Claude session. Check logs.")
if self._attach_only:
self.send_message(
chat_id,
f"⚠️ No canonical session '{self._shared_session_name}' found.\n"
"Start it with the launch wrapper first.",
)
else:
self.send_message(chat_id, "Failed to start Claude session. Check logs.")
return
# Send processing indicator
@@ -746,7 +767,14 @@ class BaseBot:
# Ensure tmux session
if not self.ensure_tmux_session():
logger.error("Cannot process file - tmux session unavailable")
self.send_message(chat_id, "Failed to start Claude session. Check logs.")
if self._attach_only:
self.send_message(
chat_id,
f"⚠️ No canonical session '{self._shared_session_name}' found.\n"
"Start it with the launch wrapper first.",
)
else:
self.send_message(chat_id, "Failed to start Claude session. Check logs.")
if file_type == "text":
file_path.unlink(missing_ok=True)
return
@@ -905,8 +933,10 @@ class BaseBot:
self.send_message(
chat_id,
f"Branch @{branch_name} found at {branch_path}.\n\n"
"Now paste the BotFather token for the new bot.\n"
"(Get one from @BotFather -> /newbot)\n\n"
f"⚠️ BotFather automation unavailable: {telethon_reason}\n"
"Falling back to manual token flow.\n\n"
"Paste the BotFather token for the new bot.\n"
"(Get one from @BotFather → /newbot)\n\n"
"/cancel to abort.",
)
logger.info(
@@ -1250,8 +1280,15 @@ class BaseBot:
"Shared session '%s' found — injecting into existing session",
self._shared_session_name,
)
self._write_mirror_mapping()
return True
else:
if self._attach_only:
logger.error(
"Attach-only: session '%s' not found — refusing to spawn",
self._shared_session_name,
)
return False
self._using_shared_session = False
self.session_name = f"telegram-{self.bot_id}"
logger.warning(
@@ -1260,9 +1297,16 @@ class BaseBot:
)
except FileNotFoundError:
logger.warning("tmux not found while checking shared session '%s'", self._shared_session_name)
if self._attach_only:
return False
self._using_shared_session = False
self.session_name = f"telegram-{self.bot_id}"
# attach-only mode: never spawn a new session
if self._attach_only:
logger.error("Attach-only: no shared session to attach to")
return False
if self._tmux_session_exists():
return True
@@ -1472,31 +1516,173 @@ class BaseBot:
except OSError as e:
logger.warning("Failed to clean stale pending file: %s", e)
def _get_transcript_line_count(self) -> int:
"""
Count lines in the Claude JSONL transcript for Layer 3 position tracking.
def _resolve_active_transcript(self) -> tuple[str | None, int]:
"""Identify the ACTIVE Claude JSONL transcript and return its path and line count.
Returns:
Line count of the JSONL transcript, or 0 if unavailable
Uses the tmux session's pane PID to walk the process tree and find which
JSONL file the Claude process has open. Falls back to most-recently-modified
JSONL that was touched in the last 5 minutes. Returns (None, 0) when no
active transcript can be identified — safe default that resets the cursor.
"""
slug = str(self.work_dir).replace("/", "-")
# Look for transcript files matching the session pattern
projects_dir = Path.home() / ".claude" / "projects" / slug
if not projects_dir.exists():
return 0
return None, 0
# Find the most recent JSONL transcript
jsonl_files = sorted(projects_dir.glob("*.jsonl"), key=lambda p: p.stat().st_mtime, reverse=True)
jsonl_files = list(projects_dir.glob("*.jsonl"))
if not jsonl_files:
return 0
return None, 0
# Strategy 1: find the JSONL open by a child of the tmux pane
pane_pid = self._get_tmux_pane_pid()
if pane_pid:
target_names = {f.name for f in jsonl_files}
found = self._find_open_jsonl(pane_pid, projects_dir, target_names)
if found:
return str(found), self._count_file_lines(found)
# Strategy 2: most recently modified JSONL, but only if touched < 5 min ago
now = time.time()
recent = sorted(
((f, f.stat().st_mtime) for f in jsonl_files),
key=lambda x: x[1],
reverse=True,
)
if recent and (now - recent[0][1]) < 300:
return str(recent[0][0]), self._count_file_lines(recent[0][0])
return None, 0
def _find_open_jsonl(self, pane_pid: int, projects_dir: Path, target_names: set[str]) -> Path | None:
"""Scan /proc fd links for descendant PIDs to find an open JSONL in *projects_dir*."""
child_pids = self._get_descendant_pids(pane_pid)
for pid in child_pids:
match = self._scan_pid_fds(pid, projects_dir, target_names)
if match:
return match
return None
@staticmethod
def _scan_pid_fds(pid: int, projects_dir: Path, target_names: set[str]) -> Path | None:
"""Check /proc/<pid>/fd/ for an open JSONL matching *target_names*."""
fd_dir = Path(f"/proc/{pid}/fd")
try:
text = jsonl_files[0].read_text(encoding="utf-8").strip()
entries = list(fd_dir.iterdir())
except OSError as e:
logger.info("Cannot list fds for pid %d: %s", pid, e)
return None
for fd in entries:
try:
target = fd.resolve()
except OSError as e:
logger.info("Cannot resolve fd %s: %s", fd.name, e)
continue
if target.parent == projects_dir and target.name in target_names:
return target
return None
def _get_tmux_pane_pid(self) -> int | None:
"""Get the shell PID of the first pane in the current tmux session."""
try:
result = subprocess.run(
["tmux", "list-panes", "-t", self.session_name, "-F", "#{pane_pid}"],
capture_output=True,
text=True,
timeout=5,
)
if result.returncode == 0 and result.stdout.strip():
return int(result.stdout.strip().split("\n")[0])
except (subprocess.TimeoutExpired, OSError, ValueError) as e:
logger.info("Could not get tmux pane PID: %s", e)
return None
@staticmethod
def _get_descendant_pids(parent_pid: int) -> list[int]:
"""Walk /proc to collect all descendant PIDs of *parent_pid*."""
children: dict[int, list[int]] = {}
try:
for entry in Path("/proc").iterdir():
if not entry.name.isdigit():
continue
try:
stat_text = (entry / "stat").read_text(encoding="utf-8")
ppid = int(stat_text.split(") ")[1].split()[1])
children.setdefault(ppid, []).append(int(entry.name))
except (OSError, IndexError, ValueError) as e:
logger.info("Skipping /proc/%s/stat: %s", entry.name, e)
continue
except OSError as e:
logger.info("Cannot scan /proc for descendant PIDs: %s", e)
return []
result: list[int] = []
queue = children.get(parent_pid, [])[:]
while queue:
pid = queue.pop()
result.append(pid)
queue.extend(children.get(pid, []))
return result
@staticmethod
def _count_file_lines(path: Path) -> int:
"""Count newline-delimited lines in *path*."""
try:
text = path.read_text(encoding="utf-8").strip()
return len(text.split("\n")) if text else 0
except OSError as e:
logger.warning("Could not read transcript for line count: %s", e)
return 0
def _get_transcript_line_count(self) -> int:
"""Count lines in the active JSONL transcript (compat shim)."""
_, count = self._resolve_active_transcript()
return count
def _write_mirror_mapping(self) -> None:
"""Write persistent mirror mapping file per THE CONTRACT (TDPLAN-0009).
Written at first successful attach AND rewritten whenever the active
transcript changes (session restart). @hooks reads this every turn
to find the chat_id and cursor for mirror delivery.
"""
if not self._attach_only:
return
chat_id = self._active_chat_id or getattr(self, "_config_chat_id", None)
if chat_id is None:
return
transcript_path, line_count = self._resolve_active_transcript()
if self._mirror_mapping_written:
if transcript_path == self._last_transcript_path:
return
logger.info(
"Transcript changed (%s → %s) — rewriting mirror mapping",
self._last_transcript_path,
transcript_path,
)
mapping_dir = Path.home() / ".aipass" / "telegram_bots"
mapping_dir.mkdir(parents=True, exist_ok=True)
mapping_file = mapping_dir / f"bot-{self.bot_id}.json"
mapping = {
"chat_id": chat_id,
"bot_token": self.bot_token,
"session_name": self.session_name,
"work_dir": str(self.work_dir),
"mirror": True,
"transcript_line_after": line_count,
}
try:
mapping_file.write_text(json.dumps(mapping, indent=2), encoding="utf-8")
self._mirror_mapping_written = True
self._last_transcript_path = transcript_path
logger.info("Mirror mapping written: %s (cursor=%d)", mapping_file, line_count)
except OSError as e:
logger.error("Failed to write mirror mapping: %s", e)
# =============================================
# HEARTBEAT THREAD
# =============================================
@@ -1520,8 +1706,8 @@ class BaseBot:
if self._heartbeat_stop.is_set():
break
# Only update if pending file still exists and tmux alive
if not self.pending_file.exists():
# Stop if response has been delivered or tmux died
if self._is_pending_delivered():
break
if not self._tmux_session_exists():
break
@@ -1540,6 +1726,17 @@ class BaseBot:
self._heartbeat_thread.join(timeout=5)
self._heartbeat_thread = None
def _is_pending_delivered(self) -> bool:
"""Check if the pending file has been marked as delivered by @hooks."""
if not self.pending_file.exists():
return True
try:
data = json.loads(self.pending_file.read_text(encoding="utf-8"))
return bool(data.get("delivered"))
except (json.JSONDecodeError, OSError) as e:
logger.warning("Failed to read pending file %s: %s", self.pending_file, e)
return False
@staticmethod
def _format_elapsed(seconds: float) -> str:
"""
@@ -1715,32 +1912,45 @@ class BaseBot:
self._monitor_streamer.start()
logger.info("Boot-started monitor streamer (chat_id=%s, mode=%s)", chat_id, mode)
def _monitor_subscription_file(self) -> Path:
"""Return path to the local monitor subscription file."""
return Path.home() / ".aipass" / "telegram_bots" / f".{self.bot_id}_monitor.json"
def _load_monitor_subscription(self) -> dict | None:
"""Load monitor subscription from the @api secrets store."""
try:
result = _api_get_secret("telegram", "monitor", as_json=True)
if isinstance(result, dict) and result.get("chat_id"):
return result
"""Load monitor subscription from local state file."""
sub_file = self._monitor_subscription_file()
if not sub_file.exists():
return None
except Exception as e:
try:
data = json.loads(sub_file.read_text(encoding="utf-8"))
if isinstance(data, dict) and data.get("chat_id"):
return data
return None
except (json.JSONDecodeError, OSError) as e:
logger.warning("Failed to load monitor subscription: %s", e)
return None
def _save_monitor_subscription(self, chat_id: int, mode: str) -> bool:
"""Persist monitor subscription to the @api secrets store."""
"""Persist monitor subscription to local state file."""
sub_file = self._monitor_subscription_file()
try:
_api_set_secret("telegram", "monitor", {"chat_id": chat_id, "mode": mode}, as_json=True)
sub_file.parent.mkdir(parents=True, exist_ok=True)
sub_file.write_text(
json.dumps({"chat_id": chat_id, "mode": mode}, indent=2),
encoding="utf-8",
)
return True
except Exception as e:
except OSError as e:
logger.error("Failed to save monitor subscription: %s", e)
return False
def _clear_monitor_subscription(self) -> bool:
"""Clear persisted monitor subscription."""
sub_file = self._monitor_subscription_file()
try:
_api_set_secret("telegram", "monitor", {}, as_json=True)
sub_file.unlink(missing_ok=True)
return True
except Exception as e:
except OSError as e:
logger.error("Failed to clear monitor subscription: %s", e)
return False
@@ -1935,9 +2145,10 @@ if __name__ == "__main__":
allowed_user_ids=config.get("allowed_user_ids", []),
branch_name=config.get("branch_name"),
shared_session=config.get("shared_session"),
attach_only=config.get("attach_only", False),
)
if hasattr(bot, "_config_chat_id") and config.get("chat_id"):
if config.get("chat_id"):
bot._config_chat_id = config["chat_id"]
sys.exit(bot.run())
@@ -66,6 +66,7 @@ from .bot_registry import (
TELEGRAM_API = "https://api.telegram.org/bot{token}"
SYSTEMD_DIR = Path.home() / ".config" / "systemd" / "user"
_BOT_CONFIG_DIR = Path.home() / ".aipass" / "telegram_bots"
CLAUDE_BIN = str(Path.home() / ".local" / "bin" / "claude")
# Command menu built from telegram_standards (single source of truth)
@@ -324,6 +325,84 @@ def start_bot_process(bot_id: str) -> bool:
return False
def launch_mirror_session(
bot_id: str,
work_dir: str,
session_name: str,
) -> bool:
"""Launch the canonical tmux mirror session with --dangerously-skip-permissions.
Creates a detached tmux session running Claude Code in autonomous mirror mode.
The session uses --dangerously-skip-permissions because a detached tmux session
has no operator to approve permission prompts — without it, the session hangs.
Args:
bot_id: Bot identifier (set as AIPASS_BOT_ID env var in the session).
work_dir: Working directory for the Claude session.
session_name: tmux session name (e.g. "telegram-api").
Returns:
True if the session was launched successfully.
"""
import os
try:
result = subprocess.run(
["tmux", "has-session", "-t", session_name],
capture_output=True,
)
if result.returncode == 0:
logger.info("Mirror session '%s' already running — skipping launch", session_name)
return True
except FileNotFoundError:
logger.error("tmux not found — cannot launch mirror session")
return False
env = os.environ.copy()
env.pop("CLAUDECODE", None)
try:
subprocess.run(
["tmux", "new-session", "-d", "-s", session_name, "-c", work_dir],
check=True,
capture_output=True,
env=env,
)
except subprocess.CalledProcessError as e:
logger.error("Failed to create tmux session '%s': %s", session_name, e)
return False
subprocess.run(
[
"tmux",
"send-keys",
"-t",
session_name,
f"export AIPASS_BOT_ID={bot_id}",
"Enter",
],
capture_output=True,
)
import time
time.sleep(0.3)
claude_cmd = f"AIPASS_SESSION_TYPE=interactive-mirror {CLAUDE_BIN} --dangerously-skip-permissions"
subprocess.run(
["tmux", "send-keys", "-t", session_name, claude_cmd, "Enter"],
capture_output=True,
)
logger.info(
"Mirror session '%s' launched (bot_id=%s, work_dir=%s, skip-perms=true)",
session_name,
bot_id,
work_dir,
)
return True
def stop_service(bot_id: str) -> bool:
"""
Stop the systemd user service for a bot.
@@ -359,6 +438,41 @@ def stop_service(bot_id: str) -> bool:
return False
def start_service(bot_id: str) -> bool:
"""
Start the systemd user service for a bot.
Runs: systemctl --user start telegram-bot@{bot_id}
Args:
bot_id: Bot identifier used in the service template.
Returns:
True if the service was started successfully, False otherwise.
"""
SERVICE_NAME = f"telegram-bot@{bot_id}"
try:
result = subprocess.run(
["systemctl", "--user", "start", SERVICE_NAME],
capture_output=True,
text=True,
timeout=10,
)
if result.returncode == 0:
logger.info("Started systemd service: %s", SERVICE_NAME)
return True
logger.warning("Failed to start service %s: %s", SERVICE_NAME, result.stderr.strip())
return False
except subprocess.TimeoutExpired:
logger.warning("Timeout starting service: %s", SERVICE_NAME)
return False
except OSError as e:
logger.warning("Error starting service %s: %s", SERVICE_NAME, e)
return False
# =============================================
# BOT LIFECYCLE
# =============================================
@@ -371,6 +485,9 @@ def create_bot(
work_dir: Optional[str] = None,
bot_name: Optional[str] = None,
allowed_user_ids: Optional[list[int]] = None,
shared_session: Optional[str] = None,
attach_only: bool = False,
chat_id: Optional[int] = None,
) -> Optional[dict]:
"""
Create a new bot: validate, write config, register, setup systemd.
@@ -383,7 +500,8 @@ def create_bot(
5. Register in bot registry
6. Set BotFather commands via setMyCommands API
7. Enable systemd service
8. Auto-start the bot process
7.5. Mirror mode: launch canonical tmux session with skip-perms
8. Auto-start the bot (systemd for mirror, Popen for standard)
Args:
bot_id: Unique identifier for this bot (e.g., "dev_central", "base").
@@ -392,6 +510,9 @@ def create_bot(
work_dir: Working directory for Claude sessions. Defaults to home dir if None.
bot_name: Human-readable bot name. Auto-generated if None.
allowed_user_ids: List of Telegram user IDs allowed to use this bot.
shared_session: tmux session name for mirror mode (attach to existing session).
attach_only: When True, bot only attaches — never spawns its own session.
chat_id: Pre-configured Telegram chat ID for mirror mapping.
Returns:
Bot info dict on success, None on any failure.
@@ -454,6 +575,9 @@ def create_bot(
"work_dir": RESOLVED_WORK_DIR,
"allowed_user_ids": allowed_user_ids or [],
"created_at": datetime.now(timezone.utc).isoformat(),
"shared_session": shared_session,
"attach_only": attach_only,
"chat_id": chat_id,
}
# Step 4a: Write to @api secrets store (load_bot_config reads from here)
@@ -494,16 +618,30 @@ def create_bot(
# Step 7: Enable systemd service
enable_service(bot_id)
# Step 8: Auto-start the bot process
started = start_bot_process(bot_id)
# Step 7.5: Mirror mode — launch canonical tmux session
mirror_launched = False
if shared_session and attach_only:
mirror_launched = launch_mirror_session(session_name=shared_session, bot_id=bot_id, work_dir=RESOLVED_WORK_DIR)
if not mirror_launched:
logger.warning(
"Mirror session launch failed for '%s' — bot may need manual session start",
bot_id,
)
# Step 8: Auto-start the bot poller
if shared_session and attach_only:
started = start_service(bot_id)
else:
started = start_bot_process(bot_id)
logger.info(
"Bot created: %s (@%s, branch=%s, work_dir=%s, started=%s)",
"Bot created: %s (@%s, branch=%s, work_dir=%s, started=%s, mirror=%s)",
bot_id,
BOT_USERNAME,
branch_name,
RESOLVED_WORK_DIR,
started,
mirror_launched,
)
return {
@@ -515,6 +653,9 @@ def create_bot(
"config_path": str(CONFIG_PATH),
"service_name": f"telegram-bot@{bot_id}",
"auto_started": started,
"shared_session": shared_session,
"attach_only": attach_only,
"mirror_launched": mirror_launched,
}
@@ -1,17 +1,9 @@
# =================== AIPass ====================
# Name: bot_operations.py - Bot operation handlers for multi-bot module
# Date: 2026-02-24
# Version: 1.0.0
# Category: api/handlers/telegram
#
# CHANGELOG (Max 5 entries):
# - v1.0.0 (2026-02-24): Initial - start, stop, status, list operations for multi-bot system
#
# CODE STANDARDS:
# - Pure functions with proper error handling (graceful - never raise)
# - No Prax imports (handler tier 3)
# - Stdlib only (subprocess for systemd)
# - Returns values for caller to log/display - no handler-level logging
# Name: bot_operations.py
# Description: Bot lifecycle operation handlers — start, stop, status, list
# Version: 1.0.1
# Created: 2026-02-24
# Modified: 2026-06-29
# =============================================
"""
@@ -32,6 +24,12 @@ All functions return values - the module layer handles logging and display.
import subprocess
from pathlib import Path
# Logging
from aipass.prax import logger
# JSON handler (seedgo standard)
from aipass.skills.apps.handlers.json import json_handler # noqa: F401
# Internal handler imports
from .base_bot import BaseBot
from .branch_plugin import BranchPlugin
@@ -57,6 +55,7 @@ def start_bot(bot_id: str) -> int | None:
Returns:
Bot exit code, or None if config loading failed.
"""
json_handler.log_operation("start_bot", {"bot_id": bot_id})
config = load_bot_config(bot_id)
if not config:
return None
@@ -69,6 +68,9 @@ def start_bot(bot_id: str) -> int | None:
bot_name = config.get("bot_name", f"AIPass {bot_id} Bot")
allowed_user_ids = config.get("allowed_user_ids", [])
branch_name = config.get("branch_name")
shared_session = config.get("shared_session")
attach_only = config.get("attach_only", False)
chat_id = config.get("chat_id")
if branch_name:
bot = BranchPlugin(
@@ -78,6 +80,8 @@ def start_bot(bot_id: str) -> int | None:
work_dir=work_dir,
bot_name=bot_name,
allowed_user_ids=allowed_user_ids,
shared_session=shared_session,
attach_only=attach_only,
)
else:
bot = BaseBot(
@@ -86,8 +90,13 @@ def start_bot(bot_id: str) -> int | None:
work_dir=work_dir,
bot_name=bot_name,
allowed_user_ids=allowed_user_ids,
shared_session=shared_session,
attach_only=attach_only,
)
if chat_id is not None:
bot._config_chat_id = chat_id
return bot.run()
@@ -117,8 +126,10 @@ def stop_bot(bot_id: str) -> tuple[bool, str]:
return False, f"Failed to stop {service_name}: {result.stderr.strip()}"
except subprocess.TimeoutExpired:
logger.warning("Timeout stopping %s", service_name)
return False, f"Timeout stopping {service_name}"
except OSError as e:
logger.warning("Error stopping %s: %s", service_name, e)
return False, f"Error stopping {service_name}: {e}"
@@ -1,6 +1,29 @@
# =================== AIPass ====================
# Name: telegram_standards.py
# Description: Standard command registry, response builders, and helpers for Telegram bots
# Version: 1.0.0
# Created: 2026-02-24
# Modified: 2026-06-29
# =============================================
"""
Telegram Standards — shared command registry and response builders.
Provides the canonical STANDARD_COMMANDS dict, text builder functions for
/start, /help, /status, and /new responses, the BotFather setMyCommands
payload builder, and the parse_command / handle_standard_command dispatcher
used by BaseBot and all subclasses.
All functions are pure (no side effects) except _tmux_session_exists which
calls subprocess to check tmux state.
"""
import subprocess
from typing import Optional
from aipass.skills.apps.handlers.json import json_handler # noqa: F401
from aipass.prax import logger
# =============================================
# STANDARD COMMAND REGISTRY
@@ -136,6 +159,7 @@ def build_status_text(
uptime: Optional[str] = None,
message_count: Optional[int] = None,
chat_id: Optional[str | int] = None,
daemon_uptime: Optional[str] = None,
) -> str:
"""
Build the /status response.
@@ -146,9 +170,10 @@ def build_status_text(
Args:
session_name: tmux session name (e.g., "telegram-assistant").
branch_name: Branch name (e.g., "assistant").
uptime: Optional human-readable uptime string.
message_count: Optional count of messages processed.
uptime: Optional conversation uptime (resets on /new).
message_count: Optional count of messages in current conversation.
chat_id: Optional Telegram chat ID to display.
daemon_uptime: Optional daemon process uptime (since boot).
Returns:
Formatted status text string.
@@ -165,6 +190,8 @@ def build_status_text(
lines.append(f"Uptime: {uptime}")
if message_count is not None:
lines.append(f"Messages: {message_count}")
if daemon_uptime:
lines.append(f"Daemon up: {daemon_uptime}")
return "\n".join(lines)
@@ -279,6 +306,8 @@ def handle_standard_command(
- tuple[str, str]: ("new", response_text) for /new command
- None: Command is not a standard command
"""
json_handler.log_operation("standard_command", {"command": command, "branch": branch_name})
if command == "start":
return build_welcome_text(
bot_name=bot_name,
@@ -327,5 +356,5 @@ def _tmux_session_exists(session_name: str) -> bool:
)
return result.returncode == 0
except FileNotFoundError:
# tmux not installed
logger.warning("tmux not found while checking session '%s'", session_name)
return False
@@ -1,18 +1,19 @@
# ===================AIPASS====================
# META DATA HEADER
# Name: conftest.py - Telegram skill test configuration
# Date: 2026-06-15
# =================== AIPass ====================
# Name: conftest.py
# Description: Telegram skill test configuration — path setup and shared fixtures
# Version: 1.0.0
# Category: skills/telegram/tests
#
# CHANGELOG (Max 5 entries):
# - v1.0.0 (2026-06-15): Initial implementation — prax log redirect + path setup
#
# CODE STANDARDS:
# - Adds src/ and skill root to sys.path for test imports
# Created: 2026-06-15
# Modified: 2026-06-29
# =============================================
"""Telegram skill test configuration."""
"""
Telegram skill test configuration.
Sets up sys.path so that both aipass.* (installed package) and the local
apps.handlers.* namespace are importable from tests without a full pip install.
Also stubs the optional telethon dependency and redirects Prax logger output
to a temp dir so test runs don't pollute production log files.
"""
import os
import shutil
@@ -27,26 +28,24 @@ if "AIPASS_TEST_LOG_DIR" not in os.environ:
import pytest
# Add src/ to path so aipass.* is importable
_src_root = Path(__file__).resolve().parents[5] # noqa: E402
# sys.path setup is intentional test infrastructure — both entries are needed:
# _src_root → resolves aipass.* installed-package imports
# _skill_root → resolves the local apps.handlers.* namespace used by all tests
_src_root = Path(__file__).resolve().parents[5]
if str(_src_root) not in sys.path:
sys.path.insert(0, str(_src_root))
# Add telegram skill root so apps.handlers.* is importable
_skill_root = Path(__file__).resolve().parents[1] # noqa: E402
_skill_root = Path(__file__).resolve().parents[1]
if str(_skill_root) not in sys.path:
sys.path.insert(0, str(_skill_root))
# Telethon stub — telethon is an OPTIONAL runtime dependency (MTProto client),
# deliberately NOT in pyproject so the core stays lightweight (botfather_client.py
# guards it with TELETHON_AVAILABLE). The botfather_client tests mock all Telethon
# classes (patch("telethon.TelegramClient"), etc.), but unittest.mock.patch must
# IMPORT the target's parent module to set the attribute — which raises
# ModuleNotFoundError when telethon isn't installed (e.g. in CI). Register a minimal
# stub so those patch targets resolve. The guard never clobbers a real telethon if
# one is installed. Real FloodWaitError/RPCError classes are required for the
# success/timeout tests, where _send_and_wait imports them but does not patch them.
# Telethon stub — telethon is an optional dependency (pyproject [telegram] extra).
# In CI or minimal installs it may not be present. botfather_client.py guards with
# TELETHON_AVAILABLE. The tests mock Telethon classes (patch("telethon.TelegramClient")),
# but unittest.mock.patch must IMPORT the parent module — which raises
# ModuleNotFoundError when telethon isn't installed. Register a minimal stub so those
# patch targets resolve. Never clobbers a real telethon if one is installed.
if "telethon" not in sys.modules:
_telethon_stub = types.ModuleType("telethon")
_telethon_errors = types.ModuleType("telethon.errors")
@@ -0,0 +1,297 @@
# =================== AIPass ====================
# Name: test_attach_only.py
# Description: Tests for attach-only mode (TDPLAN-0009 Stage 1)
# Version: 1.0.0
# Created: 2026-06-29
# Modified: 2026-06-29
# =============================================
"""
Tests for attach-only mode (TDPLAN-0009 Stage 1).
The bot attaches to a pre-existing tmux session and never spawns its own.
When attach_only=True + shared_session is set:
- Attaches to existing named tmux session (no spawn)
- Missing session → loud error, NO spawn of telegram-{bot_id}
- Persistent mapping file written with all CONTRACT fields + seeded cursor
- No lock created for the shared session
- inject_message reaches the session (unchanged, tested via existing tests)
"""
import json
from unittest.mock import patch, MagicMock
import pytest
from apps.handlers.base_bot import BaseBot # type: ignore[import-not-found]
@pytest.fixture
def _patch_base_bot_deps(tmp_path):
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_bot(tmp_path, _patch_base_bot_deps, attach_only=False, shared_session=None):
workdir = tmp_path / "workdir"
workdir.mkdir(exist_ok=True)
with patch("apps.handlers.base_bot.PENDING_DIR", tmp_path):
bot = BaseBot(
bot_id="mirror_test",
bot_token="123:FAKETOKEN",
work_dir=workdir,
bot_name="Mirror Test Bot",
allowed_user_ids=[111],
branch_name="devpulse",
shared_session=shared_session,
attach_only=attach_only,
)
bot.send_message = MagicMock(return_value={"ok": True, "message_id": 1})
return bot
# =============================================
# 1. ensure_tmux_session — attach-only behavior
# =============================================
class TestAttachOnly:
"""Attach-only mode: attach to existing session, never spawn."""
def test_attaches_to_existing_session(self, tmp_path, _patch_base_bot_deps):
"""When shared session exists, bot attaches and returns True."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
result_obj = MagicMock()
result_obj.returncode = 0
with patch("subprocess.run", return_value=result_obj):
assert bot.ensure_tmux_session() is True
assert bot.session_name == "devpulse"
assert bot._using_shared_session is True
def test_missing_session_returns_false(self, tmp_path, _patch_base_bot_deps):
"""When shared session is missing, attach-only returns False (no spawn)."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
result_obj = MagicMock()
result_obj.returncode = 1
with patch("subprocess.run", return_value=result_obj):
assert bot.ensure_tmux_session() is False
# Session name should NOT fall back to telegram-{bot_id}
assert bot._using_shared_session is False
def test_missing_session_does_not_spawn(self, tmp_path, _patch_base_bot_deps):
"""Attach-only never creates a telegram-{bot_id} session."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
result_obj = MagicMock()
result_obj.returncode = 1
with patch("subprocess.run", return_value=result_obj) as mock_run:
bot.ensure_tmux_session()
# Should only have checked has-session, never new-session
calls = [str(c) for c in mock_run.call_args_list]
for call in calls:
assert "new-session" not in call
def test_no_shared_session_config_returns_false(self, tmp_path, _patch_base_bot_deps):
"""Attach-only without shared_session configured returns False."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session=None)
assert bot.ensure_tmux_session() is False
def test_non_attach_mode_falls_back_on_missing(self, tmp_path, _patch_base_bot_deps):
"""Without attach_only, missing shared session falls back to own session."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=False, shared_session="devpulse")
has_session = MagicMock(returncode=1)
with patch("subprocess.run", return_value=has_session):
# Falls back to telegram-{bot_id} and tries to spawn — we just check the name
bot.ensure_tmux_session()
assert bot.session_name == "telegram-mirror_test"
def test_handle_message_shows_error_on_attach_fail(self, tmp_path, _patch_base_bot_deps):
"""handle_message shows specific error when attach-only fails."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
result_obj = MagicMock()
result_obj.returncode = 1
with patch("subprocess.run", return_value=result_obj):
bot.handle_message(42, "hello", {"message_id": 1})
msg = bot.send_message.call_args[0][1]
assert "No canonical session" in msg
assert "devpulse" in msg
# =============================================
# 2. Mirror mapping file (THE CONTRACT)
# =============================================
class TestMirrorMapping:
"""Persistent mapping file written at attach with CONTRACT fields."""
def test_mapping_written_on_attach(self, tmp_path, _patch_base_bot_deps):
"""Mapping file written when attach-only bot attaches to existing session."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
bot._active_chat_id = 42
mapping_dir = tmp_path / ".aipass" / "telegram_bots"
result_obj = MagicMock()
result_obj.returncode = 0
with (
patch("subprocess.run", return_value=result_obj),
patch("pathlib.Path.home", return_value=tmp_path),
):
bot.ensure_tmux_session()
mapping_file = mapping_dir / "bot-mirror_test.json"
assert mapping_file.exists()
data = json.loads(mapping_file.read_text())
assert data["chat_id"] == 42
assert data["bot_token"] == "123:FAKETOKEN"
assert data["session_name"] == "devpulse"
assert data["mirror"] is True
assert "transcript_line_after" in data
assert data["work_dir"] == str(tmp_path / "workdir")
def test_mapping_uses_config_chat_id(self, tmp_path, _patch_base_bot_deps):
"""Mapping uses _config_chat_id when _active_chat_id is not set."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
bot._config_chat_id = 99
mapping_dir = tmp_path / ".aipass" / "telegram_bots"
result_obj = MagicMock()
result_obj.returncode = 0
with (
patch("subprocess.run", return_value=result_obj),
patch("pathlib.Path.home", return_value=tmp_path),
):
bot.ensure_tmux_session()
mapping_file = mapping_dir / "bot-mirror_test.json"
assert mapping_file.exists()
data = json.loads(mapping_file.read_text())
assert data["chat_id"] == 99
def test_mapping_deferred_when_no_chat_id(self, tmp_path, _patch_base_bot_deps):
"""Mapping deferred if no chat_id at attach time."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
mapping_dir = tmp_path / ".aipass" / "telegram_bots"
result_obj = MagicMock()
result_obj.returncode = 0
with (
patch("subprocess.run", return_value=result_obj),
patch("pathlib.Path.home", return_value=tmp_path),
):
bot.ensure_tmux_session()
mapping_file = mapping_dir / "bot-mirror_test.json"
assert not mapping_file.exists()
assert bot._mirror_mapping_written is False
def test_mapping_written_once(self, tmp_path, _patch_base_bot_deps):
"""Mapping file only written once (idempotent)."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
bot._active_chat_id = 42
result_obj = MagicMock()
result_obj.returncode = 0
with (
patch("subprocess.run", return_value=result_obj),
patch("pathlib.Path.home", return_value=tmp_path),
):
bot.ensure_tmux_session()
bot.ensure_tmux_session()
assert bot._mirror_mapping_written is True
def test_mapping_not_written_for_non_attach(self, tmp_path, _patch_base_bot_deps):
"""Non-attach-only bots don't write mirror mapping."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=False, shared_session="devpulse")
bot._active_chat_id = 42
mapping_dir = tmp_path / ".aipass" / "telegram_bots"
result_obj = MagicMock()
result_obj.returncode = 0
with (
patch("subprocess.run", return_value=result_obj),
patch("pathlib.Path.home", return_value=tmp_path),
):
bot.ensure_tmux_session()
mapping_file = mapping_dir / "bot-mirror_test.json"
assert not mapping_file.exists()
# =============================================
# 3. Lock file skipped in attach-only
# =============================================
class TestAttachOnlyLock:
"""Lock file still created in attach-only (prevents duplicate pollers)."""
def test_lock_still_created_in_attach_mode(self, tmp_path, _patch_base_bot_deps):
"""run() still creates lock in attach-only mode (one poller per bot_id)."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="devpulse")
with (
patch.object(bot, "_check_lock", return_value=False) as mock_check,
patch.object(bot, "_create_lock") as mock_create,
patch.object(bot, "verify_connection", return_value=True),
patch.object(bot, "_set_command_menu"),
patch.object(bot, "_boot_monitor"),
patch.object(bot, "clean_stale_pending"),
patch.object(bot, "_load_offset", return_value=0),
patch.object(bot, "poll_updates", side_effect=KeyboardInterrupt),
patch("apps.handlers.base_bot.PENDING_DIR", tmp_path),
):
bot.run()
mock_check.assert_called_once()
mock_create.assert_called_once()
def test_lock_created_in_normal_mode(self, tmp_path, _patch_base_bot_deps):
"""run() creates lock in normal (non-attach) mode."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=False)
with (
patch.object(bot, "_check_lock", return_value=False),
patch.object(bot, "_create_lock") as mock_create,
patch.object(bot, "verify_connection", return_value=True),
patch.object(bot, "_set_command_menu"),
patch.object(bot, "_boot_monitor"),
patch.object(bot, "clean_stale_pending"),
patch.object(bot, "_load_offset", return_value=0),
patch.object(bot, "poll_updates", side_effect=KeyboardInterrupt),
patch("apps.handlers.base_bot.PENDING_DIR", tmp_path),
):
bot.run()
mock_create.assert_called_once()
# =============================================
# 4. Config loading
# =============================================
class TestAttachOnlyConfig:
"""attach_only flag passed through config."""
def test_attach_only_defaults_false(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
assert bot._attach_only is False
def test_attach_only_set_true(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True)
assert bot._attach_only is True
@@ -0,0 +1,246 @@
# =================== AIPass ====================
# Name: test_heartbeat_delivered.py
# Description: Tests for heartbeat delivered-flag fix (DPLAN-0223)
# Version: 1.0.0
# Created: 2026-06-29
# Modified: 2026-06-29
# =============================================
"""
Tests for heartbeat delivered-flag fix (DPLAN-0223).
The heartbeat now stops when the pending file contains 'delivered': true
(set by @hooks _advance_pending), instead of waiting for file deletion.
Tests cover:
- Heartbeat stops when pending file has delivered=true
- Heartbeat continues when pending file exists but not delivered
- Heartbeat stops when pending file is absent (backward compat)
- Multi-Stop keeps file alive (delivered flag, not deletion)
- Reply not clobbered: heartbeat does NOT edit after delivery
"""
import json
import time
import pytest
from unittest.mock import patch
from apps.handlers.base_bot import BaseBot # type: ignore[import-not-found]
@pytest.fixture
def _patch_base_bot_deps(tmp_path):
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_bot(tmp_path, _patch_base_bot_deps, pending_dir=None):
pdir = pending_dir or tmp_path
with patch("apps.handlers.base_bot.PENDING_DIR", pdir):
workdir = tmp_path / "workdir"
workdir.mkdir(exist_ok=True)
bot = BaseBot(
bot_id="heartbeat_test",
bot_token="123:FAKETOKEN",
work_dir=workdir,
bot_name="Heartbeat Test Bot",
allowed_user_ids=[111],
branch_name=None,
)
bot.pending_file = pdir / "bot-heartbeat_test.json"
return bot
# =============================================
# 1. _is_pending_delivered
# =============================================
class TestIsPendingDelivered:
"""Verify _is_pending_delivered reads the delivered flag correctly."""
def test_returns_true_when_delivered(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
bot.pending_file.write_text(json.dumps({"chat_id": 42, "delivered": True}), encoding="utf-8")
assert bot._is_pending_delivered() is True
def test_returns_false_when_not_delivered(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
bot.pending_file.write_text(json.dumps({"chat_id": 42}), encoding="utf-8")
assert bot._is_pending_delivered() is False
def test_returns_true_when_file_absent(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
assert not bot.pending_file.exists()
assert bot._is_pending_delivered() is True
def test_returns_false_on_corrupt_json(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
bot.pending_file.write_text("not json", encoding="utf-8")
assert bot._is_pending_delivered() is False
def test_returns_false_when_delivered_is_false(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
bot.pending_file.write_text(json.dumps({"chat_id": 42, "delivered": False}), encoding="utf-8")
assert bot._is_pending_delivered() is False
# =============================================
# 2. HEARTBEAT STOPS ON DELIVERED
# =============================================
class TestHeartbeatStopsOnDelivered:
"""Heartbeat thread exits when pending file has delivered=true."""
def test_heartbeat_stops_when_delivered_mid_loop(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
bot.pending_file.write_text(json.dumps({"chat_id": 42}), encoding="utf-8")
call_count = 0
def fake_edit(chat_id, msg_id, text):
nonlocal call_count
call_count += 1
# Simulate delivery on first heartbeat edit
bot.pending_file.write_text(json.dumps({"chat_id": 42, "delivered": True}), encoding="utf-8")
with (
patch.object(bot, "edit_message", side_effect=fake_edit),
patch.object(bot, "_tmux_session_exists", return_value=True),
patch("apps.handlers.base_bot.HEARTBEAT_INTERVAL", 0.1),
):
bot._start_heartbeat(42, 999)
time.sleep(0.5)
bot._stop_heartbeat()
# Should have stopped after 1 edit (when delivered was set)
assert call_count <= 2
def test_heartbeat_continues_when_not_delivered(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
bot.pending_file.write_text(json.dumps({"chat_id": 42}), encoding="utf-8")
call_count = 0
def fake_edit(chat_id, msg_id, text):
nonlocal call_count
call_count += 1
with (
patch.object(bot, "edit_message", side_effect=fake_edit),
patch.object(bot, "_tmux_session_exists", return_value=True),
patch("apps.handlers.base_bot.HEARTBEAT_INTERVAL", 0.1),
):
bot._start_heartbeat(42, 999)
time.sleep(0.5)
bot._stop_heartbeat()
# Should have ticked multiple times since never delivered
assert call_count >= 2
def test_heartbeat_stops_when_file_absent(self, tmp_path, _patch_base_bot_deps):
"""Backward compat: heartbeat still stops when file is gone."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
assert not bot.pending_file.exists()
with (
patch.object(bot, "edit_message") as mock_edit,
patch.object(bot, "_tmux_session_exists", return_value=True),
patch("apps.handlers.base_bot.HEARTBEAT_INTERVAL", 0.1),
):
bot._start_heartbeat(42, 999)
time.sleep(0.4)
bot._stop_heartbeat()
# File never existed, so heartbeat breaks immediately — no edits
mock_edit.assert_not_called()
# =============================================
# 3. MULTI-STOP KEEPS FILE ALIVE
# =============================================
class TestMultiStopFileAlive:
"""Pending file survives delivery (advanced, not deleted)."""
def test_pending_file_survives_with_delivered_flag(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
initial_data = {"chat_id": 42, "transcript_line_after": 100}
bot.pending_file.write_text(json.dumps(initial_data), encoding="utf-8")
# Simulate _advance_pending behavior (what @hooks does)
data = json.loads(bot.pending_file.read_text(encoding="utf-8"))
data["delivered"] = True
data["transcript_line_after"] = 200
bot.pending_file.write_text(json.dumps(data), encoding="utf-8")
# File still exists
assert bot.pending_file.exists()
# But heartbeat sees it as delivered
assert bot._is_pending_delivered() is True
# Cursor advanced
reloaded = json.loads(bot.pending_file.read_text(encoding="utf-8"))
assert reloaded["transcript_line_after"] == 200
def test_new_message_overwrites_delivered(self, tmp_path, _patch_base_bot_deps):
"""A new write_pending_file clears the delivered flag."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
# Simulate previous delivery
bot.pending_file.write_text(json.dumps({"chat_id": 42, "delivered": True}), encoding="utf-8")
assert bot._is_pending_delivered() is True
# Now write a new pending (as handle_message does)
with patch.object(bot, "_get_transcript_line_count", return_value=300):
bot.write_pending_file(42, 999, 1000)
# delivered flag should be gone
assert bot._is_pending_delivered() is False
data = json.loads(bot.pending_file.read_text(encoding="utf-8"))
assert "delivered" not in data
# =============================================
# 4. REPLY NOT CLOBBERED
# =============================================
class TestReplyNotClobbered:
"""After delivery, heartbeat must NOT edit the message (no clobber)."""
def test_no_edit_after_delivered(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.pending_file.parent.mkdir(parents=True, exist_ok=True)
# Start with delivered=true (simulates hooks already delivered before heartbeat tick)
bot.pending_file.write_text(json.dumps({"chat_id": 42, "delivered": True}), encoding="utf-8")
with (
patch.object(bot, "edit_message") as mock_edit,
patch.object(bot, "_tmux_session_exists", return_value=True),
patch("apps.handlers.base_bot.HEARTBEAT_INTERVAL", 0.1),
):
bot._start_heartbeat(42, 999)
time.sleep(0.4)
bot._stop_heartbeat()
# Heartbeat saw delivered immediately, never edited
mock_edit.assert_not_called()
@@ -0,0 +1,456 @@
# =================== AIPass ====================
# Name: test_mirror_session.py
# Description: Tests for TDPLAN-0009 FINISH — mirror session launch, transcript resolver, create→attach
# Version: 1.0.0
# Created: 2026-06-29
# Modified: 2026-06-29
# =============================================
"""
Tests for the mirror session system (TDPLAN-0009 FINISH).
Covers:
- launch_mirror_session: tmux session with --dangerously-skip-permissions
- start_service: systemd user service start
- create_bot mirror params: shared_session, attach_only, chat_id in config
- _resolve_active_transcript: PID-based transcript detection
- _write_mirror_mapping: transcript-change detection and rewrite
- _config_chat_id initialization
"""
import json
from pathlib import Path
from unittest.mock import MagicMock, patch
import pytest
from apps.handlers.base_bot import BaseBot # type: ignore[import-not-found]
from apps.handlers.bot_factory import ( # type: ignore[import-not-found]
launch_mirror_session,
start_service,
)
# =============================================
# Fixtures
# =============================================
@pytest.fixture
def _patch_base_bot_deps(tmp_path):
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_bot(tmp_path, _patch_base_bot_deps, attach_only=False, shared_session=None):
workdir = tmp_path / "workdir"
workdir.mkdir(exist_ok=True)
with patch("apps.handlers.base_bot.PENDING_DIR", tmp_path):
bot = BaseBot(
bot_id="mirror_test",
bot_token="123:FAKETOKEN",
work_dir=workdir,
bot_name="Mirror Test Bot",
allowed_user_ids=[111],
branch_name="api",
shared_session=shared_session,
attach_only=attach_only,
)
bot.send_message = MagicMock(return_value={"ok": True, "message_id": 1})
return bot
# =============================================
# 1. launch_mirror_session
# =============================================
class TestLaunchMirrorSession:
"""Tests for launch_mirror_session() in bot_factory."""
def test_creates_tmux_session_with_skip_perms(self):
"""tmux new-session created and claude launched with skip-permissions."""
no_session = MagicMock(returncode=1)
ok = MagicMock(returncode=0)
side = [no_session, ok, ok, ok]
with (
patch("apps.handlers.bot_factory.subprocess.run", side_effect=side),
patch("time.sleep"),
):
result = launch_mirror_session(
session_name="telegram-api",
bot_id="api",
work_dir="/tmp/test",
)
assert result is True
def test_claude_launched_with_correct_flags(self):
"""AIPASS_SESSION_TYPE=interactive-mirror and skip-perms flag set."""
calls_made = []
first_call = True
def _track(*args, **kwargs):
"""Record subprocess.run calls for assertion."""
nonlocal first_call
calls_made.append(args[0] if args else kwargs.get("args", []))
if first_call:
first_call = False
return MagicMock(returncode=1)
result = MagicMock(returncode=0)
return result
with (
patch("apps.handlers.bot_factory.subprocess.run", side_effect=_track),
patch("time.sleep"),
):
launch_mirror_session(
session_name="telegram-api",
bot_id="api",
work_dir="/tmp/test",
)
send_keys_calls = [c for c in calls_made if "send-keys" in c]
claude_cmd = send_keys_calls[-1][-2]
assert "interactive-mirror" in claude_cmd
def test_idempotent_if_session_exists(self):
"""Returns True without creating if session already exists."""
has_session = MagicMock(returncode=0)
with patch("apps.handlers.bot_factory.subprocess.run", return_value=has_session) as mock_run:
result = launch_mirror_session(session_name="telegram-api", bot_id="api", work_dir="/tmp/test")
assert result is True
assert mock_run.call_count == 1
def test_returns_false_when_tmux_not_found(self):
"""Returns False when tmux is not installed."""
with patch("apps.handlers.bot_factory.subprocess.run", side_effect=FileNotFoundError):
result = launch_mirror_session(session_name="telegram-api", bot_id="api", work_dir="/tmp/test")
assert result is False
# =============================================
# 2. start_service
# =============================================
class TestStartService:
"""Tests for start_service() in bot_factory."""
def test_starts_systemd_service(self):
"""Calls systemctl --user start telegram-bot@{bot_id}."""
mock_result = MagicMock(returncode=0)
with patch("apps.handlers.bot_factory.subprocess.run", return_value=mock_result) as mock_run:
result = start_service("api")
assert result is True
mock_run.assert_called_once()
cmd = mock_run.call_args[0][0]
assert cmd == ["systemctl", "--user", "start", "telegram-bot@api"]
def test_returns_false_on_failure(self):
"""Returns False when systemctl returns non-zero."""
mock_result = MagicMock(returncode=1, stderr="unit not found")
with patch("apps.handlers.bot_factory.subprocess.run", return_value=mock_result):
assert start_service("api") is False
def test_returns_false_on_timeout(self):
"""Returns False on TimeoutExpired."""
import subprocess
err = subprocess.TimeoutExpired(cmd="", timeout=10)
with patch("apps.handlers.bot_factory.subprocess.run", side_effect=err):
assert start_service("api") is False
# =============================================
# 3. create_bot mirror params
# =============================================
class TestCreateBotMirror:
"""Tests for create_bot() with mirror params."""
@pytest.fixture
def _mock_create_deps(self):
"""Mock all external deps of create_bot."""
bot_info = {"username": "test_bot", "id": 123}
branch_info = {"name": "api", "path": "/home/test/api"}
patches = [
patch("apps.handlers.bot_factory.validate_token", return_value=bot_info),
patch("apps.handlers.bot_factory.validate_branch", return_value=branch_info),
patch("apps.handlers.bot_factory.get_bot", return_value=None),
patch("apps.handlers.bot_factory.get_bot_by_branch", return_value=None),
patch("apps.handlers.bot_factory.ensure_registry"),
patch("apps.handlers.bot_factory._api_set_secret"),
patch("apps.handlers.bot_factory.register_bot", return_value=True),
patch("apps.handlers.bot_factory.set_bot_commands"),
patch("apps.handlers.bot_factory.build_botfather_commands", return_value=[]),
patch("apps.handlers.bot_factory.enable_service", return_value=True),
patch("apps.handlers.bot_factory.start_bot_process", return_value=True),
patch("apps.handlers.bot_factory.launch_mirror_session", return_value=True),
patch("apps.handlers.bot_factory.start_service", return_value=True),
patch("apps.handlers.bot_factory._BOT_CONFIG_DIR", Path("/tmp/test_bots")),
]
mocks = {}
started = []
for p in patches:
m = p.start()
started.append(p)
name = p.attribute if hasattr(p, "attribute") and p.attribute else str(p).split(".")[-1].rstrip("'>)")
mocks[name] = m
yield mocks
for p in started:
p.stop()
def test_config_includes_mirror_fields(self, _mock_create_deps, tmp_path):
"""Config written with shared_session, attach_only, chat_id."""
from apps.handlers.bot_factory import create_bot # type: ignore[import-not-found]
with patch("apps.handlers.bot_factory._BOT_CONFIG_DIR", tmp_path):
result = create_bot(
bot_id="api",
bot_token="123:FAKE",
branch_name="api",
shared_session="telegram-api",
attach_only=True,
chat_id=42,
)
assert result is not None
config_file = tmp_path / "api.json"
assert config_file.exists()
config = json.loads(config_file.read_text())
assert config["shared_session"] == "telegram-api"
assert config["attach_only"] is True
assert config["chat_id"] == 42
def test_launches_mirror_session_when_attach_only(self, _mock_create_deps, tmp_path):
"""launch_mirror_session called when shared_session + attach_only."""
from apps.handlers.bot_factory import create_bot # type: ignore[import-not-found]
with patch("apps.handlers.bot_factory._BOT_CONFIG_DIR", tmp_path):
create_bot(
bot_id="api",
bot_token="123:FAKE",
branch_name="api",
shared_session="telegram-api",
attach_only=True,
)
_mock_create_deps["launch_mirror_session"].assert_called_once()
def test_starts_via_systemd_when_mirror(self, _mock_create_deps, tmp_path):
"""Mirror bot started via start_service, not start_bot_process."""
from apps.handlers.bot_factory import create_bot # type: ignore[import-not-found]
with patch("apps.handlers.bot_factory._BOT_CONFIG_DIR", tmp_path):
create_bot(
bot_id="api",
bot_token="123:FAKE",
branch_name="api",
shared_session="telegram-api",
attach_only=True,
)
_mock_create_deps["start_service"].assert_called_once()
_mock_create_deps["start_bot_process"].assert_not_called()
def test_starts_via_popen_when_not_mirror(self, _mock_create_deps, tmp_path):
"""Non-mirror bot still uses start_bot_process."""
from apps.handlers.bot_factory import create_bot # type: ignore[import-not-found]
with patch("apps.handlers.bot_factory._BOT_CONFIG_DIR", tmp_path):
create_bot(
bot_id="api",
bot_token="123:FAKE",
branch_name="api",
)
_mock_create_deps["start_bot_process"].assert_called_once()
_mock_create_deps["launch_mirror_session"].assert_not_called()
# =============================================
# 4. Transcript resolution and mapping rewrite
# =============================================
class TestMirrorMappingRewrite:
"""Tests for transcript-change detection in _write_mirror_mapping."""
def test_mapping_rewrites_on_transcript_change(self, tmp_path, _patch_base_bot_deps):
"""When transcript path changes, mapping is rewritten with new cursor."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="api")
bot._active_chat_id = 42
bot.session_name = "api"
mapping_dir = tmp_path / ".aipass" / "telegram_bots"
with patch("pathlib.Path.home", return_value=tmp_path):
with patch.object(bot, "_resolve_active_transcript", return_value=("/path/a.jsonl", 10)):
bot._write_mirror_mapping()
assert bot._mirror_mapping_written is True
assert bot._last_transcript_path == "/path/a.jsonl"
data1 = json.loads((mapping_dir / "bot-mirror_test.json").read_text())
assert data1["transcript_line_after"] == 10
with patch.object(bot, "_resolve_active_transcript", return_value=("/path/b.jsonl", 5)):
bot._write_mirror_mapping()
data2 = json.loads((mapping_dir / "bot-mirror_test.json").read_text())
assert data2["transcript_line_after"] == 5
assert bot._last_transcript_path == "/path/b.jsonl"
def test_mapping_not_rewritten_same_transcript(self, tmp_path, _patch_base_bot_deps):
"""When transcript path unchanged, mapping is not rewritten."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="api")
bot._active_chat_id = 42
bot.session_name = "api"
mapping_dir = tmp_path / ".aipass" / "telegram_bots"
with patch("pathlib.Path.home", return_value=tmp_path):
with patch.object(bot, "_resolve_active_transcript", return_value=("/path/a.jsonl", 10)):
bot._write_mirror_mapping()
mtime1 = (mapping_dir / "bot-mirror_test.json").stat().st_mtime
import time
time.sleep(0.05)
with patch.object(bot, "_resolve_active_transcript", return_value=("/path/a.jsonl", 20)):
bot._write_mirror_mapping()
mtime2 = (mapping_dir / "bot-mirror_test.json").stat().st_mtime
assert mtime1 == mtime2
def test_mapping_rewrites_when_transcript_none(self, tmp_path, _patch_base_bot_deps):
"""When transcript is None both times, still written only once."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="api")
bot._active_chat_id = 42
bot.session_name = "api"
with patch("pathlib.Path.home", return_value=tmp_path):
with patch.object(bot, "_resolve_active_transcript", return_value=(None, 0)):
bot._write_mirror_mapping()
assert bot._mirror_mapping_written is True
bot._write_mirror_mapping()
def test_last_transcript_path_updated(self, tmp_path, _patch_base_bot_deps):
"""_last_transcript_path updated after successful write."""
bot = _make_bot(tmp_path, _patch_base_bot_deps, attach_only=True, shared_session="api")
bot._active_chat_id = 42
bot.session_name = "api"
assert bot._last_transcript_path is None
with patch("pathlib.Path.home", return_value=tmp_path):
with patch.object(bot, "_resolve_active_transcript", return_value=("/path/transcript.jsonl", 15)):
bot._write_mirror_mapping()
assert bot._last_transcript_path == "/path/transcript.jsonl"
# =============================================
# 5. Transcript resolver
# =============================================
class TestResolveActiveTranscript:
"""Tests for _resolve_active_transcript."""
def test_returns_none_when_no_projects_dir(self, tmp_path, _patch_base_bot_deps):
"""Returns (None, 0) when projects dir does not exist."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with patch("pathlib.Path.home", return_value=tmp_path):
path, count = bot._resolve_active_transcript()
assert path is None
assert count == 0
def test_returns_none_when_no_jsonl_files(self, tmp_path, _patch_base_bot_deps):
"""Returns (None, 0) when projects dir exists but no JSONL files."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
slug = str(bot.work_dir).replace("/", "-")
projects_dir = tmp_path / ".claude" / "projects" / slug
projects_dir.mkdir(parents=True)
with patch("pathlib.Path.home", return_value=tmp_path):
path, count = bot._resolve_active_transcript()
assert path is None
assert count == 0
def test_falls_back_to_recent_mtime(self, tmp_path, _patch_base_bot_deps):
"""Uses most recent JSONL when PID check fails and file is < 5min old."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
slug = str(bot.work_dir).replace("/", "-")
projects_dir = tmp_path / ".claude" / "projects" / slug
projects_dir.mkdir(parents=True)
transcript = projects_dir / "abc123.jsonl"
transcript.write_text('{"type":"message"}\n{"type":"response"}\n')
with (
patch("pathlib.Path.home", return_value=tmp_path),
patch.object(bot, "_get_tmux_pane_pid", return_value=None),
):
path, count = bot._resolve_active_transcript()
assert path == str(transcript)
assert count == 2
def test_ignores_old_files(self, tmp_path, _patch_base_bot_deps):
"""JSONL files older than 5 minutes are not selected."""
import os
bot = _make_bot(tmp_path, _patch_base_bot_deps)
slug = str(bot.work_dir).replace("/", "-")
projects_dir = tmp_path / ".claude" / "projects" / slug
projects_dir.mkdir(parents=True)
transcript = projects_dir / "old.jsonl"
transcript.write_text('{"type":"message"}\n')
old_time = 1000000.0
os.utime(transcript, (old_time, old_time))
with (
patch("pathlib.Path.home", return_value=tmp_path),
patch.object(bot, "_get_tmux_pane_pid", return_value=None),
):
path, count = bot._resolve_active_transcript()
assert path is None
assert count == 0
# =============================================
# 6. _config_chat_id init
# =============================================
class TestConfigChatId:
"""Tests for _config_chat_id initialization."""
def test_config_chat_id_initialized_to_none(self, tmp_path, _patch_base_bot_deps):
"""BaseBot.__init__ initializes _config_chat_id to None."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
assert bot._config_chat_id is None
def test_config_chat_id_settable(self, tmp_path, _patch_base_bot_deps):
"""_config_chat_id can be set after construction."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot._config_chat_id = 42
assert bot._config_chat_id == 42
@@ -1,8 +1,16 @@
# =================== AIPass ====================
# Name: test_monitor.py
# Description: Tests for /monitor command — system-wide log subscription (DPLAN-0221)
# Version: 1.0.0
# Created: 2026-06-29
# Modified: 2026-06-29
# =============================================
"""
Tests for /monitor command — system-wide log subscription feature (DPLAN-0221).
Tests cover:
- Subscribe persists {chat_id, mode} via @api set_secret
- Subscribe persists {chat_id, mode} to local file
- _boot_monitor reads persisted subscription and starts the streamer
- LogStreamer level_filter: default keeps WARNING/ERROR/CRITICAL, drops INFO
- LogStreamer level_filter: 'all' keeps everything
@@ -11,6 +19,8 @@ Tests cover:
- /monitor command routing (on, all, off, status, bare)
"""
from pathlib import Path
import pytest
from unittest.mock import patch, MagicMock
@@ -25,6 +35,7 @@ from apps.handlers.log_streamer import LogStreamer # type: ignore[import-not-fo
@pytest.fixture
def _patch_base_bot_deps(tmp_path):
"""Patch heavy BaseBot dependencies to allow lightweight instantiation."""
sub_file = tmp_path / "monitor_sub.json"
patches = [
patch("apps.handlers.base_bot.PENDING_DIR", tmp_path),
patch("apps.handlers.base_bot.signal.signal"),
@@ -32,13 +43,13 @@ def _patch_base_bot_deps(tmp_path):
]
for p in patches:
p.start()
yield
yield sub_file
for p in patches:
p.stop()
def _make_bot(tmp_path, _patch_base_bot_deps):
"""Create a BaseBot with monitor subscription mocked."""
"""Create a BaseBot with monitor subscription redirected to tmp_path."""
from apps.handlers.base_bot import BaseBot # type: ignore[import-not-found]
workdir = tmp_path / "workdir"
@@ -51,6 +62,9 @@ def _make_bot(tmp_path, _patch_base_bot_deps):
allowed_user_ids=[111],
branch_name=None,
)
# Redirect subscription file to tmp_path so tests don't touch real HOME
sub_file: Path = _patch_base_bot_deps
bot._monitor_subscription_file = lambda: sub_file # type: ignore[assignment]
return bot
@@ -62,35 +76,38 @@ def _make_bot(tmp_path, _patch_base_bot_deps):
class TestSubscribePersists:
"""Verify _monitor_subscribe persists {chat_id, mode} and can be reloaded."""
def test_subscribe_calls_set_secret(self, tmp_path, _patch_base_bot_deps):
def test_subscribe_writes_file(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
sub_file: Path = _patch_base_bot_deps
with (
patch("apps.handlers.base_bot._api_set_secret") as mock_set,
patch("apps.handlers.base_bot._api_get_secret", return_value=None),
patch.object(bot, "send_message"),
patch("apps.handlers.base_bot.LogStreamer") as MockStreamer,
):
MockStreamer.return_value = MagicMock()
bot._monitor_subscribe(42, "default")
mock_set.assert_called_once_with("telegram", "monitor", {"chat_id": 42, "mode": "default"}, as_json=True)
import json
data = json.loads(sub_file.read_text())
assert data == {"chat_id": 42, "mode": "default"}
def test_subscribe_roundtrip_reload(self, tmp_path, _patch_base_bot_deps):
"""set_secret data can be read back by _load_monitor_subscription."""
"""Written file can be read back by _load_monitor_subscription."""
bot = _make_bot(tmp_path, _patch_base_bot_deps)
stored = {"chat_id": 42, "mode": "all"}
sub_file: Path = _patch_base_bot_deps
import json
with patch("apps.handlers.base_bot._api_get_secret", return_value=stored):
result = bot._load_monitor_subscription()
sub_file.write_text(json.dumps({"chat_id": 42, "mode": "all"}))
assert result == stored
result = bot._load_monitor_subscription()
assert result == {"chat_id": 42, "mode": "all"}
assert result["chat_id"] == 42
assert result["mode"] == "all"
def test_subscribe_starts_streamer(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_set_secret"),
patch.object(bot, "send_message"),
patch("apps.handlers.base_bot.LogStreamer") as MockStreamer,
):
@@ -112,7 +129,6 @@ class TestSubscribePersists:
def test_subscribe_sends_confirmation(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_set_secret"),
patch.object(bot, "send_message") as mock_send,
patch("apps.handlers.base_bot.LogStreamer", return_value=MagicMock()),
):
@@ -127,7 +143,6 @@ class TestSubscribePersists:
bot._monitor_streamer = old_streamer
with (
patch("apps.handlers.base_bot._api_set_secret"),
patch.object(bot, "send_message"),
patch("apps.handlers.base_bot.LogStreamer", return_value=MagicMock()),
):
@@ -136,10 +151,9 @@ class TestSubscribePersists:
def test_subscribe_aborts_on_save_failure(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_set_secret", side_effect=RuntimeError("boom")),
patch.object(bot, "send_message") as mock_send,
):
# Point subscription file at an unwritable path to trigger OSError
bot._monitor_subscription_file = lambda: Path("/dev/null/impossible/sub.json") # type: ignore[assignment]
with patch.object(bot, "send_message") as mock_send:
bot._monitor_subscribe(42, "default")
msg = mock_send.call_args[0][1]
assert "Failed" in msg
@@ -156,12 +170,12 @@ class TestBootMonitor:
def test_boot_starts_streamer_from_persisted(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
stored = {"chat_id": 42, "mode": "default"}
sub_file: Path = _patch_base_bot_deps
import json
with (
patch("apps.handlers.base_bot._api_get_secret", return_value=stored),
patch("apps.handlers.base_bot.LogStreamer") as MockStreamer,
):
sub_file.write_text(json.dumps({"chat_id": 42, "mode": "default"}))
with patch("apps.handlers.base_bot.LogStreamer") as MockStreamer:
mock_instance = MagicMock()
MockStreamer.return_value = mock_instance
@@ -179,31 +193,29 @@ class TestBootMonitor:
def test_boot_noop_when_no_subscription(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_get_secret", return_value=None),
patch("apps.handlers.base_bot.LogStreamer") as MockStreamer,
):
# No file written — subscription absent
with patch("apps.handlers.base_bot.LogStreamer") as MockStreamer:
bot._boot_monitor()
MockStreamer.assert_not_called()
assert bot._monitor_streamer is None
def test_boot_noop_when_empty_subscription(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_get_secret", return_value={}),
patch("apps.handlers.base_bot.LogStreamer") as MockStreamer,
):
sub_file: Path = _patch_base_bot_deps
sub_file.write_text("{}")
with patch("apps.handlers.base_bot.LogStreamer") as MockStreamer:
bot._boot_monitor()
MockStreamer.assert_not_called()
def test_boot_respects_mode_all(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
stored = {"chat_id": 99, "mode": "all"}
sub_file: Path = _patch_base_bot_deps
import json
with (
patch("apps.handlers.base_bot._api_get_secret", return_value=stored),
patch("apps.handlers.base_bot.LogStreamer") as MockStreamer,
):
sub_file.write_text(json.dumps({"chat_id": 99, "mode": "all"}))
with patch("apps.handlers.base_bot.LogStreamer") as MockStreamer:
MockStreamer.return_value = MagicMock()
bot._boot_monitor()
MockStreamer.assert_called_once_with(
@@ -344,10 +356,7 @@ class TestMonitorOff:
mock_streamer = MagicMock()
bot._monitor_streamer = mock_streamer
with (
patch("apps.handlers.base_bot._api_set_secret"),
patch.object(bot, "send_message"),
):
with patch.object(bot, "send_message"):
bot._monitor_unsubscribe(42)
mock_streamer.stop.assert_called_once()
@@ -355,20 +364,19 @@ class TestMonitorOff:
def test_off_clears_subscription(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_set_secret") as mock_set,
patch.object(bot, "send_message"),
):
sub_file: Path = _patch_base_bot_deps
import json
sub_file.write_text(json.dumps({"chat_id": 42, "mode": "default"}))
with patch.object(bot, "send_message"):
bot._monitor_unsubscribe(42)
mock_set.assert_called_once_with("telegram", "monitor", {}, as_json=True)
assert not sub_file.exists()
def test_off_sends_confirmation(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_set_secret"),
patch.object(bot, "send_message") as mock_send,
):
with patch.object(bot, "send_message") as mock_send:
bot._monitor_unsubscribe(42)
mock_send.assert_called_once()
@@ -378,10 +386,7 @@ class TestMonitorOff:
bot = _make_bot(tmp_path, _patch_base_bot_deps)
assert bot._monitor_streamer is None
with (
patch("apps.handlers.base_bot._api_set_secret"),
patch.object(bot, "send_message"),
):
with patch.object(bot, "send_message"):
bot._monitor_unsubscribe(42)
assert bot._monitor_streamer is None
@@ -445,24 +450,23 @@ class TestMonitorStatus:
def test_status_when_not_subscribed(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
with (
patch("apps.handlers.base_bot._api_get_secret", return_value=None),
patch.object(bot, "send_message") as mock_send,
):
# No file written — no subscription
with patch.object(bot, "send_message") as mock_send:
bot._monitor_status(42)
msg = mock_send.call_args[0][1]
assert "not subscribed" in msg
def test_status_when_subscribed_and_running(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
sub_file: Path = _patch_base_bot_deps
import json
sub_file.write_text(json.dumps({"chat_id": 42, "mode": "default"}))
mock_streamer = MagicMock()
mock_streamer._running = True
bot._monitor_streamer = mock_streamer
with (
patch("apps.handlers.base_bot._api_get_secret", return_value={"chat_id": 42, "mode": "default"}),
patch.object(bot, "send_message") as mock_send,
):
with patch.object(bot, "send_message") as mock_send:
bot._monitor_status(42)
msg = mock_send.call_args[0][1]
assert "streaming" in msg
@@ -470,12 +474,13 @@ class TestMonitorStatus:
def test_status_shows_mode_label(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
sub_file: Path = _patch_base_bot_deps
import json
sub_file.write_text(json.dumps({"chat_id": 42, "mode": "all"}))
bot._monitor_streamer = MagicMock(_running=True)
with (
patch("apps.handlers.base_bot._api_get_secret", return_value={"chat_id": 42, "mode": "all"}),
patch.object(bot, "send_message") as mock_send,
):
with patch.object(bot, "send_message") as mock_send:
bot._monitor_status(42)
msg = mock_send.call_args[0][1]
assert "firehose" in msg
@@ -1,3 +1,11 @@
# =================== AIPass ====================
# Name: test_multi_bot.py
# Description: Comprehensive tests for BaseBot and BranchPlugin
# Version: 1.0.0
# Created: 2026-06-15
# Modified: 2026-06-29
# =============================================
"""
Comprehensive pytest tests for BaseBot and BranchPlugin.
@@ -25,8 +33,8 @@ import time
import pytest
from unittest.mock import patch, MagicMock
from apps.handlers.base_bot import BaseBot
from apps.handlers.branch_plugin import BranchPlugin
from apps.handlers.base_bot import BaseBot # type: ignore[import-not-found]
from apps.handlers.branch_plugin import BranchPlugin # type: ignore[import-not-found]
# =============================================
@@ -1553,6 +1561,20 @@ class TestCreateAutomated:
state = self.bot._create_state[self.chat_id]
assert state["branch_name"] == "flow"
@patch("apps.handlers.base_bot.get_bot_by_branch", return_value=None)
@patch("apps.handlers.base_bot.validate_branch")
@patch("apps.handlers.base_bot.check_telethon_setup")
def test_manual_fallback_message_shows_reason(self, mock_check, mock_validate, mock_get_bot):
"""Manual fallback message includes the reason automation is unavailable."""
mock_check.return_value = (False, "Telethon library not installed. Run: pip install telethon")
mock_validate.return_value = {"name": "flow", "path": "/home/aipass/flow"}
self.bot._handle_create_command(self.chat_id, "chat flow")
msg = self.bot.send_message.call_args[0][1]
assert "Telethon library not installed" in msg
assert "Falling back to manual token flow" in msg
@patch("apps.handlers.base_bot.create_bot")
@patch("apps.handlers.base_bot.create_bot_via_botfather")
def test_automated_success_message_contains_service_info(self, mock_bf_create, mock_create_bot):
@@ -0,0 +1,234 @@
# =================== AIPass ====================
# Name: test_status_reset.py
# Description: Tests for /status conversation vs daemon uptime + /new counter reset
# Version: 1.0.0
# Created: 2026-06-29
# Modified: 2026-06-29
# =============================================
"""
Tests for /status conversation vs daemon uptime + /new counter reset.
Tests cover:
- /new resets message_count to 0
- /new resets conversation_start
- /status shows conversation uptime (not daemon uptime) as primary
- /status shows daemon uptime separately
- build_status_text includes daemon_uptime line when provided
- build_status_text omits daemon_uptime line when not provided
"""
import time
import pytest
from unittest.mock import patch
from apps.handlers.base_bot import BaseBot # type: ignore[import-not-found]
from apps.handlers.telegram_standards import build_status_text # type: ignore[import-not-found]
@pytest.fixture
def _patch_base_bot_deps(tmp_path):
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_bot(tmp_path, _patch_base_bot_deps):
workdir = tmp_path / "workdir"
workdir.mkdir()
return BaseBot(
bot_id="status_test",
bot_token="123:FAKETOKEN",
work_dir=workdir,
bot_name="Status Test Bot",
allowed_user_ids=[111],
branch_name=None,
)
# =============================================
# 1. /new RESETS COUNTERS
# =============================================
class TestNewResetsCounters:
"""/new resets message_count and conversation_start."""
def test_new_resets_message_count(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.state["message_count"] = 15
with (
patch.object(bot, "send_message"),
patch.object(bot, "_kill_tmux_session"),
):
bot._dispatch_command(42, ("new", ""))
assert bot.state["message_count"] == 0
def test_new_resets_conversation_start(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.state["conversation_start"] = time.time() - 3600 # 1 hour ago
with (
patch.object(bot, "send_message"),
patch.object(bot, "_kill_tmux_session"),
):
before = time.time()
bot._dispatch_command(42, ("new", ""))
after = time.time()
assert before <= bot.state["conversation_start"] <= after
def test_new_does_not_reset_daemon_start(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
original_start = bot.state["start_time"]
bot.state["message_count"] = 5
with (
patch.object(bot, "send_message"),
patch.object(bot, "_kill_tmux_session"),
):
bot._dispatch_command(42, ("new", ""))
assert bot.state["start_time"] == original_start
# =============================================
# 2. /status SHOWS CORRECT UPTIMES
# =============================================
class TestStatusUptimes:
"""/status shows conversation uptime as primary and daemon uptime separately."""
def test_status_passes_daemon_uptime(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.state["start_time"] = time.time() - 7200 # daemon up 2h
bot.state["conversation_start"] = time.time() - 300 # conv 5m
with (
patch.object(bot, "send_message"),
patch("apps.handlers.base_bot.build_status_text", wraps=build_status_text) as mock_build,
patch("apps.handlers.telegram_standards._tmux_session_exists", return_value=True),
):
bot._dispatch_command(42, ("status", ""))
call_kwargs = mock_build.call_args
args = call_kwargs[1] if call_kwargs[1] else {}
if not args:
_, kwargs = mock_build.call_args
args = kwargs
assert "daemon_uptime" in args
assert "2h" in args["daemon_uptime"]
def test_status_conversation_uptime_after_new(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.state["start_time"] = time.time() - 7200 # daemon up 2h
with (
patch.object(bot, "send_message"),
patch.object(bot, "_kill_tmux_session"),
):
bot._dispatch_command(42, ("new", ""))
# Now check /status — conversation uptime should be near 0
with (
patch.object(bot, "send_message") as mock_send,
patch("apps.handlers.telegram_standards._tmux_session_exists", return_value=True),
):
bot._dispatch_command(42, ("status", ""))
msg = mock_send.call_args[0][1]
assert "Uptime: 0h 0m" in msg
assert "Daemon up: 2h" in msg
def test_status_messages_zero_after_new(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
bot.state["message_count"] = 42
with (
patch.object(bot, "send_message"),
patch.object(bot, "_kill_tmux_session"),
):
bot._dispatch_command(42, ("new", ""))
with (
patch.object(bot, "send_message") as mock_send,
patch("apps.handlers.telegram_standards._tmux_session_exists", return_value=True),
):
bot._dispatch_command(42, ("status", ""))
msg = mock_send.call_args[0][1]
assert "Messages: 0" in msg
# =============================================
# 3. build_status_text
# =============================================
class TestBuildStatusText:
"""build_status_text renders daemon_uptime when provided."""
def test_includes_daemon_uptime(self):
with patch("apps.handlers.telegram_standards._tmux_session_exists", return_value=True):
text = build_status_text(
session_name="telegram-base",
branch_name="base",
uptime="0h 5m 0s",
message_count=3,
daemon_uptime="12h 0m 0s",
)
assert "Daemon up: 12h 0m 0s" in text
assert "Uptime: 0h 5m 0s" in text
def test_omits_daemon_uptime_when_none(self):
with patch("apps.handlers.telegram_standards._tmux_session_exists", return_value=True):
text = build_status_text(
session_name="telegram-base",
branch_name="base",
uptime="1h 0m 0s",
message_count=5,
)
assert "Daemon up" not in text
assert "Uptime: 1h 0m 0s" in text
def test_uptime_before_daemon_uptime(self):
with patch("apps.handlers.telegram_standards._tmux_session_exists", return_value=True):
text = build_status_text(
session_name="telegram-base",
branch_name="base",
uptime="0h 1m 0s",
daemon_uptime="5h 0m 0s",
)
uptime_pos = text.index("Uptime:")
daemon_pos = text.index("Daemon up:")
assert uptime_pos < daemon_pos
# =============================================
# 4. CONVERSATION_START IN STATE
# =============================================
class TestConversationStartState:
"""conversation_start is initialized and tracked."""
def test_init_sets_conversation_start(self, tmp_path, _patch_base_bot_deps):
before = time.time()
bot = _make_bot(tmp_path, _patch_base_bot_deps)
after = time.time()
assert before <= bot.state["conversation_start"] <= after
def test_conversation_start_equals_start_time_at_boot(self, tmp_path, _patch_base_bot_deps):
bot = _make_bot(tmp_path, _patch_base_bot_deps)
assert abs(bot.state["conversation_start"] - bot.state["start_time"]) < 0.1