fix: TG mirror live-test round — agent_type filter unblocked main chats (daemon-backed sessions carry agent_type=claude; subagents never fire UserPromptSubmit; skip is now agent_id-based) + TG-echo gate (bot stores injected_prompt in pending file, relay skips fresh undelivered text-match; raw injections carry no marker). Found live by Patrick's morning door-tests; mirror proven both directions. 791 TG tests green (incl. cross-file mock fix @skills missed).
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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"])
|
||||
|
||||
|
||||
@@ -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"),
|
||||
|
||||
@@ -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
|
||||
# =============================================
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user