"""M3: channel-directory frame handlers (``channel.*`` + directory queries). Mixin for ``adapter.IrisAdapter``. Each request is answered by broadcasting the matching ``channel.*`` event carrying the request ``id``: the requester's pending request completes on the id, and every other device reconciles its local copy from the same frame (single broadcast serves as event + response). """ import logging from typing import Any from hermes_constants import get_hermes_home from . import protocol from . import purge as purge_bridge from .defaults import DEFAULT_HOME_CHANNEL from .mixin_base import IrisAdapterBase logger = logging.getLogger(__name__) class ChannelFrameHandlers(IrisAdapterBase): """Channel directory management (see module docstring).""" async def on_channel_create(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload name = payload.get("name") if not isinstance(name, str) or not name.strip(): await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "channel.create requires a name", id=frame.id ), ) return kind = payload.get("kind") kind = kind if kind in ("channel", "thread") else "channel" parent_chat_id = payload.get("parent_chat_id") if not isinstance(parent_chat_id, str) or not parent_chat_id.strip(): parent_chat_id = None if kind == "thread" and not parent_chat_id: parent_chat_id = frame.chat_id or self.home_channel try: entry = self._channels.create(name=name, kind=kind, parent_chat_id=parent_chat_id) except ValueError as e: await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id) ) return resp = protocol.channel_created(entry) resp.id = frame.id await self._http_server.fanout(resp) # M5: banner + push mirror (parked in the outbox when offline). await self._broadcast_or_log( entry["chat_id"], protocol.notification( entry["chat_id"], protocol.NOTIF_CHANNEL_CREATED, "Channels", f"New channel: {name}", ), ) async def on_channel_rename(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.rename requires chat_id", id=frame.id ), ) return name = frame.payload.get("name") if not isinstance(name, str) or not name.strip(): await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "channel.rename requires a name", id=frame.id ), ) return try: entry = self._channels.rename(chat_id, name) except ValueError as e: await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id) ) return if entry is None: await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) return resp = protocol.channel_renamed(entry) resp.id = frame.id await self._http_server.fanout(resp) # M5: banner + push mirror (parked in the outbox when offline). await self._broadcast_or_log( chat_id, protocol.notification( chat_id, protocol.NOTIF_CHANNEL_RENAMED, "Channels", f"Renamed to {name}", ), ) async def on_channel_set_default(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.set_default requires chat_id", id=frame.id ), ) return entry = self._channels.set_default(chat_id) if entry is None: await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) return # Reuse the renamed event shape: it carries the full entry (incl. the # new is_default flag) so every device reconciles the default change. resp = protocol.channel_renamed(entry) resp.id = frame.id await self._http_server.fanout(resp) async def on_channel_favorite(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.favorite requires chat_id", id=frame.id ), ) return on = bool(frame.payload.get("on")) entry = self._channels.set_favorite(chat_id, on) if entry is None: await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) return # Reuse the renamed event shape: it carries the full entry (incl. the # new favorite flag) so every device reconciles the change. resp = protocol.channel_renamed(entry) resp.id = frame.id await self._http_server.fanout(resp) async def on_channel_icon(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.icon requires chat_id", id=frame.id ), ) return payload = frame.payload icon = payload.get("icon") icon = icon if isinstance(icon, str) and icon else None color = payload.get("color") color = color if isinstance(color, str) and color else None # Guard against a runaway base64 blob (a channel icon is small). if icon is not None and len(icon) > 512 * 1024: await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "channel icon too large", id=frame.id), ) return entry = self._channels.set_icon(chat_id, icon, color) if entry is None: await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) return resp = protocol.channel_renamed(entry) resp.id = frame.id await self._http_server.fanout(resp) async def on_channel_set_automation(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.set_automation requires chat_id", id=frame.id ), ) return on = bool(frame.payload.get("on")) entry = self._channels.set_automation(chat_id, on) if entry is None: await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, f"cannot set automation on {chat_id} (unknown or default)", id=frame.id, ), ) return # Reuse the renamed event shape: it carries the full entry (incl. the # new automation flag) so every device reconciles the change. resp = protocol.channel_renamed(entry) resp.id = frame.id await self._http_server.fanout(resp) async def on_channel_delete(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.delete requires chat_id", id=frame.id ), ) return entry = self._channels.delete(chat_id) if entry is None: await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, f"cannot delete {chat_id} (unknown or default)", id=frame.id, ), ) return # Complete deletion: wipe the lane's history from the outbox (so # ``history`` / ``sync`` can't resurrect it) and from the hermes # session store (so no search trace survives). A channel delete takes # its threads with it (thread_id=None); a thread delete is scoped to # its parent channel + thread_id. if entry.get("kind") == "thread": lane_chat_id = entry.get("parent_chat_id") or chat_id thread_id = chat_id else: lane_chat_id = chat_id thread_id = None removed_frames = self._outbox.delete_lane(lane_chat_id, thread_id=thread_id) removed_msgs = purge_bridge.delete_lane( get_hermes_home() / "state.db", lane_chat_id, thread_id=thread_id ) logger.info( "iris: channel.delete %s kind=%s outbox_frames=%s session_msgs=%s", chat_id, entry.get("kind"), removed_frames, removed_msgs, ) resp = protocol.channel_deleted(chat_id) resp.id = frame.id await self._http_server.fanout(resp) # M5: banner + push mirror (parked in the outbox when offline). await self._broadcast_or_log( chat_id, protocol.notification( chat_id, protocol.NOTIF_CHANNEL_DELETED, "Channels", f"{entry.get('name') or chat_id} deleted", ), ) async def on_channel_list(self, frame: protocol.Frame, device_id: str) -> None: channels = self._channels.list(include_archived=False) resp = protocol.channel_list(channels) resp.id = frame.id await self._reply(device_id, resp) # ── Chat info ───────────────────────────────────────────────────────── def _channel_name(self, chat_id: str) -> str: """Channel display name (M3: from the channel directory).""" if not chat_id: return "chat" entry = self._channels.get(chat_id) if entry is not None: return entry["name"] if chat_id in (self.home_channel, DEFAULT_HOME_CHANNEL): return self.home_channel_name return chat_id async def get_chat_info(self, chat_id: str) -> dict[str, Any]: """Return ``{name, type, chat_id}`` for a chat (M3: directory-backed).""" entry = self._channels.get(chat_id) kind = entry["kind"] if entry else "channel" return { "name": self._channel_name(chat_id), "type": "dm" if kind == "default" else "channel", "chat_id": chat_id, } def channel_list(self) -> list[dict[str, Any]]: """Channel directory for ``hello.ack`` (M3: full non-archived list).""" return self._channels.list(include_archived=False) # ── M3: core channel-directory hook (cron / send_message name resolution) async def list_channels(self) -> list[dict[str, Any]]: """Expose the directory to the gateway's core channel directory. ``gateway/channel_directory.build_channel_directory`` calls this to populate ``channel_directory.json``, which ``resolve_channel_name`` reads for friendly-name -> chat_id resolution (cron + send_message). Threads are addressed via the explicit ``iris::`` syntax (see ``_parse_target_ref``), so only channels are listed here. """ out: list[dict[str, Any]] = [] for entry in self._channels.list(include_archived=False): if entry["kind"] == "thread": continue out.append( { "id": entry["chat_id"], "name": entry["name"], "type": "dm" if entry["kind"] == "default" else "channel", } ) return out # ── M3: thread handoff (gateway create_handoff_thread) ──────────────── async def create_handoff_thread(self, parent_chat_id: str, name: str) -> str | None: """Mint a named thread under *parent_chat_id* (gateway handoff path). Returns the new ``thread_id`` (``t_``) so the handed-off session is isolated in its own lane, or ``None`` when the parent is unknown. """ parent = self._channels.get(parent_chat_id) if parent is None: # Unknown parent: still mint a thread under it so the handoff has a # lane (the directory row is created lazily on first use). parent_chat_id = parent_chat_id or self.home_channel try: entry = self._channels.create( name=name or "Handoff", kind="thread", parent_chat_id=parent_chat_id ) except Exception: logger.warning("iris: create_handoff_thread failed", exc_info=True) return None await self._broadcast_both(protocol.channel_created(entry)) return entry["chat_id"] # --------------------------------------------------------------------------- # Plugin entry point # ---------------------------------------------------------------------------