Wire Workstream K's summon-pipeline integration into G's ingestion handoff
Both G and K were built independently against the paper contract and correctly left this connection point for the integrator (documented in both their reports). Calls process_device_reading_for_summon from _process_reading, skipping array-valued readings (no defined single scalar baseline for those). Adds an end-to-end integration test that connects a real /ws/session, posts device telemetry through it, and confirms an anomalous reading actually reaches the live session as an anomaly_ack — proving the two independently-built pieces genuinely connect, not just that each compiles. 165/165 backend tests pass.
This commit is contained in:
@@ -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
|
Scope boundary (per the spec): this module owns pairing, ingestion
|
||||||
auth/validation, `last_seen_at` bookkeeping, and the in-process dashboard
|
auth/validation, `last_seen_at` bookkeeping, and the in-process dashboard
|
||||||
pub/sub broadcast. It does NOT feed readings into the séance/summon
|
pub/sub broadcast. Summon-pipeline anomaly detection (Workstream K) is
|
||||||
pipeline — that's Workstream K, which extends `_process_reading` below.
|
wired in at `_process_reading` below, calling into
|
||||||
|
`app.device_anomaly.process_device_reading_for_summon`.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import asyncio
|
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 async_session_maker as _default_session_maker
|
||||||
from app.db import get_db
|
from app.db import get_db
|
||||||
from app.deps import SESSION_COOKIE_NAME, get_current_user
|
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.auth_session import AuthSession, generate_session_token, hash_token
|
||||||
from app.models.device import Device
|
from app.models.device import Device
|
||||||
from app.models.user import User
|
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:
|
async def _process_reading(device: Device, reading: Reading) -> None:
|
||||||
"""Per-reading processing handoff point.
|
"""Per-reading processing handoff point.
|
||||||
|
|
||||||
Workstream G (this file) intentionally leaves this a no-op: pairing,
|
Runs the reading through Workstream K's summon-pipeline integration
|
||||||
bearer auth, payload validation/size caps, per-device rate limiting, and
|
(`process_device_reading_for_summon`): per-(user_id, device_id,
|
||||||
the live dashboard broadcast (`_broadcast_reading`, above) are this
|
sensor_type) rolling-baseline anomaly detection, pushing flagged
|
||||||
workstream's full scope.
|
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
|
Array-valued readings (`reading.value: list[float]`) are skipped here —
|
||||||
function to run per-(user_id, device_id, sensor_type) rolling-baseline
|
K's detector operates on a single scalar signal, and the contract
|
||||||
anomaly detection and, when a reading is flagged, push
|
doesn't define how a multi-element reading collapses to one baseline
|
||||||
`{"type": "anomaly", "source": <sensor_type>, "frequency": ...,
|
value. They still reach the live dashboard via `_broadcast_reading`
|
||||||
"magnitude": ...}` into the owning user's active `SeanceState.anomalies`
|
below; only summon-pipeline anomaly detection is skipped for them.
|
||||||
— 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
|
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)
|
@router.post("/api/device/telemetry", status_code=status.HTTP_202_ACCEPTED)
|
||||||
|
|||||||
@@ -1,11 +1,14 @@
|
|||||||
import hashlib
|
import hashlib
|
||||||
|
import uuid
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
from sqlalchemy import select
|
from sqlalchemy import select
|
||||||
|
|
||||||
import app.routes.device as device_module
|
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.models.device import Device
|
||||||
from app.rate_limit import RateLimiter
|
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
|
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
|
# Rate limiting
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|||||||
Reference in New Issue
Block a user