diff --git a/backend/app/routes/device.py b/backend/app/routes/device.py index 6149c17..6319f17 100644 --- a/backend/app/routes/device.py +++ b/backend/app/routes/device.py @@ -4,8 +4,9 @@ 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. +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 @@ -32,6 +33,7 @@ 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 @@ -198,22 +200,28 @@ async def _broadcast_reading(device: Device, reading: Reading, at: datetime) -> 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. + 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. - 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. + 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. """ - return None + 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) diff --git a/backend/tests/test_device.py b/backend/tests/test_device.py index 75f924b..95cb269 100644 --- a/backend/tests/test_device.py +++ b/backend/tests/test_device.py @@ -1,11 +1,14 @@ import hashlib +import uuid import pytest from sqlalchemy import select import app.routes.device as device_module +from app.device_anomaly import reset_state as reset_device_anomaly_state from app.models.device import Device from app.rate_limit import RateLimiter +from app.ws import get_active_session # --------------------------------------------------------------------------- @@ -267,6 +270,111 @@ async def test_telemetry_updates_last_seen_at(client): assert after.json()["devices"][0]["last_seen_at"] is not None +# --------------------------------------------------------------------------- +# Integration: a device-triggered anomaly actually reaches an active séance +# (Workstream K's process_device_reading_for_summon, wired in at +# _process_reading — proves the two workstreams' independently-built pieces +# actually connect end to end through the real ingestion endpoint). +# +# Connects a real /ws/session (which creates a genuine ContactSession row +# and self-registers via app.ws's active-session registry) rather than +# hand-constructing a SeanceState — a manually-built one would have no +# backing ContactSession row, and _handle_anomaly's _record_event would hit +# a foreign-key violation. Uses sync_client throughout: _handle_anomaly +# writes via app.ws's own module-level session_maker, which only +# sync_client's fixture swaps to the NullPool test engine (see +# conftest.py) — the same reason test_ws_session.py uses it too. +# --------------------------------------------------------------------------- + + +def _ws_session_connect(sync_client, token): + return sync_client.websocket_connect( + "/ws/session", headers={"cookie": f"qm_session={token}"} + ) + + +def _read_until(ws, msg_type, max_frames=30, **match): + for _ in range(max_frames): + frame = ws.receive_json() + if frame.get("type") != msg_type: + continue + if all(frame.get(key) == value for key, value in match.items()): + return frame + raise AssertionError(f"never saw frame of type {msg_type!r} matching {match!r}") + + +def test_anomalous_reading_reaches_active_seance_session(sync_client, monkeypatch): + reset_device_anomaly_state() + # This test posts twice in quick succession — the real per-device + # telemetry_limiter (~1/sec) would reject the second post as 429 before + # it's even processed, since both land in the same window. + monkeypatch.setattr( + device_module, "telemetry_limiter", RateLimiter(max_requests=100, window_seconds=60) + ) + + register_resp = sync_client.post( + "/auth/register", json={"username": "hwseeker1", "password": "spookyspooky"} + ) + sync_client.post( + "/auth/login", json={"username": "hwseeker1", "password": "spookyspooky"} + ) + token = sync_client.cookies.get("qm_session") + user_id = uuid.UUID(register_resp.json()["id"]) + pair_resp = sync_client.post("/api/device", json={"name": "Presence Node"}) + device = pair_resp.json() + headers = {"Authorization": f"Bearer {device['token']}"} + + with _ws_session_connect(sync_client, token) as ws: + _read_until(ws, "session") + assert get_active_session(user_id) is not None + + # First reading establishes the "nothing detected" baseline — not + # itself anomalous (see detect_boolean_transition's None-previous + # guard), so nothing new should arrive on the socket for it. Confirm + # via a ping/pong round-trip instead of a fixed sleep. + sync_client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading("presence", 0, "bool")]}, + headers=headers, + ) + ws.send_json({"type": "ping"}) + _read_until(ws, "pong") + + # Second reading transitions false->true: a genuine anomaly, must + # reach the active session and surface as an anomaly_ack — exactly + # like the browser-based modes' own anomaly frames do. + sync_client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading("presence", 1, "bool")]}, + headers=headers, + ) + ack = _read_until(ws, "anomaly_ack") + assert ack["count"] == 1 + + reset_device_anomaly_state() + + +def test_reading_for_user_with_no_active_session_is_a_quiet_noop(sync_client): + reset_device_anomaly_state() + sync_client.post( + "/auth/register", json={"username": "hwseeker2", "password": "spookyspooky"} + ) + sync_client.post( + "/auth/login", json={"username": "hwseeker2", "password": "spookyspooky"} + ) + pair_resp = sync_client.post("/api/device", json={"name": "Idle Node"}) + device = pair_resp.json() + headers = {"Authorization": f"Bearer {device['token']}"} + + response = sync_client.post( + "/api/device/telemetry", + json={"readings": [_valid_reading("presence", 1, "bool")]}, + headers=headers, + ) + assert response.status_code == 202 + reset_device_anomaly_state() + + # --------------------------------------------------------------------------- # Rate limiting # ---------------------------------------------------------------------------