feat: device pairing, telemetry ingestion, live dashboard WS (Workstream G)
Implements the backend half of the ESP32-P4 sensor node spec's pairing, ingestion, and live-broadcast contract: - New Device model (backend/app/models/device.py): id, user_id FK, name, token_hash (unique+indexed), created_at, last_seen_at. Reuses generate_session_token()/hash_token() from auth_session.py verbatim for the one-time raw pairing token / stored hash. - POST /api/device, GET /api/device (session-cookie authenticated REST pairing endpoints) and POST /api/device/telemetry (device bearer-token authenticated ingestion, per-device rate limited, 16KB body cap, 64 reading cap, strict shape validation — never a 500 on garbage input) in backend/app/routes/device.py. - /ws/device-feed live dashboard WS (qm_session cookie authenticated), fanning out ingested readings to the owning user's connected dashboard sockets via an in-process dict[user_id, connections] registry, each with its own send-queue + single sender task (mirrors app.ws's SeanceState/_sender convention). - last_seen_at updates on every successful ingestion. - _process_reading(device, reading) left as an explicit no-op handoff point for Workstream K's summon-pipeline integration. Backend suite: 102 passed. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
321
backend/app/routes/device.py
Normal file
321
backend/app/routes/device.py
Normal file
@@ -0,0 +1,321 @@
|
||||
"""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. It does NOT feed readings into the séance/summon
|
||||
pipeline — that's Workstream K, which extends `_process_reading` below.
|
||||
"""
|
||||
|
||||
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.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.
|
||||
|
||||
Workstream G (this file) intentionally leaves this a no-op: pairing,
|
||||
bearer auth, payload validation/size caps, per-device rate limiting, and
|
||||
the live dashboard broadcast (`_broadcast_reading`, above) are this
|
||||
workstream's full scope.
|
||||
|
||||
Workstream K's summon-pipeline integration extends/replaces this
|
||||
function to run per-(user_id, device_id, sensor_type) rolling-baseline
|
||||
anomaly detection and, when a reading is flagged, push
|
||||
`{"type": "anomaly", "source": <sensor_type>, "frequency": ...,
|
||||
"magnitude": ...}` into the owning user's active `SeanceState.anomalies`
|
||||
— the same path `_handle_anomaly` in `app/ws.py` already uses for the
|
||||
browser-based modes (`wire`/`evp`/`radio`/`emf`). It is async so that
|
||||
integration can await DB/session-registry work; called once per reading
|
||||
in `ingest_telemetry` below, in submission order.
|
||||
"""
|
||||
return None
|
||||
|
||||
|
||||
@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)
|
||||
Reference in New Issue
Block a user