feat: TG user-comment mirror — user messages from ALL doors (terminal/remote) now mirror to the branch TG chat with origin tag. UserPromptSubmit relay handler (self-contained in telegram skill, crash-isolated last entry in hooks.json), human-only noise fences (system/task notifications, slash-command output, dispatch wakes, subagents, TG-echo, dupes), inbound hardening (stale-pending clean before write, undelivered-overwrite warning). 47 new TG tests, execution proven via engine.jsonl, positive path live-verified to real TG chat.
This commit is contained in:
@@ -38,6 +38,11 @@
|
||||
"handler": "aipass.hooks.apps.handlers.lifecycle.auto_process.handle",
|
||||
"matcher": "",
|
||||
"timeout": 120
|
||||
},
|
||||
"user_message_relay": {
|
||||
"enabled": true,
|
||||
"handler": "aipass.skills.lib.telegram.apps.handlers.user_message_relay.handle",
|
||||
"matcher": ""
|
||||
}
|
||||
},
|
||||
|
||||
|
||||
@@ -11,6 +11,26 @@ PyPI version — not the changelog header.
|
||||
|
||||
## [2026-07-14]
|
||||
|
||||
### Added
|
||||
|
||||
- **Telegram user-comment mirror: the TG chat now shows the whole conversation,
|
||||
whichever door you speak through.** Patrick's spec from the live cross-door
|
||||
drill: his own messages typed in the terminal or claude.ai remote never
|
||||
appeared in TG — only the replies did. New `user_message_relay` UserPromptSubmit
|
||||
handler (@skills-built, self-contained in the telegram skill, registered by
|
||||
@hooks as the last, crash-isolated entry) posts genuine user messages to the
|
||||
branch's TG chat with an origin tag, silently (`disable_notification`). Noise
|
||||
fences keep it human-only: system/task notifications, slash-command output,
|
||||
dispatch wake prompts, sub-agent prompts, TG-origin echoes, and consecutive
|
||||
dupes are all skipped (structural session-type detection was investigated and
|
||||
rejected — it's session-wide, would eat genuine mid-flight messages). Inbound
|
||||
hardening rides along: stale pending files cleaned before each write, and an
|
||||
undelivered-response overwrite now logs a warning instead of silently losing
|
||||
the reply. 47 new TG tests; registration execution-proven via engine.jsonl and
|
||||
the positive path live-verified — a terminal-door message delivered to the
|
||||
real TG chat. TG dormancy/proactive push deliberately untouched (design chat
|
||||
with Patrick pending).
|
||||
|
||||
### Fixed
|
||||
|
||||
- **DPLAN-0241 round 4 (night shift): user flags survive every launch path, and
|
||||
|
||||
@@ -696,6 +696,21 @@ class BaseBot:
|
||||
)
|
||||
return
|
||||
|
||||
# Inbound reliability: clean stale pending + warn on in-flight overwrite
|
||||
self.clean_stale_pending()
|
||||
if self.pending_file.exists():
|
||||
try:
|
||||
prev = json.loads(self.pending_file.read_text(encoding="utf-8"))
|
||||
if not prev.get("delivered"):
|
||||
prev_id = prev.get("message_id", "?")
|
||||
logger.warning(
|
||||
"Overwriting undelivered pending (msg_id=%s) with new message %d",
|
||||
prev_id,
|
||||
message_id,
|
||||
)
|
||||
except (json.JSONDecodeError, OSError):
|
||||
pass
|
||||
|
||||
# Send processing indicator
|
||||
processing_result = self.send_message(chat_id, PROCESSING_MSG)
|
||||
processing_msg_id = processing_result.get("message_id") if processing_result else None
|
||||
|
||||
@@ -0,0 +1,163 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: user_message_relay.py
|
||||
# Description: Relay user messages from non-TG doors to the branch TG chat
|
||||
# Version: 1.1.0
|
||||
# Created: 2026-07-14
|
||||
# Modified: 2026-07-14
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
User message relay — posts user messages from non-TG doors to the branch TG chat.
|
||||
|
||||
UserPromptSubmit hook handler. When a user types in terminal or remote, their
|
||||
message is posted to the branch's Telegram chat so the chat reads like the full
|
||||
conversation. TG-origin messages are skipped (already visible in chat).
|
||||
|
||||
Registration: @hooks adds this to .aipass/hooks.json + ~/.claude/settings.json.
|
||||
"""
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
from urllib.error import URLError
|
||||
from urllib.request import Request, urlopen
|
||||
|
||||
from aipass.prax import logger
|
||||
|
||||
|
||||
MIRROR_DIR = Path.home() / ".aipass" / "telegram_bots"
|
||||
PENDING_DIR = Path.home() / ".aipass" / "telegram_pending"
|
||||
TG_ORIGIN_MARKER = "via Telegram:"
|
||||
TELEGRAM_MAX_LENGTH = 4096
|
||||
|
||||
# Dispatch-wake detection: AIPASS_SESSION_TYPE env var ("dispatched"/"daemon") is
|
||||
# session-wide, not per-prompt — a dispatched session can still receive genuine
|
||||
# user input mid-flight. No per-prompt structural indicator exists in hook_data
|
||||
# (confirmed by inspecting all handlers + hook_test.py mock payloads). Fallback:
|
||||
# match the known automated wake-prompt prefix.
|
||||
_DISPATCH_WAKE_PREFIX = "Hi. Check inbox, process new emails"
|
||||
|
||||
_last_relay_hash: str = ""
|
||||
|
||||
|
||||
def _try_load_bot(path: Path) -> dict | None:
|
||||
if not path.exists():
|
||||
return None
|
||||
try:
|
||||
data = json.loads(path.read_text(encoding="utf-8"))
|
||||
if isinstance(data, dict) and data.get("chat_id") and data.get("bot_token"):
|
||||
return data
|
||||
return None
|
||||
except (json.JSONDecodeError, OSError):
|
||||
return None
|
||||
|
||||
|
||||
def find_bot_for_cwd(cwd: str) -> dict | None:
|
||||
"""Find a mirror/pending bot file whose work_dir contains the given CWD."""
|
||||
cwd_path = Path(cwd)
|
||||
|
||||
env_bot_id = os.environ.get("AIPASS_BOT_ID")
|
||||
if env_bot_id:
|
||||
for search_dir in [MIRROR_DIR, PENDING_DIR]:
|
||||
data = _try_load_bot(search_dir / f"bot-{env_bot_id}.json")
|
||||
if data:
|
||||
return data
|
||||
|
||||
for search_dir in [MIRROR_DIR, PENDING_DIR]:
|
||||
if not search_dir.exists():
|
||||
continue
|
||||
for bot_file in sorted(search_dir.glob("bot-*.json")):
|
||||
data = _try_load_bot(bot_file)
|
||||
if not data or not data.get("work_dir"):
|
||||
continue
|
||||
try:
|
||||
cwd_path.relative_to(Path(data["work_dir"]))
|
||||
return data
|
||||
except ValueError:
|
||||
continue
|
||||
return None
|
||||
|
||||
|
||||
def send_user_message(bot_token: str, chat_id: int, text: str, origin: str = "\U0001f5a5️") -> bool:
|
||||
"""Post a user message to the TG chat with origin tag."""
|
||||
formatted = f"{origin}\n{text}"
|
||||
if len(formatted) > TELEGRAM_MAX_LENGTH:
|
||||
formatted = formatted[:TELEGRAM_MAX_LENGTH]
|
||||
|
||||
url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
|
||||
payload = json.dumps(
|
||||
{
|
||||
"chat_id": chat_id,
|
||||
"text": formatted,
|
||||
"disable_notification": True,
|
||||
}
|
||||
).encode("utf-8")
|
||||
req = Request(url, data=payload, headers={"Content-Type": "application/json"})
|
||||
|
||||
try:
|
||||
with urlopen(req, timeout=10) as resp:
|
||||
result = json.loads(resp.read())
|
||||
return result.get("ok", False)
|
||||
except (URLError, Exception) as e:
|
||||
logger.warning("[TG] user message relay send failed: %s", e)
|
||||
return False
|
||||
|
||||
|
||||
def _is_system_noise(prompt: str) -> bool:
|
||||
"""Detect non-human system noise that should not be mirrored to TG."""
|
||||
if prompt.startswith("[SYSTEM NOTIFICATION"):
|
||||
return True
|
||||
if "<task-notification>" in prompt:
|
||||
return True
|
||||
if "<command-name>" in prompt or "<local-command-stdout>" in prompt:
|
||||
return True
|
||||
if "messages below were generated by the user while running local commands" in prompt:
|
||||
return True
|
||||
if prompt.startswith(_DISPATCH_WAKE_PREFIX):
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def handle(hook_data: dict) -> dict:
|
||||
"""UserPromptSubmit hook handler — relay user message to branch TG chat.
|
||||
|
||||
Skips: subagent prompts, system noise (notifications, local-command output,
|
||||
dispatch wakes), TG-origin messages, duplicate consecutive messages, and
|
||||
branches with no TG bot configured.
|
||||
"""
|
||||
global _last_relay_hash # noqa: PLW0603
|
||||
try:
|
||||
if hook_data.get("agent_type", ""):
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
|
||||
prompt = hook_data.get("prompt", "")
|
||||
if not prompt or not prompt.strip():
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
|
||||
if _is_system_noise(prompt):
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
|
||||
if TG_ORIGIN_MARKER in prompt:
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
|
||||
msg_hash = hashlib.md5(prompt.encode()).hexdigest()
|
||||
if msg_hash == _last_relay_hash:
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
|
||||
cwd = hook_data.get("cwd", "") or str(Path.cwd())
|
||||
bot_data = find_bot_for_cwd(cwd)
|
||||
if not bot_data:
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
|
||||
bot_token = bot_data["bot_token"]
|
||||
chat_id = int(bot_data["chat_id"])
|
||||
|
||||
if send_user_message(bot_token, chat_id, prompt):
|
||||
_last_relay_hash = msg_hash
|
||||
logger.info("[TG] user message relayed to chat_id=%s", chat_id)
|
||||
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
except Exception as e:
|
||||
logger.warning("[TG] user message relay error: %s", e)
|
||||
return {"stdout": "", "exit_code": 0}
|
||||
@@ -0,0 +1,185 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: test_inbound_reliability.py
|
||||
# Description: Tests for inbound message reliability hardening in BaseBot
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-07-14
|
||||
# Modified: 2026-07-14
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
Tests for inbound reliability hardening in handle_message.
|
||||
|
||||
Tests cover:
|
||||
- Stale pending cleaned before writing new pending
|
||||
- Warning logged when overwriting undelivered pending
|
||||
- No warning when previous pending was delivered
|
||||
- Corrupt pending file doesn't crash
|
||||
"""
|
||||
|
||||
import json
|
||||
import time
|
||||
|
||||
import pytest
|
||||
from unittest.mock import patch
|
||||
|
||||
from aipass.skills.lib.telegram.apps.handlers.base_bot import BaseBot
|
||||
|
||||
|
||||
# =============================================
|
||||
# HELPERS
|
||||
# =============================================
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def _patch_base_bot_deps(tmp_path):
|
||||
"""Patch heavy BaseBot dependencies for lightweight instantiation."""
|
||||
patches = [
|
||||
patch(
|
||||
"aipass.skills.lib.telegram.apps.handlers.base_bot.PENDING_DIR",
|
||||
tmp_path,
|
||||
),
|
||||
patch("aipass.skills.lib.telegram.apps.handlers.base_bot.signal.signal"),
|
||||
patch("aipass.skills.lib.telegram.apps.handlers.base_bot.atexit.register"),
|
||||
]
|
||||
for p in patches:
|
||||
p.start()
|
||||
yield
|
||||
for p in patches:
|
||||
p.stop()
|
||||
|
||||
|
||||
def _make_bot(tmp_path, _patch_base_bot_deps):
|
||||
workdir = tmp_path / "workdir"
|
||||
workdir.mkdir(exist_ok=True)
|
||||
return BaseBot(
|
||||
bot_id="test_bot",
|
||||
bot_token="123:FAKETOKEN",
|
||||
work_dir=workdir,
|
||||
bot_name="Test Bot",
|
||||
allowed_user_ids=[111],
|
||||
branch_name="testbranch",
|
||||
)
|
||||
|
||||
|
||||
# =============================================
|
||||
# TESTS
|
||||
# =============================================
|
||||
|
||||
|
||||
class TestInboundReliability:
|
||||
def test_clean_stale_called_before_write(self, tmp_path, _patch_base_bot_deps):
|
||||
"""handle_message calls clean_stale_pending before write_pending_file."""
|
||||
bot = _make_bot(tmp_path, _patch_base_bot_deps)
|
||||
call_order = []
|
||||
|
||||
with (
|
||||
patch.object(
|
||||
bot,
|
||||
"ensure_tmux_session",
|
||||
return_value=True,
|
||||
),
|
||||
patch.object(bot, "send_message", return_value={"message_id": 1}),
|
||||
patch.object(
|
||||
bot,
|
||||
"clean_stale_pending",
|
||||
side_effect=lambda: call_order.append("clean"),
|
||||
),
|
||||
patch.object(
|
||||
bot,
|
||||
"write_pending_file",
|
||||
side_effect=lambda *a: (call_order.append("write"), True)[1],
|
||||
),
|
||||
patch.object(bot, "inject_message", return_value=True),
|
||||
patch.object(bot, "_start_heartbeat"),
|
||||
):
|
||||
bot.handle_message(42, "hello", {"message_id": 100})
|
||||
|
||||
assert call_order == ["clean", "write"]
|
||||
|
||||
def test_warns_on_undelivered_overwrite(self, tmp_path, _patch_base_bot_deps, caplog):
|
||||
"""Warning logged when overwriting a pending file that wasn't delivered."""
|
||||
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,
|
||||
"message_id": 50,
|
||||
"delivered": False,
|
||||
"timestamp": time.time(),
|
||||
}
|
||||
)
|
||||
)
|
||||
|
||||
with (
|
||||
patch.object(bot, "ensure_tmux_session", return_value=True),
|
||||
patch.object(bot, "send_message", return_value={"message_id": 1}),
|
||||
patch.object(bot, "clean_stale_pending"),
|
||||
patch.object(bot, "write_pending_file", return_value=True),
|
||||
patch.object(bot, "inject_message", return_value=True),
|
||||
patch.object(bot, "_start_heartbeat"),
|
||||
):
|
||||
bot.handle_message(42, "new msg", {"message_id": 200})
|
||||
|
||||
assert any("Overwriting undelivered pending" in r.message for r in caplog.records)
|
||||
assert any("msg_id=50" in r.message for r in caplog.records)
|
||||
|
||||
def test_no_warn_when_delivered(self, tmp_path, _patch_base_bot_deps, caplog):
|
||||
"""No warning when previous pending was already delivered."""
|
||||
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,
|
||||
"message_id": 50,
|
||||
"delivered": True,
|
||||
"timestamp": time.time(),
|
||||
}
|
||||
)
|
||||
)
|
||||
|
||||
with (
|
||||
patch.object(bot, "ensure_tmux_session", return_value=True),
|
||||
patch.object(bot, "send_message", return_value={"message_id": 1}),
|
||||
patch.object(bot, "clean_stale_pending"),
|
||||
patch.object(bot, "write_pending_file", return_value=True),
|
||||
patch.object(bot, "inject_message", return_value=True),
|
||||
patch.object(bot, "_start_heartbeat"),
|
||||
):
|
||||
bot.handle_message(42, "new msg", {"message_id": 200})
|
||||
|
||||
assert not any("Overwriting undelivered pending" in r.message for r in caplog.records)
|
||||
|
||||
def test_corrupt_pending_no_crash(self, tmp_path, _patch_base_bot_deps):
|
||||
"""Corrupt pending file doesn't crash handle_message."""
|
||||
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{{{")
|
||||
|
||||
with (
|
||||
patch.object(bot, "ensure_tmux_session", return_value=True),
|
||||
patch.object(bot, "send_message", return_value={"message_id": 1}),
|
||||
patch.object(bot, "clean_stale_pending"),
|
||||
patch.object(bot, "write_pending_file", return_value=True),
|
||||
patch.object(bot, "inject_message", return_value=True),
|
||||
patch.object(bot, "_start_heartbeat"),
|
||||
):
|
||||
bot.handle_message(42, "hello", {"message_id": 100})
|
||||
|
||||
def test_no_pending_file_no_crash(self, tmp_path, _patch_base_bot_deps):
|
||||
"""Missing pending file doesn't crash handle_message."""
|
||||
bot = _make_bot(tmp_path, _patch_base_bot_deps)
|
||||
|
||||
with (
|
||||
patch.object(bot, "ensure_tmux_session", return_value=True),
|
||||
patch.object(bot, "send_message", return_value={"message_id": 1}),
|
||||
patch.object(bot, "clean_stale_pending"),
|
||||
patch.object(bot, "write_pending_file", return_value=True),
|
||||
patch.object(bot, "inject_message", return_value=True),
|
||||
patch.object(bot, "_start_heartbeat"),
|
||||
):
|
||||
bot.handle_message(42, "hello", {"message_id": 100})
|
||||
@@ -0,0 +1,355 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: test_user_message_relay.py
|
||||
# Description: Tests for user message relay — UserPromptSubmit hook handler
|
||||
# Version: 1.1.0
|
||||
# Created: 2026-07-14
|
||||
# Modified: 2026-07-14
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
Tests for user_message_relay — the UserPromptSubmit hook handler that posts
|
||||
user messages from non-TG doors to the branch TG chat.
|
||||
|
||||
Tests cover:
|
||||
- find_bot_for_cwd: env var priority, CWD matching, missing dirs, no match
|
||||
- send_user_message: formatting, truncation, network errors
|
||||
- _is_system_noise: system notifications, task notifications, local-command
|
||||
output, Caveat line, dispatch wake prompts
|
||||
- handle(): skip subagent, skip empty, skip system noise, skip TG-origin,
|
||||
skip dupe, happy path
|
||||
- Dedup: consecutive identical messages skipped, different messages pass
|
||||
"""
|
||||
|
||||
import json
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
import aipass.skills.lib.telegram.apps.handlers.user_message_relay as relay_mod
|
||||
from aipass.skills.lib.telegram.apps.handlers.user_message_relay import (
|
||||
_is_system_noise,
|
||||
find_bot_for_cwd,
|
||||
handle,
|
||||
send_user_message,
|
||||
)
|
||||
|
||||
|
||||
# =============================================
|
||||
# FIXTURES
|
||||
# =============================================
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _reset_dedup():
|
||||
"""Reset the dedup hash between tests."""
|
||||
relay_mod._last_relay_hash = ""
|
||||
yield
|
||||
relay_mod._last_relay_hash = ""
|
||||
|
||||
|
||||
@pytest.fixture()
|
||||
def bot_dirs(tmp_path):
|
||||
"""Create mirror + pending dirs with a test bot file."""
|
||||
mirror = tmp_path / "mirror"
|
||||
pending = tmp_path / "pending"
|
||||
mirror.mkdir()
|
||||
pending.mkdir()
|
||||
|
||||
work = tmp_path / "branch_workdir"
|
||||
work.mkdir()
|
||||
|
||||
bot_data = {
|
||||
"chat_id": 42,
|
||||
"bot_token": "123:FAKETOKEN",
|
||||
"work_dir": str(work),
|
||||
"bot_id": "test_bot",
|
||||
}
|
||||
(mirror / "bot-test_bot.json").write_text(json.dumps(bot_data))
|
||||
|
||||
with (
|
||||
patch.object(relay_mod, "MIRROR_DIR", mirror),
|
||||
patch.object(relay_mod, "PENDING_DIR", pending),
|
||||
):
|
||||
yield {"mirror": mirror, "pending": pending, "work": work, "bot_data": bot_data}
|
||||
|
||||
|
||||
# =============================================
|
||||
# 1. find_bot_for_cwd
|
||||
# =============================================
|
||||
|
||||
|
||||
class TestFindBotForCwd:
|
||||
def test_finds_by_cwd_match(self, bot_dirs):
|
||||
result = find_bot_for_cwd(str(bot_dirs["work"]))
|
||||
assert result is not None
|
||||
assert result["chat_id"] == 42
|
||||
assert result["bot_token"] == "123:FAKETOKEN"
|
||||
|
||||
def test_finds_by_cwd_subdirectory(self, bot_dirs):
|
||||
sub = bot_dirs["work"] / "some" / "subdir"
|
||||
sub.mkdir(parents=True)
|
||||
result = find_bot_for_cwd(str(sub))
|
||||
assert result is not None
|
||||
assert result["chat_id"] == 42
|
||||
|
||||
def test_returns_none_no_match(self, bot_dirs):
|
||||
result = find_bot_for_cwd("/tmp/nowhere")
|
||||
assert result is None
|
||||
|
||||
def test_env_var_priority(self, bot_dirs):
|
||||
with patch.dict("os.environ", {"AIPASS_BOT_ID": "test_bot"}):
|
||||
result = find_bot_for_cwd("/tmp/anywhere")
|
||||
assert result is not None
|
||||
assert result["bot_id"] == "test_bot"
|
||||
|
||||
def test_env_var_missing_bot(self, tmp_path):
|
||||
mirror = tmp_path / "mirror"
|
||||
mirror.mkdir()
|
||||
with (
|
||||
patch.object(relay_mod, "MIRROR_DIR", mirror),
|
||||
patch.object(relay_mod, "PENDING_DIR", tmp_path / "pending"),
|
||||
patch.dict("os.environ", {"AIPASS_BOT_ID": "nonexistent"}),
|
||||
):
|
||||
result = find_bot_for_cwd("/tmp")
|
||||
assert result is None
|
||||
|
||||
def test_skips_bot_without_chat_id(self, tmp_path):
|
||||
mirror = tmp_path / "mirror"
|
||||
mirror.mkdir()
|
||||
(mirror / "bot-bad.json").write_text(json.dumps({"bot_token": "tok", "work_dir": "/tmp"}))
|
||||
with (
|
||||
patch.object(relay_mod, "MIRROR_DIR", mirror),
|
||||
patch.object(relay_mod, "PENDING_DIR", tmp_path / "pending"),
|
||||
):
|
||||
result = find_bot_for_cwd("/tmp")
|
||||
assert result is None
|
||||
|
||||
def test_skips_corrupt_json(self, tmp_path):
|
||||
mirror = tmp_path / "mirror"
|
||||
mirror.mkdir()
|
||||
(mirror / "bot-bad.json").write_text("not json{{{")
|
||||
with (
|
||||
patch.object(relay_mod, "MIRROR_DIR", mirror),
|
||||
patch.object(relay_mod, "PENDING_DIR", tmp_path / "pending"),
|
||||
):
|
||||
result = find_bot_for_cwd("/tmp")
|
||||
assert result is None
|
||||
|
||||
def test_missing_dirs_no_crash(self, tmp_path):
|
||||
with (
|
||||
patch.object(relay_mod, "MIRROR_DIR", tmp_path / "nope1"),
|
||||
patch.object(relay_mod, "PENDING_DIR", tmp_path / "nope2"),
|
||||
):
|
||||
result = find_bot_for_cwd("/tmp")
|
||||
assert result is None
|
||||
|
||||
def test_pending_dir_fallback(self, tmp_path):
|
||||
mirror = tmp_path / "mirror"
|
||||
pending = tmp_path / "pending"
|
||||
mirror.mkdir()
|
||||
pending.mkdir()
|
||||
work = tmp_path / "work"
|
||||
work.mkdir()
|
||||
bot_data = {"chat_id": 99, "bot_token": "tok", "work_dir": str(work)}
|
||||
(pending / "bot-pend.json").write_text(json.dumps(bot_data))
|
||||
with (
|
||||
patch.object(relay_mod, "MIRROR_DIR", mirror),
|
||||
patch.object(relay_mod, "PENDING_DIR", pending),
|
||||
):
|
||||
result = find_bot_for_cwd(str(work))
|
||||
assert result is not None
|
||||
assert result["chat_id"] == 99
|
||||
|
||||
|
||||
# =============================================
|
||||
# 2. send_user_message
|
||||
# =============================================
|
||||
|
||||
|
||||
class TestSendUserMessage:
|
||||
def test_sends_formatted_message(self):
|
||||
mock_resp = MagicMock()
|
||||
mock_resp.read.return_value = json.dumps({"ok": True}).encode()
|
||||
mock_resp.__enter__ = lambda s: s
|
||||
mock_resp.__exit__ = MagicMock(return_value=False)
|
||||
|
||||
_urlopen = "aipass.skills.lib.telegram.apps.handlers.user_message_relay.urlopen"
|
||||
with patch(_urlopen, return_value=mock_resp) as mock_url:
|
||||
result = send_user_message("tok", 42, "hello world")
|
||||
assert result is True
|
||||
call_args = mock_url.call_args
|
||||
req = call_args[0][0]
|
||||
body = json.loads(req.data)
|
||||
assert body["chat_id"] == 42
|
||||
assert "hello world" in body["text"]
|
||||
assert body["disable_notification"] is True
|
||||
|
||||
def test_origin_tag_in_message(self):
|
||||
mock_resp = MagicMock()
|
||||
mock_resp.read.return_value = json.dumps({"ok": True}).encode()
|
||||
mock_resp.__enter__ = lambda s: s
|
||||
mock_resp.__exit__ = MagicMock(return_value=False)
|
||||
|
||||
_urlopen = "aipass.skills.lib.telegram.apps.handlers.user_message_relay.urlopen"
|
||||
with patch(_urlopen, return_value=mock_resp) as mock_url:
|
||||
send_user_message("tok", 42, "test", origin="TERM")
|
||||
body = json.loads(mock_url.call_args[0][0].data)
|
||||
assert body["text"].startswith("TERM\n")
|
||||
|
||||
def test_truncation_at_4096(self):
|
||||
mock_resp = MagicMock()
|
||||
mock_resp.read.return_value = json.dumps({"ok": True}).encode()
|
||||
mock_resp.__enter__ = lambda s: s
|
||||
mock_resp.__exit__ = MagicMock(return_value=False)
|
||||
|
||||
_urlopen = "aipass.skills.lib.telegram.apps.handlers.user_message_relay.urlopen"
|
||||
with patch(_urlopen, return_value=mock_resp) as mock_url:
|
||||
send_user_message("tok", 42, "x" * 5000)
|
||||
body = json.loads(mock_url.call_args[0][0].data)
|
||||
assert len(body["text"]) <= 4096
|
||||
|
||||
def test_network_error_returns_false(self):
|
||||
with patch(
|
||||
"aipass.skills.lib.telegram.apps.handlers.user_message_relay.urlopen",
|
||||
side_effect=Exception("network down"),
|
||||
):
|
||||
result = send_user_message("tok", 42, "test")
|
||||
assert result is False
|
||||
|
||||
|
||||
# =============================================
|
||||
# 3. handle() — hook handler
|
||||
# =============================================
|
||||
|
||||
|
||||
class TestHandle:
|
||||
def test_skips_subagent(self):
|
||||
result = handle({"agent_type": "subagent", "prompt": "hello"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_empty_prompt(self):
|
||||
result = handle({"prompt": ""})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_whitespace_only(self):
|
||||
result = handle({"prompt": " "})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_missing_prompt(self):
|
||||
result = handle({})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_tg_origin(self):
|
||||
result = handle({"prompt": "Patrick via Telegram: hello"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_no_bot_found(self, tmp_path):
|
||||
with (
|
||||
patch.object(relay_mod, "MIRROR_DIR", tmp_path / "nope1"),
|
||||
patch.object(relay_mod, "PENDING_DIR", tmp_path / "nope2"),
|
||||
):
|
||||
result = handle({"prompt": "hello", "cwd": "/tmp/nowhere"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_happy_path_relays(self, bot_dirs):
|
||||
with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send:
|
||||
result = handle({"prompt": "hello world", "cwd": str(bot_dirs["work"])})
|
||||
assert result["exit_code"] == 0
|
||||
mock_send.assert_called_once_with("123:FAKETOKEN", 42, "hello world")
|
||||
|
||||
def test_updates_dedup_hash_on_success(self, bot_dirs):
|
||||
with patch.object(relay_mod, "send_user_message", return_value=True):
|
||||
handle({"prompt": "hello", "cwd": str(bot_dirs["work"])})
|
||||
assert relay_mod._last_relay_hash != ""
|
||||
|
||||
def test_skips_consecutive_duplicate(self, bot_dirs):
|
||||
with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send:
|
||||
handle({"prompt": "hello", "cwd": str(bot_dirs["work"])})
|
||||
handle({"prompt": "hello", "cwd": str(bot_dirs["work"])})
|
||||
mock_send.assert_called_once()
|
||||
|
||||
def test_allows_different_after_dupe(self, bot_dirs):
|
||||
with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send:
|
||||
handle({"prompt": "hello", "cwd": str(bot_dirs["work"])})
|
||||
handle({"prompt": "world", "cwd": str(bot_dirs["work"])})
|
||||
assert mock_send.call_count == 2
|
||||
|
||||
def test_no_dedup_on_send_failure(self, bot_dirs):
|
||||
with patch.object(relay_mod, "send_user_message", return_value=False) as mock_send:
|
||||
handle({"prompt": "hello", "cwd": str(bot_dirs["work"])})
|
||||
handle({"prompt": "hello", "cwd": str(bot_dirs["work"])})
|
||||
assert mock_send.call_count == 2
|
||||
|
||||
def test_exception_returns_clean(self, bot_dirs):
|
||||
with patch.object(relay_mod, "find_bot_for_cwd", side_effect=RuntimeError("boom")):
|
||||
result = handle({"prompt": "hello", "cwd": "/tmp"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_system_notification(self):
|
||||
result = handle({"prompt": "[SYSTEM NOTIFICATION - NOT USER INPUT]\nSome event happened"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_task_notification(self):
|
||||
result = handle({"prompt": "Some text\n<task-notification>\n<task-id>abc</task-id>\n</task-notification>"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_local_command_output(self):
|
||||
result = handle({"prompt": "output\n<command-name>/help</command-name>"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_local_command_stdout(self):
|
||||
result = handle({"prompt": "ran a command\n<local-command-stdout>stuff</local-command-stdout>"})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_caveat_line(self):
|
||||
prompt = "Caveat: messages below were generated by the user while running local commands\nsome output"
|
||||
result = handle({"prompt": prompt})
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_skips_dispatch_wake(self):
|
||||
result = handle(
|
||||
{"prompt": ("Hi. Check inbox, process new emails, update memories when done. IMPORTANT: delete lock file")}
|
||||
)
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
def test_genuine_message_still_relays(self, bot_dirs):
|
||||
with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send:
|
||||
handle({"prompt": "Can you fix that bug?", "cwd": str(bot_dirs["work"])})
|
||||
mock_send.assert_called_once()
|
||||
|
||||
|
||||
# =============================================
|
||||
# 4. _is_system_noise — unit tests
|
||||
# =============================================
|
||||
|
||||
|
||||
class TestIsSystemNoise:
|
||||
def test_system_notification_prefix(self):
|
||||
assert _is_system_noise("[SYSTEM NOTIFICATION - NOT USER INPUT]\nblah") is True
|
||||
|
||||
def test_task_notification_tag(self):
|
||||
assert _is_system_noise("prefix\n<task-notification>\nstuff") is True
|
||||
|
||||
def test_command_name_tag(self):
|
||||
assert _is_system_noise("<command-name>/foo</command-name>") is True
|
||||
|
||||
def test_local_command_stdout_tag(self):
|
||||
assert _is_system_noise("<local-command-stdout>output</local-command-stdout>") is True
|
||||
|
||||
def test_caveat_line(self):
|
||||
assert _is_system_noise("messages below were generated by the user while running local commands") is True
|
||||
|
||||
def test_dispatch_wake_prefix(self):
|
||||
assert _is_system_noise("Hi. Check inbox, process new emails, update memories when done.") is True
|
||||
|
||||
def test_genuine_message_passes(self):
|
||||
assert _is_system_noise("Can you fix that bug?") is False
|
||||
|
||||
def test_empty_passes(self):
|
||||
assert _is_system_noise("") is False
|
||||
|
||||
def test_partial_match_not_triggered(self):
|
||||
assert _is_system_noise("I got a SYSTEM NOTIFICATION today") is False
|
||||
|
||||
def test_dispatch_prefix_mid_message_not_triggered(self):
|
||||
assert _is_system_noise("He said Hi. Check inbox, process new emails") is False
|
||||
Reference in New Issue
Block a user