A client that completes TCP but vanishes mid-TLS-handshake (e.g. a phone losing its network/VPN while traveling) blocked ssl.SSLSocket.accept() inside serve_forever forever: the gateway stopped accepting any new device connections (the app could not reconnect), and on the next restart httpd.shutdown() froze the whole event loop until the shutdown watchdog killed the process (ARIA journal 2026-09-11 / 2026-09-23). - Move the TLS handshake out of the accept loop: it now runs in the per-connection thread under a hard timeout (HANDSHAKE_TIMEOUT_S, 10 s); a failed/timed-out handshake just closes the socket. - stop() no longer blocks the event loop: shutdown()/server_close()/ join run in an executor under asyncio.wait_for(10 s); if the bound expires the daemon threads are abandoned. - Regression test: a silent half-open TCP connection must not stop fresh TLS connections from being served, and stop() must stay bounded.
989 lines
43 KiB
Python
989 lines
43 KiB
Python
"""HTTP transport (docs/19): the gateway's device-facing server.
|
|
|
|
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
|
|
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
|
|
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 (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).
|
|
|
|
HTTP is the ONLY transport: a bind failure is FATAL (the app has no other
|
|
way to reach the gateway).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import json
|
|
import logging
|
|
import queue
|
|
import socket
|
|
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 dispatch, protocol
|
|
from . import media as media_bridge
|
|
from .pairing import verify_token
|
|
|
|
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 transport (the WS-era default was 8790).
|
|
DEFAULT_HTTP_PORT = 8791
|
|
|
|
# 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.
|
|
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()
|
|
|
|
# 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:
|
|
"""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 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
|
|
``dispatch_frame`` (``dispatch.py``).
|
|
"""
|
|
|
|
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, dispatch._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 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()
|
|
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}"
|
|
# acquire_scoped_lock returns (acquired, existing_record); the
|
|
# tuple is always truthy, so test the first element (matching
|
|
# gateway/platforms/base.py's canonical usage).
|
|
acquired, _ = acquire_scoped_lock("iris", lock_key)
|
|
if not acquired:
|
|
logger.warning(
|
|
"iris: HTTP port %s:%s in use by another profile; server 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)
|
|
# The handshake runs in the per-connection thread with a
|
|
# hard timeout (see _ThreadingHTTPD.process_request).
|
|
# Wrapping the *listening* socket here instead would make
|
|
# serve_forever's accept() block inside do_handshake() on
|
|
# a half-open connection (TCP established, client gone
|
|
# mid-handshake), wedging ALL new device connections until
|
|
# the gateway is restarted.
|
|
httpd.set_tls(ctx)
|
|
except Exception as e:
|
|
logger.warning("iris: HTTP server 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="iris-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("iris: HTTP server 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
|
|
t = self._thread
|
|
self._thread = None
|
|
if httpd is not None or t is not None:
|
|
# shutdown() blocks until the serve_forever loop exits and
|
|
# server_close() may join handler threads — both must run on a
|
|
# worker thread (never the asyncio loop thread) with a hard
|
|
# timeout, or a wedged server would freeze the whole gateway.
|
|
# The threads are daemons: if the bounded wait expires they die
|
|
# with the process and there is nothing left to do.
|
|
loop = asyncio.get_running_loop()
|
|
|
|
def _stop_httpd() -> None:
|
|
if httpd is not None:
|
|
with contextlib.suppress(Exception):
|
|
httpd.shutdown()
|
|
with contextlib.suppress(Exception):
|
|
httpd.server_close()
|
|
if t is not None and t is not threading.current_thread():
|
|
t.join(timeout=5.0)
|
|
|
|
try:
|
|
await asyncio.wait_for(
|
|
loop.run_in_executor(None, _stop_httpd),
|
|
timeout=10.0,
|
|
)
|
|
except Exception:
|
|
logger.warning(
|
|
"iris: HTTP server teardown did not finish in time; abandoning daemon threads"
|
|
)
|
|
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("iris", 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)
|
|
# M5: a live subscriber will sync the outbox -- tell the adapter to
|
|
# drop any held-back (deferred) pushes so the turn-end flush doesn't
|
|
# duplicate what the app already shows. getattr-guard: test doubles
|
|
# may use a bare adapter stub.
|
|
on_online = getattr(self._adapter, "on_device_online", None)
|
|
if on_online is not None:
|
|
on_online()
|
|
|
|
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("iris: 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.
|
|
|
|
Token model (docs/09 §9.3): a REVOKED device_id is rejected no matter
|
|
which token it presents (per-device isolation). Otherwise the shared
|
|
``IRIS_TOKEN`` (bootstrap / legacy) or the device's own per-device
|
|
token (minted at pairing, returned in ``hello.ack.device_token``)
|
|
both authenticate — each compared in constant time."""
|
|
auth = handler.headers.get("Authorization") or ""
|
|
token = auth[len("Bearer ") :] if auth.startswith("Bearer ") else None
|
|
device_id = (handler.headers.get("X-Iris-Device") or "").strip()
|
|
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
|
|
with contextlib.suppress(Exception):
|
|
if self._devices.is_revoked(device_id):
|
|
logger.warning("iris: http rejected: device %s is revoked", device_id)
|
|
_send_json(handler, 401, {"error": "device revoked"})
|
|
return None
|
|
if not verify_token(token, self._adapter.token):
|
|
# Not the shared token: try the device's own per-device token.
|
|
device_token = None
|
|
with contextlib.suppress(Exception):
|
|
device_token = self._devices.token_for(device_id)
|
|
if not (device_token and verify_token(token, device_token)):
|
|
_send_json(handler, 401, {"error": "unauthorized"})
|
|
return None
|
|
if (
|
|
not self._adapter.allow_all
|
|
and self._adapter.allowed_users
|
|
and device_id not in self._adapter.allowed_users
|
|
):
|
|
logger.warning("iris: 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 frame limit."""
|
|
with self._buckets_lock:
|
|
b = self._buckets.get(device_id)
|
|
if b is None:
|
|
b = self._buckets[device_id] = dispatch._TokenBucket(
|
|
dispatch.INBOUND_RATE_PER_S, dispatch.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:
|
|
# 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 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)
|
|
)
|
|
_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.dispatch_frame(self._adapter, frame, device_id)
|
|
except Exception:
|
|
logger.warning("iris: 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.fanout(frame, cursor=None)
|
|
|
|
# ── GET /v1/events (SSE) ──────────────────────────────────────────────
|
|
|
|
def _handle_sse(self, handler: BaseHTTPRequestHandler, device_id: str, parsed: Any) -> None: # noqa: PLR0912,PLR0915
|
|
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
|
|
# App release version (the repo-root VERSION baked into the build);
|
|
# stored in the device registry's caps JSON so `hermes` can see which
|
|
# app version each device runs (old-version awareness). A missing
|
|
# header (old app build) must not wipe a previously stored version,
|
|
# so merge over the existing caps instead of replacing them.
|
|
app_version = (handler.headers.get("X-Iris-App-Version") or "").strip()[:40]
|
|
try:
|
|
existing_caps = dict(self._devices.get(device_id) or {}).get("caps") or {}
|
|
if app_version:
|
|
existing_caps["app_version"] = app_version
|
|
self._devices.upsert(
|
|
device_id,
|
|
device_name or device_id,
|
|
existing_caps or None,
|
|
fcm_token,
|
|
ntfy_topic,
|
|
)
|
|
except Exception:
|
|
logger.warning("iris: device registry upsert failed", exc_info=True)
|
|
# Per-device token (docs/09 §9.3): minted once at pairing (idempotent
|
|
# across (re)connects) and returned in the hello below; the app
|
|
# stores it and presents it instead of the shared token from then on.
|
|
device_token = ""
|
|
try:
|
|
device_token = self._devices.issue_token(device_id)
|
|
except Exception:
|
|
logger.warning("iris: device token issuance 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.
|
|
self._add_sub(sub)
|
|
reason = "eof"
|
|
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()
|
|
# INFO (not DEBUG like the per-request log): the stream
|
|
# lifecycle is the primary "is the device connected?" signal
|
|
# for debugging flaky links — a gap here is invisible at the
|
|
# gateway's default log level.
|
|
logger.info(
|
|
"iris: SSE stream opened: %s (cursor=%d, app_version=%s)",
|
|
device_id,
|
|
cursor,
|
|
app_version or "?",
|
|
)
|
|
# 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),
|
|
device_token=device_token,
|
|
)
|
|
self._write_sse(handler, "hello", None, hello.to_json())
|
|
self._write_sse(
|
|
handler, "frame", None, protocol.status(self._adapter.gateway_status()).to_json()
|
|
)
|
|
# 2b. Live todo-list snapshot (ephemeral state, never outboxed):
|
|
# a (re)connecting device re-learns the agent's current plan here.
|
|
for snap in self._adapter.todo_snapshot_frames():
|
|
self._write_sse(handler, "frame", None, snap.to_json())
|
|
# 3. Live frames (cursor=None frames have no id).
|
|
while True:
|
|
if sub.closed.is_set():
|
|
# stop() can land between the initial writes above and
|
|
# this loop (the handler thread is descheduled under
|
|
# load): drain the frames queued before the close — e.g.
|
|
# the status{restarting} teardown broadcast — so the
|
|
# client sees them before EOF instead of losing them to
|
|
# the closed check.
|
|
try:
|
|
item = sub.q.get_nowait()
|
|
except queue.Empty:
|
|
reason = "stopped"
|
|
break
|
|
else:
|
|
try:
|
|
item = sub.q.get(timeout=SSE_HEARTBEAT_S)
|
|
except queue.Empty:
|
|
self._write_raw(handler, ": hb\n\n")
|
|
continue
|
|
if item is _STOP:
|
|
reason = "stopped"
|
|
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):
|
|
# Client went away mid-stream: normal (the app reconnects with
|
|
# Last-Event-ID and catches up from the outbox).
|
|
reason = "client-gone"
|
|
finally:
|
|
logger.info("iris: SSE stream closed: %s (%s)", device_id, reason)
|
|
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()
|
|
|
|
# ── POST /v1/media (upload, docs/19 §19.15) ───────────────────────────
|
|
|
|
def _handle_media_upload(self, handler: BaseHTTPRequestHandler, device_id: str) -> None: # noqa: PLR0911
|
|
"""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:
|
|
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):
|
|
# Client went away while we held the poll: normal.
|
|
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.
|
|
|
|
With TLS enabled (``set_tls``) the handshake runs in the
|
|
per-connection thread under a hard timeout — never in the
|
|
``serve_forever`` accept loop. ``ssl.SSLSocket.accept()`` would
|
|
otherwise block that loop inside ``do_handshake()`` on a half-open
|
|
connection (TCP established but the client vanished mid-handshake,
|
|
e.g. a phone losing its network/VPN), and the gateway would stop
|
|
accepting any new device connections until it is restarted.
|
|
"""
|
|
|
|
daemon_threads = True
|
|
allow_reuse_address = True
|
|
|
|
# A client that completes TCP but never finishes the TLS handshake
|
|
# must not hold the connection open indefinitely.
|
|
HANDSHAKE_TIMEOUT_S = 10.0
|
|
|
|
def __init__(self, addr: tuple[str, int], http_server: HttpServer):
|
|
super().__init__(addr, _Handler)
|
|
self.http_server = http_server
|
|
self._tls_ctx: ssl.SSLContext | None = None
|
|
|
|
def set_tls(self, ctx: ssl.SSLContext) -> None:
|
|
self._tls_ctx = ctx
|
|
|
|
def process_request( # noqa: A003 # type: ignore[override]
|
|
self, request: socket.socket, client_address: Any
|
|
) -> None:
|
|
"""Spawn the handler thread; with TLS, the handshake happens in
|
|
that thread first, under ``HANDSHAKE_TIMEOUT_S`` (see class
|
|
docstring). A failed/timed-out handshake just closes the socket —
|
|
the accept loop is never blocked by it."""
|
|
if self._tls_ctx is None:
|
|
super().process_request(request, client_address)
|
|
return
|
|
tls_ctx = self._tls_ctx
|
|
|
|
def _handshake_then_handle() -> None:
|
|
try:
|
|
request.settimeout(self.HANDSHAKE_TIMEOUT_S)
|
|
# wrap_socket() performs the handshake (default
|
|
# do_handshake_on_connect=True); restore blocking mode for
|
|
# the request handler afterwards.
|
|
tls_sock = tls_ctx.wrap_socket(request, server_side=True)
|
|
tls_sock.settimeout(None)
|
|
except OSError as e: # ssl.SSLError, timeout, reset, ...
|
|
with contextlib.suppress(OSError):
|
|
request.close()
|
|
logger.debug("iris http: TLS handshake failed (%s): %s", client_address, e)
|
|
return
|
|
super(_ThreadingHTTPD, self).process_request(tls_sock, client_address)
|
|
|
|
threading.Thread(target=_handshake_then_handle, name="iris-tls", daemon=True).start()
|
|
|
|
|
|
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("iris 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
|
|
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, PLR0911, PLR0912
|
|
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/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
|
|
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:
|
|
# 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,
|
|
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
|
|
hs._handle_frame(self, device_id, frame)
|