""" Core Zulip adapter for Hermes agents (Tanko, Mumuni, Koonimo, Koby). Architecture (mirrors pi-zulip-extension): Zulip event queue → poll loop → parse message → spawn agent subprocess → capture stdout → post response to Zulip The Hermes agent is invoked as a subprocess (configurable command) so the plugin remains decoupled from the agent runtime. In a future iteration, this could switch to a direct IPC/socket if the agent exposes one. See ADR-005 (@mention detection), ADR-006 (@all-bots), ADR-009 (error handling). """ import asyncio import json import logging import re import shlex import subprocess import time from datetime import datetime from typing import Any, Dict, Optional logger = logging.getLogger(__name__) # Regex to strip Zulip mention artifacts from message bodies (ADR-008) MENTION_CLEANER = re.compile(r'@\*\*[^*]+\*\*') class ZulipAdapter: """Zulip adapter for Hermes agents. Connects to Zulip via event queues (not the deprecated call_on_each_message), polls for new events, routes @mentions and DMs to the Hermes agent via subprocess, and posts responses back to Zulip. """ def __init__(self, config: Dict[str, Any]) -> None: self.config = config self.connected = False self._client = None self._queue_id: Optional[str] = None self._last_event_id: int = -1 self._poll_task: Optional[asyncio.Task] = None # Agent subprocess config self._agent_command: str = config.get("agent", {}).get( "command", "hermes chat" ) self._agent_timeout: int = config.get("error_handling", {}).get( "timeout_seconds", 60 ) # Bot identity for filtering own messages self._bot_email: str = config.get("zulip", {}).get("email", "") self._bot_id: int = config.get("swarm", {}).get("bot_id", 0) # Stream config self._stream: str = config.get("zulip", {}).get("stream", "agent-hub") self._all_bots_user_id: int = config.get("zulip", {}).get( "all_bots_user_id", 1 ) # ------------------------------------------------------------------ # Connection lifecycle # ------------------------------------------------------------------ async def connect(self) -> None: """Connect to Zulip and register an event queue.""" if self.connected: logger.info("Already connected to Zulip.") return try: import zulip server_url = self.config["zulip"]["server_url"] email = self.config["zulip"]["email"] api_key = self.config["zulip"]["api_key"] self._client = zulip.Client( server_url=server_url, email=email, api_key=api_key ) # Register an event queue (like the pi extension) queue_res = self._client.register( event_types=["message"] ) self._queue_id = queue_res.get("queue_id") self._last_event_id = queue_res.get("last_event_id", -1) self.connected = True logger.info( f"Connected to Zulip: {server_url} as {email} " f"(queue={self._queue_id})" ) except Exception as e: logger.error(f"Failed to connect to Zulip: {e}") self.connected = False raise async def disconnect(self) -> None: """Disconnect and cancel the poll loop.""" if self._poll_task: self._poll_task.cancel() try: await self._poll_task except asyncio.CancelledError: pass self._poll_task = None self.connected = False logger.info("Disconnected from Zulip.") # ------------------------------------------------------------------ # Event polling (async loop, like pi extension's setInterval) # ------------------------------------------------------------------ async def poll_forever(self, poll_interval: float = 3.0) -> None: """Continuously poll the Zulip event queue and process messages. This is the main event loop. Runs until cancelled. """ while self.connected and self._client and self._queue_id: try: events = await self._poll_events() for event in events: await self._process_event(event) except asyncio.CancelledError: break except Exception as e: logger.error(f"Poll error: {e}") # Check for queue expiry if "BAD_EVENT_QUEUE_ID" in str(e): logger.info("Queue expired, reconnecting...") self.connected = False await self._reconnect_with_retry() continue await asyncio.sleep(poll_interval) async def _poll_events(self) -> list: """Fetch events from the Zulip event queue.""" import zulip if not self._client or not self._queue_id: return [] response = self._client.get_events( queue_id=self._queue_id, last_event_id=self._last_event_id, dont_block=True, ) if response.get("result") != "success": raise RuntimeError( f"Events API error: {response.get('msg', 'unknown')}" ) events = response.get("events", []) for event in events: if event.get("id", 0) > self._last_event_id: self._last_event_id = event["id"] return [e for e in events if e.get("type") == "message"] async def _reconnect_with_retry(self) -> None: """Retry connection with backoff.""" retries = self.config.get("error_handling", {}).get("retry_count", 3) delay = self.config.get("error_handling", {}).get( "retry_delay_seconds", 5 ) for attempt in range(retries): logger.info( f"Reconnect attempt {attempt + 1}/{retries} " f"in {delay}s..." ) await asyncio.sleep(delay) try: await self.connect() if self.connected: logger.info("Reconnected successfully.") return except Exception as e: logger.error(f"Reconnect attempt {attempt + 1} failed: {e}") delay *= 2 # Exponential backoff logger.error("Max reconnection retries reached. Giving up.") # ------------------------------------------------------------------ # Message processing # ------------------------------------------------------------------ async def _process_event(self, event: Dict[str, Any]) -> None: """Route a Zulip message event to the Hermes agent.""" msg = event.get("message", {}) msg_type = msg.get("type", "") # "private" or "stream" sender_email = msg.get("sender_email", "") sender_name = msg.get("sender_full_name", "Unknown") content = msg.get("content", "") # Ignore own messages if sender_email == self._bot_email: return # Determine if this message targets this bot mentioned_users = msg.get("mentioned_users", []) mentioned_user_ids = [u.get("user_id") for u in mentioned_users] is_dm = msg_type == "private" is_mention = self._bot_id in mentioned_user_ids is_all_bots = self._all_bots_user_id in mentioned_user_ids # DM-first: process all private messages (like ADR-001/ADR-002) if is_dm: logger.info(f"DM from {sender_name}: {content[:80]}...") await self._route_to_agent( message_type="dm", sender_name=sender_name, content=content, recipient=sender_email, ) return # Stream: only respond to @mentions and @all-bots (ADR-005, ADR-006) if msg_type == "stream" and (is_mention or is_all_bots): stream_name = msg.get("display_recipient", "unknown") topic = msg.get("subject", "general") logger.info( f"{'@mention' if is_mention else '@all-bots'} " f"in #{stream_name} > {topic}" ) # Clean @mention artifacts from content (ADR-008) clean_content = MENTION_CLEANER.sub("", content).strip() await self._route_to_agent( message_type="mention" if is_mention else "all-bots", sender_name=sender_name, content=clean_content, recipient=stream_name, topic=topic, stream_id=msg.get("stream_id"), ) async def _route_to_agent( self, message_type: str, sender_name: str, content: str, recipient: str, topic: Optional[str] = None, stream_id: Optional[int] = None, ) -> None: """Send the message to the Hermes agent subprocess and relay the response back to Zulip. Architecture note: This uses a subprocess (configurable via agent.command in config.yaml). This mirrors the initial pi extension approach. A future iteration could use IPC or a socket if the Hermes agent exposes one. """ # 1. Send typing indicator await self._send_typing_indicator(recipient, "start") # 2. Send a "Thinking..." placeholder to Zulip immediately. # On agent_end, this message will be edited with the final response. placeholder = ( f":robot: _Processing your message..._" ) placeholder_msg_id = None try: placeholder_msg_id = await self._send_message( message_type=message_type, recipient=recipient, content=placeholder, topic=topic, stream_id=stream_id, ) except Exception as e: logger.warning(f"Failed to send placeholder: {e}") # 3. Invoke the Hermes agent subprocess response_text = await self._invoke_agent(content, sender_name) # 4. Stop typing indicator await self._send_typing_indicator(recipient, "stop") # 5. Post (or edit) the response if not response_text.strip(): logger.warning( f"Hermes agent returned empty response for " f"{message_type} from {sender_name}" ) return # Truncate to Zulip's 10K char limit MAX_ZULIP_MSG = 10000 truncated = ( response_text[:MAX_ZULIP_MSG] + "\n\n[...truncated at Zulip limit]" if len(response_text) > MAX_ZULIP_MSG else response_text ) try: if placeholder_msg_id: await self._edit_message(placeholder_msg_id, truncated) logger.info( f"Finalized response to {sender_name} " f"({len(truncated)} chars)" ) else: await self._send_message( message_type=message_type, recipient=recipient, content=truncated, topic=topic, stream_id=stream_id, ) logger.info( f"Sent response to {sender_name} " f"({len(truncated)} chars)" ) except Exception as e: logger.error(f"Failed to post response: {e}") # ------------------------------------------------------------------ # Hermes agent subprocess invocation # ------------------------------------------------------------------ async def _invoke_agent( self, message: str, sender_name: str ) -> str: """Spawn the Hermes agent subprocess with the message. The agent command is configurable (agent.command in config.yaml). Default: "hermes chat" The message is passed via stdin (pipe). The agent's stdout is captured as the response. This is synchronous (run in executor) to avoid blocking the event loop during subprocess execution. """ cmd_str = self._agent_command cmd = shlex.split(cmd_str) logger.info( f"Invoking agent: {' '.join(cmd)} " f"(timeout={self._agent_timeout}s)" ) def _run() -> str: try: result = subprocess.run( cmd, input=message, capture_output=True, text=True, timeout=self._agent_timeout, ) if result.returncode != 0: stderr = result.stderr.strip() logger.error( f"Agent exited with code {result.returncode}: " f"{stderr}" ) return ( f":warning: Agent error (exit {result.returncode}). " f"Please try again later." ) return result.stdout.strip() except subprocess.TimeoutExpired: logger.error( f"Agent timed out after {self._agent_timeout}s" ) return ( f":hourglass: Agent timed out after " f"{self._agent_timeout}s. Please try again." ) except FileNotFoundError: logger.error(f"Agent command not found: {cmd_str}") return ( f":warning: Agent command not found: `{cmd_str}`. " f"Check configuration." ) except Exception as e: logger.error(f"Agent invocation error: {e}") return ( f":warning: Failed to invoke agent: {e}" ) return await asyncio.get_event_loop().run_in_executor(None, _run) # ------------------------------------------------------------------ # Zulip API helpers # ------------------------------------------------------------------ async def _send_message( self, message_type: str, recipient: str, content: str, topic: Optional[str] = None, stream_id: Optional[int] = None, ) -> Optional[int]: """Send a message to Zulip. Returns the message ID if successful.""" if not self._client: return None def _send() -> Optional[int]: if message_type in ("dm", "private"): # Private message: recipient is email payload = { "type": "private", "to": [recipient], "content": content, } else: # Stream message payload = { "type": "stream", "to": stream_id or recipient, "subject": topic or "general", "content": content, } try: result = self._client.send_message(payload) if result.get("result") == "success": return result.get("id") logger.error( f"Send message failed: {result.get('msg', 'unknown')}" ) return None except Exception as e: logger.error(f"Send message error: {e}") return None return await asyncio.get_event_loop().run_in_executor(None, _send) async def _edit_message( self, message_id: int, content: str ) -> bool: """Edit a previously sent Zulip message (for streaming updates).""" if not self._client: return False def _edit() -> bool: try: result = self._client.update_message( {"message_id": message_id, "content": content} ) if result.get("result") != "success": logger.warning( f"Edit message {message_id}: " f"{result.get('msg', 'unknown')}" ) return result.get("result") == "success" except Exception as e: logger.warning(f"Edit message {message_id} error: {e}") return False return await asyncio.get_event_loop().run_in_executor(None, _edit) async def _send_typing_indicator( self, recipient: str, operation: str ) -> None: """Send a typing indicator (start/stop).""" if not self._client: return def _typing() -> None: try: # The zulip Python library doesn't have a direct typing API # Use the REST endpoint directly via requests import requests server_url = self.config["zulip"]["server_url"] email = self.config["zulip"]["email"] api_key = self.config["zulip"]["api_key"] # For private messages, recipient is the user's email # We need their user_id. Use the API directly. to_data = json.dumps([recipient]) requests.post( f"{server_url}/api/v1/typing", auth=requests.auth.HTTPBasicAuth(email, api_key), data={"to": to_data, "op": operation}, timeout=5, ) except Exception as e: # Non-critical pass await asyncio.get_event_loop().run_in_executor(None, _typing) # ------------------------------------------------------------------ # Synchronous entry point for the Hermes runner # ------------------------------------------------------------------ def run(self) -> None: """Synchronous entry point that starts the async poll loop. Called by the Hermes runner (e.g., from a systemd service). """ asyncio.run(self._run_async()) async def _run_async(self) -> None: """Async entry point.""" try: await self.connect() if self.connected: logger.info("Starting event poll loop...") await self.poll_forever() except asyncio.CancelledError: logger.info("Poll loop cancelled.") except Exception as e: logger.error(f"Fatal error in event loop: {e}") finally: await self.disconnect()