852 lines
36 KiB
Python
852 lines
36 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}"
|
|
if not acquire_scoped_lock("android", lock_key):
|
|
logger.warning(
|
|
"android: 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("android: 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="android-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("android: 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("android", 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)
|
|
|
|
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("android: 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."""
|
|
auth = handler.headers.get("Authorization") or ""
|
|
token = auth[len("Bearer ") :] if auth.startswith("Bearer ") else None
|
|
if not verify_token(token, self._adapter.token):
|
|
_send_json(handler, 401, {"error": "unauthorized"})
|
|
return None
|
|
device_id = (handler.headers.get("X-Iris-Device") or "").strip()
|
|
if not device_id or len(device_id) > dispatch.MAX_DEVICE_ID_LEN:
|
|
_send_json(handler, 401, {"error": "X-Iris-Device header required"})
|
|
return None
|
|
if (
|
|
not self._adapter.allow_all
|
|
and self._adapter.allowed_users
|
|
and device_id not in self._adapter.allowed_users
|
|
):
|
|
logger.warning("android: 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("android: HTTP dispatch failed for %s", frame.type, exc_info=True)
|
|
finally:
|
|
# Pop our sink entry (a newer request from the same device may
|
|
# have replaced it). If the HTTP response was already sent
|
|
# (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:
|
|
qs = parse_qs(parsed.query)
|
|
cursor = _parse_cursor(qs.get("cursor", [None])[0], handler.headers.get("Last-Event-ID"))
|
|
# Device registration (the HTTP equivalent of the WS hello upsert):
|
|
# the SSE open carries the device name + push tokens as optional
|
|
# headers; upsert is idempotent and COALESCEs absent tokens, so a
|
|
# re-open never clobbers a newer fcm.register value.
|
|
device_name = (handler.headers.get("X-Iris-Device-Name") or "").strip()[:120]
|
|
fcm_token = handler.headers.get("X-Iris-Fcm-Token") or None
|
|
ntfy_topic = handler.headers.get("X-Iris-Ntfy-Topic") or None
|
|
try:
|
|
self._devices.upsert(
|
|
device_id,
|
|
device_name or device_id,
|
|
None,
|
|
fcm_token,
|
|
ntfy_topic,
|
|
)
|
|
except Exception:
|
|
logger.warning("android: device registry upsert failed", exc_info=True)
|
|
sub = _Subscriber(device_id=device_id, kind="sse")
|
|
# Register BEFORE the replay so a frame appended in between is
|
|
# fanned out to us (and de-duped by cursor below) instead of lost.
|
|
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("android: 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),
|
|
)
|
|
self._write_sse(handler, "hello", None, hello.to_json())
|
|
self._write_sse(
|
|
handler, "frame", None, protocol.status(self._adapter.gateway_status()).to_json()
|
|
)
|
|
# 3. Live frames (cursor=None frames have no id).
|
|
while not sub.closed.is_set():
|
|
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("android: 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:
|
|
"""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("android 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
|
|
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)
|