steward_mirror/steward/bot/telegram.py
Andrew Ridgway f79214bb22
fix: address pr_reviewer findings (data loss, memory, prompt injection)
- CI/CD: stop deleting the steward namespace on every deploy; use
  kubectl apply (dry-run -> apply) so the PVC and conversation memory
  survive deployments.
- core: bound _histories with an LRU eviction (max 1000 active threads)
  to prevent unbounded memory growth.
- core: wrap knowledge-base context in KB START/END delimiters and
  instruct the LLM to treat it as data, mitigating indirect prompt
  injection.
- matrix: wrap message processing in try/except so failures are logged
  instead of silently dropped.
- telegram: remove now-dead flush/tags prompt constants (centralized in
  core).

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
2026-08-18 22:01:59 +10:00

540 lines
19 KiB
Python

"""Telegram bot interface for Steward."""
import logging
import re
from dataclasses import dataclass, field
from typing import Any
from telegram import Update
from telegram.constants import ParseMode
from telegram.ext import (
Application,
CommandHandler,
ContextTypes,
MessageHandler,
filters,
)
from steward.bot.core import ConversationService
from steward.bot.thread_key import ThreadKey
from steward.config import Settings
from steward.proposals.generator import Proposal, ProposalGenerator
logger = logging.getLogger(__name__)
_CONVERSATION_SYSTEM_APPENDIX = """
Telegram conversation guidance:
- Reply like a thoughtful software engineer in a chat, not like a one-shot FAQ bot.
- Ask a brief follow-up question when the requested outcome, constraints, or preferred
option is unclear.
- You may send multiple short Telegram messages when that makes the conversation easier to follow.
- To send multiple messages, separate each message with a line containing only [MESSAGE].
- When the user needs to choose from 2-10 concise options, use Telegram's native poll format:
[POLL]
question: The decision to make?
- First option
- Second option
[/POLL]
- Use polls only for discrete choices. Use a normal clarifying question for open-ended input.
""".strip()
_MESSAGE_MARKER_RE = re.compile(r"^\s*\[MESSAGE\]\s*$", re.IGNORECASE)
_POLL_START_RE = re.compile(r"^\s*\[POLL\]\s*$", re.IGNORECASE)
_POLL_END_RE = re.compile(r"^\s*\[/POLL\]\s*$", re.IGNORECASE)
_OPTION_RE = re.compile(r"^\s*(?:[-*•]|\d+[.)]|option\s*:?)\s*(.+?)\s*$", re.IGNORECASE)
_QUESTION_RE = re.compile(r"^\s*question\s*:\s*(.+?)\s*$", re.IGNORECASE)
_MAX_POLL_OPTIONS = 10
_MAX_POLL_QUESTION_LENGTH = 300
_MAX_POLL_OPTION_LENGTH = 100
@dataclass(frozen=True)
class PollRequest:
"""A Telegram poll requested by the assistant response."""
question: str
options: list[str]
TelegramResponseAction = str | PollRequest
@dataclass(frozen=True)
class TelegramResponsePlan:
"""Telegram actions derived from a single assistant response."""
actions: list[TelegramResponseAction] = field(default_factory=list)
@property
def messages(self) -> list[str]:
"""Text messages in this response plan."""
return [action for action in self.actions if isinstance(action, str)]
@property
def polls(self) -> list[PollRequest]:
"""Polls in this response plan."""
return [action for action in self.actions if isinstance(action, PollRequest)]
def _thread_key(update: Update) -> ThreadKey | None:
"""Return the ThreadKey if the message is part of a Telegram thread, else None."""
msg = update.message
chat = update.effective_chat
if msg is None or chat is None:
return None
thread_id = msg.message_thread_id
if thread_id is None:
return None
return ThreadKey(platform="telegram", scope=str(chat.id), thread=str(thread_id))
def _is_allowed(user_id: int, settings: Settings) -> bool:
"""Return True if the user is in the allow-list (or no list is configured)."""
if not settings.telegram_allowed_user_ids:
return True
allowed = user_id in settings.telegram_allowed_user_ids
if not allowed:
logger.info("Ignoring update from unauthorized user %s", user_id)
return allowed
def _is_group_enabled(chat_id: int, settings: Settings) -> bool:
"""Return True if the group/channel is in the configured list (or no list is configured)."""
if not settings.telegram_group_ids:
return True
return chat_id in settings.telegram_group_ids
def _is_chat_enabled(chat: Any, settings: Settings) -> bool:
"""Return True if the chat is allowed.
Private chats are governed only by the user allow-list. Group/channel allow-listing
applies only to non-private chats.
"""
if getattr(chat, "type", None) == "private":
return True
allowed = _is_group_enabled(chat.id, settings)
if not allowed:
logger.info(
"Ignoring update in unauthorized chat %s (type=%s); configured group IDs: %s",
chat.id,
getattr(chat, "type", None),
settings.telegram_group_ids,
)
return allowed
def _telegram_system_prompt(settings: Settings) -> str:
"""Return the conversational Telegram system prompt for normal chat turns."""
base_prompt = settings.openai_system_prompt.strip()
if not base_prompt:
return _CONVERSATION_SYSTEM_APPENDIX
return f"{base_prompt}\n\n{_CONVERSATION_SYSTEM_APPENDIX}"
def _append_message_action(actions: list[TelegramResponseAction], lines: list[str]) -> None:
message = "\n".join(lines).strip()
if message:
actions.append(message)
lines.clear()
def _compact_telegram_text(text: str, max_length: int) -> str:
return re.sub(r"\s+", " ", text).strip()[:max_length]
def _parse_poll(lines: list[str]) -> PollRequest | None:
question = ""
options: list[str] = []
for line in lines:
stripped = line.strip()
if not stripped:
continue
question_match = _QUESTION_RE.match(stripped)
if question_match and not question:
question = question_match.group(1)
continue
option_match = _OPTION_RE.match(stripped)
if option_match:
options.append(option_match.group(1))
continue
if not question:
question = stripped
question = _compact_telegram_text(question, _MAX_POLL_QUESTION_LENGTH)
unique_options: list[str] = []
for option in options:
compact_option = _compact_telegram_text(option, _MAX_POLL_OPTION_LENGTH)
if compact_option and compact_option not in unique_options:
unique_options.append(compact_option)
unique_options = unique_options[:_MAX_POLL_OPTIONS]
if not question or len(unique_options) < 2:
return None
return PollRequest(question=question, options=unique_options)
def _parse_telegram_response(reply: str) -> TelegramResponsePlan:
"""Parse assistant response directives into ordered Telegram actions."""
actions: list[TelegramResponseAction] = []
message_lines: list[str] = []
poll_lines: list[str] = []
in_poll = False
for line in reply.splitlines():
if in_poll:
if _POLL_END_RE.match(line):
poll = _parse_poll(poll_lines)
if poll is None:
message_lines.extend(["[POLL]", *poll_lines, "[/POLL]"])
else:
actions.append(poll)
poll_lines.clear()
in_poll = False
else:
poll_lines.append(line)
continue
if _MESSAGE_MARKER_RE.match(line):
_append_message_action(actions, message_lines)
continue
if _POLL_START_RE.match(line):
_append_message_action(actions, message_lines)
in_poll = True
continue
message_lines.append(line)
if in_poll:
message_lines.extend(["[POLL]", *poll_lines])
_append_message_action(actions, message_lines)
if not actions and reply.strip():
actions.append(reply.strip())
return TelegramResponsePlan(actions=actions)
def _format_response_for_history(plan: TelegramResponsePlan) -> str:
"""Render Telegram actions back into natural assistant history text."""
history_parts: list[str] = []
for action in plan.actions:
if isinstance(action, str):
history_parts.append(action)
else:
options = "; ".join(action.options)
history_parts.append(f"Poll: {action.question} ({options})")
return "\n\n".join(history_parts)
async def _send_telegram_response(update: Update, plan: TelegramResponsePlan) -> None:
"""Send all Telegram actions for a parsed assistant response."""
for action in plan.actions:
if isinstance(action, str):
await _send_long(update, action)
else:
await update.message.reply_poll( # type: ignore[union-attr]
question=action.question,
options=action.options,
is_anonymous=False,
)
async def _send_long(update: Update, text: str) -> None:
"""Send text, splitting across messages if it exceeds Telegram's 4096-char limit."""
limit = 4096
for i in range(0, len(text), limit):
await update.message.reply_text( # type: ignore[union-attr]
text[i : i + limit],
parse_mode=ParseMode.MARKDOWN,
)
async def start_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
"""Handle /start."""
settings: Settings = context.bot_data["settings"]
user = update.effective_user
if user is None or not _is_allowed(user.id, settings):
return
chat = update.effective_chat
if chat is None or not _is_chat_enabled(chat, settings):
return
await update.message.reply_text( # type: ignore[union-attr]
"Hello, I'm *Steward* \U0001f916\n\n"
"I'm your AI-assisted personal operations platform.\n"
"Talk to me naturally, or use:\n"
"/help \u2013 show available commands\n"
"/clear \u2013 reset conversation history\n"
"/flush \u2013 summarise and archive this thread's memory\n"
"/recall [query] \u2013 retrieve archived thread summaries\n"
"/analyse \u2013 run a manual API analysis right now",
parse_mode=ParseMode.MARKDOWN,
)
async def help_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
"""Handle /help."""
settings: Settings = context.bot_data["settings"]
user = update.effective_user
if user is None or not _is_allowed(user.id, settings):
return
chat = update.effective_chat
if chat is None or not _is_chat_enabled(chat, settings):
return
await update.message.reply_text( # type: ignore[union-attr]
"*Steward commands*\n\n"
"/start \u2013 greeting\n"
"/help \u2013 this message\n"
"/clear \u2013 reset conversation history for this context\n"
"/flush \u2013 summarise the current thread, store the summary, and compress memory\n"
" _(only available inside a message thread)_\n"
"/recall [query] \u2013 show this thread's summary, list all summaries, or search\n"
" by keyword\n"
"/analyse \u2013 trigger an immediate API analysis and proposal",
parse_mode=ParseMode.MARKDOWN,
)
async def clear_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
"""Handle /clear \u2013 wipe conversation history for this context.
Inside a message thread: clears the thread's unbounded history.
Outside a thread: clears the per-user capped history.
"""
settings: Settings = context.bot_data["settings"]
service: ConversationService = context.bot_data["service"]
user = update.effective_user
if user is None or not _is_allowed(user.id, settings):
return
chat = update.effective_chat
if chat is None or not _is_chat_enabled(chat, settings):
return
key = _thread_key(update)
if key is not None:
service.clear_history(key)
await update.message.reply_text( # type: ignore[union-attr]
"Thread conversation history cleared."
)
else:
service.clear_history(ThreadKey(platform="telegram", scope=str(user.id)))
await update.message.reply_text( # type: ignore[union-attr]
"Conversation history cleared."
)
async def flush_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
"""Handle /flush \u2013 summarise thread memory, persist it, and compress in-memory history."""
settings: Settings = context.bot_data["settings"]
service: ConversationService = context.bot_data["service"]
user = update.effective_user
if user is None or not _is_allowed(user.id, settings):
return
chat = update.effective_chat
if chat is None or not _is_chat_enabled(chat, settings):
return
key = _thread_key(update)
if key is None:
await update.message.reply_text( # type: ignore[union-attr]
"\u26a0\ufe0f /flush can only be used inside a message thread."
)
return
if not service.has_history(key):
await update.message.reply_text( # type: ignore[union-attr]
"This thread has no conversation history to flush."
)
return
await update.message.reply_text( # type: ignore[union-attr]
"\U0001f4be Summarising thread memory\u2026"
)
thread_summary = await service.flush_history(key)
if thread_summary is None:
await update.message.reply_text( # type: ignore[union-attr]
"This thread has no conversation history to flush."
)
return
await update.message.reply_text( # type: ignore[union-attr]
f"\u2705 Thread memory flushed and stored "
f"(thread `{thread_summary.thread_id}`, "
f"{thread_summary.message_count} messages summarised).",
parse_mode=ParseMode.MARKDOWN,
)
async def recall_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
"""Handle /recall [query] \u2013 retrieve stored thread summaries."""
settings: Settings = context.bot_data["settings"]
service: ConversationService = context.bot_data["service"]
user = update.effective_user
if user is None or not _is_allowed(user.id, settings):
return
chat = update.effective_chat
if chat is None or not _is_chat_enabled(chat, settings):
return
args: list[str] = context.args or []
if args:
query = " ".join(args).strip()
results = service.search_summaries(query)
if not results:
await update.message.reply_text( # type: ignore[union-attr]
f"No memories found matching *{query}*. "
"Try a different keyword or use /flush to add more summaries.",
parse_mode=ParseMode.MARKDOWN,
)
return
lines = [f"\U0001f50d *Knowledge base search: {query}*\n"]
for s in results:
date = s.flushed_at[:10]
first_line = s.summary.split("\n")[0][:80]
lines.append(f"\u2022 Thread `{s.thread_id}` ({date}): {first_line}\u2026")
if s.tags:
lines.append(f" \U0001f3f7 {', '.join(s.tags)}")
await _send_long(update, "\n".join(lines))
return
key = _thread_key(update)
if key is not None:
stored = service.get_summary(key)
if stored is None:
await update.message.reply_text( # type: ignore[union-attr]
"No stored summary for this thread yet. Use /flush to create one."
)
else:
await _send_long(update, stored.format_for_telegram())
return
all_summaries = service.all_summaries()
if not all_summaries:
await update.message.reply_text( # type: ignore[union-attr]
"No thread summaries stored yet. Use /flush inside a message thread."
)
return
lines = ["\U0001f4da *Stored thread summaries*\n"]
for s in all_summaries:
date = s.flushed_at[:10]
first_line = s.summary.split("\n")[0][:80]
lines.append(f"\u2022 Thread `{s.thread_id}` ({date}): {first_line}\u2026")
if s.tags:
lines.append(f" \U0001f3f7 {', '.join(s.tags)}")
await _send_long(update, "\n".join(lines))
async def analyse_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
"""Handle /analyse \u2013 run the proposal generator on demand."""
settings: Settings = context.bot_data["settings"]
user = update.effective_user
if user is None or not _is_allowed(user.id, settings):
return
chat = update.effective_chat
if chat is None or not _is_chat_enabled(chat, settings):
return
await update.message.reply_text("Running analysis, please wait...") # type: ignore[union-attr]
generator = ProposalGenerator(settings, context.bot_data["llm"])
proposal = await generator.run()
if proposal is None:
await update.message.reply_text( # type: ignore[union-attr]
"Analysis could not be completed. Check that `ANALYSIS_TARGET_URL` is configured."
)
return
await _send_long(update, proposal.format_for_telegram())
async def message_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
"""Handle plain text messages \u2013 forward to the shared pipeline and reply."""
settings: Settings = context.bot_data["settings"]
service: ConversationService = context.bot_data["service"]
user = update.effective_user
chat = update.effective_chat
if user is None or not _is_allowed(user.id, settings):
return
if chat is None or not _is_chat_enabled(chat, settings):
return
text = update.message.text # type: ignore[union-attr]
if not text:
return
system_prompt = _telegram_system_prompt(settings)
key = _thread_key(update)
if key is None:
key = ThreadKey(platform="telegram", scope=str(user.id))
history_cap = 40
else:
history_cap = None
reply = await service.process_message(
key,
str(user.id),
text,
system_prompt,
history_cap=history_cap,
history_formatter=lambda raw: _format_response_for_history(_parse_telegram_response(raw)),
)
response_plan = _parse_telegram_response(reply)
await _send_telegram_response(update, response_plan)
def build_application(
settings: Settings,
service: ConversationService,
generator: ProposalGenerator | None = None,
) -> Application: # type: ignore[type-arg]
"""Build and return the Telegram Application."""
app = Application.builder().token(settings.telegram_bot_token).build()
app.bot_data["settings"] = settings
app.bot_data["service"] = service
app.bot_data["llm"] = service.llm
app.bot_data["generator"] = generator
logger.info("Configured allowed users: %s", settings.telegram_allowed_user_ids)
logger.info("Configured group IDs: %s", settings.telegram_group_ids)
app.add_handler(CommandHandler("start", start_handler))
app.add_handler(CommandHandler("help", help_handler))
app.add_handler(CommandHandler("clear", clear_handler))
app.add_handler(CommandHandler("flush", flush_handler))
app.add_handler(CommandHandler("recall", recall_handler))
app.add_handler(CommandHandler("analyse", analyse_handler))
app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, message_handler))
return app
async def send_proposal(app: Application, proposal: Proposal, user_ids: list[int]) -> None: # type: ignore[type-arg]
"""Send a proposal message to all configured user IDs."""
text = proposal.format_for_telegram()
for uid in user_ids:
try:
await app.bot.send_message(
chat_id=uid,
text=text,
parse_mode=ParseMode.MARKDOWN,
)
except Exception:
logger.exception("Failed to send proposal to user %d", uid)