From e5c7d690b896b334a788bf8d6d7b11fd7039ff13 Mon Sep 17 00:00:00 2001 From: ARIA Date: Sat, 22 Aug 2026 14:10:14 +0200 Subject: [PATCH] HTTP fallback leg (docs/19): POST /v1/frame + SSE /v1/events + long-poll Second, short-lived-connection transport next to the WS: same frames, same outbox/cursor, same token, served over plain HTTP (stdlib ThreadingHTTPServer bridged into the asyncio loop; zero new deps). - http_server.py: /v1/health (unauthenticated), POST /v1/frame (accept-and-ack; validation rejections as 4xx error frames), SSE /v1/events (outbox catch-up with id=cursor, event: hello, 15s heartbeat, bounded-queue backpressure), long-poll /v1/poll (25s hold). Bearer token + X-Iris-Device (same allowlist as WS hello), 64 KiB body cap, per-device rate limit, optional TLS, non-fatal bind failure. - ws_server.py: inbound dispatch chain extracted to shared dispatch_frame() used by both transports. - adapter.py: ANDROID_HTTP_PORT/CERT/KEY config; start/stop next to the WS; delivery counting in _broadcast_or_log (an SSE subscriber is a live subscriber -> no push, docs/19 19.8); _reply() routes point-to-point replies into the in-flight HTTP response (reply sink) or broadcasts when the device has no live WS (19.7); status/typing/ channel events fan out to both transports. - ws_probe.py: --http mode (health + POST + SSE turn drive, same assertion flags); tests/README updated. - Tests: hermes-agent/tests/gateway/test_android_http.py (23 tests, incl. the 19.8 delivery-counting regression); test_android.py (74) still green. --- docs/19-http-fallback-transport.md | 318 ++++++++++++++ docs/README.md | 7 +- gateway-plugin/adapter.py | 176 +++++--- gateway-plugin/http_server.py | 666 +++++++++++++++++++++++++++++ gateway-plugin/tests/README.md | 10 +- gateway-plugin/tests/ws_probe.py | 144 +++++++ gateway-plugin/ws_server.py | 83 ++-- 7 files changed, 1315 insertions(+), 89 deletions(-) create mode 100644 docs/19-http-fallback-transport.md create mode 100644 gateway-plugin/http_server.py diff --git a/docs/19-http-fallback-transport.md b/docs/19-http-fallback-transport.md new file mode 100644 index 0000000..0e30ce0 --- /dev/null +++ b/docs/19-http-fallback-transport.md @@ -0,0 +1,318 @@ +# 19 — HTTP Fallback Transport (the "HTTP leg") + +A second, **short-lived-connection** transport next to the WebSocket: the same +JSON frames, the same outbox/cursor, the same token — served over plain HTTP +by the gateway. When the WS is down (flaky network, NAT timeout, app just +relaunched), the app **sends over `POST` and receives over SSE** instead of +waiting 2–20 s for a WS redial. + +Status: **design (proposed, not built)**. Complements — does not replace — +`04-wire-protocol.md` (frames), `08-push.md` (outbox/sync/push), and +`09-pairing-security.md` (auth model). + +## 19.1 Problem + +Today the WS is the *only* transport, and the app hard-gates sending on a live +socket (`ChatScreen.doSend()` no-ops unless `State.Connected`; +`GatewayClient.sendMessage()` drops when `socket == null`). Consequences: + +- **App killed → reopened:** full cold dial (TCP + TLS + `hello`/`hello.ack`, + 15 s dial timeout) before the user can send. On a flaky network the first + dial often fails → backoff → second dial. Observed: 2–20 s of "can't send". +- **Long-lived WS is the most fragile connection type on mobile:** idle + sockets expire in router/CGNAT NAT tables, die on WiFi↔cellular handover, + and are killed aggressively by OEM power management (MIUI on the test + device). There is no foreground service holding the WS. +- **Stale detection is slow:** 20 s ping interval, 60 s reap — a dead-but- + unclosed socket can sit for up to a minute before redial. + +## 19.2 Why HTTP (and why not the alternatives) + +Short-lived HTTP requests are dramatically more resilient on mobile networks +than a long-lived socket: no NAT table entry to expire, no proxy idle-kill, +each request is a fresh connection (fast with TLS resumption), and they work +through the restrictive proxies that mangle WebSockets. Sending a message +becomes a single `POST` that completes in well under a second on a LAN — +independent of whether the WS is up. + +Alternatives considered and rejected (research, 2026-08): + +| Option | Verdict | +| --- | --- | +| **MQTT broker** (QoS 1, persistent sessions) | Best protocol for flaky links, but new infra (broker process) + new Python dep (`paho-mqtt`, breaks the zero-new-deps rule) + new Kotlin dep + frame↔topic bridge. Overkill for a 1-user agent. | +| **ntfy as the send path** (app publishes to a topic the gateway subscribes to) | Adds a third party to the critical send path; public ntfy.sh is already known-flaky. Not worth it. | +| **WebTransport / QUIC** | The *real* fix for handover flakiness (connection migration), but no OkHttp support and `aioquic` is a new Python dep. Future option if this doc's approach is still not enough. | +| **gRPC** | New deps both sides; no advantage over WS+SSE here. | +| **Inverted connection** (app runs a local HTTP server, gateway pushes to the phone) | LAN-only, breaks on cellular/remote, security mess. Rejected. | +| **Matrix / full chat server** | Massive overkill for a personal agent. | + +**Zero new Python dependencies is preserved:** the HTTP leg is stdlib +`http.server` (a `ThreadingHTTPServer` in a daemon thread) bridged into the +gateway's asyncio loop. The app side uses the OkHttp it already depends on +(hand-rolled SSE reader — the format is trivial; `okhttp-eventsource` is an +acceptable alternative if preferred). + +## 19.3 Shape + +``` + ┌──────────────────────── hermes gateway process ───────────────────────┐ + │ AndroidAdapter │ + │ │ frames (same protocol.Frame objects) │ + │ ▼ │ + │ _broadcast_or_log ──► outbox.append(cursor) ──► push (if no live) │ + │ │ │ │ + │ ▼ ▼ │ + │ WsServer (asyncio, :8790) HttpServer (stdlib thread, :8791) │ + │ primary: full protocol fallback: POST /v1/frame, │ + │ incl. binary media GET /v1/events (SSE), /v1/poll │ + └───────────────┬──────────────────────────────┬────────────────────────┘ + │ WS (primary) │ HTTP (fallback) + ┌─────────┴──────────────────────────────┴────────┐ + │ APP: transport state machine │ + │ WS up → WS only (media works, lowest latency)│ + │ WS down → send via POST, receive via SSE/poll │ + └─────────────────────────────────────────────────┘ +``` + +**v1 scope** + +| Over HTTP (v1) | WS-only (v1) | +| --- | --- | +| All JSON request frames (`message.send`, `search`, `channel.*`, `commands.catalog`, `agent.stop`/`agent.steer`, …) via one generic endpoint | Binary media upload (chunked binary frames) | +| All event/response frames via SSE (or long-poll) | Binary media pull stream | +| `sync` catch-up (same outbox, same cursor) | — | + +Media stays WS-only in v1: it is the one part of the protocol that is +inherently binary/streaming, and attachments are a rarer action than sending +text. While in HTTP-fallback mode the composer disables the attach button +("media needs the live connection"). HTTP media endpoints are a v2 item +(§19.13). + +## 19.4 Gateway: `gateway-plugin/http_server.py` + +New module, started/stopped by `AndroidAdapter.connect()`/`disconnect()` next +to the WS server. + +- **Server:** `http.server.ThreadingHTTPServer` + `BaseHTTPRequestHandler`, + run in a **daemon thread** (one thread per connection — fine at single-user + scale). The handler thread never touches adapter state directly; it bridges + into the gateway's asyncio loop with + `asyncio.run_coroutine_threadsafe(coro, loop)` (the loop is captured at + start, same loop the WS server runs on). +- **Config:** `ANDROID_HTTP_PORT` (default **8791**), same bind host as the WS + (`ANDROID_WS_HOST`). Optional TLS via `ANDROID_HTTP_CERT`/`ANDROID_HTTP_KEY` + (`ssl.SSLContext` on the server) — same posture as the WS: plaintext on a + trusted LAN by default, TLS for remote/Tailscale setups. +- **Bind failure is NON-fatal** (unlike the WS): log a warning, disable the + HTTP leg, show it in the inspector. The plugin must keep working WS-only. +- Port-conflict lock: same flock pattern the WS uses (`host:port` key). + +### Endpoints + +| Endpoint | Auth | Purpose | +| --- | --- | --- | +| `GET /v1/health` | none | Liveness probe → `200 {"ok": true}`. Leaks nothing (no token echo, no device info). The app races this against the WS dial at startup. | +| `POST /v1/frame` | Bearer token | Accept **any** JSON frame the WS accepts (except binary media). Body = one frame envelope (`04-wire-protocol.md`). Dispatched through the *same* adapter handlers as WS (`on_message_send`, `on_search`, …). | +| `GET /v1/events?cursor=N` | Bearer token | **SSE** stream: catch-up from the outbox, then live frames (§19.5). | +| `GET /v1/poll?cursor=N` | Bearer token | **Long-poll** fallback where SSE is blocked (§19.6). | + +### Auth & limits + +- `Authorization: Bearer `; verified with the existing constant-time + `verify_token()`; `401` on failure. Device identity via `X-Iris-Device` + header (same `device_id` the app uses for `hello`; same allowlist check). +- Request body cap **64 KiB**, `Content-Type: application/json` enforced + (frames are small; media never travels here in v1). +- Rate limit: token bucket per device, same parameters as the WS inbound + limit (`INBOUND_RATE_PER_S` / `INBOUND_BURST`); `429` on exceed. +- No CORS headers (app clients only); unknown paths → `404`. + +## 19.5 SSE stream design (`GET /v1/events`) + +Wire format (standard SSE, three fields): + +``` +id: 1043 +event: frame +data: {"v":1,"type":"message","chat_id":"android:default",...} + +: hb ← comment heartbeat every 15 s (keeps proxies alive) +``` + +- **`id` = outbox cursor.** This is what makes resume trivial: on reconnect + the client sends `Last-Event-ID` (or `?cursor=`) and the server replays + `outbox.replay(cursor)` — exactly the `sync` semantics, no new machinery. +- **Stream open sequence:** + 1. Replay outbox rows with `cursor > N` (bounded by the existing + `_REPLAY_LIMIT`), each as an `event: frame` with its `id`. + 2. One `event: hello` carrying the `hello.ack` payload + (`server_caps`, `sync_cursor`, `last_pushed_cursor`, `channels`) — the + HTTP equivalent of pairing-ack; the app treats it like `hello.ack`. + 3. Live frames as they are produced. +- **Live fan-out hook:** in `adapter._broadcast_or_log`, after + `outbox.append()` returns the cursor, push `(cursor, frame_json)` into every + live HTTP subscriber's queue. The direct `status` broadcasts + (`ws_server.broadcast(protocol.status(...))`) get a second fan-out call + with `cursor = None` (SSE event without `id`). +- **Thread model:** each SSE connection owns its handler thread, which blocks + on a cross-thread `queue.get()` (via `run_coroutine_threadsafe`, 30 s + timeout → write `: hb` and loop) and writes to `wfile` + `flush()`. +- **Backpressure:** bounded queue (256). A subscriber that can't keep up is + dropped; the client reconnects with `Last-Event-ID` and catches up from the + outbox. Single-user scale makes this a non-event in practice. +- **App-side reader:** hand-rolled over OkHttp's streaming `ResponseBody` + (read lines; `id:` / `event:` / `data:`; blank line = dispatch). ~100 lines, + no new dependency. Reconnect with exponential backoff + `Last-Event-ID`. + +## 19.6 Long-poll fallback (`GET /v1/poll`) + +For networks/proxies that buffer or kill SSE: + +- `GET /v1/poll?cursor=N` → server holds the request (asyncio waiter on the + subscriber queue) until a frame with `cursor > N` exists or **25 s** pass. +- Response: `200 {"cursor": , "frames": [ ... ]}` (frames may + be empty on timeout; the app immediately re-polls with the new cursor). +- The app switches to long-poll automatically after **two consecutive SSE + open failures**, and back to SSE on the next full (re)connect. + +## 19.7 Request/response over HTTP + +`POST /v1/frame` is accept-and-ack: + +- `202 {"ok": true}` — frame accepted and dispatched. +- `4xx` with an **error frame as the JSON body** for validation rejections + (empty message, automation-channel read-only, rate limit → `429`, bad JSON + → `400`). These are the same `error` frames the WS path sends via + `send_to`; over HTTP they double as the HTTP response. +- **Async responses** (user echo, `search` results, `channel.list`, the agent + reply, streaming updates) arrive on the **event stream** carrying the same + `id` — the app's existing request-id correlation works unchanged. +- Consequence: `send_to(device_id, …)` error replies for HTTP-originated + requests are instead **broadcast** (single-user model; the SSE stream + delivers them). The dispatch refactor must tag the origin so WS-originated + requests keep point-to-point errors. + +## 19.8 Delivery counting & push interaction (critical) + +`_broadcast_or_log` fires push when `delivered == 0`. With the HTTP leg, a +device reading SSE **is** a live subscriber: + +```python +delivered = await self._ws_server.broadcast(frame) +delivered += await self._http_server.fanout(frame, cursor) # live SSE/poll subs +... +if delivered == 0: # → outbox + push (unchanged) +``` + +If this is forgotten, every message would push *and* stream to a device that +is already receiving it. Related bookkeeping: + +- `has_devices()` / `status` must count HTTP subscribers as connected devices + (mark the device's transport `ws` | `http` in the connection registry). +- `last_pushed_cursor` / notification dedupe (`08-push.md` §8.8) is + unchanged — SSE-replayed frames carry the same `cursor` envelope as + `sync`-replayed ones, so the app's existing dedupe applies. + +## 19.9 App side + +New `iris/net/HttpGateway.kt` (OkHttp) + a transport state machine inside +`GatewayClient` (or a thin `Transport` wrapper around it): + +- **API:** `health(timeoutMs)`, `postFrame(json): Result`, + `events(cursor, onFrame, onHello): Job` (SSE reader), `poll(cursor): Result`. +- **State machine:** + + | State | Send path | Receive path | + | --- | --- | --- | + | `WS_CONNECTED` | WS frame | WS | + | `HTTP_FALLBACK` | `POST /v1/frame` | SSE (or long-poll) | + | `CONNECTING` / `RECONNECTING` | queued/dropped as today | — | + +- **On WS loss:** switch to `HTTP_FALLBACK` **immediately** — open the SSE + stream (catch-up from the local cursor is free) and route sends to POST. + No backoff gate on the send path; the WS redial loop keeps running in the + background. +- **At startup (the key UX fix):** race the WS dial against + `GET /v1/health` (2 s timeout). WS dial fails + health OK → straight into + `HTTP_FALLBACK`: the user can send in **< 1 s** after opening the app, + instead of waiting out dial timeouts and backoff. +- **On WS reconnect:** close the SSE stream, resume WS-only (lowest latency, + media available again). +- **Send path:** `sendMessage()` builds the same `message.send` frame JSON and + writes it to WS or POST depending on state. The `State.Connected` gate in + `ChatScreen.doSend()` becomes `state is Connected || state is HttpFallback`. +- **Media:** disabled in the composer while in `HTTP_FALLBACK` (v1). +- **UI:** status pill shows "connected" (WS) or "connected · http" (fallback) + — both green; the fallback is a healthy state, not an error. + +## 19.10 Security + +- Same token, constant-time verify, same bind host, same device allowlist as + the WS (`09-pairing-security.md` threat model unchanged — the HTTP leg adds + no new trust boundary, only a second door with the same lock). +- `/v1/health` is unauthenticated by design (it answers "is the gateway + alive?"); it must not reflect tokens, device ids, or version strings. +- TLS: optional, same cert pattern as the WS; plaintext is a LAN-only + default, identical to today's WS posture. +- New attack-surface items to keep small: 64 KiB body cap, strict + content-type, per-device rate limit, no directory listing, no CORS. + +## 19.11 Failure modes + +| Failure | Behavior | +| --- | --- | +| Gateway fully down | Both legs dead → app shows offline; sends queue (app-side outbox, follow-up work) or are dropped with a visible "not sent" state. Push is the wake path when the gateway comes back (`08-push.md`). | +| WS down, HTTP up | Normal `HTTP_FALLBACK` operation — text chat fully functional, media paused. | +| SSE blocked by a proxy | Two failures → long-poll loop (§19.6). | +| HTTP port firewalled, WS up | WS-only operation (today's behavior); `health` fails at startup, no fallback attempted. | +| Both flaky | Existing WS backoff + SSE/poll backoff run independently; outbox + cursor keep both paths idempotent. | +| Slow SSE subscriber | Dropped at queue overflow; reconnects with `Last-Event-ID`, catches up from outbox. | + +## 19.12 Testing + +- **Python** (`hermes-agent/tests/gateway/test_android_http.py`, run via + `scripts/run_tests.sh`): + - auth: bad/missing token → 401; allowlist rejection; constant-time verify reused. + - `POST /v1/frame`: valid `message.send` dispatches (agent turn fires); + empty text → 400 error frame; automation channel → 409/400; rate limit → 429. + - SSE: catch-up rows carry correct `id`s; `event: hello` present; a live + frame appended after connect arrives on the stream; `Last-Event-ID` + resume replays exactly the delta; heartbeat observed within 15 s. + - long-poll: returns on new frame; empty 200 at timeout with advanced cursor. + - **delivery counting:** frame with only an SSE subscriber → `delivered ≥ 1` + → **no push fired** (the critical regression test for §19.8). +- **Probe:** `ws_probe.py` gains an `--http` mode (health, post, SSE read with + assertion flags, per `gateway-plugin/tests/README.md`). +- **Kotlin** (`:shared` commonTest): SSE parser (multi-line data, comments, + `Last-Event-ID` bookkeeping); transport state machine transitions (fake + clock: WS-loss → immediate fallback; startup race → fallback in < 1 s). +- **E2E** (`e2e.py`, new scenario): point the app at a dead WS port with the + HTTP leg live → send a message → assert user echo + agent reply arrive via + SSE; timing assertion: send → user echo < 1 s on LAN. Live-verify on the + device via ADB (screenshot of the "connected · http" pill). + +## 19.13 Non-goals (v1) / future + +- **Media over HTTP** (v2): `POST /v1/media` (chunked, same sha256 contract as + `07-media.md`) + `GET /v1/media/{id}` for pull/playback. Unblocks + attachments in fallback mode. +- **App-side send outbox** (companion work, separate doc): queue sends locally + when *both* legs are down; drains over whichever leg recovers. This doc + removes the 2–20 s wait; the outbox removes the last "gateway was down for + 30 s" data-loss case. +- **QUIC / WebTransport** if handover flakiness persists after this + the + outbox (connection migration would make the fallback rare). +- Per-device tokens (`16-open-questions.md` #3) apply to both legs identically + when implemented. + +## 19.14 Effort & change list + +| Slice | Files | Est. | +| --- | --- | --- | +| Gateway leg | new `gateway-plugin/http_server.py` (~450 lines); `adapter.py` hooks (start/stop, fan-out in `_broadcast_or_log` + status path, delivery counting, dispatch-origin tag); `protocol.py` unchanged | 2–3 d | +| App leg | new `app/shared/.../net/HttpGateway.kt` (SSE reader + poll); `GatewayClient.kt` state machine + startup race; `ChatScreen.kt` gate + status pill; composer media-disable in fallback | 2–3 d | +| Tests + e2e + docs | per §19.12; `frames.schema.json` unchanged (no new frame types); `09-pairing-security.md` + `13-testing.md` cross-references | 1–2 d | + +Total: **~1 week**, each slice independently shippable (gateway leg is +inert until the app uses it; app leg degrades to today's behavior if the +HTTP port is closed). diff --git a/docs/README.md b/docs/README.md index e12ab85..5fcd02b 100644 --- a/docs/README.md +++ b/docs/README.md @@ -8,6 +8,7 @@ This folder is the single source of truth for *what to build and why*. Read it top-to-bottom once, then use the numbered docs as a lookup while implementing. > ⚠️ **READ FIRST — two hard rules** +> > 1. **`hermes-agent/` (sibling of this folder) is a read-only research > reference. It must NEVER be committed, pushed, or shipped.** It is > git-ignored at the repo root. We only *install* our plugin into a live @@ -20,7 +21,7 @@ top-to-bottom once, then use the numbered docs as a lookup while implementing. ## Reading order | # | File | When to read | -|---|------|--------------| +| --- | ------ | -------------- | | 0 | [`00-overview.md`](00-overview.md) | Always first. Vision, scope, disclaimers, locked decisions. | | 1 | [`01-architecture.md`](01-architecture.md) | Before touching code. System shape + rationale. | | 2 | [`02-monorepo.md`](02-monorepo.md) | When scaffolding the repo. | @@ -39,8 +40,10 @@ top-to-bottom once, then use the numbered docs as a lookup while implementing. | 15 | [`15-hermes-reference.md`](15-hermes-reference.md) | **Cheat-sheet** of hermes-agent source to read. | | 16 | [`16-open-questions.md`](16-open-questions.md) | Decisions made + open items. | | 17 | [`17-future-control-surface.md`](17-future-control-surface.md) | **Backlog** — what the app could control beyond chat (cron, kanban, models, …). | +| 19 | [`19-http-fallback-transport.md`](19-http-fallback-transport.md) | **Design** — HTTP fallback leg (POST + SSE/long-poll) so the app can send/receive when the WS is down. | Machine-readable / diagrams: + - [`protocol/frames.schema.json`](protocol/frames.schema.json) — wire-frame schema. - [`diagrams/architecture.mmd`](diagrams/architecture.mmd) — mermaid architecture. @@ -65,4 +68,4 @@ project** (`app/`) with a shared KMP module (`app/shared`). - **Phase:** M0–M6 complete; M7 (polish + E2E + docs) in progress. - **Owner decisions locked:** see [`16-open-questions.md`](16-open-questions.md). -- **Last updated:** 2026-08-19. \ No newline at end of file +- **Last updated:** 2026-08-19. diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index 2293f68..66b69f4 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -62,6 +62,7 @@ import contextlib import json import logging import os +import queue import re import threading import time @@ -117,6 +118,7 @@ from . import protocol # noqa: E402 from . import purge as purge_bridge # noqa: E402 from . import search as search_bridge # noqa: E402 from .channels import get_directory # noqa: E402 +from .http_server import HttpServer # noqa: E402 from .outbox import Outbox # noqa: E402 from .pairing import ( # noqa: E402 DeviceRegistry, @@ -450,6 +452,7 @@ async def _build_runtime_footer(meta: dict[str, Any]) -> dict[str, Any]: DEFAULT_HOST = "127.0.0.1" DEFAULT_PORT = 8790 +DEFAULT_HTTP_PORT = 8791 # docs/19: HTTP fallback leg DEFAULT_HOME_CHANNEL = "android:default" DEFAULT_HOME_CHANNEL_NAME = "Default" DEFAULT_PUSH_BACKEND = "fcm" @@ -1090,6 +1093,10 @@ class AndroidAdapter(BasePlatformAdapter): self.port = _parse_port( os.getenv("ANDROID_WS_PORT", "") or str(extra.get("port", DEFAULT_PORT)) ) + # docs/19: HTTP fallback leg (same bind host as the WS; optional TLS). + self.http_port = _parse_port( + os.getenv("ANDROID_HTTP_PORT", "") or str(extra.get("http_port", DEFAULT_HTTP_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 @@ -1122,6 +1129,8 @@ class AndroidAdapter(BasePlatformAdapter): # 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", "") + self.http_cert = _get_scoped_secret("ANDROID_HTTP_CERT") or extra.get("http_cert", "") + self.http_key = _get_scoped_secret("ANDROID_HTTP_KEY") or extra.get("http_key", "") # Auth allowed = os.getenv("ANDROID_ALLOWED_USERS", "").strip() @@ -1133,6 +1142,13 @@ class AndroidAdapter(BasePlatformAdapter): # Runtime state self._devices = DeviceRegistry(get_hermes_home() / "android" / "devices.db") self._ws_server = WsServer(self, self._devices) + # docs/19: HTTP fallback leg (inert until the app uses it; a bind + # failure disables it without affecting the WS). + self._http_server = HttpServer(self, self._devices) + # docs/19 §19.7: reply sinks for in-flight HTTP requests — while a + # POST /v1/frame is being dispatched, the handler's point-to-point + # replies are captured here and returned as the HTTP response. + self._http_reply_sinks: dict[str, tuple[queue.Queue, threading.Event]] = {} self._connected = False # M2: per-chat turn state for outbound frame classification. self._turns: dict[str, _TurnState] = {} @@ -1210,12 +1226,17 @@ class AndroidAdapter(BasePlatformAdapter): self._connected = False return False + # docs/19: start the HTTP fallback leg next to the WS. Bind failure + # is NON-fatal (unlike the WS): the plugin keeps working WS-only. + await self._http_server.start() + # M5: announce gateway health to connected clients (none yet at # startup; the frame + plumbing exist for future transitions). # Reset in case this adapter instance previously went down (the # gateway may reconnect the same adapter after a fatal error). self._gateway_status = protocol.STATUS_ONLINE await self._ws_server.broadcast(protocol.status(self._gateway_status)) + await self._http_server.fanout(protocol.status(self._gateway_status), cursor=None) # M3: ensure the default (home) channel exists in the directory so the # app's channel list and cron home delivery have a stable anchor. @@ -1248,6 +1269,7 @@ class AndroidAdapter(BasePlatformAdapter): self._gateway_status = protocol.STATUS_RESTARTING with contextlib.suppress(Exception): await self._ws_server.broadcast(protocol.status(self._gateway_status)) + await self._http_server.fanout(protocol.status(self._gateway_status), cursor=None) with contextlib.suppress(ImportError): from gateway.status import release_scoped_lock @@ -1258,6 +1280,10 @@ class AndroidAdapter(BasePlatformAdapter): await self._ws_server.stop() except Exception: logger.warning("android: WS server stop failed", exc_info=True) + try: + await self._http_server.stop() + except Exception: + logger.warning("android: HTTP server stop failed", exc_info=True) # Best-effort shutdown: a close failure on an already-closed store is # not actionable at disconnect time. with contextlib.suppress(Exception): @@ -1642,6 +1668,50 @@ class AndroidAdapter(BasePlatformAdapter): # (e.g. tool_progress off) so they can't leak into the next turn. _reset_tool_results() + # ── docs/19: HTTP-leg reply routing ─────────────────────────────────── + + def _http_register_sink( + self, device_id: str, entry: tuple[queue.Queue, threading.Event] + ) -> None: + self._http_reply_sinks[device_id] = entry + + def _http_pop_sink(self, device_id: str) -> tuple[queue.Queue, threading.Event] | None: + return self._http_reply_sinks.pop(device_id, None) + + def _http_pop_sink_if( + self, device_id: str, sink: queue.Queue + ) -> tuple[queue.Queue, threading.Event] | None: + """Pop the sink entry only if it is still ours (a newer request from + the same device may have replaced it).""" + entry = self._http_reply_sinks.get(device_id) + if entry is None or entry[0] is not sink: + return None + return self._http_reply_sinks.pop(device_id, None) + + async def _broadcast_both(self, frame: "protocol.Frame") -> None: + """Bare (non-outbox) broadcast to both transports (docs/19): the + frame reaches WS devices and live SSE/long-poll subscribers.""" + await self._ws_server.broadcast(frame) + await self._http_server.fanout(frame, cursor=None) + + async def _reply(self, device_id: str, frame: "protocol.Frame") -> None: + """Point-to-point reply with HTTP-leg fallback (docs/19 §19.7). + + WS-originated requests keep point-to-point delivery. For an + in-flight HTTP request (a reply sink is registered) the frame goes + into the HTTP response. If the device has no live WS and no sink + (e.g. it dropped mid-request), the frame is broadcast so the SSE + stream delivers it (single-user model). + """ + entry = self._http_reply_sinks.get(device_id) + if entry is not None: + entry[0].put(frame) + return + if await self._ws_server.send_to(device_id, frame): + return + await self._ws_server.broadcast(frame) + await self._http_server.fanout(frame, cursor=None) + 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 @@ -1656,6 +1726,10 @@ class AndroidAdapter(BasePlatformAdapter): except Exception: logger.warning("android: outbox append failed", exc_info=True) return + # docs/19 §19.8: a device reading SSE/long-poll IS a live subscriber + # — count it in the delivery total or every message would push AND + # stream to a device that is already receiving it. + delivered += await self._http_server.fanout(frame, cursor) if delivered == 0: logger.info( "android: no live devices for %s; %s frame parked in outbox (cursor=%s)", @@ -1819,11 +1893,15 @@ class AndroidAdapter(BasePlatformAdapter): 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)) + frame = protocol.typing(chat_id, True, thread_id=thread_id) + await self._ws_server.broadcast(frame) + await self._http_server.fanout(frame, cursor=None) 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)) + frame = protocol.typing(chat_id, False) + await self._ws_server.broadcast(frame) + await self._http_server.fanout(frame, cursor=None) # ── M4: outbound media (agent -> app) ───────────────────────────────── # @@ -1959,7 +2037,7 @@ class AndroidAdapter(BasePlatformAdapter): ) if not text.strip() and not media_refs: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "message.send requires non-empty text", id=frame.id @@ -1978,7 +2056,7 @@ class AndroidAdapter(BasePlatformAdapter): # enforcement). target = self._channels.get(chat_id) if target is not None and target.get("automation"): - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, @@ -2025,7 +2103,7 @@ class AndroidAdapter(BasePlatformAdapter): 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)) + await self._broadcast_both(protocol.channel_created(entry, auto=True)) self._schedule_thread_title_upgrade(entry["chat_id"], text) # M4: resolve media refs (single-use; unknown ref -> error). @@ -2035,7 +2113,7 @@ class AndroidAdapter(BasePlatformAdapter): for ref in media_refs: entry = self._media.get_inbound(ref) if entry is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, f"unknown media_ref {ref}", id=frame.id @@ -2120,7 +2198,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.read_receipt(chat_id, message_id), ) @@ -2157,7 +2235,7 @@ class AndroidAdapter(BasePlatformAdapter): return try: asyncio.run_coroutine_threadsafe( - self._ws_server.broadcast(protocol.channel_renamed(renamed)), + self._broadcast_both(protocol.channel_renamed(renamed)), loop, ) except Exception: @@ -2178,7 +2256,7 @@ class AndroidAdapter(BasePlatformAdapter): payload = frame.payload media_ref = str(payload.get("media_ref") or "").strip() if not media_ref or len(media_ref) > MAX_MEDIA_REF_LEN: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "media.upload.start requires media_ref", id=frame.id @@ -2187,7 +2265,7 @@ class AndroidAdapter(BasePlatformAdapter): return kind = payload.get("kind") if kind not in media_bridge.KINDS: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, f"unsupported media kind {kind!r}", id=frame.id @@ -2202,7 +2280,7 @@ class AndroidAdapter(BasePlatformAdapter): except (TypeError, ValueError): size = -1 if size <= 0: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, @@ -2212,7 +2290,7 @@ class AndroidAdapter(BasePlatformAdapter): ) return if size > self.max_upload_bytes: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_MEDIA_TOO_LARGE, @@ -2233,7 +2311,7 @@ class AndroidAdapter(BasePlatformAdapter): 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)) + await self._reply(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. @@ -2244,7 +2322,7 @@ class AndroidAdapter(BasePlatformAdapter): return # stray binary frame: ignore (forward-compat) session.feed(chunk) if session.failed: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(session.error_code, session.error_message, id=session.request_id), ) @@ -2255,7 +2333,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "media.upload.end requires media_ref", id=frame.id @@ -2265,9 +2343,9 @@ class AndroidAdapter(BasePlatformAdapter): 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)) + await self._reply(device_id, protocol.error(e.code, e.message, id=frame.id)) return - await self._ws_server.send_to( + await self._reply( device_id, protocol.media_upload_ack(True, entry.media_id, id=frame.id) ) @@ -2276,7 +2354,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, f"unknown media_id {media_id!r}", id=frame.id @@ -2287,7 +2365,7 @@ class AndroidAdapter(BasePlatformAdapter): # moved / been replaced since the offer). safe = validate_media_delivery_path(entry.path) if safe is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, "media no longer deliverable", id=frame.id), ) @@ -2299,11 +2377,11 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( 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)) + await self._reply(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.""" @@ -2320,7 +2398,7 @@ class AndroidAdapter(BasePlatformAdapter): payload = frame.payload name = payload.get("name") if not isinstance(name, str) or not name.strip(): - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "channel.create requires a name", id=frame.id @@ -2337,7 +2415,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id) ) return @@ -2358,7 +2436,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.rename requires chat_id", id=frame.id @@ -2367,7 +2445,7 @@ class AndroidAdapter(BasePlatformAdapter): return name = frame.payload.get("name") if not isinstance(name, str) or not name.strip(): - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "channel.rename requires a name", id=frame.id @@ -2377,12 +2455,12 @@ class AndroidAdapter(BasePlatformAdapter): try: entry = self._channels.rename(chat_id, name) except ValueError as e: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id) ) return if entry is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) @@ -2404,7 +2482,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.set_default requires chat_id", id=frame.id @@ -2413,7 +2491,7 @@ class AndroidAdapter(BasePlatformAdapter): return entry = self._channels.set_default(chat_id) if entry is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) @@ -2427,7 +2505,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.favorite requires chat_id", id=frame.id @@ -2437,7 +2515,7 @@ class AndroidAdapter(BasePlatformAdapter): on = bool(frame.payload.get("on")) entry = self._channels.set_favorite(chat_id, on) if entry is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) @@ -2451,7 +2529,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.icon requires chat_id", id=frame.id @@ -2465,14 +2543,14 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "channel icon too large", id=frame.id), ) return entry = self._channels.set_icon(chat_id, icon, color) if entry is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id), ) @@ -2484,7 +2562,7 @@ class AndroidAdapter(BasePlatformAdapter): async def on_channel_set_automation(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.set_automation requires chat_id", id=frame.id @@ -2494,7 +2572,7 @@ class AndroidAdapter(BasePlatformAdapter): on = bool(frame.payload.get("on")) entry = self._channels.set_automation(chat_id, on) if entry is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, @@ -2512,7 +2590,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "channel.delete requires chat_id", id=frame.id @@ -2521,7 +2599,7 @@ class AndroidAdapter(BasePlatformAdapter): return entry = self._channels.delete(chat_id) if entry is None: - await self._ws_server.send_to( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, @@ -2570,7 +2648,7 @@ class AndroidAdapter(BasePlatformAdapter): 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) + await self._reply(device_id, resp) # ── Slash-command catalog (app's "/" drawer) ────────────────────────── @@ -2581,7 +2659,7 @@ class AndroidAdapter(BasePlatformAdapter): the typed prefix client-side; the catalog is static per gateway run, so no caching is needed here.""" resp = protocol.commands_catalog(_slash_command_catalog(), id=frame.id) - await self._ws_server.send_to(device_id, resp) + await self._reply(device_id, resp) # ── M3: search (app -> agent) ───────────────────────────────────────── @@ -2589,7 +2667,7 @@ class AndroidAdapter(BasePlatformAdapter): payload = frame.payload query = payload.get("query") if not isinstance(query, str) or not query.strip(): - await self._ws_server.send_to( + await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "search requires a query", id=frame.id), ) @@ -2612,7 +2690,7 @@ class AndroidAdapter(BasePlatformAdapter): 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) + await self._reply(device_id, resp) # ── M3: sync (reconnect catch-up) ───────────────────────────────────── @@ -2639,9 +2717,9 @@ class AndroidAdapter(BasePlatformAdapter): cursor=e.get("cursor"), v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION, ) - await self._ws_server.send_to(device_id, replayed) + await self._reply(device_id, replayed) done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id) - await self._ws_server.send_to(device_id, done) + await self._reply(device_id, done) # ── Full message history (initial channel open / scroll-up) ─────────── @@ -2658,7 +2736,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "history requires a chat_id", id=frame.id), ) @@ -2689,7 +2767,7 @@ class AndroidAdapter(BasePlatformAdapter): oldest_message_id=page["oldest_message_id"], id=frame.id, ) - await self._ws_server.send_to(device_id, resp) + await self._reply(device_id, resp) # ── Message deletion (app -> agent) ─────────────────────────────────── @@ -2708,7 +2786,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "message.delete requires chat_id", id=frame.id @@ -2724,7 +2802,7 @@ class AndroidAdapter(BasePlatformAdapter): 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( + await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "message.delete requires message_ids", id=frame.id @@ -2987,7 +3065,7 @@ class AndroidAdapter(BasePlatformAdapter): except Exception: logger.warning("android: create_handoff_thread failed", exc_info=True) return None - await self._ws_server.broadcast(protocol.channel_created(entry)) + await self._broadcast_both(protocol.channel_created(entry)) return entry["chat_id"] diff --git a/gateway-plugin/http_server.py b/gateway-plugin/http_server.py new file mode 100644 index 0000000..e7435e6 --- /dev/null +++ b/gateway-plugin/http_server.py @@ -0,0 +1,666 @@ +"""HTTP fallback transport (docs/19): the "HTTP leg". + +A second, short-lived-connection transport next to the WebSocket: the same +JSON frames, the same outbox/cursor, the same token — served over plain HTTP +by the gateway. When the WS is down (flaky network, NAT timeout, app just +relaunched), the app sends over ``POST /v1/frame`` and receives over +``GET /v1/events`` (SSE) or ``GET /v1/poll`` (long-poll) instead of waiting +for a WS redial. + +Zero new Python dependencies: stdlib ``http.server`` (a +``ThreadingHTTPServer`` in a daemon thread) bridged into the gateway's +asyncio loop with ``asyncio.run_coroutine_threadsafe``. + +Endpoints (docs/19 §19.4): + * ``GET /v1/health`` — unauthenticated liveness probe. + * ``POST /v1/frame`` — accept-and-ack for any JSON frame the WS + accepts (except binary media, which stays + WS-only in v1). + * ``GET /v1/events?cursor=N`` — SSE stream: outbox catch-up, then live + frames (``id`` = outbox cursor, so resume + is just ``Last-Event-ID``). + * ``GET /v1/poll?cursor=N`` — long-poll fallback where SSE is blocked. + +Auth: ``Authorization: Bearer `` (constant-time ``verify_token``) + +``X-Iris-Device`` header (same device id / allowlist as the WS ``hello``). + +Bind failure is NON-fatal (unlike the WS): the plugin keeps working +WS-only. +""" + +from __future__ import annotations + +import asyncio +import contextlib +import json +import logging +import queue +import ssl +import threading +import time +from dataclasses import dataclass, field +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from typing import Any +from urllib.parse import parse_qs, urlparse + +from . import protocol +from .pairing import verify_token +from .ws_server import ( + INBOUND_BURST, + INBOUND_RATE_PER_S, + MAX_DEVICE_ID_LEN, + _TokenBucket, + dispatch_frame, +) + +logger = logging.getLogger(__name__) + +# Default port for the HTTP leg (WS default is 8790). +DEFAULT_HTTP_PORT = 8791 + +# Request body cap for POST /v1/frame (frames are small; media never +# travels here in v1). +MAX_BODY_BYTES = 64 * 1024 + +# Per-subscriber live-frame queue. A subscriber that can't keep up is +# dropped; it reconnects with Last-Event-ID and catches up from the outbox. +SUB_QUEUE_MAX = 256 + +# SSE comment heartbeat cadence (keeps proxies from idling the stream). +SSE_HEARTBEAT_S = 15.0 + +# Long-poll hold time (docs/19 §19.6). +POLL_TIMEOUT_S = 25.0 + +# How long POST /v1/frame waits for a synchronous validation rejection +# before acking 202 and letting the handler (e.g. the agent turn) run on. +ACCEPT_ACK_TIMEOUT_S = 5.0 + +# Sentinel pushed into subscriber queues on shutdown. +_STOP = object() + +# Frame types POST /v1/frame must not accept (docs/19 §19.3): media is +# inherently binary/streaming (WS-only in v1); hello/ping are +# transport-specific (auth is via headers, liveness via /v1/health). +HTTP_REJECTED_TYPES = frozenset( + { + protocol.TYPE_HELLO, + protocol.TYPE_PING, + protocol.TYPE_MEDIA_UPLOAD_START, + protocol.TYPE_MEDIA_UPLOAD_END, + protocol.TYPE_MEDIA_PULL, + } +) + + +def _with_cursor(frame: dict[str, Any], cursor: int) -> str: + """Re-serialize an outbox frame dict with its cursor in the envelope + (same tagging ``sync`` replay uses, docs/08 §8.7).""" + d = dict(frame) + d["cursor"] = cursor + return json.dumps(d, separators=(",", ":"), ensure_ascii=False) + + +def _parse_cursor(*raws: Any) -> int: + """First parseable non-negative int wins (``?cursor=`` beats + ``Last-Event-ID``); 0 when nothing usable.""" + for raw in raws: + if raw is None: + continue + try: + v = int(str(raw).strip()) + except (TypeError, ValueError): + continue + if v >= 0: + return v + return 0 + + +def _send_json(handler: BaseHTTPRequestHandler, status: int, obj: Any) -> None: + body = json.dumps(obj, separators=(",", ":")).encode("utf-8") + handler.send_response(status) + handler.send_header("Content-Type", "application/json") + handler.send_header("Content-Length", str(len(body))) + handler.end_headers() + with contextlib.suppress(BrokenPipeError, ConnectionResetError, OSError): + handler.wfile.write(body) + handler.wfile.flush() + + +def _send_frame_json(handler: BaseHTTPRequestHandler, status: int, frame_json: str) -> None: + """Send a protocol frame as the HTTP response body (docs/19 §19.7: + error frames double as the HTTP response).""" + body = frame_json.encode("utf-8") + handler.send_response(status) + handler.send_header("Content-Type", "application/json") + handler.send_header("Content-Length", str(len(body))) + handler.end_headers() + with contextlib.suppress(BrokenPipeError, ConnectionResetError, OSError): + handler.wfile.write(body) + handler.wfile.flush() + + +@dataclass +class _Subscriber: + """One live HTTP subscriber (SSE stream or long-poll request).""" + + device_id: str + kind: str # "sse" | "poll" + q: queue.Queue = field(default_factory=lambda: queue.Queue(maxsize=SUB_QUEUE_MAX)) + closed: threading.Event = field(default_factory=threading.Event) + + +class HttpServer: + """The plugin's HTTP fallback server + live subscriber registry. + + The handler threads never touch adapter state directly: inbound frames + are bridged into the gateway's asyncio loop (captured at ``start()``) + with ``asyncio.run_coroutine_threadsafe`` and dispatched through the + same ``dispatch_frame`` the WS server uses. + """ + + def __init__(self, adapter: Any, devices: Any): + self._adapter = adapter + self._devices = devices + self._loop: asyncio.AbstractEventLoop | None = None + self._httpd: _ThreadingHTTPD | None = None + self._thread: threading.Thread | None = None + self._subs: dict[str, list[_Subscriber]] = {} + self._subs_lock = threading.Lock() + self._buckets: dict[str, _TokenBucket] = {} + self._buckets_lock = threading.Lock() + self._lock_key: str | None = None + self.enabled = False + self.bound_port = 0 + + # ── Lifecycle ───────────────────────────────────────────────────────── + + async def start(self) -> None: + """Bind and start serving. NEVER raises: a bind failure disables the + HTTP leg (the plugin keeps working WS-only, docs/19 §19.4).""" + if self.enabled: + return + self._loop = asyncio.get_running_loop() + host = self._adapter.host + port = self._adapter.http_port + + # Port-conflict lock: same flock pattern the WS uses. + try: + from gateway.status import acquire_scoped_lock + + lock_key = f"http:{host}:{port}" + if not acquire_scoped_lock("android", lock_key): + logger.warning( + "android: HTTP port %s:%s in use by another profile; HTTP leg disabled", + host, + port, + ) + return + self._lock_key = lock_key + except ImportError: + self._lock_key = None # status module not available (e.g. tests) + + try: + httpd = _ThreadingHTTPD((host, port), self) + self.bound_port = int(httpd.server_address[1]) + if self._adapter.http_cert and self._adapter.http_key: + ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER) + ctx.load_cert_chain(self._adapter.http_cert, self._adapter.http_key) + httpd.socket = ctx.wrap_socket(httpd.socket, server_side=True) + except Exception as e: + logger.warning( + "android: HTTP fallback leg disabled (bind %s:%s failed: %s)", host, port, e + ) + self._release_lock() + return + + self._httpd = httpd + self._thread = threading.Thread( + target=httpd.serve_forever, name="android-http", daemon=True + ) + self._thread.start() + self.enabled = True + scheme = "https" if (self._adapter.http_cert and self._adapter.http_key) else "http" + logger.info( + "android: HTTP fallback leg listening on %s://%s:%s", scheme, host, self.bound_port + ) + + async def stop(self) -> None: + """Stop serving and unblock all subscribers.""" + self.enabled = False + with self._subs_lock: + subs = [s for lst in self._subs.values() for s in lst] + self._subs.clear() + for s in subs: + s.closed.set() + with contextlib.suppress(Exception): + s.q.put_nowait(_STOP) + httpd = self._httpd + self._httpd = None + if httpd is not None: + # shutdown() must be called from a thread other than the one + # running serve_forever(); we are on the asyncio loop thread. + with contextlib.suppress(Exception): + httpd.shutdown() + with contextlib.suppress(Exception): + httpd.server_close() + t = self._thread + self._thread = None + if t is not None and t is not threading.current_thread(): + t.join(timeout=5.0) + self._release_lock() + + def _release_lock(self) -> None: + with contextlib.suppress(ImportError): + from gateway.status import release_scoped_lock + + if self._lock_key: + release_scoped_lock("android", self._lock_key) + self._lock_key = None + + # ── Subscriber registry ─────────────────────────────────────────────── + + def has_devices(self) -> bool: + with self._subs_lock: + return bool(self._subs) + + def device_ids(self) -> list: + with self._subs_lock: + return list(self._subs.keys()) + + def _add_sub(self, sub: _Subscriber) -> None: + with self._subs_lock: + self._subs.setdefault(sub.device_id, []).append(sub) + + def _remove_sub(self, sub: _Subscriber) -> None: + with self._subs_lock: + lst = self._subs.get(sub.device_id) + if lst: + with contextlib.suppress(ValueError): + lst.remove(sub) + if not lst: + del self._subs[sub.device_id] + + # ── Outbound ────────────────────────────────────────────────────────── + + async def fanout(self, frame: protocol.Frame, cursor: int | None = None) -> int: + """Push a frame to every live HTTP subscriber. Returns subscribers + reached — the HTTP half of the delivery count (docs/19 §19.8): a + device reading SSE is a live subscriber, so a frame fanned out here + must not also fire a push.""" + if not self.enabled: + return 0 + data = frame.to_json() + with self._subs_lock: + subs = [s for lst in self._subs.values() for s in lst] + sent = 0 + for s in subs: + try: + s.q.put_nowait((cursor, data)) + sent += 1 + except queue.Full: + # Slow subscriber: drop it. The client reconnects with + # Last-Event-ID and catches up from the outbox. + logger.info("android: dropping slow HTTP subscriber %s", s.device_id) + s.closed.set() + self._remove_sub(s) + return sent + + # ── Auth / limits (handler threads) ─────────────────────────────────── + + def _authenticate(self, handler: BaseHTTPRequestHandler) -> str | None: + """Verify Bearer token + device identity. Returns the device_id, or + None after sending a 401.""" + auth = handler.headers.get("Authorization") or "" + token = auth[len("Bearer ") :] if auth.startswith("Bearer ") else None + if not verify_token(token, self._adapter.token): + _send_json(handler, 401, {"error": "unauthorized"}) + return None + device_id = (handler.headers.get("X-Iris-Device") or "").strip() + if not device_id or len(device_id) > MAX_DEVICE_ID_LEN: + _send_json(handler, 401, {"error": "X-Iris-Device header required"}) + return None + if ( + not self._adapter.allow_all + and self._adapter.allowed_users + and device_id not in self._adapter.allowed_users + ): + logger.warning("android: http rejected: device %s not allowlisted", device_id) + _send_json(handler, 401, {"error": "device not allowed"}) + return None + with contextlib.suppress(Exception): + self._devices.touch(device_id) + return device_id + + def _rate_limited(self, device_id: str) -> bool: + """Per-device token bucket, same parameters as the WS inbound limit.""" + with self._buckets_lock: + b = self._buckets.get(device_id) + if b is None: + b = self._buckets[device_id] = _TokenBucket(INBOUND_RATE_PER_S, INBOUND_BURST) + return not b.consume() + + # ── POST /v1/frame ──────────────────────────────────────────────────── + + def _handle_frame( + self, handler: BaseHTTPRequestHandler, device_id: str, frame: protocol.Frame + ) -> None: + """Accept-and-ack (docs/19 §19.7): 202 once the frame is dispatched; + a synchronous validation rejection comes back as the 4xx body. + Long-running handlers (the agent turn) keep running after the ack — + their async output arrives on the event stream.""" + loop = self._loop + if loop is None or not loop.is_running(): + _send_frame_json( + handler, + 503, + protocol.error(protocol.ERR_INTERNAL, "gateway loop not running").to_json(), + ) + return + # Reply sink: while this request is being dispatched, frames the + # handler would send via send_to() are captured here instead (the + # adapter's _reply routes them in). ``abandoned`` is set once the + # HTTP response has been sent without consuming the sink (the + # long-running-handler case); the dispatch's finally then delivers + # any late replies via the event stream instead of losing them. + sink: queue.Queue = queue.Queue() + abandoned = threading.Event() + self._adapter._http_register_sink(device_id, (sink, abandoned)) + try: + task = asyncio.run_coroutine_threadsafe( + self._dispatch_guarded(frame, device_id, sink, abandoned), loop + ) + except Exception: + self._adapter._http_pop_sink(device_id) + _send_frame_json( + handler, 500, protocol.error(protocol.ERR_INTERNAL, "dispatch failed").to_json() + ) + return + frames: list[protocol.Frame] = [] + deadline = time.monotonic() + ACCEPT_ACK_TIMEOUT_S + while True: + try: + frames.append(sink.get(timeout=0.05)) + break + except queue.Empty: + if task.done(): + # All replies are in the sink now (the handler finished); + # drain them all. + while True: + try: + frames.append(sink.get_nowait()) + except queue.Empty: + break + break + if time.monotonic() >= deadline: + # Long-running handler (the agent turn): ack now; late + # replies go to the event stream (the dispatch's finally + # sees ``abandoned`` and delivers them there). + abandoned.set() + break + if not frames: + _send_json(handler, 202, {"ok": True}) + elif len(frames) == 1: + f = frames[0] + status = 429 if f.payload.get("code") == protocol.ERR_RATE_LIMITED else ( + 400 if f.type == protocol.TYPE_ERROR else 200 + ) + _send_frame_json(handler, status, f.to_json()) + else: + # Multi-frame reply (e.g. a sync replay): deliver it all on the + # event stream; the ack stays plain. + for f in frames: + with contextlib.suppress(Exception): + asyncio.run_coroutine_threadsafe( + self._deliver_via_stream(f), loop + ) + _send_json(handler, 202, {"ok": True}) + + async def _dispatch_guarded( + self, + frame: protocol.Frame, + device_id: str, + sink: queue.Queue, + abandoned: threading.Event, + ) -> None: + try: + await dispatch_frame(self._adapter, frame, device_id) + except Exception: + logger.warning( + "android: HTTP dispatch failed for %s", frame.type, exc_info=True + ) + finally: + # Pop our sink entry (a newer request from the same device may + # have replaced it). If the HTTP response was already sent + # (abandoned), any replies still in the sink are delivered via + # the event stream instead of being lost. In the normal case the + # handler thread has already drained the sink, so nothing is + # left to deliver. + popped = self._adapter._http_pop_sink_if(device_id, sink) + if popped is not None and abandoned.is_set(): + while True: + try: + f = sink.get_nowait() + except queue.Empty: + break + with contextlib.suppress(Exception): + await self._deliver_via_stream(f) + + async def _deliver_via_stream(self, frame: protocol.Frame) -> None: + await self._adapter._ws_server.broadcast(frame) + await self.fanout(frame, cursor=None) + + # ── GET /v1/events (SSE) ────────────────────────────────────────────── + + def _handle_sse( + self, handler: BaseHTTPRequestHandler, device_id: str, parsed: Any + ) -> None: + qs = parse_qs(parsed.query) + cursor = _parse_cursor(qs.get("cursor", [None])[0], handler.headers.get("Last-Event-ID")) + sub = _Subscriber(device_id=device_id, kind="sse") + # Register BEFORE the replay so a frame appended in between is + # fanned out to us (and de-duped by cursor below) instead of lost. + self._add_sub(sub) + try: + handler.send_response(200) + handler.send_header("Content-Type", "text/event-stream") + handler.send_header("Cache-Control", "no-cache") + handler.send_header("X-Accel-Buffering", "no") + handler.end_headers() + # 1. Catch-up from the outbox (id = cursor; the envelope also + # carries the cursor for the app's push dedupe). + max_cursor = cursor + for e in self._adapter._outbox.replay(cursor): + c = int(e["cursor"]) + max_cursor = max(max_cursor, c) + self._write_sse(handler, "frame", c, _with_cursor(e["frame"], c)) + # 2. hello (the HTTP equivalent of hello.ack) + current status. + hello = protocol.hello_ack( + server_caps=self._adapter.server_caps(), + sync_cursor=self._adapter._outbox.latest_cursor(), + channels=self._adapter.channel_list(), + last_pushed_cursor=self._adapter._devices.last_pushed_cursor(device_id), + ) + self._write_sse(handler, "hello", None, hello.to_json()) + self._write_sse( + handler, "frame", None, protocol.status(self._adapter.gateway_status()).to_json() + ) + # 3. Live frames (cursor=None frames have no id). + while not sub.closed.is_set(): + try: + item = sub.q.get(timeout=SSE_HEARTBEAT_S) + except queue.Empty: + self._write_raw(handler, ": hb\n\n") + continue + if item is _STOP: + break + c, data = item + if c is not None and c <= max_cursor: + continue # already replayed above + self._write_sse(handler, "frame", c, data) + except (BrokenPipeError, ConnectionResetError, OSError): + pass # client went away: normal + finally: + self._remove_sub(sub) + + @staticmethod + def _write_sse( + handler: BaseHTTPRequestHandler, event: str, cursor: int | None, data: str + ) -> None: + lines = "" + if cursor is not None: + lines += f"id: {cursor}\n" + lines += f"event: {event}\ndata: {data}\n\n" + handler.wfile.write(lines.encode("utf-8")) + handler.wfile.flush() + + @staticmethod + def _write_raw(handler: BaseHTTPRequestHandler, text: str) -> None: + handler.wfile.write(text.encode("utf-8")) + handler.wfile.flush() + + # ── GET /v1/poll (long-poll) ────────────────────────────────────────── + + def _handle_poll( + self, handler: BaseHTTPRequestHandler, device_id: str, parsed: Any + ) -> None: + qs = parse_qs(parsed.query) + cursor = _parse_cursor(qs.get("cursor", [None])[0]) + sub = _Subscriber(device_id=device_id, kind="poll") + self._add_sub(sub) + try: + frames: list[str] = [] + max_cursor = cursor + for e in self._adapter._outbox.replay(cursor): + c = int(e["cursor"]) + max_cursor = max(max_cursor, c) + frames.append(_with_cursor(e["frame"], c)) + deadline = time.monotonic() + POLL_TIMEOUT_S + while not frames and not sub.closed.is_set() and time.monotonic() < deadline: + remaining = deadline - time.monotonic() + try: + item = sub.q.get(timeout=min(remaining, 5.0)) + except queue.Empty: + continue + if item is _STOP: + break + c, data = item + if c is None or c <= max_cursor: + continue + max_cursor = c + frames.append(data) + hwm = max(max_cursor, self._adapter._outbox.latest_cursor()) + _send_json(handler, 200, {"cursor": hwm, "frames": frames}) + except (BrokenPipeError, ConnectionResetError, OSError): + pass + finally: + self._remove_sub(sub) + + +class _ThreadingHTTPD(ThreadingHTTPServer): + """One thread per connection (fine at single-user scale); daemon + threads so a stuck handler can't block process exit.""" + + daemon_threads = True + allow_reuse_address = True + + def __init__(self, addr: tuple[str, int], http_server: HttpServer): + super().__init__(addr, _Handler) + self.http_server = http_server + + +class _Handler(BaseHTTPRequestHandler): + # HTTP/1.0 (default): the connection closes after each response. That + # matches the transport's design (short-lived connections) and avoids + # Content-Length bookkeeping on the streamed SSE response. + server: _ThreadingHTTPD + + def log_message(self, fmt: str, *args: Any) -> None: # noqa: A003 + logger.debug("android http: " + fmt, *args) + + # ── Routing ─────────────────────────────────────────────────────────── + + def do_GET(self) -> None: # noqa: N802 + hs = self.server.http_server + if not hs.enabled: + _send_json(self, 503, {"error": "http leg disabled"}) + return + parsed = urlparse(self.path) + if parsed.path == "/v1/health": + # Unauthenticated by design: it answers "is the gateway + # alive?" and must not reflect tokens, device ids, or versions. + _send_json(self, 200, {"ok": True}) + return + if parsed.path == "/v1/events": + device_id = hs._authenticate(self) + if device_id is not None: + hs._handle_sse(self, device_id, parsed) + return + if parsed.path == "/v1/poll": + device_id = hs._authenticate(self) + if device_id is not None: + hs._handle_poll(self, device_id, parsed) + return + _send_json(self, 404, {"error": "not found"}) + + def do_POST(self) -> None: # noqa: N802 + hs = self.server.http_server + if not hs.enabled: + _send_json(self, 503, {"error": "http leg disabled"}) + return + parsed = urlparse(self.path) + if parsed.path != "/v1/frame": + _send_json(self, 404, {"error": "not found"}) + return + device_id = hs._authenticate(self) + if device_id is None: + return + if hs._rate_limited(device_id): + _send_frame_json( + self, + 429, + protocol.error( + protocol.ERR_RATE_LIMITED, "http frame rate limit exceeded" + ).to_json(), + ) + return + ctype = (self.headers.get("Content-Type") or "").split(";")[0].strip().lower() + if ctype != "application/json": + _send_frame_json( + self, + 400, + protocol.error( + protocol.ERR_INTERNAL, "Content-Type must be application/json" + ).to_json(), + ) + return + try: + length = int(self.headers.get("Content-Length") or 0) + except ValueError: + length = 0 + if length <= 0 or length > MAX_BODY_BYTES: + _send_frame_json( + self, + 413, + protocol.error( + protocol.ERR_INTERNAL, f"body must be 1..{MAX_BODY_BYTES} bytes" + ).to_json(), + ) + return + body = self.rfile.read(length) + frame = protocol.Frame.from_json(body) + if frame is None: + _send_frame_json( + self, 400, protocol.error(protocol.ERR_INTERNAL, "invalid frame").to_json() + ) + return + if frame.type in HTTP_REJECTED_TYPES: + _send_frame_json( + self, + 400, + protocol.error( + protocol.ERR_UNSUPPORTED, f"{frame.type} requires the live connection" + ).to_json(), + ) + return + hs._handle_frame(self, device_id, frame) diff --git a/gateway-plugin/tests/README.md b/gateway-plugin/tests/README.md index ed98366..07e244a 100644 --- a/gateway-plugin/tests/README.md +++ b/gateway-plugin/tests/README.md @@ -43,6 +43,12 @@ Beyond the base modes (`--send`, `--upload`, `--pull-offer`, `--sync`, - `--offer-grace S` — with `--pull-offer`, keep listening S seconds after the final message for a `media.offer` (offers are emitted post-turn, right after the final; default 15). +- `--http [--http-url http://host:port]` — docs/19: drive the turn over + the **HTTP fallback leg** instead of WS: `GET /v1/health`, + `POST /v1/frame` (the `message.send`), receive over SSE `GET + /v1/events`. The same assertion flags apply. The base URL defaults to + the `--url` host with scheme `ws(s)` → `http(s)` and port 8791 + (`ANDROID_HTTP_PORT`). Exit codes: `0` ok (incl. SKIP for absent M7 frames), `2` connect fail, `3` no hello.ack, `4` expected hello.ack, `5` authfail expected but @@ -51,7 +57,9 @@ acked, `6` timeout, `7` no final message, `8` upload/sync fail, `12` assert-tools fail, `13` assert-commentary fail, `14` search fail (error or zero hits), `15` channel.create/list fail, `16` channel.delete fail, `17` watch timeout, `18` read.receipt arrived before the sent -message, `19` status frame with empty payload. +message, `19` status frame with empty payload, `20` `--http` health +check failed, `21` `--http` SSE open failed, `22` `--http` +`POST /v1/frame` rejected (4xx). ## E2E driver (`e2e.py`) diff --git a/gateway-plugin/tests/ws_probe.py b/gateway-plugin/tests/ws_probe.py index 6bc76c5..f68ab39 100644 --- a/gateway-plugin/tests/ws_probe.py +++ b/gateway-plugin/tests/ws_probe.py @@ -56,6 +56,13 @@ Request modes (no turn driven unless --send/--upload also given): --watch CHAT_ID wait up to --timeout for a message to land in CHAT_ID (cron delivery E2E, scenario 7) +HTTP fallback leg (docs/19): + --http drive the turn over the HTTP leg instead of WS: + GET /v1/health, POST /v1/frame (message.send), receive + over SSE /v1/events. The same assertion flags apply. + --http-url http://host:port base for --http (default: derived + from --url, ws(s) -> http(s), port 8791) + Exit codes: 0 ok (incl. SKIP for absent M7 frames) 2 connect failed @@ -76,6 +83,9 @@ Exit codes: 17 --watch timed out (no message landed in the channel) 18 --assert-read-receipt failed (frame arrived before the sent message) 19 --assert-status failed (status frame arrived with an empty payload) + 20 --http: health check failed + 21 --http: SSE open failed + 22 --http: POST /v1/frame rejected (4xx) """ import argparse @@ -690,6 +700,124 @@ async def run(args) -> int: return 0 +def run_http(args, base: str) -> int: + """docs/19: drive a turn over the HTTP fallback leg — GET /v1/health, + POST /v1/frame (message.send), receive over SSE /v1/events. Blocking + (stdlib http.client); the same assertion flags apply as the WS leg.""" + from http.client import HTTPConnection + from urllib.parse import urlparse + + u = urlparse(base) + host = u.hostname or "127.0.0.1" + port = u.port or (443 if u.scheme == "https" else 80) + headers = { + "Authorization": f"Bearer {args.token}", + "X-Iris-Device": args.device, + } + + # 1. health (unauthenticated liveness probe). + try: + conn = HTTPConnection(host, port, timeout=5) + conn.request("GET", "/v1/health") + r = conn.getresponse() + body = r.read() + conn.close() + except Exception as e: + print(f"!! health check failed: {e}") + return 20 + if r.status != 200: + print(f"!! health check failed: HTTP {r.status} {body[:200]!r}") + return 20 + print(f"== health ok: {body!r}") + + # 2. open the SSE stream. + sse = HTTPConnection(host, port, timeout=args.timeout) + sse.request("GET", "/v1/events", headers=headers) + resp = sse.getresponse() + if resp.status != 200: + print(f"!! SSE open failed: HTTP {resp.status}") + return 21 + print("== SSE open (/v1/events)") + + # 3. POST the message.send frame (accept-and-ack). + if args.send: + frame = { + "v": 1, "id": 1, "type": "message.send", + "chat_id": "android:default", "payload": {"text": args.send}, + } + conn = HTTPConnection(host, port, timeout=30) + conn.request( + "POST", "/v1/frame", body=json.dumps(frame), + headers={**headers, "Content-Type": "application/json"}, + ) + r = conn.getresponse() + body = r.read() + conn.close() + print(f"== POST /v1/frame -> {r.status} {body[:200]!r}") + if r.status >= 400: + print("!! POST /v1/frame rejected") + return 22 + + # 4. read SSE until the final assistant message (same final-detection + # logic as the WS leg). + st = _TurnState() + got_final = False + seen_final_frame = False + deadline = time.time() + args.timeout + cur_data: list[str] = [] + + def feed(line: str) -> bool: + nonlocal cur_data, got_final, seen_final_frame + line = line.rstrip("\r\n") + if line == "": + if cur_data: + data = _print_frame("\n".join(cur_data)) + if data is not None: + ftype = data.get("type") + payload = data.get("payload") or {} + st.track(ftype, payload) + if ftype == "message" and payload.get("role") == "assistant": + got_final = True + if ftype == "message.stop": + seen_final_frame = True + if ftype == "typing" and payload.get("on") is False and seen_final_frame: + got_final = True + cur_data = [] + return got_final + if line.startswith(":"): + return got_final # heartbeat comment + field, _, value = line.partition(":") + if value.startswith(" "): + value = value[1:] + if field == "data": + cur_data.append(value) + return got_final + + sock = getattr(getattr(resp.fp, "raw", None), "_sock", None) + try: + while time.time() < deadline and not got_final: + if sock is not None: + sock.settimeout(max(0.1, deadline - time.time())) + line = resp.fp.readline() + if not line: + break + if feed(line.decode("utf-8")): + break + finally: + sse.close() + + if not got_final: + print(f"!! no final assistant message (HTTP leg, {args.timeout:.0f}s)") + return 7 + print("== final assistant message received (via SSE)") + + for code, ok, msg in _evaluate_assertions(args, st): + if not ok: + print(f"!! {msg}") + return code + return 0 + + def main() -> int: p = argparse.ArgumentParser(description=__doc__) p.add_argument("--url", default=os.getenv("ANDROID_WS_URL", "ws://127.0.0.1:8790/ws")) @@ -739,12 +867,28 @@ def main() -> int: p.add_argument("--offer-grace", type=float, default=15.0, help="seconds to wait for a media.offer after the final " "message when --pull-offer (default 15)") + p.add_argument("--http", action="store_true", + help="docs/19: drive the turn over the HTTP fallback leg " + "(health + POST /v1/frame + SSE /v1/events) instead of WS") + p.add_argument("--http-url", default="", + help="docs/19: http(s)://host:port base for --http " + "(default: derived from --url, port 8791)") args = p.parse_args() if not args.token and not args.authfail: p.error("--token (or $ANDROID_TOKEN) is required") if args.assert_read_receipt and not args.send: p.error("--assert-read-receipt requires --send (the receipt must follow " "the sent message)") + if args.http: + if args.http_url: + base = args.http_url + else: + from urllib.parse import urlparse + + u = urlparse(args.url) + scheme = "https" if u.scheme == "wss" else "http" + base = f"{scheme}://{u.hostname or '127.0.0.1'}:8791" + return run_http(args, base) return asyncio.run(run(args)) diff --git a/gateway-plugin/ws_server.py b/gateway-plugin/ws_server.py index 6e31248..6d7f521 100644 --- a/gateway-plugin/ws_server.py +++ b/gateway-plugin/ws_server.py @@ -66,6 +66,50 @@ CLOSE_SHUTDOWN = 1001 MAX_DEVICE_ID_LEN = 128 +async def dispatch_frame(adapter: Any, frame: protocol.Frame, device_id: str) -> None: + """Shared inbound frame dispatch for the WS and HTTP transports + (docs/19 §19.4). Transport-specific frames (``ping``/``pong``, binary + media chunks) are handled by their own server before this is called; + unknown types are ignored (forward-compat).""" + if frame.type == protocol.TYPE_MESSAGE_SEND: + await adapter.on_message_send(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_CREATE: + await adapter.on_channel_create(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_RENAME: + await adapter.on_channel_rename(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_SET_DEFAULT: + await adapter.on_channel_set_default(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_FAVORITE: + await adapter.on_channel_favorite(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_ICON: + await adapter.on_channel_icon(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_SET_AUTOMATION: + await adapter.on_channel_set_automation(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_DELETE: + await adapter.on_channel_delete(frame, device_id) + elif frame.type == protocol.TYPE_CHANNEL_LIST: + await adapter.on_channel_list(frame, device_id) + elif frame.type == protocol.TYPE_COMMANDS_CATALOG: + await adapter.on_commands_catalog(frame, device_id) + elif frame.type == protocol.TYPE_SEARCH: + await adapter.on_search(frame, device_id) + elif frame.type == protocol.TYPE_SYNC: + await adapter.on_sync(frame, device_id) + elif frame.type == protocol.TYPE_HISTORY: + await adapter.on_history(frame, device_id) + elif frame.type == protocol.TYPE_MESSAGE_DELETE: + await adapter.on_message_delete(frame, device_id) + elif frame.type == protocol.TYPE_MEDIA_UPLOAD_START: + await adapter.on_media_upload_start(frame, device_id) + elif frame.type == protocol.TYPE_MEDIA_UPLOAD_END: + await adapter.on_media_upload_end(frame, device_id) + elif frame.type == protocol.TYPE_MEDIA_PULL: + await adapter.on_media_pull(frame, device_id) + elif frame.type == protocol.TYPE_FCM_REGISTER: + await adapter.on_fcm_register(frame, device_id) + # Unknown types are ignored (forward-compat). + + class _TokenBucket: """Minimal token bucket (stdlib only). One instance per connection.""" @@ -370,43 +414,8 @@ class WsServer: if frame.type == protocol.TYPE_PING: ts = frame.payload.get("ts") await self._send_quiet(ws, protocol.pong(ts if isinstance(ts, int) else None)) - elif frame.type == protocol.TYPE_MESSAGE_SEND: - await self._adapter.on_message_send(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_CREATE: - await self._adapter.on_channel_create(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_RENAME: - await self._adapter.on_channel_rename(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_SET_DEFAULT: - await self._adapter.on_channel_set_default(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_FAVORITE: - await self._adapter.on_channel_favorite(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_ICON: - await self._adapter.on_channel_icon(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_SET_AUTOMATION: - await self._adapter.on_channel_set_automation(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_DELETE: - await self._adapter.on_channel_delete(frame, device_id) - elif frame.type == protocol.TYPE_CHANNEL_LIST: - await self._adapter.on_channel_list(frame, device_id) - elif frame.type == protocol.TYPE_COMMANDS_CATALOG: - await self._adapter.on_commands_catalog(frame, device_id) - elif frame.type == protocol.TYPE_SEARCH: - await self._adapter.on_search(frame, device_id) - elif frame.type == protocol.TYPE_SYNC: - await self._adapter.on_sync(frame, device_id) - elif frame.type == protocol.TYPE_HISTORY: - await self._adapter.on_history(frame, device_id) - elif frame.type == protocol.TYPE_MESSAGE_DELETE: - await self._adapter.on_message_delete(frame, device_id) - elif frame.type == protocol.TYPE_MEDIA_UPLOAD_START: - await self._adapter.on_media_upload_start(frame, device_id) - elif frame.type == protocol.TYPE_MEDIA_UPLOAD_END: - await self._adapter.on_media_upload_end(frame, device_id) - elif frame.type == protocol.TYPE_MEDIA_PULL: - await self._adapter.on_media_pull(frame, device_id) - elif frame.type == protocol.TYPE_FCM_REGISTER: - await self._adapter.on_fcm_register(frame, device_id) - # Unknown types are ignored (forward-compat). + return True + await dispatch_frame(self._adapter, frame, device_id) return True # ── Helpers ───────────────────────────────────────────────────────────