"""Frame schemas -- the single source of truth for the wire protocol. Every frame the plugin sends/receives is modelled here as a dataclass + constants. ``docs/protocol/frames.schema.json`` is generated/mirrored from this module, and the Kotlin side mirrors these shapes (see ``docs/04-wire-protocol.md``). Milestone M1: hello/hello.ack, message, message.send, error, ping/pong, typing (typing is pulled forward from M2 so the app gets a live "working…" indicator during the first milestone). Milestone M2: message.start/update/stop, reasoning (on message / message.stop), tool.start/progress/end, commentary. Milestone M3: channel.*, search, sync. Milestone M4: media.*. Milestone M5: notification, fcm.register, read.receipt, status. """ import json from dataclasses import dataclass, field from typing import Any, Dict, List, Optional PROTOCOL_VERSION = 1 # --------------------------------------------------------------------------- # Frame type constants (the ``type`` field of every frame) # --------------------------------------------------------------------------- # Pairing / lifecycle TYPE_HELLO = "hello" TYPE_HELLO_ACK = "hello.ack" TYPE_ERROR = "error" TYPE_PING = "ping" TYPE_PONG = "pong" # Chat TYPE_MESSAGE = "message" TYPE_MESSAGE_SEND = "message.send" TYPE_TYPING = "typing" # Streaming (M2) TYPE_MESSAGE_START = "message.start" TYPE_MESSAGE_UPDATE = "message.update" TYPE_MESSAGE_STOP = "message.stop" # Message deletion (app requests; broadcast to all devices) TYPE_MESSAGE_DELETE = "message.delete" TYPE_MESSAGE_DELETED = "message.deleted" # Tool activity (M2) TYPE_TOOL_START = "tool.start" TYPE_TOOL_PROGRESS = "tool.progress" TYPE_TOOL_END = "tool.end" # Intermediate assistant beat (M2) TYPE_COMMENTARY = "commentary" # Channels / threads (M3) TYPE_CHANNEL_CREATE = "channel.create" TYPE_CHANNEL_RENAME = "channel.rename" TYPE_CHANNEL_SET_DEFAULT = "channel.set_default" TYPE_CHANNEL_FAVORITE = "channel.favorite" TYPE_CHANNEL_ICON = "channel.icon" TYPE_CHANNEL_DELETE = "channel.delete" TYPE_CHANNEL_CREATED = "channel.created" TYPE_CHANNEL_RENAMED = "channel.renamed" TYPE_CHANNEL_DELETED = "channel.deleted" TYPE_CHANNEL_LIST = "channel.list" # Search (M3) TYPE_SEARCH = "search" TYPE_SEARCH_RESULTS = "search.results" # Reconnect catch-up (M3 outbox; extended by M5 push) TYPE_SYNC = "sync" TYPE_SYNC_DONE = "sync.done" # Full message history (initial channel open / scroll-up pagination) TYPE_HISTORY = "history" # Media (M4) TYPE_MEDIA_UPLOAD_START = "media.upload.start" TYPE_MEDIA_UPLOAD_END = "media.upload.end" TYPE_MEDIA_UPLOAD_ACK = "media.upload.ack" TYPE_MEDIA_OFFER = "media.offer" TYPE_MEDIA_PULL = "media.pull" TYPE_MEDIA_PULL_END = "media.pull.end" # Push / notifications (M5) TYPE_NOTIFICATION = "notification" TYPE_FCM_REGISTER = "fcm.register" TYPE_READ_RECEIPT = "read.receipt" # Gateway health (M5) TYPE_STATUS = "status" # --------------------------------------------------------------------------- # Error codes (``error`` frame payload.code) # --------------------------------------------------------------------------- ERR_AUTH = "auth" ERR_NOT_FOUND = "not_found" ERR_RATE_LIMITED = "rate_limited" ERR_MEDIA_TOO_LARGE = "media_too_large" ERR_UNSUPPORTED = "unsupported" ERR_INTERNAL = "internal" # --------------------------------------------------------------------------- # Message roles (``message`` frame payload.role) # --------------------------------------------------------------------------- ROLE_USER = "user" ROLE_ASSISTANT = "assistant" ROLE_SYSTEM = "system" ROLE_CRON = "cron" # --------------------------------------------------------------------------- # Notification kinds (``notification`` frame payload.kind, M5) # --------------------------------------------------------------------------- NOTIF_CHANNEL_CREATED = "channel_created" NOTIF_CHANNEL_RENAMED = "channel_renamed" NOTIF_CHANNEL_DELETED = "channel_deleted" NOTIF_CRON = "cron" NOTIF_APPROVAL = "approval" NOTIF_CLARIFY = "clarify" NOTIF_GENERIC = "generic" # Kinds that push even when a device is live (the app may be backgrounded; # it decides whether to also show an in-app banner). HIGH_PRIORITY_NOTIF_KINDS = frozenset({NOTIF_APPROVAL, NOTIF_CLARIFY, NOTIF_CRON}) # --------------------------------------------------------------------------- # Gateway health states (``status`` frame payload.state) # --------------------------------------------------------------------------- STATUS_ONLINE = "online" STATUS_RESTARTING = "restarting" STATUS_DEGRADED = "degraded" # --------------------------------------------------------------------------- # Envelope # --------------------------------------------------------------------------- @dataclass class Frame: """One wire frame. ``v`` is always serialised; ``id``/``chat_id``/``thread_id`` are omitted when ``None`` (events carry no ``id``; chat-scoped frames carry ``chat_id``/``thread_id`` at the top level for convenience). """ type: str payload: Dict[str, Any] = field(default_factory=dict) id: Optional[int] = None chat_id: Optional[str] = None thread_id: Optional[str] = None v: int = PROTOCOL_VERSION def to_dict(self) -> Dict[str, Any]: d: Dict[str, Any] = {"v": self.v, "type": self.type} if self.id is not None: d["id"] = self.id if self.chat_id is not None: d["chat_id"] = self.chat_id if self.thread_id is not None: d["thread_id"] = self.thread_id d["payload"] = self.payload return d def to_json(self) -> str: return json.dumps(self.to_dict(), separators=(",", ":"), ensure_ascii=False) @classmethod def from_json(cls, raw: "str | bytes") -> Optional["Frame"]: """Parse a text frame. Returns ``None`` for anything not a valid frame (bad JSON, non-object, missing/invalid ``type``) so callers can ignore malformed input (forward-compat).""" try: data = json.loads(raw) except (json.JSONDecodeError, TypeError, UnicodeDecodeError, ValueError): return None if not isinstance(data, dict): return None ftype = data.get("type") if not isinstance(ftype, str) or not ftype: return None payload = data.get("payload") if not isinstance(payload, dict): payload = {} fid = data.get("id") if not isinstance(fid, int) or isinstance(fid, bool): fid = None chat_id = data.get("chat_id") if not isinstance(chat_id, str): chat_id = None thread_id = data.get("thread_id") if not isinstance(thread_id, str): thread_id = None return cls(type=ftype, payload=payload, id=fid, chat_id=chat_id, thread_id=thread_id) # --------------------------------------------------------------------------- # Frame constructors (server -> app) # --------------------------------------------------------------------------- def hello_ack( server_caps: Dict[str, Any], sync_cursor: int = 0, channels: Optional[list] = None, ) -> Frame: return Frame( type=TYPE_HELLO_ACK, payload={ "server_caps": server_caps, "sync_cursor": sync_cursor, "channels": channels or [], }, ) def message( chat_id: str, message_id: str, role: str, text: str, *, thread_id: Optional[str] = None, reasoning: Optional[str] = None, media: Optional[list] = None, reply_to: Optional[str] = None, model: Optional[str] = None, tokens: Optional[int] = None, ts: Optional[int] = None, ) -> Frame: payload: Dict[str, Any] = { "message_id": message_id, "role": role, "text": text, } if reasoning: payload["reasoning"] = reasoning if media: payload["media"] = media if reply_to: payload["reply_to"] = reply_to if model: payload["model"] = model if tokens is not None: payload["tokens"] = tokens if ts is not None: payload["ts"] = ts return Frame(type=TYPE_MESSAGE, chat_id=chat_id, thread_id=thread_id, payload=payload) def typing(chat_id: str, on: bool = True, *, thread_id: Optional[str] = None) -> Frame: return Frame( type=TYPE_TYPING, chat_id=chat_id, thread_id=thread_id, payload={"on": on}, ) # --------------------------------------------------------------------------- # Streaming frames (M2) # --------------------------------------------------------------------------- def message_start( chat_id: str, message_id: str, role: str = ROLE_ASSISTANT, *, thread_id: Optional[str] = None, ) -> Frame: """Open a streaming bubble.""" return Frame( type=TYPE_MESSAGE_START, chat_id=chat_id, thread_id=thread_id, payload={"message_id": message_id, "role": role}, ) def message_update( chat_id: str, message_id: str, text: str, *, thread_id: Optional[str] = None, ) -> Frame: """Replace the live bubble text (full snapshot).""" return Frame( type=TYPE_MESSAGE_UPDATE, chat_id=chat_id, thread_id=thread_id, payload={"message_id": message_id, "text": text}, ) def message_stop( chat_id: str, message_id: str, final_text: str, *, thread_id: Optional[str] = None, reasoning: Optional[str] = None, model: Optional[str] = None, tokens: Optional[int] = None, ts: Optional[int] = None, ) -> Frame: """Finalize a streaming bubble.""" payload: Dict[str, Any] = { "message_id": message_id, "final_text": final_text, } if reasoning: payload["reasoning"] = reasoning if model: payload["model"] = model if tokens is not None: payload["tokens"] = tokens if ts is not None: payload["ts"] = ts return Frame( type=TYPE_MESSAGE_STOP, chat_id=chat_id, thread_id=thread_id, payload=payload, ) # --------------------------------------------------------------------------- # Tool activity frames (M2) # --------------------------------------------------------------------------- def tool_start( chat_id: str, index: int, name: str, *, thread_id: Optional[str] = None, preview: Optional[str] = None, args: Optional[Dict[str, Any]] = None, ) -> Frame: payload: Dict[str, Any] = {"index": index, "name": name} if preview: payload["preview"] = preview if args: payload["args"] = args return Frame( type=TYPE_TOOL_START, chat_id=chat_id, thread_id=thread_id, payload=payload, ) def tool_progress( chat_id: str, index: int, name: str, *, thread_id: Optional[str] = None, note: Optional[str] = None, ) -> Frame: payload: Dict[str, Any] = {"index": index, "name": name} if note: payload["note"] = note return Frame( type=TYPE_TOOL_PROGRESS, chat_id=chat_id, thread_id=thread_id, payload=payload, ) def tool_end( chat_id: str, index: int, name: str, *, thread_id: Optional[str] = None, ok: bool = True, duration: Optional[float] = None, output_preview: Optional[str] = None, ) -> Frame: payload: Dict[str, Any] = {"index": index, "name": name, "ok": ok} if duration is not None: payload["duration"] = duration if output_preview: payload["output_preview"] = output_preview return Frame( type=TYPE_TOOL_END, chat_id=chat_id, thread_id=thread_id, payload=payload, ) # --------------------------------------------------------------------------- # Commentary frame (M2) # --------------------------------------------------------------------------- def commentary( chat_id: str, message_id: str, text: str, *, thread_id: Optional[str] = None, ) -> Frame: """An intermediate assistant beat (between tool iterations).""" return Frame( type=TYPE_COMMENTARY, chat_id=chat_id, thread_id=thread_id, payload={"message_id": message_id, "text": text}, ) # --------------------------------------------------------------------------- # Channel directory frames (M3) # --------------------------------------------------------------------------- def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]: """Project a directory entry onto the wire shape.""" payload: Dict[str, Any] = { "chat_id": entry.get("chat_id"), "name": entry.get("name"), "kind": entry.get("kind", "channel"), } if entry.get("parent_chat_id") is not None: payload["parent_chat_id"] = entry["parent_chat_id"] if entry.get("is_default"): payload["is_default"] = True if entry.get("archived"): payload["archived"] = True if entry.get("favorite"): payload["favorite"] = True if entry.get("icon"): payload["icon"] = entry["icon"] if entry.get("color"): payload["color"] = entry["color"] return payload def channel_created(entry: Dict[str, Any], auto: bool = False) -> Frame: """Broadcast: a channel/thread was created. ``auto=True`` marks a thread the gateway minted itself for an incoming message (auto-threading, docs/06 §6.3): the app jumps into it and the name is an instant derived title, upgraded by the LLM via a follow-up ``channel.renamed``. """ payload = _channel_payload(entry) if auto: payload["auto"] = True return Frame(type=TYPE_CHANNEL_CREATED, payload=payload) def channel_renamed(entry: Dict[str, Any]) -> Frame: """Broadcast: a channel/thread was renamed.""" return Frame(type=TYPE_CHANNEL_RENAMED, payload=_channel_payload(entry)) def channel_deleted(chat_id: str) -> Frame: """Broadcast: a channel was archived (soft-deleted).""" return Frame(type=TYPE_CHANNEL_DELETED, payload={"chat_id": chat_id}) def channel_list(channels: List[Dict[str, Any]]) -> Frame: """Full directory (response to a ``channel.list`` request).""" return Frame( type=TYPE_CHANNEL_LIST, payload={"channels": [_channel_payload(c) for c in channels]}, ) # --------------------------------------------------------------------------- # Search frames (M3) # --------------------------------------------------------------------------- def search_results( query: str, scope: str, hits: List[Dict[str, Any]], *, id: Optional[int] = None, ) -> Frame: """Response to a ``search`` request. Each hit: ``{message_id, chat_id, thread_id, role, snippet, ts}``. """ return Frame( type=TYPE_SEARCH_RESULTS, id=id, payload={"query": query, "scope": scope, "hits": hits}, ) # --------------------------------------------------------------------------- # Sync frames (M3 outbox) # --------------------------------------------------------------------------- def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame: """Terminal frame of a ``sync`` replay: the new cursor to persist.""" return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor}) # --------------------------------------------------------------------------- # History frame (full message history for a chat/thread) # --------------------------------------------------------------------------- def history( chat_id: str, messages: List[Dict[str, Any]], has_more: bool, *, thread_id: Optional[str] = None, oldest_message_id: Optional[str] = None, id: Optional[int] = None, ) -> Frame: """Response to a ``history`` request: a page of final messages for a chat/thread, ordered oldest → newest. ``has_more`` signals older pages exist; ``oldest_message_id`` is the ``before_message_id`` for the next (older) page.""" payload: Dict[str, Any] = { "messages": messages, "has_more": has_more, } if oldest_message_id: payload["oldest_message_id"] = oldest_message_id return Frame( type=TYPE_HISTORY, id=id, chat_id=chat_id, thread_id=thread_id, payload=payload, ) # --------------------------------------------------------------------------- # Message deletion frames # --------------------------------------------------------------------------- def message_deleted( chat_id: str, message_ids: List[str], *, thread_id: Optional[str] = None, id: Optional[int] = None, ) -> Frame: """Broadcast: the given message(s) were deleted from a chat/thread. Carried by the response to a ``message.delete`` request (``id`` set) and broadcast to every device so all of them drop the message(s) from their cache. Also outboxed, so a device that was offline learns of the deletion on its next ``sync``. """ return Frame( type=TYPE_MESSAGE_DELETED, id=id, chat_id=chat_id, thread_id=thread_id, payload={"message_ids": list(message_ids)}, ) # --------------------------------------------------------------------------- # Push / notification frames (M5) # --------------------------------------------------------------------------- def notification( chat_id: str, kind: str, title: str, body: str, *, thread_id: Optional[str] = None, ts: Optional[int] = None, ) -> Frame: """Event: a transient in-app banner (and a push mirror when the device is offline). ``kind`` is one of the ``NOTIF_*`` constants.""" payload: Dict[str, Any] = {"kind": kind, "title": title, "body": body} if ts is not None: payload["ts"] = ts return Frame( type=TYPE_NOTIFICATION, chat_id=chat_id, thread_id=thread_id, payload=payload ) def fcm_register( fcm_token: Optional[str] = None, ntfy_topic: Optional[str] = None ) -> Frame: """Request: update the device's push tokens (FCM rotation / ntfy topic).""" payload: Dict[str, Any] = {} if fcm_token: payload["fcm_token"] = fcm_token if ntfy_topic: payload["ntfy_topic"] = ntfy_topic return Frame(type=TYPE_FCM_REGISTER, payload=payload) def read_receipt(chat_id: str, message_id: str) -> Frame: """Ack to the originating device: the agent received and started processing the user's message (the app shows ✓✓ on the user bubble).""" return Frame( type=TYPE_READ_RECEIPT, chat_id=chat_id, payload={"message_id": message_id}, ) # --------------------------------------------------------------------------- # Gateway health frame (M5) # --------------------------------------------------------------------------- def status(state: str) -> Frame: """Gateway health state (``state`` is one of the ``STATUS_*`` constants).""" return Frame(type=TYPE_STATUS, payload={"state": state}) # --------------------------------------------------------------------------- # Media frames (M4) # --------------------------------------------------------------------------- def media_offer( media_id: str, kind: str, mime: str, size: int, filename: str, *, chat_id: Optional[str] = None, thread_id: Optional[str] = None, message_id: Optional[str] = None, ) -> Frame: """Event: the agent produced media the app can fetch via ``media.pull``. ``message_id`` (optional) associates the offer with the assistant message it belongs to (the app falls back to the lane's last assistant message). """ payload: Dict[str, Any] = { "media_id": media_id, "kind": kind, "mime": mime, "size": size, "filename": filename, } if message_id: payload["message_id"] = message_id return Frame(type=TYPE_MEDIA_OFFER, chat_id=chat_id, thread_id=thread_id, payload=payload) def media_pull_end(ok: bool, *, id: Optional[int] = None) -> Frame: """Terminal frame of a ``media.pull`` binary stream.""" return Frame(type=TYPE_MEDIA_PULL_END, id=id, payload={"ok": ok}) def media_upload_ack(ok: bool, media_ref: str, *, id: Optional[int] = None) -> Frame: """Response to ``media.upload.end``: the ref is cached and may be used in a ``message.send`` ``media_refs``. Failures use ``error`` frames instead.""" return Frame( type=TYPE_MEDIA_UPLOAD_ACK, id=id, payload={"ok": ok, "media_ref": media_ref}, ) def error(code: str, message: str, *, id: Optional[int] = None) -> Frame: return Frame(type=TYPE_ERROR, id=id, payload={"code": code, "message": message}) def pong(ts: Optional[int] = None) -> Frame: payload: Dict[str, Any] = {} if ts is not None: payload["ts"] = ts return Frame(type=TYPE_PONG, payload=payload)