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.
535 lines
19 KiB
Python
535 lines
19 KiB
Python
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
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@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)
|