""" Android Platform Adapter for Hermes Agent (Iris x Hermes). A plugin-based gateway adapter that runs a WebSocket server *inside* the ``hermes gateway`` process. The native Android / Desktop app connects to it with a pairing token and talks to the agent over a single WS transport (chat, streaming, tools, media, pairing, push-token). Zero new Python dependencies: ``websockets`` and ``httpx`` are hermes core deps. Zero hermes-core changes. Milestone M1: the gateway core loop (text round-trip). The WS server binds and authenticates devices (``hello`` with constant-time token check), the adapter emits ``message`` frames from ``send()`` and turns inbound ``message.send`` frames into ``MessageEvent``s for ``handle_message()``. Milestone M2: agent transparency. ``send()``/``edit_message()`` are mapped to ``message.start``/``message.update``/``message.stop`` (streaming), tool progress is classified into structured ``tool.start``/``tool.end`` frames, interim commentary becomes ``commentary`` frames, and the code-style reasoning prefix is split into a ``reasoning`` field. Outbox and search land in M3; media, push, and desktop land in later milestones (see ``docs/14-milestones.md``). Milestone M4: media. Inbound ``media.upload`` (chunked binary frames) is reassembled in a temp file, verified (size + sha256), re-sniffed, and cached via hermes ``cache_*_from_bytes``; the resulting refs attach to the next ``message.send`` as ``MessageEvent.media_urls``. Outbound ``send_*`` calls register the (delivery-validated) file in the media registry and emit ``media.offer``; ``media.pull`` streams the file back as chunked binary frames, re-checking ``validate_media_delivery_path`` at pull time. Milestone M5: push + offline. Frames with no live subscriber are parked in the outbox (M3) AND wake the device via the push backend (``push.py``: FCM HTTP v1 primary, ntfy fallback, selected by ``ANDROID_PUSH_BACKEND``). ``notification`` frames render in-app banners and mirror to push (channel events, cron deliveries, approvals, clarifies); high-priority kinds push even when a device is live. ``fcm.register`` rotates push tokens (registry + live connection). The outbox enforces a row cap with a throttled prune notice. Configuration in config.yaml:: gateway: platforms: android: enabled: true extra: host: 127.0.0.1 port: 8790 home_channel: android:default push_backend: fcm outbox_retention_hours: 72 max_upload_bytes: 104857600 Or via environment variables (overrides config.yaml; secrets live in .env): ANDROID_TOKEN, ANDROID_WS_HOST, ANDROID_WS_PORT, ANDROID_HOME_CHANNEL, ANDROID_PUSH_BACKEND, ANDROID_FCM_SERVICE_ACCOUNT, NTFY_TOPIC, ... """ import asyncio import json import logging import os import re import threading import time import uuid from collections import deque from dataclasses import dataclass, field from typing import Any, Dict, List, Optional, Tuple from agent.secret_scope import UnscopedSecretError as _UnscopedSecretError from agent.secret_scope import get_secret as _scoped_get_secret def _get_scoped_secret(name, default=None): """Scope-aware credential read with the default-profile startup fallback. Secondary profiles construct their adapters under a profile secret scope -- the scope is authoritative and a scoped miss returns ``default`` (no cross-profile borrow from ``os.environ``, which may hold another profile's value). The DEFAULT profile's adapter constructs and sends *unscoped* under multiplexing, where a bare ``get_secret`` would raise ``UnscopedSecretError`` and crash this path; there ``os.environ`` is that profile's own value, so fall back to it. Same pattern as the IRC ``IRC_SERVER_PASSWORD`` read (``plugins/platforms/irc/adapter.py``). """ try: val = _scoped_get_secret(name, default) except _UnscopedSecretError: val = os.getenv(name) return val if val is not None else default logger = logging.getLogger(__name__) # --------------------------------------------------------------------------- # Lazy import: BasePlatformAdapter and friends live in the main repo. # We import at module level (as the bundled plugins do) but guard the heavy # gateway imports so the plugin can be discovered before the gateway is fully # initialised. # --------------------------------------------------------------------------- from gateway.platforms.base import ( # noqa: E402 BasePlatformAdapter, SendResult, MessageEvent, MessageType, validate_media_delivery_path, ) from gateway.config import Platform # noqa: E402 from hermes_constants import get_hermes_home # noqa: E402 from . import media as media_bridge # noqa: E402 from . import protocol # noqa: E402 from . import search as search_bridge # noqa: E402 from .channels import get_directory # noqa: E402 from .outbox import Outbox # noqa: E402 from .push import NtfyBackend, PushBackend, build_push_backend # noqa: E402 from .pairing import ( # noqa: E402 DeviceRegistry, generate_token, pairing_url, qr_payload, ) from .ws_server import WsServer # noqa: E402 # --------------------------------------------------------------------------- # M2 — reasoning capture (streaming) # # The gateway streams only ``content`` to the platform and suppresses the # final send (which would carry the prepended reasoning), so the model's # separate ``reasoning_content`` is otherwise lost in the streaming case. # hermes exposes a plugin ``on_stream_delta`` hook that fires reasoning # deltas with ``kind="reasoning"`` (gated by ``plugins.stream_reasoning_deltas``). # We accumulate those deltas here and attach the result to the turn's # ``message.stop`` frame. Single-chat for now (android:default), so a # module-level buffer suffices; it is reset at each turn start. # --------------------------------------------------------------------------- _reasoning_parts: List[str] = [] _reasoning_lock = threading.Lock() # Barrier: set by the hook worker once it has processed the first content # delta (kind="text"). The worker drains a FIFO queue and reasoning deltas are # enqueued before content deltas, so at that point every reasoning delta has # already been appended -- a reliable "reasoning flushed" signal that avoids # racing message.stop against the async hook thread. _reasoning_flushed = threading.Event() def _on_stream_delta(**kwargs: Any) -> None: """Plugin hook: capture reasoning deltas (kind="reasoning").""" kind = kwargs.get("kind") if kind == "reasoning": delta = kwargs.get("delta") or "" if delta: with _reasoning_lock: _reasoning_parts.append(delta) elif kind == "text": _reasoning_flushed.set() async def _wait_for_reasoning_flushed(timeout: float = 0.3) -> None: """Wait (without blocking the event loop) until the hook worker has processed all reasoning deltas, or *timeout* seconds elapse.""" loop = asyncio.get_running_loop() deadline = loop.time() + timeout while loop.time() < deadline: if _reasoning_flushed.is_set(): return await asyncio.sleep(0.01) def _take_reasoning() -> str: """Drain and return the accumulated reasoning (empty string if none).""" with _reasoning_lock: parts = _reasoning_parts[:] _reasoning_parts.clear() _reasoning_flushed.clear() return "".join(parts).strip() def _reset_reasoning() -> None: with _reasoning_lock: _reasoning_parts.clear() _reasoning_flushed.clear() # --------------------------------------------------------------------------- # M2 — tool-result capture (post_tool_call hook) # # The gateway renders tool *progress* lines to the platform but never streams # the tool *output* (it is the agent's concern, persisted to history, not # presentation). To let the app show the full call + result on demand # (Settings → Tool detail), we capture each completed tool call via the # ``post_tool_call`` hook and attach it to the ``tool.end`` frame. # # Global FIFO (like the reasoning buffer): a personal android gateway serves # one active turn at a time, and records are matched to the open tool by name # in completion order. Bounded so a runaway turn can't grow it without limit. # --------------------------------------------------------------------------- _tool_results: "deque[Dict[str, Any]]" = deque() _tool_results_lock = threading.Lock() _MAX_TOOL_RESULTS = 200 _MAX_OUTPUT_PREVIEW = 8000 def _on_post_tool_call(**kwargs: Any) -> None: """Plugin hook: capture a completed tool call's result + timing.""" result = kwargs.get("result") record = { "tool_name": kwargs.get("tool_name") or "", "result": (str(result) if result is not None else "")[:_MAX_OUTPUT_PREVIEW], "duration_ms": kwargs.get("duration_ms") or 0, "status": kwargs.get("status") or "ok", } with _tool_results_lock: _tool_results.append(record) while len(_tool_results) > _MAX_TOOL_RESULTS: _tool_results.popleft() def _take_tool_result(tool_name: str) -> Optional[Dict[str, Any]]: """Pop the first completed record matching *tool_name* (FIFO), else None.""" if not tool_name: return None with _tool_results_lock: for i, rec in enumerate(_tool_results): if rec["tool_name"] == tool_name: del _tool_results[i] return rec return None def _reset_tool_results() -> None: """Clear captured records (turn boundary — drop anything unconsumed).""" with _tool_results_lock: _tool_results.clear() def _tool_end_fields(tool_name: str) -> Dict[str, Any]: """Build the ``tool.end`` enrichment (ok/duration/output_preview) from the captured hook record for *tool_name*; empty dict when none is available (e.g. tool_progress off, or the call came from another session).""" rec = _take_tool_result(tool_name) if rec is None: return {} fields: Dict[str, Any] = { "ok": rec["status"] == "ok", "output_preview": rec["result"] or None, } if rec["duration_ms"]: fields["duration"] = round(rec["duration_ms"] / 1000.0, 3) return fields # --------------------------------------------------------------------------- # Defaults # --------------------------------------------------------------------------- DEFAULT_HOST = "127.0.0.1" DEFAULT_PORT = 8790 DEFAULT_HOME_CHANNEL = "android:default" DEFAULT_HOME_CHANNEL_NAME = "Default" DEFAULT_PUSH_BACKEND = "fcm" DEFAULT_OUTBOX_RETENTION_HOURS = 72 DEFAULT_MAX_UPLOAD_BYTES = 100 * 1024 * 1024 # 100 MB def _truthy(value: Optional[str]) -> bool: return (value or "").strip().lower() in {"1", "true", "yes", "on"} # --------------------------------------------------------------------------- # M2 — turn state + outbound classification # # The main gateway delivers through the legacy callback path: the stream # consumer calls ``send()`` (first bubble of a segment) and ``edit_message()`` # (updates), tool progress flows through ``send()``/``edit_message()`` of an # accumulated line buffer, and interim commentary arrives as a plain # ``send()``. We classify each outbound call into a structured frame using a # per-chat turn state machine + content markers: # # * ``metadata["expect_edits"] is True`` -> streaming segment start # * ``metadata["notify"] is True`` -> final message (or fallback final) # * tool-progress line format -> tool.start / tool.end # * anything else -> commentary # # Verified empirically against the live gateway with ``tests/ws_probe.py``. # --------------------------------------------------------------------------- # Streaming cursor the gateway appends to in-progress edits (" ▉"). Stripped # before we forward text to the app (the app renders its own live indicator). _STREAMING_CURSOR = " ▉" # Code-style reasoning prefix (gateway/run.py, reasoning_style="code"): # "💭 **Reasoning:**\n```\n\n```\n\n" _REASONING_PREFIX = "💭 **Reasoning:**\n```\n" _REASONING_CLOSE = "\n```\n\n" # A gateway tool-progress line begins with a (non-ASCII) tool emoji. _TOOL_LINE_RE = re.compile(r"^(\S+)\s+(.+)$") _TOOL_NAME_PREVIEW_RE = re.compile(r'^(\S+):\s*"(.*)"\s*$') _TOOL_NAME_BARE_RE = re.compile(r"^(\S+)\.\.\.\s*$") _TOOL_NAME_ARGS_RE = re.compile(r"^(\S+)\(([^)]*)\)\s*$") # Terminal code block: "💻 terminal\n```\n\n```" _TOOL_CODEBLOCK_HEAD_RE = re.compile(r"^(\S+)\s+(\S+)\s*$") # Reverse map of the gateway's friendly tool verbs (agent/display.py # _TOOL_VERBS) so a verb-form line ("🔍 Searching the web for …") can be # recovered to a structured (tool_name, preview). Longest-first matching is # done at parse time. Verbs shared by several tools map to the most common. _VERB_TO_TOOL: Dict[str, str] = { "Searching the web": "web_search", "Searching files": "search_files", "Searching past sessions": "session_search", "Running code": "execute_code", "Running": "terminal", "Reading skill": "skill_view", "Reading": "read_file", "Writing": "write_file", "Editing": "patch", "Browsing": "browser_navigate", "Clicking": "browser_click", "Typing": "browser_type", "Generating image": "image_generate", "Generating video": "video_generate", "Generating speech": "text_to_speech", "Looking at the image": "vision_analyze", "Listing skills": "skills_list", "Updating skill": "skill_manage", "Updating memory": "memory", "Updating tasks": "todo", "Delegating": "delegate_task", "Scheduling": "cronjob", "Asking": "clarify", } # Verbs that take a " for " connector before the preview. _VERB_FOR_CONNECTOR = {"web_search", "search_files"} def _mint_message_id() -> str: return f"m_{uuid.uuid4().hex[:16]}" def _thread_id_from_metadata(metadata: Optional[Dict[str, Any]]) -> Optional[str]: if not metadata: return None tid = metadata.get("thread_id") if isinstance(tid, str) and tid: return tid return None def _derive_thread_name(text: str) -> str: """Instant auto-thread name from the user's opening message (no model). Reuses hermes' session-title derivation (``agent/title_generator.py``): a deterministic slice of the user's own words, so the thread is named the moment it is created. The LLM upgrade (``_schedule_thread_title_upgrade``) replaces it moments later — the same two-stage titling hermes uses for sessions (derived < llm < user). """ try: from agent.title_generator import derive_title title = derive_title(text) except Exception: logger.debug("Thread name derivation failed", exc_info=True) title = None return (title or "").strip() or "New thread" def _strip_streaming_cursor(text: str) -> str: if text and text.endswith(_STREAMING_CURSOR): return text[: -len(_STREAMING_CURSOR)] return text def _push_preview(text: Any, limit: int = 120) -> str: """Short single-line preview for push bodies (lock-screen privacy: no secrets, no full bodies -- full content arrives via ``sync``).""" s = " ".join(str(text or "").split()) if len(s) > limit: s = s[: limit - 1] + "…" return s # Cron delivery wrap (cron/scheduler.py ``_deliver_result``, # cron.wrap_response: true): # "Cronjob Response: \n(job_id: )\n-------------\n\n\n\n # To stop or manage this job, send me a new message (e.g. ...)." _CRON_WRAP_RE = re.compile( r"^Cronjob Response: (.+?)\n\(job_id: [^)]*\)\n-+\n\n" ) _CRON_FOOTER = "\n\nTo stop or manage this job" def _cron_brief(content: str, job_id: str) -> Tuple[str, str]: """Parse a cron delivery into ``(job_name, inner_text)``. Falls back to ``(job_id, content)`` when the wrap is disabled (``cron.wrap_response: false``) or unrecognised. """ m = _CRON_WRAP_RE.match(content or "") if not m: return str(job_id or "cron"), (content or "").strip() name = m.group(1).strip() body = content[m.end():] idx = body.rfind(_CRON_FOOTER) if idx != -1: body = body[:idx] return name, body.strip() def _split_reasoning(text: str) -> Tuple[Optional[str], str]: """Split a code-style reasoning prefix off the front of *text*. Returns ``(reasoning, body)``; ``reasoning`` is ``None`` when no prefix is present (reasoning off / no reasoning / non-code style). Best-effort parse of a stable, gateway-owned format: on any mismatch the fallback is ``(None, full text)`` so the answer still renders. """ if not text or not text.startswith(_REASONING_PREFIX): return None, text close_idx = text.find(_REASONING_CLOSE, len(_REASONING_PREFIX)) if close_idx == -1: return None, text reasoning = text[len(_REASONING_PREFIX):close_idx] body = text[close_idx + len(_REASONING_CLOSE):] return reasoning, body def _parse_tool_line(line: str) -> Optional[Tuple[str, Optional[str]]]: """Parse a single gateway tool-progress line into ``(name, preview)``. Returns ``None`` when the line is not a tool line. The gateway formats tool lines as `` : ""``, `` ...``, `` (keys)``, or a friendly verb phrase (`` …``). The verb form is lossy (no tool name), so we surface the verb as the name. """ line = line.strip() if not line: return None m = _TOOL_LINE_RE.match(line) if not m: return None emoji, rest = m.group(1), m.group(2) if emoji.isascii(): return None # a tool line always leads with a non-ASCII emoji mp = _TOOL_NAME_PREVIEW_RE.match(rest) if mp: return mp.group(1), mp.group(2) mb = _TOOL_NAME_BARE_RE.match(rest) if mb: return mb.group(1), None ma = _TOOL_NAME_ARGS_RE.match(rest) if ma: return ma.group(1), None # Friendly verb phrase: reverse-map to (tool_name, preview). verb_parsed = _parse_verb_phrase(rest) if verb_parsed is not None: return verb_parsed # Unrecognised: use the phrase as the label. return rest, None def _parse_verb_phrase(phrase: str) -> Optional[Tuple[str, Optional[str]]]: """Reverse-map a friendly verb phrase to ``(tool_name, preview)``. Matches the longest verb first so "Running code" wins over "Running". Returns ``None`` when no known verb leads the phrase. """ for verb in sorted(_VERB_TO_TOOL, key=len, reverse=True): tool = _VERB_TO_TOOL[verb] if phrase == verb: return tool, None if tool in _VERB_FOR_CONNECTOR and phrase.startswith(verb + " for "): return tool, phrase[len(verb) + len(" for "):].strip() or None if phrase.startswith(verb + " "): return tool, phrase[len(verb) + 1:].strip() or None return None def _extract_code_block(content: str) -> Optional[str]: """Return the first fenced code block's body in *content*, else ``None``. Used to recover the terminal command from a tool-progress code block (`` terminal`` head line + fenced command). """ m = re.search(r"```[^\n]*\n(.*?)\n```", content, re.DOTALL) if m: return m.group(1).strip() or None return None def _extract_verbose_args(line: str, content: str) -> Optional[Dict[str, Any]]: """Recover the full args dict from a verbose tool line, else ``None``. In verbose mode the gateway renders `` (keys)`` on one line and the full args JSON on the line that follows. When *line* is such a header, return the parsed JSON object from the following line. """ parts = line.strip().split(None, 1) if len(parts) < 2 or not _TOOL_NAME_ARGS_RE.match(parts[1]): return None lines = content.splitlines() for i, ln in enumerate(lines): if ln.strip() != line.strip(): continue for follow in lines[i + 1:]: follow = follow.strip() if not follow: continue if follow.startswith("{"): try: obj = json.loads(follow) return obj if isinstance(obj, dict) else None except Exception: return None return None # next non-empty line is not the args JSON return None def _short_preview_from_args(args: Dict[str, Any], cap: int = 60) -> Optional[str]: """Derive a short one-line preview from a verbose args dict. The verbose line carries no explicit preview, so the Truncated display would otherwise lose its one-liner. Use the first non-empty string value (whitespace-collapsed, capped) as a stand-in. """ if not isinstance(args, dict): return None for value in args.values(): if isinstance(value, str) and value.strip(): s = " ".join(value.split()) return s[: cap - 1] + "…" if len(s) > cap else s return None def _is_tool_progress(content: str) -> bool: """Heuristic: does *content* look like gateway tool-progress line(s)? Tool progress is delivered as one or more lines, each led by a tool emoji (or a terminal code block). Commentary is free-form prose. We classify on the first non-empty line; subsequent lines of the same bubble are tracked by message id, not re-classified. """ if not content: return False lines = [ln for ln in content.splitlines() if ln.strip()] if not lines: return False first = lines[0].strip() # Terminal code block: " terminal" then a fenced command. if len(lines) > 1 and lines[1].strip().startswith("```"): return _TOOL_CODEBLOCK_HEAD_RE.match(first) is not None return _parse_tool_line(first) is not None def _is_gateway_lifecycle_notice(content: str) -> bool: """True for hermes gateway lifecycle notices (restart / shutdown / online). These are system notices, not tool progress. Their leading ⚠️/♻️ emoji would otherwise trip the tool-line heuristic and render them as a never-completing tool card (an endless spinner, since no ``tool.end`` ever arrives for a notice that is not a real tool). """ if not content: return False c = content.strip() return any( marker in c for marker in ("Gateway restarting", "Gateway shutting down", "Gateway online") ) @dataclass class _TurnState: """Per-chat turn state for outbound frame classification (M2).""" active: bool = False # message_id of the currently streaming segment (message.start open). stream_id: Optional[str] = None # message_id of the current tool-progress bubble (editable line buffer). tool_msg_id: Optional[str] = None # Monotonic per-turn tool counter (start -> end correlation). tool_index: int = 0 # Index of the most recently started tool (awaiting tool.end). open_tool_index: Optional[int] = None # Name of the most recently started tool (matches the post_tool_call # record when the tool completes, so tool.end can carry its output). open_tool_name: Optional[str] = None # Tool lines already emitted as tool.start (dedup across edits). seen_tool_lines: set = field(default_factory=set) # --------------------------------------------------------------------------- # Passive / config probes (called from status displays -- no side effects) # --------------------------------------------------------------------------- def check_requirements() -> bool: """PASSIVE dependency probe: ``websockets`` importable + token set. Must be side-effect free (called from ``hermes setup`` / ``status`` / dashboard readiness). Never installs. """ try: import websockets # noqa: F401 (core dep) except Exception: return False return bool(_get_scoped_secret("ANDROID_TOKEN")) def validate_config(config) -> bool: """Given a PlatformConfig, is the platform properly configured?""" extra = getattr(config, "extra", {}) or {} token = _get_scoped_secret("ANDROID_TOKEN") or extra.get("token", "") return bool(token) def is_connected(config) -> bool: """Is the platform configured (env or config.yaml)?""" return validate_config(config) # --------------------------------------------------------------------------- # Env-driven auto-configuration (seeds PlatformConfig.extra pre-adapter) # --------------------------------------------------------------------------- def _env_enablement() -> Optional[dict]: """Seed ``PlatformConfig.extra`` from env vars during gateway config load. Called by the platform registry's env-enablement hook BEFORE adapter construction, so ``gateway status`` and ``get_connected_platforms()`` reflect env-only configuration without instantiating the adapter. Returns ``None`` when the platform isn't minimally configured (no token); the caller then skips auto-enabling. The special ``home_channel`` key in the returned dict is handled by the core hook -- it becomes a proper ``HomeChannel`` dataclass on the ``PlatformConfig`` rather than being merged into ``extra``. """ token = _get_scoped_secret("ANDROID_TOKEN", "") if not token: return None # Seed ONLY explicitly-set env vars: the core commits this seed on top of # config.yaml (``extra.update(seed)``), so default values here would # clobber user YAML. Unset keys fall through to config.yaml / adapter # defaults. seed: Dict[str, Any] = {} host = os.getenv("ANDROID_WS_HOST", "").strip() if host: seed["host"] = host port_raw = os.getenv("ANDROID_WS_PORT", "").strip() if port_raw: seed["port"] = _parse_port(port_raw) push = os.getenv("ANDROID_PUSH_BACKEND", "").strip().lower() if push: seed["push_backend"] = push home = os.getenv("ANDROID_HOME_CHANNEL", "").strip() if home: seed["home_channel"] = { "chat_id": home, "name": os.getenv("ANDROID_HOME_CHANNEL_NAME", "").strip() or DEFAULT_HOME_CHANNEL_NAME, } return seed def _parse_port(raw: str) -> int: try: return int((raw or "").strip()) except (ValueError, TypeError): return DEFAULT_PORT # --------------------------------------------------------------------------- # Target parsing: "android:[:]" # --------------------------------------------------------------------------- def _parse_target_ref(target_ref: str) -> Optional[tuple]: """Parse a raw target string into ``(chat_id, thread_id)`` or ``None``. Recognises the native syntax ``android:[:]`` where the chat_id itself carries the ``android:`` prefix (e.g. ``android:chan_7``) and an optional thread is a trailing ``:t_``. A bare friendly name (e.g. ``Cron Reports``) is resolved against the channel directory so cron / ``send_message`` can target a channel by name immediately, without waiting for the core directory's refresh timer. Returns ``None`` for anything unrecognised so the target proceeds to the core channel-directory resolution. """ if not target_ref: return None t = target_ref.strip() if not t: return None if t.startswith("android:"): body = t[len("android:"):].strip() if not body: return None thread_id: Optional[str] = None if ":" in body: head, tail = body.rsplit(":", 1) if tail and tail.startswith("t_"): thread_id = tail body = head return (f"android:{body}", thread_id) # Bare friendly name -> resolve via the channel directory. A thread resolves # to its session lane (parent_chat_id + thread_id); a channel/default to # its chat_id. try: entry = get_directory().resolve_entry(t) except Exception: entry = None if entry is not None: if entry["kind"] == "thread": return (entry["parent_chat_id"], entry["chat_id"]) return (entry["chat_id"], None) return None # --------------------------------------------------------------------------- # Standalone (out-of-process) send -- best-effort, stretch for v1 # --------------------------------------------------------------------------- async def _standalone_send( pconfig, chat_id: str, message: str, *, thread_id: Optional[str] = None, media_files: Optional[List[str]] = None, force_document: bool = False, ) -> Dict[str, Any]: """Out-of-process delivery for cron jobs that run separately from the gateway. The outbox is served by the *running* gateway, so standalone delivery while the gateway process is fully down is best-effort only (see ``docs/00-overview.md`` "Out of scope"). For M1 this is a stub that reports the gateway is required; the real implementation lands with the outbox (M3/M5). """ return { "error": ( "android standalone send: the running gateway is required to serve " "the outbox (standalone delivery is best-effort only)" ) } # --------------------------------------------------------------------------- # Verbose tool progress (full args on the progress line) # --------------------------------------------------------------------------- def _ensure_verbose_tool_progress() -> None: """Ensure the android platform renders tool progress in ``verbose`` mode. Verbose mode makes the gateway's tool-progress line carry the FULL argument JSON (not just a ~40-char preview), which the adapter parses into the ``tool.start`` frame's ``args`` field; the app then decides how much to show (Settings → Tool detail). The tool *output* is captured separately via the ``post_tool_call`` hook (verbose mode does not stream it). Best-effort and idempotent: writes ``display.platforms.android.tool_progress: verbose`` to config.yaml only when it isn't already set. The gateway's config cache is mtime-keyed, so the write takes effect on the next turn without a restart. Never raises. """ try: from hermes_cli.config import load_config_readonly cfg = load_config_readonly() or {} display = cfg.get("display") or {} platforms = display.get("platforms") or {} android = platforms.get("android") or {} if android.get("tool_progress") == "verbose": return # already set from utils import atomic_roundtrip_yaml_update atomic_roundtrip_yaml_update( get_hermes_home() / "config.yaml", "display.platforms.android.tool_progress", "verbose", ) logger.info("android: set display.platforms.android.tool_progress=verbose") except Exception: logger.debug("android: could not ensure verbose tool_progress", exc_info=True) # --------------------------------------------------------------------------- # Interactive setup (hermes gateway setup flow) # --------------------------------------------------------------------------- def interactive_setup() -> None: """Prompt for the pairing token / host / port / push backend. M1: token generation, host/port/push prompts, and the pairing QR payload (``iris://pair?...``) + app URL printed for the Connect screen. """ try: from hermes_cli.config import ( get_env_value, save_env_value, prompt, print_info, print_success, print_warning, ) except Exception: print("android: setup helpers unavailable; set ANDROID_TOKEN in ~/.hermes/.env") return print_info("📱 Android / Desktop (Iris x Hermes)") token = get_env_value("ANDROID_TOKEN") or "" if not token: generated = generate_token() save_env_value("ANDROID_TOKEN", generated) print_success(f"Generated pairing token: {generated}") print_warning("Keep this secret -- the app presents it on connect.") else: print_info("Existing ANDROID_TOKEN found (not shown).") host = prompt("WS bind host", default=get_env_value("ANDROID_WS_HOST") or DEFAULT_HOST) save_env_value("ANDROID_WS_HOST", host or DEFAULT_HOST) port = prompt("WS port", default=str(_parse_port(get_env_value("ANDROID_WS_PORT") or ""))) save_env_value("ANDROID_WS_PORT", str(_parse_port(port))) backend = prompt("Push backend (fcm/ntfy)", default=get_env_value("ANDROID_PUSH_BACKEND") or DEFAULT_PUSH_BACKEND) save_env_value("ANDROID_PUSH_BACKEND", (backend or DEFAULT_PUSH_BACKEND).strip().lower()) # Pairing payload for the app's Connect screen (QR / manual entry). try: from hermes_cli.config import print_code url = pairing_url(host or DEFAULT_HOST, _parse_port(port)) payload = qr_payload(host or DEFAULT_HOST, _parse_port(port), token) print_info("Pair your device (scan with the app or enter on the Connect screen):") print_code(payload) print_info(f"Server URL: {url}") except Exception: url = pairing_url(host or DEFAULT_HOST, _parse_port(port)) print_info(f"Pairing URL: {qr_payload(host or DEFAULT_HOST, _parse_port(port), token)}") print_info(f"Server URL: {url}") # Always render tool progress verbosely so the app receives the full tool # call args (it decides how much to show via Settings → Tool detail). _ensure_verbose_tool_progress() print_success("Android configuration saved to ~/.hermes/.env") print_info("Restart the gateway for changes to take effect: hermes gateway restart") # --------------------------------------------------------------------------- # Android Adapter # --------------------------------------------------------------------------- class AndroidAdapter(BasePlatformAdapter): """WebSocket-backed adapter for the native Iris Android / Desktop app. M1: the WS server (``ws_server.WsServer``) authenticates devices with the pairing token, the connection registry tracks live sockets, ``send()`` emits ``message`` frames, and inbound ``message.send`` frames become ``MessageEvent``s for ``handle_message()``. """ def __init__(self, config, **kwargs): platform = Platform("android") super().__init__(config=config, platform=platform) # Ensure verbose tool progress (full args on the progress line) so the # app can show the full tool call on demand. Best-effort; idempotent. _ensure_verbose_tool_progress() extra = getattr(config, "extra", {}) or {} # Connection settings (env vars override config.yaml) self.host = os.getenv("ANDROID_WS_HOST", "").strip() or extra.get("host", DEFAULT_HOST) self.port = _parse_port(os.getenv("ANDROID_WS_PORT", "") or str(extra.get("port", DEFAULT_PORT))) self.token = _get_scoped_secret("ANDROID_TOKEN") or extra.get("token", "") self.push_backend = ( os.getenv("ANDROID_PUSH_BACKEND", "").strip().lower() or extra.get("push_backend", DEFAULT_PUSH_BACKEND) ) self.outbox_retention_hours = int( extra.get("outbox_retention_hours", DEFAULT_OUTBOX_RETENTION_HOURS) ) self.max_upload_bytes = int( extra.get("max_upload_bytes", DEFAULT_MAX_UPLOAD_BYTES) ) self._gateway_status = protocol.STATUS_ONLINE # Home channel: the core hook turns the env-seeded ``home_channel`` # dict into a HomeChannel dataclass on the config; config.yaml may # also put it in extra (dict or bare string). home = getattr(config, "home_channel", None) if home is not None and getattr(home, "chat_id", None): self.home_channel = str(home.chat_id) self.home_channel_name = str(getattr(home, "name", "") or DEFAULT_HOME_CHANNEL_NAME) else: hc = extra.get("home_channel") if isinstance(hc, dict) and hc.get("chat_id"): self.home_channel = str(hc["chat_id"]) self.home_channel_name = str(hc.get("name") or DEFAULT_HOME_CHANNEL_NAME) elif isinstance(hc, str) and hc.strip(): self.home_channel = hc.strip() self.home_channel_name = DEFAULT_HOME_CHANNEL_NAME else: self.home_channel = DEFAULT_HOME_CHANNEL self.home_channel_name = DEFAULT_HOME_CHANNEL_NAME # TLS (optional) self.ws_cert = _get_scoped_secret("ANDROID_WS_CERT") or extra.get("ws_cert", "") self.ws_key = _get_scoped_secret("ANDROID_WS_KEY") or extra.get("ws_key", "") # Auth allowed = os.getenv("ANDROID_ALLOWED_USERS", "").strip() self.allowed_users: List[str] = ( [u.strip() for u in allowed.split(",") if u.strip()] if allowed else [] ) self.allow_all = _truthy(os.getenv("ANDROID_ALLOW_ALL_USERS")) # Runtime state self._devices = DeviceRegistry(get_hermes_home() / "android" / "devices.db") self._ws_server = WsServer(self, self._devices) self._connected = False # M2: per-chat turn state for outbound frame classification. self._turns: Dict[str, _TurnState] = {} # M3: channel directory (shared singleton) + offline outbox. self._channels = get_directory() self._outbox = Outbox( get_hermes_home() / "android" / "outbox.db", retention_hours=self.outbox_retention_hours, ) # M4: media registry (inbound upload refs + outbound offers) and the # last finalized assistant message id per chat (offer association). self._media = media_bridge.MediaStore(get_hermes_home()) self._last_message_id: Dict[str, str] = {} # M5: push backend (FCM primary, ntfy fallback) + the throttle for # the outbox-prune banner. self._push: PushBackend = build_push_backend( self.push_backend, fcm_service_account=_get_scoped_secret("ANDROID_FCM_SERVICE_ACCOUNT"), fcm_server_key=_get_scoped_secret("ANDROID_FCM_SERVER_KEY"), ntfy_topic=_get_scoped_secret("NTFY_TOPIC"), ntfy_server_url=os.getenv("NTFY_SERVER_URL", "").strip() or None, ntfy_auth_token=_get_scoped_secret("NTFY_AUTH_TOKEN"), ) self._prune_notified_at = 0.0 def _turn_state(self, chat_id: str) -> _TurnState: st = self._turns.get(chat_id) if st is None: st = _TurnState() self._turns[chat_id] = st return st @property def name(self) -> str: return "Android" # ── Connection lifecycle ────────────────────────────────────────────── async def connect(self, *, is_reconnect: bool = False) -> bool: """Bring the platform up: bind the WS server on host:port.""" if not self.token: logger.error("android: ANDROID_TOKEN must be set") self._set_fatal_error( "config_missing", "ANDROID_TOKEN must be set", retryable=False, ) return False # Prevent two profiles from binding the same port/identity. try: from gateway.status import acquire_scoped_lock lock_key = f"{self.host}:{self.port}" if not acquire_scoped_lock("android", lock_key): logger.error("android: %s:%s already in use by another profile", self.host, self.port) self._set_fatal_error( "lock_conflict", "WS port in use by another profile", retryable=False, ) return False self._lock_key = lock_key except ImportError: self._lock_key = None # status module not available (e.g. tests) try: await self._ws_server.start() except Exception: self._connected = False return False # M5: announce gateway health to connected clients (none yet at # startup; the frame + plumbing exist for future transitions). await self._ws_server.broadcast(protocol.status(self._gateway_status)) # M3: ensure the default (home) channel exists in the directory so the # app's channel list and cron home delivery have a stable anchor. try: self._channels.ensure_default(self.home_channel, self.home_channel_name) except Exception: logger.warning("android: ensure_default failed", exc_info=True) # M5: push backend status (degrade gracefully when unconfigured). if not self._push.configured(): logger.warning( "android: push backend %r not configured (no credentials) -- " "offline devices will not be woken; outbox + sync still apply", self.push_backend, ) else: logger.info("android: push backend: %s", self._push.name) self._connected = True self._mark_connected() logger.info("android: connected; WS server on %s:%s", self.host, self.port) return True async def disconnect(self) -> None: """Tear down the platform: stop the server, close device sockets.""" try: from gateway.status import release_scoped_lock if getattr(self, "_lock_key", None): release_scoped_lock("android", self._lock_key) except ImportError: pass try: await self._ws_server.stop() except Exception: logger.warning("android: WS server stop failed", exc_info=True) try: self._devices.close() except Exception: pass try: self._outbox.close() except Exception: pass self._connected = False self._mark_disconnected() logger.info("android: disconnected") # ── Outbound (agent -> app) ─────────────────────────────────────────── async def send( self, chat_id: str, content: str, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, ) -> SendResult: """Send a message to a chat. M2: classify the outbound call into a structured frame using the per-chat turn state machine (see module docstring): * ``metadata["expect_edits"]`` -> ``message.start`` (streaming segment) * ``metadata["notify"]`` -> ``message`` / ``message.stop`` (final) * tool-progress line format -> ``tool.start`` (first tool bubble) * anything else -> ``commentary`` With no live devices the frame is dropped here (the outbox + push replay lands in M3/M5). """ content = content or "" meta = metadata or {} thread_id = _thread_id_from_metadata(meta) state = self._turn_state(chat_id) # 1. Streaming segment start (stream consumer first send). if meta.get("expect_edits") is True: # A new content segment means the tool the model was waiting on # has returned -> close it before the segment opens. await self._close_open_tool(chat_id, state, thread_id) message_id = _mint_message_id() state.active = True state.stream_id = message_id await self._broadcast_or_log( chat_id, protocol.message_start(chat_id, message_id, protocol.ROLE_ASSISTANT, thread_id=thread_id), ) return SendResult(success=True, message_id=message_id) # 2. Cron delivery (the cron scheduler passes ``job_id`` in metadata): # a final assistant message + a high-priority notification banner # (pushed even when a device is live, docs/08 §8.1). if meta.get("job_id"): reasoning, _ = _split_reasoning(content) if reasoning: _reset_reasoning() else: await _wait_for_reasoning_flushed() reasoning = _take_reasoning() or None name, inner = _cron_brief(content, str(meta.get("job_id"))) message_id = _mint_message_id() await self._broadcast_or_log( chat_id, protocol.notification( chat_id, protocol.NOTIF_CRON, f"Cron: {name}", _push_preview(inner), thread_id=thread_id, ), ) await self._broadcast_or_log( chat_id, protocol.message( chat_id=chat_id, message_id=message_id, role=protocol.ROLE_ASSISTANT, text=inner, thread_id=thread_id, reasoning=reasoning, ts=int(time.time() * 1000), ), ) self._last_message_id[chat_id] = message_id state.active = False return SendResult(success=True, message_id=message_id) # 3. Final message (non-streaming final, or streaming fallback final). if meta.get("notify") is True: reasoning, body = _split_reasoning(content) # Non-streaming: reasoning is prepended to content (split above). # Streaming fallback: content has no reasoning, so use the # reasoning captured via the on_stream_delta hook (wait for the # async hook worker to flush it first). if not reasoning: await _wait_for_reasoning_flushed() reasoning = _take_reasoning() or None else: _reset_reasoning() if state.stream_id: # Fallback final: close the open streaming segment in place. message_id = state.stream_id state.stream_id = None await self._broadcast_or_log( chat_id, protocol.message_stop( chat_id, message_id, body, reasoning=reasoning, thread_id=thread_id, ts=int(time.time() * 1000), ), ) else: message_id = _mint_message_id() await self._broadcast_or_log( chat_id, protocol.message( chat_id=chat_id, message_id=message_id, role=protocol.ROLE_ASSISTANT, text=body, thread_id=thread_id, reasoning=reasoning, reply_to=reply_to, ts=int(time.time() * 1000), ), ) # M4: media offers emitted after this final associate with it. self._last_message_id[chat_id] = message_id await self._close_open_tool(chat_id, state, thread_id) self._reset_tool_state(state) state.active = False return SendResult(success=True, message_id=message_id) # 4. Tool progress (first tool bubble of an editable line buffer). # Gateway lifecycle notices (restart/shutdown/online) are system # notices, not tool progress — skip the tool-line heuristic so they # fall through to commentary instead of a never-completing card. if not _is_gateway_lifecycle_notice(content) and _is_tool_progress(content): return await self._emit_tool_lines(chat_id, content, state, thread_id, is_edit=False) # 5. Commentary (interim assistant beat). message_id = _mint_message_id() state.active = True await self._broadcast_or_log( chat_id, protocol.commentary(chat_id, message_id, content, thread_id=thread_id), ) return SendResult(success=True, message_id=message_id) async def edit_message( self, chat_id: str, message_id: str, content: str, *, finalize: bool = False, metadata: Optional[Dict[str, Any]] = None, ) -> SendResult: """Edit a previously sent message (M2: drives streaming + tool updates). * ``message_id == state.stream_id`` -> ``message.update`` (``finalize=True`` -> ``message.stop``). * ``message_id == state.tool_msg_id`` -> tool-progress update (new lines -> ``tool.start``). * unknown id -> best-effort ``message.update``. """ content = content or "" thread_id = _thread_id_from_metadata(metadata) state = self._turn_state(chat_id) if message_id and message_id == state.stream_id: if finalize: reasoning, body = _split_reasoning(_strip_streaming_cursor(content)) # Streaming: the gateway drops the model's separate # reasoning_content (final send suppressed), so attach the # reasoning we captured via the on_stream_delta hook (wait # for the async hook worker to flush it first). if not reasoning: await _wait_for_reasoning_flushed() reasoning = _take_reasoning() or None else: _reset_reasoning() state.stream_id = None await self._broadcast_or_log( chat_id, protocol.message_stop( chat_id, message_id, body, reasoning=reasoning, thread_id=thread_id, ts=int(time.time() * 1000), ), ) # M4: media offers emitted after this final associate with it. self._last_message_id[chat_id] = message_id await self._close_open_tool(chat_id, state, thread_id) self._reset_tool_state(state) state.active = False else: await self._broadcast_or_log( chat_id, protocol.message_update( chat_id, message_id, _strip_streaming_cursor(content), thread_id=thread_id, ), ) return SendResult(success=True, message_id=message_id) if message_id and message_id == state.tool_msg_id: return await self._emit_tool_lines(chat_id, content, state, thread_id, is_edit=True) # Unknown id: treat as a streaming update (best effort). if finalize: await self._broadcast_or_log( chat_id, protocol.message_stop( chat_id, message_id, _strip_streaming_cursor(content), thread_id=thread_id, ts=int(time.time() * 1000), ), ) else: await self._broadcast_or_log( chat_id, protocol.message_update( chat_id, message_id, _strip_streaming_cursor(content), thread_id=thread_id, ), ) return SendResult(success=True, message_id=message_id) # ── M2: tool-progress helpers ───────────────────────────────────────── async def _emit_tool_lines( self, chat_id: str, content: str, state: _TurnState, thread_id: Optional[str], *, is_edit: bool, ) -> SendResult: """Emit ``tool.start`` for each NEW tool line in *content*. The gateway accumulates tool lines in one editable bubble; on an edit the full buffer is re-sent, so we diff against ``seen_tool_lines`` to emit only the new ones. A new tool closes the previously-open tool. """ message_id = state.tool_msg_id or _mint_message_id() state.tool_msg_id = message_id state.active = True lines = [ln for ln in content.splitlines() if ln.strip()] for line in lines: key = line.strip() if key in state.seen_tool_lines: continue state.seen_tool_lines.add(key) parsed = self._parse_tool_line_or_block(line, content) if parsed is None: continue name, preview, args = parsed # A new tool begins: close the previously-open one, attaching the # output/duration/ok captured by the post_tool_call hook. if state.open_tool_index is not None: extra = _tool_end_fields(state.open_tool_name or "") await self._broadcast_or_log( chat_id, protocol.tool_end( chat_id, state.open_tool_index, state.open_tool_name or "", ok=extra.get("ok", True), duration=extra.get("duration"), output_preview=extra.get("output_preview"), thread_id=thread_id, ), ) state.tool_index += 1 state.open_tool_index = state.tool_index state.open_tool_name = name await self._broadcast_or_log( chat_id, protocol.tool_start( chat_id, state.tool_index, name, preview=preview, args=args, thread_id=thread_id, ), ) return SendResult(success=True, message_id=message_id) @staticmethod def _parse_tool_line_or_block( line: str, content: str ) -> Optional[Tuple[str, Optional[str], Optional[Dict[str, Any]]]]: """Parse a tool line into ``(name, preview, args)``. Expands a terminal code block to its command, and a verbose header (`` (keys)``) to its full args JSON (the JSON sits on the following line). ``args`` is ``None`` unless the line is a verbose header with a parseable JSON body. """ parsed = _parse_tool_line(line) if parsed is None: return None name, preview = parsed # Terminal code block: the command lives in the fenced lines that # follow the " terminal" head line. if name == "terminal" and preview is None and "```" in content: cmd = _extract_code_block(content) if cmd: return name, cmd, None # Verbose mode: recover the full args from the following JSON line. args = _extract_verbose_args(line, content) if args is not None and preview is None: preview = _short_preview_from_args(args) return name, preview, args async def _close_open_tool( self, chat_id: str, state: _TurnState, thread_id: Optional[str] ) -> None: """Emit ``tool.end`` for the currently-open tool, if any. A tool is considered complete when the next tool starts OR a new content segment begins (the model only produces content after the tool it was waiting on has returned). """ if state.open_tool_index is not None: extra = _tool_end_fields(state.open_tool_name or "") await self._broadcast_or_log( chat_id, protocol.tool_end( chat_id, state.open_tool_index, state.open_tool_name or "", ok=extra.get("ok", True), duration=extra.get("duration"), output_preview=extra.get("output_preview"), thread_id=thread_id, ), ) state.open_tool_index = None state.open_tool_name = None def _reset_tool_state(self, state: _TurnState) -> None: """Clear per-turn tool bookkeeping (called at turn finalization).""" state.tool_msg_id = None state.seen_tool_lines = set() state.tool_index = 0 state.open_tool_index = None state.open_tool_name = None # Drop any captured tool results not consumed by a tool.end this turn # (e.g. tool_progress off) so they can't leak into the next turn. _reset_tool_results() async def _broadcast_or_log(self, chat_id: str, frame: "protocol.Frame") -> None: delivered = await self._ws_server.broadcast(frame) # M3/M5: always append to the outbox so a reconnecting app can catch # up on *all* recent frames, not just the ones that were parked. This # covers the case where the app's in-memory ChatStore is reset (e.g. # activity recreation / composition recompose) while the process stays # alive: the recreated controller re-syncs from its stale cursor and # replays the frames it missed. The push below still only fires when # there is no live subscriber (or for high-priority events). try: cursor = self._outbox.append(chat_id, frame.to_json()) except Exception: logger.warning("android: outbox append failed", exc_info=True) return if delivered == 0: logger.info( "android: no live devices for %s; %s frame parked in outbox (cursor=%s)", chat_id, frame.type, cursor, ) # M5: wake the offline device(s) via the push backend. await self._maybe_push(chat_id, frame, cursor) await self._maybe_notify_outbox_prune(chat_id) elif ( frame.type == protocol.TYPE_NOTIFICATION and frame.payload.get("kind") in protocol.HIGH_PRIORITY_NOTIF_KINDS ): # M5: high-priority events (approval/clarify/cron) push even when # a device is live -- the app may be backgrounded and decides # whether to also show an in-app banner (docs/08 §8.1). await self._maybe_push(chat_id, frame, cursor) # ── M5: push ─────────────────────────────────────────────────────────── def _push_summary( self, frame: "protocol.Frame" ) -> Optional[Tuple[str, str, str, str]]: """``(title, body, kind, priority)`` for a pushable frame, else None. Only terminal/interesting frames wake a device: intermediate streaming and tool frames are replayed by ``sync`` without a push (no notification spam per turn). """ t = frame.type p = frame.payload if t == protocol.TYPE_MESSAGE: return ( self._channel_name(frame.chat_id or ""), _push_preview(p.get("text")), "message", "normal", ) if t == protocol.TYPE_MESSAGE_STOP: return ( self._channel_name(frame.chat_id or ""), _push_preview(p.get("final_text")), "message", "normal", ) if t == protocol.TYPE_NOTIFICATION: kind = str(p.get("kind") or protocol.NOTIF_GENERIC) priority = "high" if kind in protocol.HIGH_PRIORITY_NOTIF_KINDS else "normal" return ( str(p.get("title") or "Iris"), str(p.get("body") or ""), kind, priority, ) if t == protocol.TYPE_MEDIA_OFFER: return ( self._channel_name(frame.chat_id or ""), f"New {p.get('kind') or 'media'}: {p.get('filename') or ''}".strip(), "media", "normal", ) return None async def _maybe_push(self, chat_id: str, frame: "protocol.Frame", cursor: int) -> None: """Fire the configured push backend for a parked (or high-priority) frame. Best-effort: failures are logged, never raised.""" summary = self._push_summary(frame) if summary is None: return title, body, kind, priority = summary backend = self._push if backend is None or not backend.token_field: return devices = self._devices.list() if not backend.configured() and not any( d.get(backend.token_field) for d in devices ): return data: Dict[str, Any] = {"chat_id": chat_id, "kind": kind, "cursor": str(cursor)} if frame.thread_id: data["thread_id"] = frame.thread_id message_id = frame.payload.get("message_id") if isinstance(message_id, str) and message_id: data["message_id"] = message_id for device in devices: device_id = device.get("device_id") if not device_id: continue # Prefer the live connection's token (fcm.register refreshes it # in memory) over the possibly-stale registry row. conn = self._ws_server.connection(device_id) token = getattr(conn, backend.token_field, None) if conn is not None else None if not token: token = device.get(backend.token_field) if not token: continue try: ok = await backend.send( device_id=device_id, chat_id=chat_id, title=title, body=body, data=data, token=token, priority=priority, ) except Exception: logger.warning("android: push via %s failed", backend.name, exc_info=True) continue if ok: logger.info( "android: push via %s -> %s (%s, chat=%s)", backend.name, device_id, frame.type, chat_id, ) async def _maybe_notify_outbox_prune(self, chat_id: str) -> None: """When the outbox row cap pruned old frames, tell the app (throttled to once per hour so a full box doesn't banner per frame).""" pruned = self._outbox.take_overflow_pruned() if pruned <= 0: return now = time.time() if now - self._prune_notified_at < 3600.0: return self._prune_notified_at = now await self._broadcast_or_log( chat_id, protocol.notification( chat_id, protocol.NOTIF_GENERIC, "Outbox", f"{pruned} older message(s) pruned", ), ) async def send_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]] = None) -> None: """Send a typing indicator (``typing`` frame, on=true).""" thread_id = None if metadata: tid = metadata.get("thread_id") if isinstance(tid, str) and tid: thread_id = tid await self._ws_server.broadcast(protocol.typing(chat_id, True, thread_id=thread_id)) async def stop_typing(self, chat_id: str) -> None: """Clear the typing indicator (``typing`` frame, on=false).""" await self._ws_server.broadcast(protocol.typing(chat_id, False)) # ── M4: outbound media (agent -> app) ───────────────────────────────── # # The gateway's dispatch partition (gateway/run.py) extracts MEDIA: tags / # image URLs from the final response, filters them through # ``filter_media_delivery_paths``, then calls these ``send_*`` overrides # with local file paths. We re-validate each path (defense in depth), # register it in the media registry, mint a ``media_id``, and emit # ``media.offer``; the app fetches the bytes via ``media.pull``. async def _offer_media( self, chat_id: str, path: str, kind: str, filename: Optional[str], metadata: Optional[Dict[str, Any]], ) -> SendResult: safe = validate_media_delivery_path(path) if safe is None: logger.warning("android: media path failed delivery validation: %s", path) return SendResult(success=False, error="android: media path not deliverable") try: size = os.path.getsize(safe) except OSError as e: logger.warning("android: media file unreadable %s: %s", safe, e) return SendResult(success=False, error="android: media file unreadable") entry = self._media.register_outbound( safe, kind, media_bridge.mime_for_path(safe), filename or os.path.basename(safe), size ) thread_id = _thread_id_from_metadata(metadata) frame = protocol.media_offer( entry.media_id, entry.kind, entry.mime, entry.size, entry.filename, chat_id=chat_id, thread_id=thread_id, message_id=self._last_message_id.get(chat_id), ) await self._broadcast_or_log(chat_id, frame) return SendResult(success=True, message_id=entry.media_id) async def send_image( self, chat_id: str, image_url: str, caption: Optional[str] = None, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, ) -> SendResult: """Send an image (M4: local files offered over WS; remote URLs fall back to the base text rendering).""" if image_url.startswith("file://"): from urllib.parse import unquote return await self._offer_media(chat_id, unquote(image_url[7:]), "image", None, metadata) return await super().send_image( chat_id, image_url, caption=caption, reply_to=reply_to, metadata=metadata ) async def send_image_file( self, chat_id: str, image_path: str, caption: Optional[str] = None, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, **kwargs: Any, ) -> SendResult: """Send a local image file (M4).""" return await self._offer_media(chat_id, image_path, "image", None, metadata) async def send_video( self, chat_id: str, video_path: str, caption: Optional[str] = None, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, **kwargs: Any, ) -> SendResult: """Send a video (M4).""" return await self._offer_media(chat_id, video_path, "video", None, metadata) async def send_voice( self, chat_id: str, audio_path: str, caption: Optional[str] = None, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, **kwargs: Any, ) -> SendResult: """Send a voice note / audio file (M4).""" return await self._offer_media(chat_id, audio_path, "voice", None, metadata) async def send_document( self, chat_id: str, file_path: str, caption: Optional[str] = None, file_name: Optional[str] = None, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, **kwargs: Any, ) -> SendResult: """Send a document (M4).""" return await self._offer_media(chat_id, file_path, "document", file_name, metadata) # ── Inbound (app -> agent) ──────────────────────────────────────────── async def on_message_send(self, frame: protocol.Frame, device_id: str) -> None: """Handle an inbound ``message.send`` frame. Echoes the user message to all devices (multi-device sync + ack), then builds a ``MessageEvent`` and hands it to ``handle_message()`` (the gateway's command pipeline + agent turn). M4: ``media_refs`` reference completed ``media.upload``s; they are resolved to ``MessageEvent.media_urls``/``media_types`` (local paths the agent's vision/audio tools can read) and echoed in the user message's ``media[]`` so every device renders the attachments. """ payload = frame.payload text = payload.get("text") text = text if isinstance(text, str) else "" refs_raw = payload.get("media_refs") media_refs = ( [r for r in refs_raw if isinstance(r, str) and r] if isinstance(refs_raw, list) else [] ) if not text.strip() and not media_refs: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "message.send requires non-empty text", id=frame.id), ) return chat_id = frame.chat_id or payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): chat_id = self.home_channel chat_id = chat_id.strip() thread_id = frame.thread_id or payload.get("thread_id") if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None reply_to = payload.get("reply_to") if not isinstance(reply_to, str) or not reply_to.strip(): reply_to = None # Auto-threading (the app's Threads setting, docs/06 §6.3): a message # in a channel's flat lane gets its own fresh thread, the way Telegram # topic mode mints a topic per new conversation. The thread is named # instantly from the user's opening message (derived title) and the # LLM upgrades the name in the background. The user echo, the agent # turn, and all streaming frames then carry the new thread_id. # Skipped for slash commands (session-scoped, not conversation # starters) and replies (they continue where the user is). auto_thread = payload.get("auto_thread") is True if ( auto_thread and thread_id is None and text.strip() and not text.lstrip().startswith("/") and reply_to is None ): entry = self._channels.create( name=_derive_thread_name(text), kind="thread", parent_chat_id=chat_id, ) thread_id = entry["chat_id"] # Bare broadcast (like channel.create): the directory is # re-served on hello.ack, so no outbox entry is needed. await self._ws_server.broadcast(protocol.channel_created(entry, auto=True)) self._schedule_thread_title_upgrade(entry["chat_id"], text) # M4: resolve media refs (single-use; unknown ref -> error). media_urls: List[str] = [] media_types: List[str] = [] media_wire: List[Dict[str, Any]] = [] for ref in media_refs: entry = self._media.get_inbound(ref) if entry is None: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, f"unknown media_ref {ref}", id=frame.id), ) return media_urls.append(entry.path) media_types.append(entry.mime) media_wire.append( { "media_id": entry.media_id, "kind": entry.kind, "mime": entry.mime, "size": entry.size, "filename": entry.filename, } ) device = self._devices.get(device_id) or {} user_name = device.get("name") or device_id # Echo to all devices: the sender confirms (server-assigned id), # other devices see the message too (single-user, multi-device). # Routed through _broadcast_or_log (not a bare broadcast) so the echo # is appended to the outbox: the app's ChatStore is in-memory only, so # after a process death / activity recreation the only way the user's # own message is restored is via the sync replay. Without this, user # messages vanish on reconnect while bot messages (already parked) # survive. message_id = f"m_{uuid.uuid4().hex[:16]}" echo = protocol.message( chat_id=chat_id, message_id=message_id, role=protocol.ROLE_USER, text=text, thread_id=thread_id, media=media_wire or None, reply_to=reply_to, ts=int(time.time() * 1000), ) await self._broadcast_or_log(chat_id, echo) # Refs are consumed by this message (no replay). for ref in media_refs: self._media.pop_inbound(ref) # M4: a new user turn starts -- stale offer association is dropped. self._last_message_id.pop(chat_id, None) kind = media_wire[0]["kind"] if media_wire else None if kind == "image": message_type = MessageType.PHOTO elif kind == "video": message_type = MessageType.VIDEO elif kind == "audio": message_type = MessageType.AUDIO elif kind == "voice": message_type = MessageType.VOICE elif kind == "document": message_type = MessageType.DOCUMENT else: message_type = MessageType.TEXT source = self.build_source( chat_id=chat_id, chat_name=self._channel_name(chat_id), chat_type="dm", user_id=device_id, user_name=user_name, thread_id=thread_id, ) event = MessageEvent( text=text, message_type=message_type, user_id=device_id, user_name=user_name, source=source, message_id=message_id, reply_to_message_id=reply_to, media_urls=media_urls, media_types=media_types, ) await self.handle_message(event) # M5: acknowledge the user message to the originating device (the # app shows ✓✓) at the moment it is handed to the agent. await self._ws_server.send_to( device_id, protocol.read_receipt(chat_id, message_id), ) def _schedule_thread_title_upgrade(self, thread_id: str, text: str) -> None: """Upgrade an auto-created thread's name with the model's title. Stage 2 of hermes' two-stage session titling (``agent/title_generator .py``): the thread was created with an instant derived name; this background call on the ``title_generation`` auxiliary task replaces it with the model's title and broadcasts ``channel.renamed``. Best-effort — any failure (config, model, network) leaves the derived name in place, and a thread the user already renamed or archived is untouched. """ loop = asyncio.get_running_loop() def _work() -> None: try: from agent.title_generator import generate_title title = generate_title(text) except Exception: logger.debug("Thread title upgrade failed", exc_info=True) return if not title: return entry = self._channels.get(thread_id) if entry is None or entry.get("archived"): return if (entry.get("name") or "") == title: return renamed = self._channels.rename(thread_id, title) if renamed is None: return try: asyncio.run_coroutine_threadsafe( self._ws_server.broadcast(protocol.channel_renamed(renamed)), loop, ) except Exception: logger.debug("Thread title rename broadcast failed", exc_info=True) threading.Thread( target=_work, daemon=True, name="android-thread-title" ).start() # ── M4: inbound media (app -> agent) ────────────────────────────────── # # ``media.upload.start`` -> raw binary frames (one at a time per # connection) -> ``media.upload.end``. The session streams to a temp # file (bounded RAM); on end we verify size + sha256, re-sniff the kind, # and cache via hermes ``cache_*_from_bytes``. ``media.pull`` serves an # outbound offer as chunked binary frames, re-checking the delivery-path # validation at pull time. async def on_media_upload_start(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload media_ref = str(payload.get("media_ref") or "").strip() if not media_ref or len(media_ref) > 64: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "media.upload.start requires media_ref", id=frame.id), ) return kind = payload.get("kind") if kind not in media_bridge.KINDS: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, f"unsupported media kind {kind!r}", id=frame.id), ) return mime = str(payload.get("mime") or "application/octet-stream")[:128] filename = str(payload.get("filename") or "upload")[:255] size = payload.get("size") try: size = int(size) if size is not None else -1 except (TypeError, ValueError): size = -1 if size <= 0: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "media.upload.start requires a positive size", id=frame.id), ) return if size > self.max_upload_bytes: await self._ws_server.send_to( device_id, protocol.error( protocol.ERR_MEDIA_TOO_LARGE, f"upload of {size} bytes exceeds limit ({self.max_upload_bytes})", id=frame.id, ), ) return try: self._media.create_upload( device_id, media_ref, kind, mime, filename, size, frame.id, self.max_upload_bytes, ) except media_bridge.MediaError as e: await self._ws_server.send_to(device_id, protocol.error(e.code, e.message, id=frame.id)) return # No ack: WS ordering guarantees the server processes this before the # first binary chunk; failures arrive as ``error`` frames. async def on_media_chunk(self, device_id: str, chunk: bytes) -> None: session = self._media.get_upload(device_id) if session is None: return # stray binary frame: ignore (forward-compat) session.feed(chunk) if session.failed: await self._ws_server.send_to( device_id, protocol.error(session.error_code, session.error_message, id=session.request_id), ) self._media.discard_upload(device_id, session.media_ref) async def on_media_upload_end(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload media_ref = str(payload.get("media_ref") or "").strip() sha256 = str(payload.get("sha256") or "").strip().lower() if not media_ref: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "media.upload.end requires media_ref", id=frame.id), ) return try: entry = self._media.complete_upload(device_id, media_ref, sha256) except media_bridge.MediaError as e: await self._ws_server.send_to(device_id, protocol.error(e.code, e.message, id=frame.id)) return await self._ws_server.send_to( device_id, protocol.media_upload_ack(True, entry.media_id, id=frame.id) ) async def on_media_pull(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload media_id = str(payload.get("media_id") or "").strip() entry = self._media.get_outbound(media_id) if media_id else None if entry is None: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown media_id {media_id!r}", id=frame.id), ) return # Delivery-path security: re-validate at pull time (the file may have # moved / been replaced since the offer). safe = validate_media_delivery_path(entry.path) if safe is None: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_NOT_FOUND, "media no longer deliverable", id=frame.id), ) return conn = self._ws_server.connection(device_id) if conn is None: return try: await media_bridge.stream_file(conn.ws, safe, media_bridge.DEFAULT_CHUNK_BYTES) except Exception as e: logger.warning("android: media.pull stream failed for %s: %s", media_id, e) await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_INTERNAL, f"pull failed: {e}", id=frame.id) ) return await self._ws_server.send_to(device_id, protocol.media_pull_end(True, id=frame.id)) def on_connection_closed(self, device_id: str) -> None: """M4: drop in-flight upload temp files for a disconnected device.""" self._media.discard_device(device_id) # ── M3: channel directory management (app -> agent) ─────────────────── # # 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). 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._ws_server.send_to( 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._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id) ) return resp = protocol.channel_created(entry) resp.id = frame.id await self._ws_server.broadcast(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._ws_server.send_to( 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._ws_server.send_to( 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._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id) ) return if entry is None: await self._ws_server.send_to( 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._ws_server.broadcast(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._ws_server.send_to( 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._ws_server.send_to( 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._ws_server.broadcast(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._ws_server.send_to( 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._ws_server.send_to( 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._ws_server.broadcast(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._ws_server.send_to( 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._ws_server.send_to( 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._ws_server.send_to( 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._ws_server.broadcast(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._ws_server.send_to( 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._ws_server.send_to( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"cannot delete {chat_id} (unknown or default)", id=frame.id), ) return resp = protocol.channel_deleted(chat_id) resp.id = frame.id await self._ws_server.broadcast(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._ws_server.send_to(device_id, resp) # ── M3: search (app -> agent) ───────────────────────────────────────── async def on_search(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload query = payload.get("query") if not isinstance(query, str) or not query.strip(): await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "search requires a query", id=frame.id) ) return scope = payload.get("scope") scope = scope if scope in ("all", "chat") else "all" chat_id = payload.get("chat_id") or frame.chat_id if not isinstance(chat_id, str) or not chat_id.strip(): chat_id = None thread_id = payload.get("thread_id") or frame.thread_id if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None limit = payload.get("limit") try: limit = int(limit) if limit is not None else 20 except (TypeError, ValueError): limit = 20 db_path = get_hermes_home() / "state.db" hits = search_bridge.search( db_path, query, scope=scope, chat_id=chat_id, thread_id=thread_id, limit=limit ) resp = protocol.search_results(query, scope, hits, id=frame.id) await self._ws_server.send_to(device_id, resp) # ── M3: sync (reconnect catch-up) ───────────────────────────────────── async def on_sync(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload cursor = payload.get("cursor") try: cursor = int(cursor) if cursor is not None else 0 except (TypeError, ValueError): cursor = 0 for e in self._outbox.replay(cursor): raw = e["frame"] replayed = protocol.Frame( type=raw.get("type", ""), payload=raw.get("payload", {}) if isinstance(raw.get("payload"), dict) else {}, id=raw.get("id") if isinstance(raw.get("id"), int) else None, chat_id=raw.get("chat_id") if isinstance(raw.get("chat_id"), str) else e.get("chat_id"), thread_id=raw.get("thread_id") if isinstance(raw.get("thread_id"), str) else None, v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION, ) await self._ws_server.send_to(device_id, replayed) done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id) await self._ws_server.send_to(device_id, done) # ── Full message history (initial channel open / scroll-up) ─────────── async def on_history(self, frame: protocol.Frame, device_id: str) -> None: """Handle an inbound ``history`` request. ``sync`` only replays the outbox delta since the device's cursor, so after a process death the app's in-memory ChatStore is empty and the delta does not cover older messages. ``history`` loads the full message list for a chat/thread (reconstructed from the outbox log) so the app can populate the view on first open / restart. """ payload = frame.payload chat_id = frame.chat_id or payload.get("chat_id") logger.info("android: history request from %s chat_id=%r", device_id, chat_id) if not isinstance(chat_id, str) or not chat_id.strip(): await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "history requires a chat_id", id=frame.id), ) return chat_id = chat_id.strip() thread_id = frame.thread_id or payload.get("thread_id") if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None before = payload.get("before_message_id") if not isinstance(before, str) or not before.strip(): before = None limit_raw = payload.get("limit") try: limit = int(limit_raw) if limit_raw is not None else 50 except (TypeError, ValueError): limit = 50 page = self._outbox.history( chat_id, thread_id=thread_id, before_message_id=before, limit=limit, ) resp = protocol.history( chat_id, page["messages"], page["has_more"], thread_id=thread_id, oldest_message_id=page["oldest_message_id"], id=frame.id, ) await self._ws_server.send_to(device_id, resp) # ── Message deletion (app -> agent) ─────────────────────────────────── async def on_message_delete(self, frame: protocol.Frame, device_id: str) -> None: """Handle an inbound ``message.delete`` request. Removes the requested message(s) from the outbox (so ``history`` and ``sync`` no longer return them) and broadcasts ``message.deleted`` to every device (outboxed too, so an offline device learns of the deletion on its next ``sync``). Deleting is idempotent: a message that is already gone (pruned by retention) simply yields 0 removed rows, and the ``message.deleted`` broadcast is still emitted so live caches drop it. """ payload = frame.payload chat_id = frame.chat_id or payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_NOT_FOUND, "message.delete requires chat_id", id=frame.id), ) return chat_id = chat_id.strip() thread_id = frame.thread_id or payload.get("thread_id") if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None message_ids = payload.get("message_ids") if not isinstance(message_ids, list): message_ids = [payload.get("message_id")] if payload.get("message_id") else [] message_ids = [m for m in message_ids if isinstance(m, str) and m.strip()] if not message_ids: await self._ws_server.send_to( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "message.delete requires message_ids", id=frame.id), ) return removed = 0 for mid in message_ids: removed += self._outbox.delete_message(chat_id, mid, thread_id=thread_id) logger.info( "android: message.delete from %s chat_id=%r thread_id=%r ids=%s removed=%s", device_id, chat_id, thread_id, message_ids, removed, ) resp = protocol.message_deleted(chat_id, message_ids, thread_id=thread_id) resp.id = frame.id await self._broadcast_or_log(chat_id, resp) # ── M5: push token registration ─────────────────────────────────────── async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None: """Update the device's push tokens (FCM rotation / ntfy topic). Persists to the device registry AND refreshes the live connection so the next push targets the current token without a stale read. """ fcm_token = frame.payload.get("fcm_token") ntfy_topic = frame.payload.get("ntfy_topic") fcm_token = fcm_token if isinstance(fcm_token, str) and fcm_token else None ntfy_topic = ntfy_topic if isinstance(ntfy_topic, str) and ntfy_topic else None if fcm_token is None and ntfy_topic is None: return try: self._devices.update_push_tokens( device_id, fcm_token=fcm_token, ntfy_topic=ntfy_topic ) except Exception: logger.warning("android: fcm.register update failed", exc_info=True) return conn = self._ws_server.connection(device_id) if conn is not None: if fcm_token is not None: conn.fcm_token = fcm_token if ntfy_topic is not None: conn.ntfy_topic = ntfy_topic logger.info("android: push tokens updated for %s", device_id) # ── M5: approval / clarify banners ──────────────────────────────────── async def send_slash_confirm( self, chat_id: str, title: str, message: str, session_key: str, confirm_id: str, metadata: Optional[Dict[str, Any]] = None, ) -> SendResult: """Banner + push for a slash-command approval prompt. The gateway's text fallback still renders the actionable prompt (the app has no inline buttons yet); the notification is the push-visible signal (high priority: pushed even when a device is live). """ thread_id = _thread_id_from_metadata(metadata) await self._broadcast_or_log( chat_id, protocol.notification( chat_id, protocol.NOTIF_APPROVAL, title or "Approval needed", _push_preview(message), thread_id=thread_id, ), ) return await super().send_slash_confirm( chat_id, title, message, session_key, confirm_id, metadata=metadata ) async def send_clarify( self, chat_id: str, question: str, choices: Optional[list], clarify_id: str, session_key: str, metadata: Optional[Dict[str, Any]] = None, ) -> SendResult: """Banner + push for a clarify prompt. Renders the prompt as a proper ``message`` frame (the base text fallback would route through ``send()`` and be misclassified as commentary/tool progress) and keeps the gateway's text intercept working via ``mark_awaiting_text``. """ thread_id = _thread_id_from_metadata(metadata) await self._broadcast_or_log( chat_id, protocol.notification( chat_id, protocol.NOTIF_CLARIFY, "Question", _push_preview(question), thread_id=thread_id, ), ) if choices: # Multi-select clarifies register their flag on the pending entry; # look it up by id (mirrors the base text fallback). _is_multi = False try: from tools import clarify_gateway as _cg with _cg._lock: _entry = _cg._entries.get(clarify_id) _is_multi = bool(_entry and getattr(_entry, "multi_select", False)) except Exception: _is_multi = False lines = [f"❓ {question}", ""] for i, choice in enumerate(choices, start=1): lines.append(f" {i}. {choice}") lines.append("") if _is_multi: lines.append( "Multiple selections allowed — reply with the numbers " "separated by commas or spaces (e.g. \"1, 3\"), the option " "text, or your own answer." ) else: lines.append("Reply with the number, the option text, or your own answer.") text = "\n".join(lines) # Text fallback: enable text-capture so the gateway intercept # picks up the user's typed reply (e.g. "2" or choice text). from tools.clarify_gateway import mark_awaiting_text mark_awaiting_text(clarify_id) else: text = f"❓ {question}" message_id = _mint_message_id() await self._broadcast_or_log( chat_id, protocol.message( chat_id=chat_id, message_id=message_id, role=protocol.ROLE_ASSISTANT, text=text, thread_id=thread_id, ts=int(time.time() * 1000), ), ) return SendResult(success=True, message_id=message_id) # ── 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, } # ── hello.ack helpers ───────────────────────────────────────────────── def gateway_status(self) -> str: """Current gateway health state (sent to each pairing connection).""" return self._gateway_status def server_caps(self) -> Dict[str, Any]: """Capability flags advertised in ``hello.ack`` (M4 surface).""" return { "streaming": True, # M2: message.start/update/stop "reasoning": True, # M2: reasoning field on message / message.stop "tools": True, # M2: tool.start/progress/end "media": True, # M4: media.upload/offer/pull "search": True, # M3: search frame "push": self.push_backend, # M5: ntfy server URL (app listener discovery; "" when not ntfy). "push_ntfy_server": ( self._push.server_url if isinstance(self._push, NtfyBackend) else "" ), "pickers": False, # M2+ } 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 ``android::`` 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 ) -> Optional[str]: """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("android: create_handoff_thread failed", exc_info=True) return None await self._ws_server.broadcast(protocol.channel_created(entry)) return entry["chat_id"] # --------------------------------------------------------------------------- # Plugin entry point # --------------------------------------------------------------------------- def register(ctx): """Plugin entry point: called by the Hermes plugin system.""" # M2: capture the model's separate reasoning_content during streaming so # it can be attached to the turn's message.stop frame (the gateway # otherwise drops it when streaming suppresses the final send). try: ctx.register_hook("on_stream_delta", _on_stream_delta) except Exception: logger.debug("android: on_stream_delta hook registration failed", exc_info=True) # M2: capture each completed tool call's result + timing so the tool.end # frame can carry the output (the gateway never streams tool output to # platforms). The app shows it on demand (Settings → Tool detail). try: ctx.register_hook("post_tool_call", _on_post_tool_call) except Exception: logger.debug("android: post_tool_call hook registration failed", exc_info=True) ctx.register_platform( name="android", label="Android", adapter_factory=lambda cfg: AndroidAdapter(cfg), check_fn=check_requirements, validate_config=validate_config, is_connected=is_connected, required_env=["ANDROID_TOKEN"], install_hint="No extra packages needed (websockets + httpx are core deps)", setup_fn=interactive_setup, # Env-driven auto-configuration: seeds PlatformConfig.extra with # host/port/push_backend + home_channel so env-only setups show up in # gateway status without instantiating the adapter. env_enablement_fn=_env_enablement, # Cron home-channel delivery support (deliver=android:[:]). cron_deliver_env_var="ANDROID_HOME_CHANNEL", # Out-of-process cron delivery (best-effort; outbox is gateway-served). standalone_sender_fn=_standalone_send, # Native target syntax: "android:[:]". parse_target_ref_fn=_parse_target_ref, # Auth env vars for _is_user_authorized() integration. allowed_users_env="ANDROID_ALLOWED_USERS", allow_all_env="ANDROID_ALLOW_ALL_USERS", # WS has no message-size limit. max_message_length=0, # Display. emoji="📱", pii_safe=False, allow_update_command=True, # LLM guidance. platform_hint=( "You are chatting with the user through their native Iris app " "(Android/Desktop). It renders Markdown, inline code, images, " "audio and video, and shows your reasoning and tool activity. " "Conversations are organized into channels and optional threads. " "Keep formatting rich but readable. " "You can send media files natively: to deliver a file to the user, " "include MEDIA:/absolute/path/to/file in your response. Images " "(.png, .jpg, .webp) appear as photos, audio (.ogg, .mp3, .m4a) " "plays inline, videos (.mp4, .webm, .mov) play inline, and other " "files arrive as downloadable documents." ), )