Files
ARIA 2c20b1c8a5
CI / Kotlin tests (android host + desktop) (pull_request) Successful in 8m18s
CI / Gateway plugin tests (pull_request) Failing after 15m7s
fix(http): don't let a half-open TLS connection wedge the accept loop
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.
2026-09-23 15:10:16 +02:00

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)