Second, short-lived-connection transport next to the WS: same frames, same outbox/cursor, same token, served over plain HTTP (stdlib ThreadingHTTPServer bridged into the asyncio loop; zero new deps). - http_server.py: /v1/health (unauthenticated), POST /v1/frame (accept-and-ack; validation rejections as 4xx error frames), SSE /v1/events (outbox catch-up with id=cursor, event: hello, 15s heartbeat, bounded-queue backpressure), long-poll /v1/poll (25s hold). Bearer token + X-Iris-Device (same allowlist as WS hello), 64 KiB body cap, per-device rate limit, optional TLS, non-fatal bind failure. - ws_server.py: inbound dispatch chain extracted to shared dispatch_frame() used by both transports. - adapter.py: ANDROID_HTTP_PORT/CERT/KEY config; start/stop next to the WS; delivery counting in _broadcast_or_log (an SSE subscriber is a live subscriber -> no push, docs/19 19.8); _reply() routes point-to-point replies into the in-flight HTTP response (reply sink) or broadcasts when the device has no live WS (19.7); status/typing/ channel events fan out to both transports. - ws_probe.py: --http mode (health + POST + SSE turn drive, same assertion flags); tests/README updated. - Tests: hermes-agent/tests/gateway/test_android_http.py (23 tests, incl. the 19.8 delivery-counting regression); test_android.py (74) still green.
667 lines
27 KiB
Python
667 lines
27 KiB
Python
"""HTTP fallback transport (docs/19): the "HTTP leg".
|
|
|
|
A second, short-lived-connection transport next to the WebSocket: the same
|
|
JSON frames, the same outbox/cursor, the same token — served over plain HTTP
|
|
by the gateway. When the WS is down (flaky network, NAT timeout, app just
|
|
relaunched), the app sends over ``POST /v1/frame`` and receives over
|
|
``GET /v1/events`` (SSE) or ``GET /v1/poll`` (long-poll) instead of waiting
|
|
for a WS redial.
|
|
|
|
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 WS
|
|
accepts (except binary media, which stays
|
|
WS-only in v1).
|
|
* ``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.
|
|
|
|
Auth: ``Authorization: Bearer <token>`` (constant-time ``verify_token``) +
|
|
``X-Iris-Device`` header (same device id / allowlist as the WS ``hello``).
|
|
|
|
Bind failure is NON-fatal (unlike the WS): the plugin keeps working
|
|
WS-only.
|
|
"""
|
|
|
|
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 protocol
|
|
from .pairing import verify_token
|
|
from .ws_server import (
|
|
INBOUND_BURST,
|
|
INBOUND_RATE_PER_S,
|
|
MAX_DEVICE_ID_LEN,
|
|
_TokenBucket,
|
|
dispatch_frame,
|
|
)
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Default port for the HTTP leg (WS default is 8790).
|
|
DEFAULT_HTTP_PORT = 8791
|
|
|
|
# Request body cap for POST /v1/frame (frames are small; media never
|
|
# travels here in v1).
|
|
MAX_BODY_BYTES = 64 * 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()
|
|
|
|
# Frame types POST /v1/frame must not accept (docs/19 §19.3): media is
|
|
# inherently binary/streaming (WS-only in v1); hello/ping are
|
|
# transport-specific (auth is via headers, liveness via /v1/health).
|
|
HTTP_REJECTED_TYPES = frozenset(
|
|
{
|
|
protocol.TYPE_HELLO,
|
|
protocol.TYPE_PING,
|
|
protocol.TYPE_MEDIA_UPLOAD_START,
|
|
protocol.TYPE_MEDIA_UPLOAD_END,
|
|
protocol.TYPE_MEDIA_PULL,
|
|
}
|
|
)
|
|
|
|
|
|
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 fallback 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 the
|
|
same ``dispatch_frame`` the WS server uses.
|
|
"""
|
|
|
|
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, _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 disables the
|
|
HTTP leg (the plugin keeps working WS-only, 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; HTTP leg 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 fallback leg 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 fallback leg 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) > 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 WS inbound limit."""
|
|
with self._buckets_lock:
|
|
b = self._buckets.get(device_id)
|
|
if b is None:
|
|
b = self._buckets[device_id] = _TokenBucket(INBOUND_RATE_PER_S, 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:
|
|
try:
|
|
frames.append(sink.get(timeout=0.05))
|
|
break
|
|
except queue.Empty:
|
|
if task.done():
|
|
# All replies are in the sink now (the handler finished);
|
|
# drain them all.
|
|
while True:
|
|
try:
|
|
frames.append(sink.get_nowait())
|
|
except queue.Empty:
|
|
break
|
|
break
|
|
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
|
|
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_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._adapter._ws_server.broadcast(frame)
|
|
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"))
|
|
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)
|
|
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()
|
|
# 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:
|
|
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):
|
|
pass # client went away: normal
|
|
finally:
|
|
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()
|
|
|
|
# ── 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):
|
|
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
|
|
_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/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:
|
|
_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
|
|
if frame.type in HTTP_REJECTED_TYPES:
|
|
_send_frame_json(
|
|
self,
|
|
400,
|
|
protocol.error(
|
|
protocol.ERR_UNSUPPORTED, f"{frame.type} requires the live connection"
|
|
).to_json(),
|
|
)
|
|
return
|
|
hs._handle_frame(self, device_id, frame)
|