Merge Workstream G: device pairing + ingestion + live broadcast

Resolved conflict in main.py: combined both workstreams' router imports
and registrations (device_router from G, inventory_router from C).
143/143 backend tests pass after cleaning stray pollution from an
earlier parallel workstream run against the shared test DB.
This commit is contained in:
Indiana
2026-07-24 14:57:10 +00:00
6 changed files with 780 additions and 0 deletions

View File

@@ -13,6 +13,7 @@ from app.config import settings
from app.db import Base, async_session_maker, engine
from app.routes.auth import router as auth_router
from app.routes.codex import router as codex_router
from app.routes.device import router as device_router
from app.routes.inventory import router as inventory_router
from app.routes.shop import router as shop_router
from app.session_cleanup import delete_expired_sessions
@@ -65,6 +66,7 @@ async def lifespan(app: FastAPI):
app = FastAPI(title="Quantumancy", lifespan=lifespan)
app.include_router(auth_router)
app.include_router(codex_router)
app.include_router(device_router)
app.include_router(inventory_router)
app.include_router(shop_router)
app.include_router(ws_router)

View File

@@ -1,5 +1,6 @@
from app.models.auth_session import AuthSession
from app.models.contact_session import ContactSession
from app.models.device import Device
from app.models.entity import Entity
from app.models.entity_sighting import EntitySighting
from app.models.event import Event
@@ -13,6 +14,7 @@ __all__ = [
"User",
"AuthSession",
"ContactSession",
"Device",
"Entity",
"EntitySighting",
"Event",

View File

@@ -0,0 +1,27 @@
import uuid
from datetime import datetime, timezone
from sqlalchemy import DateTime, ForeignKey, String
from sqlalchemy.orm import Mapped, mapped_column
from app.db import Base
class Device(Base):
"""A paired hardware sensor node (ESP32-P4 etc.), belonging to exactly
one user. Authenticates telemetry via a bearer token whose hash — never
the raw value — is stored here, identical convention to
`AuthSession.token_hash` (see `app/models/auth_session.py`)."""
__tablename__ = "devices"
id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4)
user_id: Mapped[uuid.UUID] = mapped_column(ForeignKey("users.id"), index=True)
name: Mapped[str] = mapped_column(String(64))
token_hash: Mapped[str] = mapped_column(String(64), unique=True, index=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), default=lambda: datetime.now(timezone.utc)
)
last_seen_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), nullable=True
)

View 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)