"""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 `` (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}" if not acquire_scoped_lock("iris", lock_key): 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 try: self._devices.upsert( device_id, device_name or device_id, 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)", device_id, cursor) # 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)