From c88fbc843a0fa47a1bb2eabb3b34c71ef9949283 Mon Sep 17 00:00:00 2001 From: Indiana Date: Fri, 24 Jul 2026 01:12:21 +0000 Subject: [PATCH] feat: device pairing, telemetry ingestion, live dashboard WS (Workstream G) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- backend/app/main.py | 2 + backend/app/models/__init__.py | 2 + backend/app/models/device.py | 27 +++ backend/app/routes/device.py | 321 +++++++++++++++++++++++++ backend/tests/conftest.py | 2 + backend/tests/test_device.py | 426 +++++++++++++++++++++++++++++++++ 6 files changed, 780 insertions(+) create mode 100644 backend/app/models/device.py create mode 100644 backend/app/routes/device.py create mode 100644 backend/tests/test_device.py diff --git a/backend/app/main.py b/backend/app/main.py index 5f9a429..38a7081 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -12,6 +12,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.shop import router as shop_router from app.session_cleanup import delete_expired_sessions from app.ws import AUDIO_DIR @@ -57,6 +58,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(shop_router) app.include_router(ws_router) diff --git a/backend/app/models/__init__.py b/backend/app/models/__init__.py index 02b6c67..46cd5ff 100644 --- a/backend/app/models/__init__.py +++ b/backend/app/models/__init__.py @@ -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 @@ -10,6 +11,7 @@ __all__ = [ "User", "AuthSession", "ContactSession", + "Device", "Entity", "EntitySighting", "Event", diff --git a/backend/app/models/device.py b/backend/app/models/device.py new file mode 100644 index 0000000..9454802 --- /dev/null +++ b/backend/app/models/device.py @@ -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 + ) diff --git a/backend/app/routes/device.py b/backend/app/routes/device.py new file mode 100644 index 0000000..6149c17 --- /dev/null +++ b/backend/app/routes/device.py @@ -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": , "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) diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 186b277..62e00b1 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -49,9 +49,11 @@ def sync_client(monkeypatch): swapped: TestClient runs the lifespan on its own portal loop, and the pooled production engine would carry connections across loops.""" import app.main as main_module + import app.routes.device as device_module import app.ws as ws_module monkeypatch.setattr(ws_module, "session_maker", TestSessionLocal) + monkeypatch.setattr(device_module, "session_maker", TestSessionLocal) monkeypatch.setattr(main_module, "engine", test_engine) with TestClient(app, base_url="https://testserver") as tc: yield tc diff --git a/backend/tests/test_device.py b/backend/tests/test_device.py new file mode 100644 index 0000000..75f924b --- /dev/null +++ b/backend/tests/test_device.py @@ -0,0 +1,426 @@ +import hashlib + +import pytest +from sqlalchemy import select + +import app.routes.device as device_module +from app.models.device import Device +from app.rate_limit import RateLimiter + + +# --------------------------------------------------------------------------- +# Async-client helpers (REST-only tests) +# --------------------------------------------------------------------------- + + +async def _register_and_login(client, username="devowner"): + await client.post("/auth/register", json={"username": username, "password": "spookyspooky"}) + await client.post("/auth/login", json={"username": username, "password": "spookyspooky"}) + + +async def _pair_device(client, name="Sensor Node 1") -> dict: + response = await client.post("/api/device", json={"name": name}) + assert response.status_code == 201 + return response.json() + + +def _valid_reading(sensor_type="temperature", value=21.4, unit="c"): + return {"sensor_type": sensor_type, "value": value, "unit": unit, "metadata": {}} + + +# --------------------------------------------------------------------------- +# Pairing REST endpoints +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_create_device_requires_session_cookie(client): + response = await client.post("/api/device", json={"name": "Sensor Node"}) + assert response.status_code == 401 + + +@pytest.mark.asyncio +async def test_create_device_returns_raw_token_once(client, db_session): + await _register_and_login(client, "pairer1") + body = await _pair_device(client, "Attic Node") + + assert body["name"] == "Attic Node" + assert "token" in body and len(body["token"]) > 20 + assert "id" in body and "created_at" in body + assert "token_hash" not in body + + # The stored hash is the sha256 of the raw token — same convention as + # AuthSession.token_hash / hash_token(). + device = await db_session.scalar(select(Device).where(Device.id == body["id"])) + assert device.token_hash == hashlib.sha256(body["token"].encode()).hexdigest() + + +@pytest.mark.asyncio +async def test_list_devices_never_exposes_token_or_hash(client): + await _register_and_login(client, "pairer2") + await _pair_device(client, "Basement Node") + + response = await client.get("/api/device") + assert response.status_code == 200 + devices = response.json()["devices"] + assert len(devices) == 1 + assert devices[0]["name"] == "Basement Node" + assert devices[0]["last_seen_at"] is None + assert "token" not in devices[0] + assert "token_hash" not in devices[0] + + +@pytest.mark.asyncio +async def test_list_devices_requires_session_cookie(client): + response = await client.get("/api/device") + assert response.status_code == 401 + + +@pytest.mark.asyncio +async def test_list_devices_scoped_to_owner(client): + await _register_and_login(client, "pairer3a") + await _pair_device(client, "Owner A Node") + await client.post("/auth/logout") + + await _register_and_login(client, "pairer3b") + response = await client.get("/api/device") + assert response.json()["devices"] == [] + + +# --------------------------------------------------------------------------- +# Telemetry ingestion — auth +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_telemetry_rejects_missing_bearer_token(client): + response = await client.post( + "/api/device/telemetry", json={"readings": [_valid_reading()]} + ) + assert response.status_code == 401 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_malformed_authorization_header(client): + response = await client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": "Token not-a-bearer-scheme"}, + ) + assert response.status_code == 401 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_unknown_token(client): + response = await client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": "Bearer totally-made-up-token"}, + ) + assert response.status_code == 401 + + +@pytest.mark.asyncio +async def test_telemetry_accepts_valid_bearer_token(client): + await _register_and_login(client, "ingest1") + device = await _pair_device(client, "Living Room Node") + + response = await client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 202 + assert response.json()["count"] == 1 + + +# --------------------------------------------------------------------------- +# Telemetry ingestion — payload validation (never a 500) +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_telemetry_rejects_oversized_body(client): + await _register_and_login(client, "ingest2") + device = await _pair_device(client, "Garage Node") + + huge_reading = _valid_reading(sensor_type="temperature") + huge_reading["metadata"] = {"blob": "x" * 20_000} + response = await client.post( + "/api/device/telemetry", + json={"readings": [huge_reading]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 413 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_oversized_readings_array(client): + await _register_and_login(client, "ingest3") + device = await _pair_device(client, "Hallway Node") + + readings = [_valid_reading(sensor_type=f"sensor{i}") for i in range(65)] + response = await client.post( + "/api/device/telemetry", + json={"readings": readings}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 422 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_garbage_sensor_type_type(client): + await _register_and_login(client, "ingest4") + device = await _pair_device(client, "Cellar Node") + + bad = _valid_reading() + bad["sensor_type"] = 12345 # must be a string + response = await client.post( + "/api/device/telemetry", + json={"readings": [bad]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 422 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_garbage_value_type(client): + await _register_and_login(client, "ingest5") + device = await _pair_device(client, "Attic Node 2") + + bad = _valid_reading() + bad["value"] = "not-a-number" + response = await client.post( + "/api/device/telemetry", + json={"readings": [bad]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 422 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_garbage_unit_type(client): + await _register_and_login(client, "ingest6") + device = await _pair_device(client, "Loft Node") + + bad = _valid_reading() + bad["unit"] = {"nested": "dict"} + response = await client.post( + "/api/device/telemetry", + json={"readings": [bad]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 422 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_non_dict_metadata(client): + await _register_and_login(client, "ingest7") + device = await _pair_device(client, "Porch Node") + + bad = _valid_reading() + bad["metadata"] = ["not", "a", "dict"] + response = await client.post( + "/api/device/telemetry", + json={"readings": [bad]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 422 + + +@pytest.mark.asyncio +async def test_telemetry_rejects_malformed_json_body(client): + await _register_and_login(client, "ingest8") + device = await _pair_device(client, "Yard Node") + + response = await client.post( + "/api/device/telemetry", + content=b"{not valid json", + headers={ + "Authorization": f"Bearer {device['token']}", + "Content-Type": "application/json", + }, + ) + assert response.status_code == 422 + + +# --------------------------------------------------------------------------- +# last_seen_at bookkeeping +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_telemetry_updates_last_seen_at(client): + await _register_and_login(client, "seenat1") + device = await _pair_device(client, "Cave Node") + + before = await client.get("/api/device") + assert before.json()["devices"][0]["last_seen_at"] is None + + await client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + + after = await client.get("/api/device") + assert after.json()["devices"][0]["last_seen_at"] is not None + + +# --------------------------------------------------------------------------- +# Rate limiting +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_telemetry_rate_limited_per_device(client, monkeypatch): + monkeypatch.setattr( + device_module, "telemetry_limiter", RateLimiter(max_requests=1, window_seconds=60) + ) + await _register_and_login(client, "ratelimited1") + device = await _pair_device(client, "Rate Limited Node") + headers = {"Authorization": f"Bearer {device['token']}"} + + first = await client.post( + "/api/device/telemetry", json={"readings": [_valid_reading()]}, headers=headers + ) + assert first.status_code == 202 + + second = await client.post( + "/api/device/telemetry", json={"readings": [_valid_reading()]}, headers=headers + ) + assert second.status_code == 429 + + +@pytest.mark.asyncio +async def test_telemetry_rate_limit_is_per_device_not_global(client, monkeypatch): + monkeypatch.setattr( + device_module, "telemetry_limiter", RateLimiter(max_requests=1, window_seconds=60) + ) + await _register_and_login(client, "ratelimited2") + device_a = await _pair_device(client, "Node A") + device_b = await _pair_device(client, "Node B") + + resp_a = await client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": f"Bearer {device_a['token']}"}, + ) + assert resp_a.status_code == 202 + + # Device B has spent nothing yet — its own bucket is untouched. + resp_b = await client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": f"Bearer {device_b['token']}"}, + ) + assert resp_b.status_code == 202 + + +# --------------------------------------------------------------------------- +# Live dashboard WS — pub/sub fan-out (needs a synchronous client that can +# hold a WS connection open alongside plain HTTP calls, same as ws.py's +# tests use `sync_client` for). +# --------------------------------------------------------------------------- + + +def _sync_register_and_login(sync_client, username="dashuser"): + sync_client.post("/auth/register", json={"username": username, "password": "spookyspooky"}) + sync_client.post("/auth/login", json={"username": username, "password": "spookyspooky"}) + return sync_client.cookies.get("qm_session") + + +def _sync_pair_device(sync_client, name="Sync Node") -> dict: + response = sync_client.post("/api/device", json={"name": name}) + assert response.status_code == 201 + return response.json() + + +def _ws_feed_connect(sync_client, token): + # Same rationale as app.ws's tests: TestClient upgrades over ws://, so + # the jar withholds the Secure qm_session cookie — pass it explicitly. + return sync_client.websocket_connect( + "/ws/device-feed", headers={"cookie": f"qm_session={token}"} + ) + + +def test_device_feed_requires_session_cookie(sync_client): + with pytest.raises(Exception): + with sync_client.websocket_connect("/ws/device-feed"): + pass + + +def test_device_feed_sends_device_list_on_connect(sync_client): + token = _sync_register_and_login(sync_client, "dashuser1") + device = _sync_pair_device(sync_client, "Feed Node") + + with _ws_feed_connect(sync_client, token) as ws: + frame = ws.receive_json() + assert frame["type"] == "devices" + assert len(frame["devices"]) == 1 + assert frame["devices"][0]["id"] == device["id"] + assert frame["devices"][0]["name"] == "Feed Node" + assert "token" not in frame["devices"][0] + + +def test_reading_posted_while_dashboard_connected_arrives_on_socket(sync_client): + token = _sync_register_and_login(sync_client, "dashuser2") + device = _sync_pair_device(sync_client, "Live Node") + + with _ws_feed_connect(sync_client, token) as ws: + ws.receive_json() # initial "devices" frame + + response = sync_client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading(sensor_type="presence", value=1, unit="bool")]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 202 + + reading_frame = ws.receive_json() + assert reading_frame["type"] == "reading" + assert reading_frame["device_id"] == device["id"] + assert reading_frame["sensor_type"] == "presence" + assert reading_frame["value"] == 1 + assert reading_frame["unit"] == "bool" + assert reading_frame["metadata"] == {} + assert "at" in reading_frame + + +def test_reading_posted_for_device_with_no_dashboard_owner_does_not_error(sync_client): + _sync_register_and_login(sync_client, "dashuser3") + device = _sync_pair_device(sync_client, "Lonely Node") + + # No /ws/device-feed connection is open for this user at all. + response = sync_client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": f"Bearer {device['token']}"}, + ) + assert response.status_code == 202 + + +def test_device_feed_only_broadcasts_to_the_owning_user(sync_client): + token_a = _sync_register_and_login(sync_client, "dashuser4a") + _sync_pair_device(sync_client, "User A Node") + + _sync_register_and_login(sync_client, "dashuser4b") + device_b = _sync_pair_device(sync_client, "User B Node") + + with _ws_feed_connect(sync_client, token_a) as ws_a: + ws_a.receive_json() # devices frame for user A (empty-ish/own list) + + sync_client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading()]}, + headers={"Authorization": f"Bearer {device_b['token']}"}, + ) + + # User A's socket must not receive user B's device reading. + # WebSocketTestSession.receive_json() has no timeout param, so poll + # its underlying queue directly with a short timeout instead of + # blocking forever waiting for a frame that must never arrive. + import queue + + with pytest.raises(queue.Empty): + ws_a._send_queue.get(timeout=0.3)