Files
iris_x_hermes/gateway-plugin/http_server.py
T
ARIA 8657e6afc6
CI / Kotlin tests (android host + desktop) (push) Successful in 5m48s
CI / Gateway plugin tests (push) Successful in 7m40s
fix(http): test acquire_scoped_lock's bool, not the always-truthy tuple
acquire_scoped_lock returns (acquired, existing_record); the old
'if not acquire_scoped_lock(...)' tested the tuple, which is always
truthy, so the 'port in use by another profile' pre-check never fired
and a conflict surfaced as a generic bind failure. Unpack and test the
first element, matching gateway/platforms/base.py's canonical usage.

Bump VERSION / plugin.yaml to 0.1.3.
2026-08-31 23:17:30 +02:00

917 lines
40 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 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)
httpd.socket = ctx.wrap_socket(httpd.socket, server_side=True)
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
if httpd is not None:
# shutdown() must be called from a thread other than the one
# running serve_forever(); we are on the asyncio loop thread.
with contextlib.suppress(Exception):
httpd.shutdown()
with contextlib.suppress(Exception):
httpd.server_close()
t = self._thread
self._thread = None
if t is not None and t is not threading.current_thread():
t.join(timeout=5.0)
self._release_lock()
def _release_lock(self) -> None:
with contextlib.suppress(ImportError):
from gateway.status import release_scoped_lock
if self._lock_key:
release_scoped_lock("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."""
daemon_threads = True
allow_reuse_address = True
def __init__(self, addr: tuple[str, int], http_server: HttpServer):
super().__init__(addr, _Handler)
self.http_server = http_server
class _Handler(BaseHTTPRequestHandler):
# HTTP/1.0 (default): the connection closes after each response. That
# matches the transport's design (short-lived connections) and avoids
# Content-Length bookkeeping on the streamed SSE response.
server: _ThreadingHTTPD
def log_message(self, fmt: str, *args: Any) -> None: # noqa: A003
logger.debug("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)