"""Device pairing, telemetry ingestion, and the live device-feed dashboard WebSocket (Workstream G of docs/superpowers/specs/2026-07-23-esp32-sensor-node-design.md). Scope boundary (per the spec): this module owns pairing, ingestion auth/validation, `last_seen_at` bookkeeping, and the in-process dashboard pub/sub broadcast. Summon-pipeline anomaly detection (Workstream K) is wired in at `_process_reading` below, calling into `app.device_anomaly.process_device_reading_for_summon`. """ import asyncio import contextlib import uuid from collections import defaultdict from dataclasses import dataclass, field from datetime import datetime, timezone from fastapi import ( APIRouter, Depends, Header, HTTPException, Request, WebSocket, WebSocketDisconnect, status, ) from pydantic import BaseModel, Field, ValidationError from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.db import async_session_maker as _default_session_maker from app.db import get_db from app.deps import SESSION_COOKIE_NAME, get_current_user from app.device_anomaly import process_device_reading_for_summon from app.models.auth_session import AuthSession, generate_session_token, hash_token from app.models.device import Device from app.models.user import User from app.rate_limit import RateLimiter router = APIRouter(tags=["device"]) # Alias so tests can swap in the NullPool test session maker — same pattern # app.ws uses for its own session_maker (see backend/tests/conftest.py). session_maker = _default_session_maker MAX_TELEMETRY_BODY_BYTES = 16 * 1024 MAX_READINGS_PER_REQUEST = 64 # ~1 request/second sustained is generous for sensor telemetry; keyed # per-device so one chatty/misbehaving node can't starve another. telemetry_limiter = RateLimiter(max_requests=1, window_seconds=1) # --------------------------------------------------------------------------- # Schemas # --------------------------------------------------------------------------- class DeviceCreateIn(BaseModel): name: str = Field(min_length=1, max_length=64) class Reading(BaseModel): # sensor_type is intentionally free-form (not an enum) — "add all sorts # of sensors, anything you can think of" per the project brief. Only # shape is validated here, never a fixed sensor list. sensor_type: str = Field(min_length=1, max_length=32) value: float | list[float] unit: str = Field(min_length=1, max_length=16) metadata: dict = Field(default_factory=dict) class TelemetryIn(BaseModel): readings: list[Reading] = Field(min_length=1, max_length=MAX_READINGS_PER_REQUEST) # --------------------------------------------------------------------------- # REST: pairing (session-cookie authenticated) # --------------------------------------------------------------------------- @router.post("/api/device", status_code=status.HTTP_201_CREATED) async def create_device( payload: DeviceCreateIn, user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): """Creates a device and returns the raw pairing token ONCE — it is never stored or retrievable again, identical convention to `generate_session_token()`/`hash_token()` in `models/auth_session.py`.""" raw_token, token_hash = generate_session_token() device = Device(user_id=user.id, name=payload.name, token_hash=token_hash) db.add(device) await db.commit() await db.refresh(device) return { "id": str(device.id), "name": device.name, "token": raw_token, "created_at": device.created_at.isoformat(), } @router.get("/api/device") async def list_devices( user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): devices = ( await db.execute( select(Device).where(Device.user_id == user.id).order_by(Device.created_at) ) ).scalars().all() return { "devices": [ { "id": str(d.id), "name": d.name, "last_seen_at": d.last_seen_at.isoformat() if d.last_seen_at else None, } for d in devices ] } # --------------------------------------------------------------------------- # Device bearer auth — headless clients only, never the qm_session cookie. # --------------------------------------------------------------------------- async def _authenticate_device( authorization: str | None = Header(default=None), db: AsyncSession = Depends(get_db), ) -> Device: if not authorization or not authorization.startswith("Bearer "): raise HTTPException(status.HTTP_401_UNAUTHORIZED, "missing device bearer token") raw_token = authorization.removeprefix("Bearer ").strip() if not raw_token: raise HTTPException(status.HTTP_401_UNAUTHORIZED, "missing device bearer token") # Look up by hash — the raw token is never compared or stored. token_hash = hash_token(raw_token) device = await db.scalar(select(Device).where(Device.token_hash == token_hash)) if device is None: raise HTTPException(status.HTTP_401_UNAUTHORIZED, "invalid device token") return device # --------------------------------------------------------------------------- # In-process dashboard pub/sub: user_id -> live /ws/device-feed connections. # Each connection gets its own send queue + single sender task, mirroring # app.ws's SeanceState/_sender convention, so concurrent ingestion requests # broadcasting to the same dashboard socket never interleave writes on it. # --------------------------------------------------------------------------- @dataclass class _DashboardConnection: user_id: uuid.UUID send_queue: asyncio.Queue = field(default_factory=asyncio.Queue) _dashboard_connections: dict[uuid.UUID, list[_DashboardConnection]] = defaultdict(list) async def _dashboard_sender(conn: _DashboardConnection, websocket: WebSocket) -> None: """The only task allowed to write to this dashboard socket.""" try: while True: message = await conn.send_queue.get() await websocket.send_json(message) except asyncio.CancelledError: pass async def _broadcast_reading(device: Device, reading: Reading, at: datetime) -> None: connections = _dashboard_connections.get(device.user_id) if not connections: return # no dashboard open for this owner right now — not an error frame = { "type": "reading", "device_id": str(device.id), "sensor_type": reading.sensor_type, "value": reading.value, "unit": reading.unit, "metadata": reading.metadata, "at": at.isoformat(), } for conn in list(connections): await conn.send_queue.put(frame) # --------------------------------------------------------------------------- # Ingestion # --------------------------------------------------------------------------- async def _process_reading(device: Device, reading: Reading) -> None: """Per-reading processing handoff point. Runs the reading through Workstream K's summon-pipeline integration (`process_device_reading_for_summon`): per-(user_id, device_id, sensor_type) rolling-baseline anomaly detection, pushing flagged anomalies into the owning user's active `SeanceState.anomalies` via the same path `_handle_anomaly` in `app/ws.py` already uses for the browser-based modes. A no-op if the user has no active séance session. Array-valued readings (`reading.value: list[float]`) are skipped here — K's detector operates on a single scalar signal, and the contract doesn't define how a multi-element reading collapses to one baseline value. They still reach the live dashboard via `_broadcast_reading` below; only summon-pipeline anomaly detection is skipped for them. """ if isinstance(reading.value, list): return await process_device_reading_for_summon( user_id=device.user_id, device_id=device.id, sensor_type=reading.sensor_type, value=reading.value, unit=reading.unit, ) @router.post("/api/device/telemetry", status_code=status.HTTP_202_ACCEPTED) async def ingest_telemetry( request: Request, device: Device = Depends(_authenticate_device), db: AsyncSession = Depends(get_db), ): if not telemetry_limiter.allow(str(device.id)): raise HTTPException( status.HTTP_429_TOO_MANY_REQUESTS, "telemetry rate limit exceeded" ) body = await request.body() if len(body) > MAX_TELEMETRY_BODY_BYTES: raise HTTPException( status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, "telemetry payload too large" ) try: payload = TelemetryIn.model_validate_json(body) except ValidationError: raise HTTPException( status.HTTP_422_UNPROCESSABLE_ENTITY, "malformed telemetry payload" ) device.last_seen_at = datetime.now(timezone.utc) await db.commit() at = datetime.now(timezone.utc) for reading in payload.readings: await _process_reading(device, reading) await _broadcast_reading(device, reading, at) return {"status": "received", "count": len(payload.readings)} # --------------------------------------------------------------------------- # Live dashboard WS — session-cookie authenticated (browser client, unlike # the ingestion endpoint above). # --------------------------------------------------------------------------- async def _authenticate_dashboard(websocket: WebSocket) -> uuid.UUID | None: raw_token = websocket.cookies.get(SESSION_COOKIE_NAME) if raw_token is None: return None token_hash = hash_token(raw_token) async with session_maker() as db: result = await db.execute( select(AuthSession).where(AuthSession.token_hash == token_hash) ) session = result.scalar_one_or_none() if session is None or session.expires_at < datetime.now(timezone.utc): return None return session.user_id @router.websocket("/ws/device-feed") async def device_feed_socket(websocket: WebSocket) -> None: user_id = await _authenticate_dashboard(websocket) if user_id is None: await websocket.close(code=4401) return await websocket.accept() conn = _DashboardConnection(user_id=user_id) _dashboard_connections[user_id].append(conn) sender = asyncio.create_task(_dashboard_sender(conn, websocket)) async with session_maker() as db: devices = ( await db.execute(select(Device).where(Device.user_id == user_id)) ).scalars().all() await conn.send_queue.put( { "type": "devices", "devices": [ { "id": str(d.id), "name": d.name, "last_seen_at": d.last_seen_at.isoformat() if d.last_seen_at else None, } for d in devices ], } ) try: while True: # The dashboard is receive-only in practice; this just blocks # until disconnect while _dashboard_sender does all the writing. await websocket.receive_text() except WebSocketDisconnect: pass finally: sender.cancel() with contextlib.suppress(asyncio.CancelledError): await sender with contextlib.suppress(ValueError): _dashboard_connections[user_id].remove(conn) if not _dashboard_connections.get(user_id): _dashboard_connections.pop(user_id, None)