"""Matrix appservice bot interface for Steward. Registers as a Synapse application service. Synapse pushes room events to the appservice HTTP server (``/_matrix/app/v1/transactions/{txnId}``); the bot replies through the client-server API using the appservice's ``as_token``. """ from __future__ import annotations import asyncio import logging from mautrix.appservice import AppService from mautrix.types import ( Event, EventType, MessageEvent, MessageType, TextMessageEventContent, ) from steward.bot.core import ConversationService from steward.bot.thread_key import ThreadKey from steward.config import Settings logger = logging.getLogger(__name__) _MATRIX_SYSTEM_APPENDIX = """ Matrix 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. - Keep replies to a single message unless splitting genuinely helps readability. """.strip() class StewardMatrixBot: """Matrix appservice bot that drives the shared conversation pipeline.""" def __init__(self, settings: Settings, service: ConversationService) -> None: self._settings = settings self._service = service self._az = AppService( server=settings.matrix.homeserver_url, domain=settings.matrix.homeserver_domain, as_token=settings.matrix.as_token, hs_token=settings.matrix.hs_token, bot_localpart=settings.matrix.bot_localpart, id=settings.matrix.appservice_id, ) self._az.matrix_event_handler(self._on_event) @property def bot_mxid(self) -> str: return self._az.bot_mxid def _is_allowed_user(self, sender: str) -> bool: allowed = self._settings.matrix.allowed_user_ids if not allowed: return True return sender in allowed def _is_allowed_room(self, room_id: str) -> bool: allowed = self._settings.matrix.allowed_room_ids if not allowed: return True return room_id in allowed def _matrix_system_prompt(self) -> str: base_prompt = self._settings.openai_system_prompt.strip() if not base_prompt: return _MATRIX_SYSTEM_APPENDIX return f"{base_prompt}\n\n{_MATRIX_SYSTEM_APPENDIX}" async def _on_event(self, evt: Event) -> None: if not isinstance(evt, MessageEvent): return if evt.sender == self.bot_mxid: return if evt.type != EventType.ROOM_MESSAGE: return if not isinstance(evt.content, TextMessageEventContent): return if not self._is_allowed_user(evt.sender): logger.info("Ignoring message from unauthorized user %s", evt.sender) return if not self._is_allowed_room(evt.room_id): logger.info("Ignoring message in unauthorized room %s", evt.room_id) return body = (evt.content.body or "").strip() if not body: return key = ThreadKey(platform="matrix", scope=evt.room_id) try: reply = await self._service.process_message( key, evt.sender, body, self._matrix_system_prompt(), history_cap=40, ) except Exception: logger.exception( "Failed to process Matrix message from %s in %s", evt.sender, evt.room_id ) return if not reply.strip(): return content = TextMessageEventContent(msgtype=MessageType.TEXT, body=reply) await self._az.intent.send_message(evt.room_id, content) async def run(self) -> None: """Start the appservice HTTP server and keep the event loop alive.""" await self._az.start( host=self._settings.matrix.listen_host, port=self._settings.matrix.listen_port, ) logger.info( "Matrix appservice listening on %s:%s (bot %s)", self._settings.matrix.listen_host, self._settings.matrix.listen_port, self.bot_mxid, ) await self._az.intent.set_displayname("Steward") while True: await asyncio.sleep(3600) async def stop(self) -> None: await self._az.stop()