HTTP transport: drop WS server, offline send queue + dead-stream watchdog
CI / Gateway plugin tests (push) Successful in 5m9s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m55s

Gateway (docs/19):
- Remove ws_server.py; frame dispatch factored into dispatch.py
- http_server: media upload/pull, pairing over HTTP
- protocol: media frames mirrored; tests + ws_probe updated for HTTP

App:
- HttpGateway: postFrame/uploadMedia/pullMedia no longer throw on
  network failure (PostResult ok=false / Result.failure) — uncaught
  SocketTimeoutException on Dispatchers.Default crashed the app
- GatewayClient: dead-stream watchdog (health probe every 10s, 2
  failures -> redial in ~20s instead of the 45s SSE read timeout);
  state flips to Reconnecting when the stream dies, restored from the
  last hello.ack on long-poll success; poke() + backoff reset on app
  resume (MainActivity.onResume)
- Offline sends: composer enabled while disconnected; a send with no
  response (status 0) stays queued (Pending) and is auto-resent on the
  next (re)connect after a 2s outbox-replay grace; gateway 4xx
  rejections fail the bubble (tap to retry, no auto-loop)
- ChatStore: echo-replace and thread-relocate also match Failed
  bubbles (POST response lost in a network drop); loadHistory dedupes
  local failed bubbles the server already has; failMessage()
- MainActivity: poke() on resume so a backgrounded app reconnects
  promptly instead of waiting out the backoff
This commit is contained in:
ARIA committed 2026-08-22 20:10:05 +02:00
1 parent 2349a95dd4
commit e6015033b6
22 files changed
+2804 -2439

No files matched your search

+274 -97
View File
@@ -1,11 +1,10 @@
"""HTTP fallback transport (docs/19): the "HTTP leg".
"""HTTP transport (docs/19): the gateway's device-facing server.
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.
Short-lived-connection transport: the same JSON frames, the same
outbox/cursor, the same token — served over plain HTTP by the gateway.
The app sends over ``POST /v1/frame`` and receives over
``GET /v1/events`` (SSE) or ``GET /v1/poll`` (long-poll); media travels
via ``POST /v1/media`` / ``GET /v1/media/{id}`` (v2, docs/19 §19.15).
Zero new Python dependencies: stdlib ``http.server`` (a
``ThreadingHTTPServer`` in a daemon thread) bridged into the gateway's
@@ -13,19 +12,30 @@ 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).
* ``POST /v1/frame`` — accept-and-ack for any JSON frame the
app sends (media uses the /v1/media
endpoints; hello/ping are
transport-specific).
* ``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.
* ``POST /v1/media`` — media upload (docs/19 §19.15, v2): the
whole file as the request body; metadata
in ``X-Iris-Media-*`` headers; sha256
contract per docs/07 §7.2.
* ``GET /v1/media/{media_id}`` — media pull (docs/19 §19.15, v2): streams
an outbound offer (``media.offer`` id)
as the response body.
Auth: ``Authorization: Bearer <token>`` (constant-time ``verify_token``) +
``X-Iris-Device`` header (same device id / allowlist as the WS ``hello``).
``X-Iris-Device`` header (device id / allowlist). Device registration
(name + push tokens) rides on the SSE open via ``X-Iris-Device-Name`` /
``X-Iris-Fcm-Token`` / ``X-Iris-Ntfy-Topic`` headers (the HTTP equivalent
of the old WS ``hello`` upsert).
Bind failure is NON-fatal (unlike the WS): the plugin keeps working
WS-only.
HTTP is the ONLY transport: a bind failure is FATAL (the app has no other
way to reach the gateway).
"""
from __future__ import annotations
@@ -43,24 +53,25 @@ from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from typing import Any
from urllib.parse import parse_qs, urlparse
from . import protocol
from . import dispatch, protocol
from . import media as media_bridge
from .pairing import verify_token
from .ws_server import (
INBOUND_BURST,
INBOUND_RATE_PER_S,
MAX_DEVICE_ID_LEN,
_TokenBucket,
dispatch_frame,
)
try: # main-repo import (same as adapter.py); absent in bare unit contexts
from gateway.platforms.base import validate_media_delivery_path
except ImportError: # pragma: no cover
validate_media_delivery_path = None # type: ignore[assignment]
logger = logging.getLogger(__name__)
# Default port for the HTTP leg (WS default is 8790).
# Default port for the HTTP transport (the WS-era default was 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
# Request body cap for POST /v1/frame. Frames are usually small, but a
# ``channel.icon`` carries a base64 blob up to 512 KiB (docs/10), so the cap
# must clear that with headroom. Media never travels here (it uses
# POST /v1/media).
MAX_BODY_BYTES = 1024 * 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.
@@ -79,18 +90,16 @@ 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,
}
)
# Max length of a client-supplied media_ref (same as the WS path).
MAX_MEDIA_REF_LEN = 64
# MediaError code -> HTTP status for the /v1/media endpoints.
_MEDIA_STATUS = {
protocol.ERR_MEDIA_TOO_LARGE: 413,
protocol.ERR_NOT_FOUND: 404,
protocol.ERR_UNSUPPORTED: 400,
protocol.ERR_INTERNAL: 500,
}
def _with_cursor(frame: dict[str, Any], cursor: int) -> str:
@@ -151,12 +160,12 @@ class _Subscriber:
class HttpServer:
"""The plugin's HTTP fallback server + live subscriber registry.
"""The plugin's HTTP 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.
with ``asyncio.run_coroutine_threadsafe`` and dispatched through
``dispatch_frame`` (``dispatch.py``).
"""
def __init__(self, adapter: Any, devices: Any):
@@ -167,7 +176,7 @@ class HttpServer:
self._thread: threading.Thread | None = None
self._subs: dict[str, list[_Subscriber]] = {}
self._subs_lock = threading.Lock()
self._buckets: dict[str, _TokenBucket] = {}
self._buckets: dict[str, dispatch._TokenBucket] = {}
self._buckets_lock = threading.Lock()
self._lock_key: str | None = None
self.enabled = False
@@ -176,8 +185,9 @@ class HttpServer:
# ── 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)."""
"""Bind and start serving. NEVER raises: a bind failure leaves
``enabled`` False, which the adapter treats as a fatal error
(HTTP is the only transport, docs/19 §19.4)."""
if self.enabled:
return
self._loop = asyncio.get_running_loop()
@@ -191,7 +201,7 @@ class HttpServer:
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",
"android: HTTP port %s:%s in use by another profile; server disabled",
host,
port,
)
@@ -208,9 +218,7 @@ class HttpServer:
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
)
logger.warning("android: HTTP server disabled (bind %s:%s failed: %s)", host, port, e)
self._release_lock()
return
@@ -221,9 +229,7 @@ class HttpServer:
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
)
logger.info("android: HTTP server listening on %s://%s:%s", scheme, host, self.bound_port)
async def stop(self) -> None:
"""Stop serving and unblock all subscribers."""
@@ -317,7 +323,7 @@ class HttpServer:
_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:
if not device_id or len(device_id) > dispatch.MAX_DEVICE_ID_LEN:
_send_json(handler, 401, {"error": "X-Iris-Device header required"})
return None
if (
@@ -333,11 +339,13 @@ class HttpServer:
return device_id
def _rate_limited(self, device_id: str) -> bool:
"""Per-device token bucket, same parameters as the WS inbound limit."""
"""Per-device token bucket, same parameters as the frame 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)
b = self._buckets[device_id] = dispatch._TokenBucket(
dispatch.INBOUND_RATE_PER_S, dispatch.INBOUND_BURST
)
return not b.consume()
# ── POST /v1/frame ────────────────────────────────────────────────────
@@ -379,31 +387,35 @@ class HttpServer:
frames: list[protocol.Frame] = []
deadline = time.monotonic() + ACCEPT_ACK_TIMEOUT_S
while True:
try:
frames.append(sink.get(timeout=0.05))
# If the handler is done, drain any replies and stop (no wait).
# This keeps fast/ignored frames from incurring the sink timeout.
if task.done():
while True:
try:
frames.append(sink.get_nowait())
except queue.Empty:
break
break
try:
frames.append(sink.get(timeout=0.01))
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
continue
# Got a frame; loop back to check task.done() (drain the rest if
# the handler finished, e.g. a sync replay).
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
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:
@@ -411,9 +423,7 @@ class HttpServer:
# event stream; the ack stays plain.
for f in frames:
with contextlib.suppress(Exception):
asyncio.run_coroutine_threadsafe(
self._deliver_via_stream(f), loop
)
asyncio.run_coroutine_threadsafe(self._deliver_via_stream(f), loop)
_send_json(handler, 202, {"ok": True})
async def _dispatch_guarded(
@@ -424,11 +434,9 @@ class HttpServer:
abandoned: threading.Event,
) -> None:
try:
await dispatch_frame(self._adapter, frame, device_id)
await dispatch.dispatch_frame(self._adapter, frame, device_id)
except Exception:
logger.warning(
"android: HTTP dispatch failed for %s", frame.type, exc_info=True
)
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
@@ -438,25 +446,39 @@ class HttpServer:
# 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)
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:
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"))
# Device registration (the HTTP equivalent of the WS hello upsert):
# the SSE open carries the device name + push tokens as optional
# headers; upsert is idempotent and COALESCEs absent tokens, so a
# re-open never clobbers a newer fcm.register value.
device_name = (handler.headers.get("X-Iris-Device-Name") or "").strip()[:120]
fcm_token = handler.headers.get("X-Iris-Fcm-Token") or None
ntfy_topic = handler.headers.get("X-Iris-Ntfy-Topic") or None
try:
self._devices.upsert(
device_id,
device_name or device_id,
None,
fcm_token,
ntfy_topic,
)
except Exception:
logger.warning("android: device registry upsert failed", exc_info=True)
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.
@@ -499,7 +521,9 @@ class HttpServer:
continue # already replayed above
self._write_sse(handler, "frame", c, data)
except (BrokenPipeError, ConnectionResetError, OSError):
pass # client went away: normal
# Client went away mid-stream: normal (the app reconnects with
# Last-Event-ID and catches up from the outbox).
pass
finally:
self._remove_sub(sub)
@@ -519,11 +543,137 @@ class HttpServer:
handler.wfile.write(text.encode("utf-8"))
handler.wfile.flush()
# ── POST /v1/media (upload, docs/19 §19.15) ───────────────────────────
def _handle_media_upload(self, handler: BaseHTTPRequestHandler, device_id: str) -> None:
"""Whole-file upload: metadata in headers, file bytes as the body.
Mirrors the WS ``media.upload`` contract (docs/07 §7.2) in one
request: the body is streamed to a temp file (bounded RAM), then
size + sha256 are verified and the file cached via the hermes
``cache_*_from_bytes`` helpers. Runs entirely on the handler thread
(plain file IO — no asyncio bridge needed)."""
media_ref = (handler.headers.get("X-Iris-Media-Ref") or "").strip()
kind = (handler.headers.get("X-Iris-Media-Kind") or "").strip()
filename = (handler.headers.get("X-Iris-Media-Filename") or "upload")[:255]
sha256 = (handler.headers.get("X-Iris-Media-Sha256") or "").strip().lower()
mime = handler.headers.get("Content-Type") or "application/octet-stream"
mime = mime.split(";")[0].strip()[:128]
try:
length = int(handler.headers.get("Content-Length") or 0)
except ValueError:
length = 0
def reject(code: str, message: str) -> None:
_send_frame_json(
handler,
_MEDIA_STATUS.get(code, 400),
protocol.error(code, message).to_json(),
)
# Same validation rules as the WS media.upload.start handler.
if not media_ref or len(media_ref) > MAX_MEDIA_REF_LEN:
reject(protocol.ERR_UNSUPPORTED, "X-Iris-Media-Ref header required")
return
if kind not in media_bridge.KINDS:
reject(protocol.ERR_UNSUPPORTED, f"unsupported media kind {kind!r}")
return
if length <= 0:
reject(protocol.ERR_UNSUPPORTED, "empty body")
return
if length > self._adapter.max_upload_bytes:
reject(
protocol.ERR_MEDIA_TOO_LARGE,
f"upload of {length} bytes exceeds limit ({self._adapter.max_upload_bytes})",
)
return
try:
sess = self._adapter._media.create_upload(
device_id,
media_ref,
kind,
mime,
filename,
length,
None,
self._adapter.max_upload_bytes,
)
except media_bridge.MediaError as e:
reject(e.code, e.message)
return
try:
remaining = length
while remaining > 0:
chunk = handler.rfile.read(min(media_bridge.DEFAULT_CHUNK_BYTES, remaining))
if not chunk:
raise media_bridge.MediaError(
protocol.ERR_INTERNAL, "client disconnected mid-upload"
)
sess.feed(chunk)
remaining -= len(chunk)
if sess.received != length:
raise media_bridge.MediaError(
protocol.ERR_INTERNAL,
f"size mismatch (declared {length}, received {sess.received})",
)
entry = self._adapter._media.complete_upload(device_id, media_ref, sha256)
except media_bridge.MediaError as e:
# complete_upload already popped the session; discard is a no-op
# in that case (feed/short-read failures leave it active).
self._adapter._media.discard_upload(device_id, media_ref)
reject(e.code, e.message)
return
except (BrokenPipeError, ConnectionResetError, OSError):
self._adapter._media.discard_upload(device_id, media_ref)
return # client went away: nothing to answer
_send_frame_json(handler, 201, protocol.media_upload_ack(True, entry.media_id).to_json())
# ── GET /v1/media/{id} (pull, docs/19 §19.15) ─────────────────────────
def _handle_media_pull(
self, handler: BaseHTTPRequestHandler, device_id: str, media_id: str
) -> None:
"""Stream an outbound offer as the response body (docs/07 §7.3).
The delivery-path validation is re-checked at pull time, exactly as
the WS ``media.pull`` handler does (the file may have moved since
the offer)."""
entry = self._adapter._media.get_outbound(media_id)
if entry is None:
_send_frame_json(
handler,
404,
protocol.error(protocol.ERR_NOT_FOUND, f"unknown media_id {media_id!r}").to_json(),
)
return
safe = validate_media_delivery_path(entry.path) if validate_media_delivery_path else None
if safe is None:
_send_frame_json(
handler,
404,
protocol.error(protocol.ERR_NOT_FOUND, "media no longer deliverable").to_json(),
)
return
filename = entry.filename.replace('"', "")
handler.send_response(200)
handler.send_header("Content-Type", entry.mime)
handler.send_header("Content-Length", str(entry.size))
handler.send_header("Content-Disposition", f'attachment; filename="{filename}"')
handler.end_headers()
try:
with open(safe, "rb") as f: # pi-lens-ignore: python-path-traversal
while True:
chunk = f.read(media_bridge.DEFAULT_CHUNK_BYTES)
if not chunk:
break
handler.wfile.write(chunk)
handler.wfile.flush()
except (BrokenPipeError, ConnectionResetError, OSError):
pass # client went away mid-pull, or the file vanished: normal
# ── GET /v1/poll (long-poll) ──────────────────────────────────────────
def _handle_poll(
self, handler: BaseHTTPRequestHandler, device_id: str, parsed: Any
) -> None:
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")
@@ -552,6 +702,7 @@ class HttpServer:
hwm = max(max_cursor, self._adapter._outbox.latest_cursor())
_send_json(handler, 200, {"cursor": hwm, "frames": frames})
except (BrokenPipeError, ConnectionResetError, OSError):
# Client went away while we held the poll: normal.
pass
finally:
self._remove_sub(sub)
@@ -601,6 +752,17 @@ class _Handler(BaseHTTPRequestHandler):
if device_id is not None:
hs._handle_poll(self, device_id, parsed)
return
if parsed.path.startswith("/v1/media/"):
media_id = parsed.path[len("/v1/media/") :]
# The id is looked up in an exact-match dict; reject anything
# path-shaped so a bad URL can't be mistaken for an id.
if media_id and "/" not in media_id:
device_id = hs._authenticate(self)
if device_id is not None:
hs._handle_media_pull(self, device_id, media_id)
else:
_send_json(self, 404, {"error": "not found"})
return
_send_json(self, 404, {"error": "not found"})
def do_POST(self) -> None: # noqa: N802
@@ -609,6 +771,21 @@ class _Handler(BaseHTTPRequestHandler):
_send_json(self, 503, {"error": "http leg disabled"})
return
parsed = urlparse(self.path)
if parsed.path == "/v1/media":
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 media rate limit exceeded"
).to_json(),
)
return
hs._handle_media_upload(self, device_id)
return
if parsed.path != "/v1/frame":
_send_json(self, 404, {"error": "not found"})
return
@@ -639,6 +816,15 @@ class _Handler(BaseHTTPRequestHandler):
except ValueError:
length = 0
if length <= 0 or length > MAX_BODY_BYTES:
# Drain the (oversize) body so the connection stays clean; cap the
# drain at MAX_BODY_BYTES so a runaway body can't wedge the thread.
if length > 0:
to_drain = min(length, MAX_BODY_BYTES)
while to_drain > 0:
chunk = self.rfile.read(min(65536, to_drain))
if not chunk:
break
to_drain -= len(chunk)
_send_frame_json(
self,
413,
@@ -654,13 +840,4 @@ class _Handler(BaseHTTPRequestHandler):
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)