From 62d047cac1c6d011b85ae62608888c970c55dda8 Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Tue, 14 Jul 2026 09:16:02 -0700 Subject: [PATCH] =?UTF-8?q?fix:=20TG=20mirror=20live-test=20round=20?= =?UTF-8?q?=E2=80=94=20agent=5Ftype=20filter=20unblocked=20main=20chats=20?= =?UTF-8?q?(daemon-backed=20sessions=20carry=20agent=5Ftype=3Dclaude;=20su?= =?UTF-8?q?bagents=20never=20fire=20UserPromptSubmit;=20skip=20is=20now=20?= =?UTF-8?q?agent=5Fid-based)=20+=20TG-echo=20gate=20(bot=20stores=20inject?= =?UTF-8?q?ed=5Fprompt=20in=20pending=20file,=20relay=20skips=20fresh=20un?= =?UTF-8?q?delivered=20text-match;=20raw=20injections=20carry=20no=20marke?= =?UTF-8?q?r).=20Found=20live=20by=20Patrick's=20morning=20door-tests;=20m?= =?UTF-8?q?irror=20proven=20both=20directions.=20791=20TG=20tests=20green?= =?UTF-8?q?=20(incl.=20cross-file=20mock=20fix=20@skills=20missed).?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 13 ++ .../lib/telegram/apps/handlers/base_bot.py | 14 +- .../apps/handlers/user_message_relay.py | 41 +++++- .../tests/test_inbound_reliability.py | 2 +- .../telegram/tests/test_user_message_relay.py | 134 +++++++++++++++++- 5 files changed, 192 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d69ec28e..effc2811 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -33,6 +33,19 @@ PyPI version — not the changelog header. ### Fixed +- **TG mirror live-test fixes: main-chat messages mirror, TG messages don't + echo.** Patrick's first morning test caught what 47 green tests missed: the + relay's sub-agent skip blocked ALL daemon-backed main chats (they run with + `--agent claude`, so `agent_type="claude"` — and real sub-agents never fire + UserPromptSubmit at all; the filter's premise was empirically wrong across the + entire engine log). Skip is now agent_id-based (defensive, never observed). + Second catch from tracing his test: TG messages inject into tmux as raw text — + no `via Telegram:` marker — so the TG-origin filter never matched and every + TG message would have echoed back once the first fix landed. New structural + gate: the bot stores the injected prompt in its pending file; the relay skips + a prompt that text-matches a fresh undelivered pending entry. Mirror proven + live by Patrick across both directions ("success :)"). 791 TG tests green. + - **DPLAN-0241 round 4 (night shift): user flags survive every launch path, and every session is born with an honest name.** R6 — the bug behind Patrick's approve-everything chat: the boot menu suppressed its bypass defaults when the diff --git a/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py b/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py index f9bf565e..db0fe9ba 100644 --- a/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py +++ b/src/aipass/skills/lib/telegram/apps/handlers/base_bot.py @@ -716,7 +716,7 @@ class BaseBot: processing_msg_id = processing_result.get("message_id") if processing_result else None # Write pending file - if not self.write_pending_file(chat_id, message_id, processing_msg_id): + if not self.write_pending_file(chat_id, message_id, processing_msg_id, injected_prompt=prompt): logger.error("Failed to write pending file") self.send_message(chat_id, "Internal error writing pending file.") return @@ -814,7 +814,7 @@ class BaseBot: processing_msg_id = processing_result.get("message_id") if processing_result else None # Write pending file - if not self.write_pending_file(chat_id, message_id, processing_msg_id): + if not self.write_pending_file(chat_id, message_id, processing_msg_id, injected_prompt=prompt): logger.error("Failed to write pending file for file upload") self.send_message(chat_id, "Internal error writing pending file.") return @@ -1440,7 +1440,13 @@ class BaseBot: # PENDING FILE MANAGEMENT # ============================================= - def write_pending_file(self, chat_id: int, message_id: int, processing_message_id: Optional[int] = None) -> bool: + def write_pending_file( + self, + chat_id: int, + message_id: int, + processing_message_id: Optional[int] = None, + injected_prompt: str = "", + ) -> bool: """ Write the pending file for Stop hook coordination. @@ -1469,6 +1475,8 @@ class BaseBot: "transcript_path": str(self._active_transcript_path) if self._active_transcript_path else None, "session_id": self._active_session_id, } + if injected_prompt: + pending_data["injected_prompt"] = injected_prompt if self._stream: pending_data["streaming"] = True diff --git a/src/aipass/skills/lib/telegram/apps/handlers/user_message_relay.py b/src/aipass/skills/lib/telegram/apps/handlers/user_message_relay.py index b03ab7ec..4fd271ca 100644 --- a/src/aipass/skills/lib/telegram/apps/handlers/user_message_relay.py +++ b/src/aipass/skills/lib/telegram/apps/handlers/user_message_relay.py @@ -1,7 +1,7 @@ # =================== AIPass ==================== # Name: user_message_relay.py # Description: Relay user messages from non-TG doors to the branch TG chat -# Version: 1.1.0 +# Version: 1.2.0 # Created: 2026-07-14 # Modified: 2026-07-14 # ============================================= @@ -19,6 +19,7 @@ Registration: @hooks adds this to .aipass/hooks.json + ~/.claude/settings.json. import hashlib import json import os +import time from pathlib import Path from urllib.error import URLError from urllib.request import Request, urlopen @@ -104,6 +105,29 @@ def send_user_message(bot_token: str, chat_id: int, text: str, origin: str = "\U return False +_PENDING_TTL = 120 + + +def _is_pending_tg_message(prompt: str, bot_data: dict) -> bool: + """Check if prompt matches a fresh pending TG injection for this bot.""" + bot_id = bot_data.get("bot_id") + if not bot_id: + return False + pending_path = PENDING_DIR / f"bot-{bot_id}.json" + if not pending_path.exists(): + return False + try: + pending = json.loads(pending_path.read_text(encoding="utf-8")) + except (json.JSONDecodeError, OSError): + return False + if pending.get("delivered"): + return False + ts = pending.get("timestamp", 0) + if time.time() - ts > _PENDING_TTL: + return False + return pending.get("injected_prompt", "") == prompt + + def _is_system_noise(prompt: str) -> bool: """Detect non-human system noise that should not be mirrored to TG.""" if prompt.startswith("[SYSTEM NOTIFICATION"): @@ -122,13 +146,17 @@ def _is_system_noise(prompt: str) -> bool: 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. + Skips: identified subagents (non-empty agent_id), 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", ""): + # agent_type is NOT a reliable subagent indicator — main branch chats run + # with agent_type="claude" (--agent claude). Subagents spawned by tools + # never fire UserPromptSubmit. Defensive: skip only if agent_id is + # non-empty, which would indicate a future CC subagent prompt route. + if hook_data.get("agent_id", ""): return {"stdout": "", "exit_code": 0} prompt = hook_data.get("prompt", "") @@ -150,6 +178,9 @@ def handle(hook_data: dict) -> dict: if not bot_data: return {"stdout": "", "exit_code": 0} + if _is_pending_tg_message(prompt, bot_data): + return {"stdout": "", "exit_code": 0} + bot_token = bot_data["bot_token"] chat_id = int(bot_data["chat_id"]) diff --git a/src/aipass/skills/lib/telegram/tests/test_inbound_reliability.py b/src/aipass/skills/lib/telegram/tests/test_inbound_reliability.py index c2d7d57b..8412be8f 100644 --- a/src/aipass/skills/lib/telegram/tests/test_inbound_reliability.py +++ b/src/aipass/skills/lib/telegram/tests/test_inbound_reliability.py @@ -87,7 +87,7 @@ class TestInboundReliability: patch.object( bot, "write_pending_file", - side_effect=lambda *a: (call_order.append("write"), True)[1], + side_effect=lambda *a, **kw: (call_order.append("write"), True)[1], ), patch.object(bot, "inject_message", return_value=True), patch.object(bot, "_start_heartbeat"), diff --git a/src/aipass/skills/lib/telegram/tests/test_user_message_relay.py b/src/aipass/skills/lib/telegram/tests/test_user_message_relay.py index 838dd242..3f4154d8 100644 --- a/src/aipass/skills/lib/telegram/tests/test_user_message_relay.py +++ b/src/aipass/skills/lib/telegram/tests/test_user_message_relay.py @@ -21,12 +21,14 @@ Tests cover: """ import json +import time 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_pending_tg_message, _is_system_noise, find_bot_for_cwd, handle, @@ -223,10 +225,22 @@ class TestSendUserMessage: class TestHandle: - def test_skips_subagent(self): - result = handle({"agent_type": "subagent", "prompt": "hello"}) + def test_skips_identified_subagent(self): + result = handle({"agent_id": "agent-123", "prompt": "hello"}) assert result["exit_code"] == 0 + def test_allows_agent_type_claude(self, bot_dirs): + with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send: + handle( + { + "agent_type": "claude", + "agent_id": "", + "prompt": "hello", + "cwd": str(bot_dirs["work"]), + } + ) + mock_send.assert_called_once() + def test_skips_empty_prompt(self): result = handle({"prompt": ""}) assert result["exit_code"] == 0 @@ -317,9 +331,123 @@ class TestHandle: handle({"prompt": "Can you fix that bug?", "cwd": str(bot_dirs["work"])}) mock_send.assert_called_once() + def test_skips_pending_tg_message(self, bot_dirs): + pending_data = { + "chat_id": 42, + "bot_token": "123:FAKETOKEN", + "bot_id": "test_bot", + "injected_prompt": "hello from TG", + "timestamp": time.time(), + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send: + result = handle({"prompt": "hello from TG", "cwd": str(bot_dirs["work"])}) + assert result["exit_code"] == 0 + mock_send.assert_not_called() + + def test_allows_message_no_pending(self, bot_dirs): + with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send: + handle({"prompt": "hello from terminal", "cwd": str(bot_dirs["work"])}) + mock_send.assert_called_once() + + def test_allows_message_stale_pending(self, bot_dirs): + pending_data = { + "chat_id": 42, + "bot_token": "123:FAKETOKEN", + "bot_id": "test_bot", + "injected_prompt": "hello from TG", + "timestamp": time.time() - 300, + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send: + handle({"prompt": "hello from TG", "cwd": str(bot_dirs["work"])}) + mock_send.assert_called_once() + + def test_allows_message_delivered_pending(self, bot_dirs): + pending_data = { + "chat_id": 42, + "bot_token": "123:FAKETOKEN", + "bot_id": "test_bot", + "injected_prompt": "hello from TG", + "timestamp": time.time(), + "delivered": True, + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send: + handle({"prompt": "hello from TG", "cwd": str(bot_dirs["work"])}) + mock_send.assert_called_once() + + def test_allows_message_different_text_pending(self, bot_dirs): + pending_data = { + "chat_id": 42, + "bot_token": "123:FAKETOKEN", + "bot_id": "test_bot", + "injected_prompt": "something else entirely", + "timestamp": time.time(), + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + with patch.object(relay_mod, "send_user_message", return_value=True) as mock_send: + handle({"prompt": "hello from terminal", "cwd": str(bot_dirs["work"])}) + mock_send.assert_called_once() + # ============================================= -# 4. _is_system_noise — unit tests +# 4. _is_pending_tg_message — unit tests +# ============================================= + + +class TestIsPendingTgMessage: + def test_matches_fresh_pending(self, bot_dirs): + pending_data = { + "injected_prompt": "test msg", + "timestamp": time.time(), + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + assert _is_pending_tg_message("test msg", {"bot_id": "test_bot"}) is True + + def test_no_match_different_text(self, bot_dirs): + pending_data = { + "injected_prompt": "other msg", + "timestamp": time.time(), + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + assert _is_pending_tg_message("test msg", {"bot_id": "test_bot"}) is False + + def test_no_match_stale(self, bot_dirs): + pending_data = { + "injected_prompt": "test msg", + "timestamp": time.time() - 300, + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + assert _is_pending_tg_message("test msg", {"bot_id": "test_bot"}) is False + + def test_no_match_delivered(self, bot_dirs): + pending_data = { + "injected_prompt": "test msg", + "timestamp": time.time(), + "delivered": True, + } + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + assert _is_pending_tg_message("test msg", {"bot_id": "test_bot"}) is False + + def test_no_bot_id(self): + assert _is_pending_tg_message("test", {}) is False + + def test_no_pending_file(self, bot_dirs): + assert _is_pending_tg_message("test", {"bot_id": "nonexistent"}) is False + + def test_corrupt_pending_file(self, bot_dirs): + (bot_dirs["pending"] / "bot-test_bot.json").write_text("not json{{{") + assert _is_pending_tg_message("test", {"bot_id": "test_bot"}) is False + + def test_no_injected_prompt_field(self, bot_dirs): + pending_data = {"timestamp": time.time()} + (bot_dirs["pending"] / "bot-test_bot.json").write_text(json.dumps(pending_data)) + assert _is_pending_tg_message("test", {"bot_id": "test_bot"}) is False + + +# ============================================= +# 5. _is_system_noise — unit tests # =============================================