Merge Workstream K: hardware anomalies feed the summon pipeline

This commit is contained in:
Indiana
2026-07-24 14:57:14 +00:00
3 changed files with 581 additions and 0 deletions

View File

@@ -0,0 +1,294 @@
"""Hardware sensor anomalies: the sixth summon source (spec §"Contract",
Workstream K). Paired ESP32 devices stream readings through Workstream G's
`POST /api/device/telemetry` ingestion endpoint; this module is the
self-contained detection+push logic that endpoint calls into per reading.
Nothing here talks HTTP or owns the `Device` model — it is pure detection
state plus a thin bridge into the existing `_handle_anomaly` séance
pipeline in `app.ws`, so hardware summons reuse 100% of the existing
signature/mint/Codex machinery.
Two detectors, one per sensor shape:
* Continuous numeric sensors (temperature, humidity, pressure, and any
future numeric sensor_type) get a rolling-baseline statistical detector,
structurally the same three guards as `telemetry.detect_wire_spike`:
minimum sample count, an absolute floor, and a 3-sigma + relative
threshold. The exact floor/threshold numbers are adapted per sensor_type
since units vary wildly (see `_deviation_floor` below) — a 20KB/s floor
makes no sense for a temperature in degrees C.
* Discrete/boolean sensors (presence, contact switches, etc.) get a simple
false->true state-transition detector — a spike in a 0/1 signal isn't
"statistically anomalous" in any meaningful sense, it's just a change of
state, and only one direction of that change (something appearing) is
the anomaly we care about.
Rolling per-key history/state lives in this module as plain in-process
dicts, matching how `detect_wire_spike`'s caller (`SeanceState.
wire_jitter_history`) manages its own history today — no external store,
doesn't need to survive a restart.
"""
import hashlib
import uuid
from statistics import pstdev
# ---------------------------------------------------------------------------
# Numeric (continuous) detector
# ---------------------------------------------------------------------------
# How many readings before a sensor's baseline is trusted. Matches
# detect_wire_spike's default (6) — devices report roughly once/second per
# the contract's rate cap, so 6 samples is a handful of seconds of warm-up,
# not a long cold-start.
DEFAULT_MIN_SAMPLES = 6
# How much history to retain per (user, device, sensor_type) key. Bounded so
# a chatty device can't grow this dict without limit.
MAX_HISTORY = 60
# Per-sensor-type absolute floors on |current - mean| deviation, in the
# sensor's own units. Unlike detect_wire_spike's floor (a floor on the raw
# value itself — near-silent byte counters should stay silent), a generic
# numeric sensor can legitimately sit at any value including zero or
# negative (temperature), so the floor here is on the *deviation*, not the
# reading. These numbers are a deliberately generous noise band per sensor
# type — reasoned from typical cheap-sensor (BME280-class) precision and
# ordinary ambient drift, not a spec:
# temperature (C): +-0.8 is well within a room's normal thermal drift
# (HVAC cycling, a door opening) and BME280 datasheet noise.
# humidity (%RH): +-3 is a normal hygrometer noise band.
# pressure (hPa): +-1.5 is routine weather drift over the timescale of
# a rolling baseline; BME280 itself is accurate to ~1hPa.
_KNOWN_DEVIATION_FLOORS: dict[str, float] = {
"temperature": 0.8,
"humidity": 3.0,
"pressure": 1.5,
}
# For an unrecognized sensor_type (contract: sensor_type is free-form, not
# an enum — "add all sorts of sensors... anything you can think of"), we
# have no unit-specific noise band to reach for. Fall back to a floor
# proportional to the rolling mean's own magnitude, with a small absolute
# minimum so a baseline near zero doesn't collapse the floor to nothing.
_DEFAULT_FLOOR_FRACTION = 0.15
_DEFAULT_FLOOR_MINIMUM = 0.01
# The relative half of the "3-sigma + relative" guard. detect_wire_spike
# requires the surge be >2.5x the mean; that multiplier doesn't translate
# cleanly to arbitrary (possibly signed, possibly near-zero-mean) sensor
# units, so here the relative check is expressed as a multiple of the
# deviation floor instead — "not just past the noise floor, comfortably
# past it" plays the same role as "not just above baseline, 2.5x above it".
RELATIVE_FLOOR_MULTIPLIER = 2.5
def _deviation_floor(sensor_type: str, mean: float) -> float:
if sensor_type in _KNOWN_DEVIATION_FLOORS:
return _KNOWN_DEVIATION_FLOORS[sensor_type]
return max(abs(mean) * _DEFAULT_FLOOR_FRACTION, _DEFAULT_FLOOR_MINIMUM)
def detect_sensor_spike(
history: list[float],
current: float,
sensor_type: str,
*,
min_samples: int = DEFAULT_MIN_SAMPLES,
) -> float | None:
"""Pure detector, modeled directly on `telemetry.detect_wire_spike`.
Returns the deviation magnitude (in the sensor's own units) when
`current` is anomalous vs. the rolling baseline, or None. Same three
guards as the wire detector: enough history, a floor so sensor noise
never cries ghost, and a combined statistical + relative threshold.
"""
if len(history) < min_samples:
return None
mean = sum(history) / len(history)
deviation = current - mean
abs_deviation = abs(deviation)
floor = _deviation_floor(sensor_type, mean)
if abs_deviation < floor:
return None
std = pstdev(history)
if abs_deviation > 3 * std and abs_deviation > floor * RELATIVE_FLOOR_MULTIPLIER:
return abs_deviation
return None
# ---------------------------------------------------------------------------
# Boolean / discrete (state-transition) detector
# ---------------------------------------------------------------------------
def detect_boolean_transition(previous: bool | None, current: bool) -> bool:
"""The anomaly is a false->true transition (e.g. `presence` going from
nothing detected to something detected). `previous=None` means we have
no prior reading for this key yet — nothing to transition *from*, so
it is never itself an anomaly (avoids flagging every device's very
first reading just because it happens to be `True`)."""
return previous is False and current is True
# ---------------------------------------------------------------------------
# Per-(user_id, device_id, sensor_type) rolling state
# ---------------------------------------------------------------------------
_HistoryKey = tuple[uuid.UUID, uuid.UUID, str]
# In-memory only, matching detect_wire_spike's caller-managed history in
# ws.py — does not need to survive a process restart.
_numeric_history: dict[_HistoryKey, list[float]] = {}
_boolean_state: dict[_HistoryKey, bool] = {}
def _key(user_id: uuid.UUID, device_id: uuid.UUID, sensor_type: str) -> _HistoryKey:
return (user_id, device_id, sensor_type)
def reset_state() -> None:
"""Test/debug helper: clear all rolling history and boolean state."""
_numeric_history.clear()
_boolean_state.clear()
def _record_numeric(
user_id: uuid.UUID, device_id: uuid.UUID, sensor_type: str, value: float
) -> float | None:
key = _key(user_id, device_id, sensor_type)
history = _numeric_history.setdefault(key, [])
deviation = detect_sensor_spike(history, value, sensor_type)
history.append(value)
del history[:-MAX_HISTORY]
return deviation
def _record_boolean(user_id: uuid.UUID, device_id: uuid.UUID, sensor_type: str, value: bool) -> bool:
key = _key(user_id, device_id, sensor_type)
previous = _boolean_state.get(key)
anomalous = detect_boolean_transition(previous, value)
_boolean_state[key] = value
return anomalous
# ---------------------------------------------------------------------------
# Numeric vs. boolean classification
# ---------------------------------------------------------------------------
# Rule for telling a boolean/discrete sensor apart from a continuous
# numeric one: the contract's own example reading marks boolean sensors
# with `unit: "bool"` (`{"sensor_type": "presence", "value": 1, "unit":
# "bool", ...}`), so that's the primary, explicit signal — a device author
# opts a sensor into transition semantics by declaring its unit as such.
# `python`/`bool` `value` types are treated the same way as a convenience
# for callers that already have a real bool in hand (e.g. Workstream G's
# ingestion handler after its own shape validation). A bare numeric 0/1
# with a *non*-bool unit (e.g. a duty-cycle percentage that happens to read
# 0 or 1) is intentionally left as numeric — guessing boolean from value
# alone would silently reinterpret real numeric sensors whenever they
# happened to read exactly 0 or 1.
_BOOL_UNITS = {"bool", "boolean"}
def is_boolean_sensor(unit: str, value: float | bool) -> bool:
if isinstance(value, bool):
return True
return unit.strip().lower() in _BOOL_UNITS
# ---------------------------------------------------------------------------
# Per-sensor-type frequency/magnitude mapping
# ---------------------------------------------------------------------------
# A stable per-sensor-type "frequency" so the same sensor_type always
# fingerprints into the same signature bucket (see
# `app.entities.signature_from_anomalies`, which buckets anomalies by the
# digit-count of their frequency). Known sensor types get a hand-picked
# value; anything else (sensor_type is free-form per the contract) gets a
# deterministic hash-derived value in the same rough order of magnitude, so
# a brand-new sensor_type nobody wrote a case for still fingerprints
# consistently every time rather than randomly.
_KNOWN_FREQUENCIES: dict[str, float] = {
"presence": 66.6,
"temperature": 111.0,
"humidity": 222.0,
"pressure": 333.0,
}
def frequency_for_sensor_type(sensor_type: str) -> float:
if sensor_type in _KNOWN_FREQUENCIES:
return _KNOWN_FREQUENCIES[sensor_type]
digest = hashlib.sha1(sensor_type.encode()).hexdigest()[:6]
return 50.0 + (int(digest, 16) % 950)
# Fixed magnitude for a boolean-transition anomaly: there is no continuous
# deviation to measure (the signal is 0/1), so a representative constant
# stands in — chosen mid-range against the magnitudes the other four modes
# typically produce (see `test_ws_session.py`'s anomaly fixtures, roughly
# single-to-low-double digits) so a presence trip reads as a normal-sized
# anomaly rather than a suspiciously flat or huge one.
BOOLEAN_ANOMALY_MAGNITUDE = 8.0
# ---------------------------------------------------------------------------
# Public entry point — called from Workstream G's ingestion handoff point
# ---------------------------------------------------------------------------
async def process_device_reading_for_summon(
user_id: uuid.UUID,
device_id: uuid.UUID,
sensor_type: str,
value: float | bool,
unit: str,
) -> None:
"""Run one ingested device reading through the appropriate anomaly
detector and, if it's anomalous, push it into the owning user's live
séance (if any) exactly the way the four browser-based modes already
do — via `app.ws._handle_anomaly`, so hardware summons reuse the whole
existing signature/mint/Codex pipeline with no new mint logic.
This is meant to be called (awaited) from Workstream G's
`POST /api/device/telemetry` handler, once per reading in the batch,
after that endpoint's own shape/auth validation — it does no HTTP or
auth work of its own. If the user has no active séance session open
(the common case — most readings arrive with no one watching), this is
a normal no-op, not an error.
Imports `app.ws` lazily-at-module-level would create no cycle today
(ws.py does not import this module), but the import is written as a
plain top-of-function local import anyway to keep this module trivially
importable/testable in isolation from the full ws.py dependency graph
(DB engine, LLM service, TTS, etc.) for anyone who only wants the pure
detector functions above.
"""
from app.ws import _handle_anomaly, get_active_session # noqa: PLC0415
if is_boolean_sensor(unit, value):
anomalous = _record_boolean(user_id, device_id, sensor_type, bool(value))
magnitude = BOOLEAN_ANOMALY_MAGNITUDE
else:
deviation = _record_numeric(user_id, device_id, sensor_type, float(value))
anomalous = deviation is not None
magnitude = round(deviation, 2) if deviation is not None else None
if not anomalous:
return
state = get_active_session(user_id)
if state is None:
return
await _handle_anomaly(
state,
{
"source": sensor_type,
"frequency": frequency_for_sensor_type(sensor_type),
"magnitude": magnitude,
},
)

View File

@@ -91,6 +91,36 @@ class SeanceState:
last_wire_anomaly_at: float = 0.0 last_wire_anomaly_at: float = 0.0
# Active-session registry (spec: ESP32 sensor node, Workstream K): maps a
# user to their live SeanceState so hardware ingestion (a plain HTTP
# request, not a WS connection) can find "does this user have a séance
# open right now" and push a hardware anomaly into it. Single
# most-recent-session mapping — if a user somehow has two `/ws/session`
# tabs open concurrently, the newer one wins the registry slot; that's a
# reasonable v1 (see spec's Contract section) since a user's attention is
# realistically in one tab at a time. Plain in-process dict, same reasoning
# as everywhere else in this file: no external store needed at this scale,
# doesn't need to survive a restart.
_active_sessions: dict[uuid.UUID, "SeanceState"] = {}
def register_active_session(user_id: uuid.UUID, state: "SeanceState") -> None:
_active_sessions[user_id] = state
def unregister_active_session(user_id: uuid.UUID, state: "SeanceState") -> None:
# Only remove if it's still *this* state — guards against a rare
# overlap where an older session's disconnect cleanup runs after a
# newer session for the same user has already registered, which would
# otherwise wipe out the newer (still-live) registry entry.
if _active_sessions.get(user_id) is state:
del _active_sessions[user_id]
def get_active_session(user_id: uuid.UUID) -> "SeanceState | None":
return _active_sessions.get(user_id)
def serialize_entity(entity: Entity) -> dict: def serialize_entity(entity: Entity) -> dict:
return { return {
"id": str(entity.id), "id": str(entity.id),
@@ -475,6 +505,7 @@ async def session_socket(websocket: WebSocket) -> None:
session_id=contact_session.id, session_id=contact_session.id,
client_ip=_client_ip(websocket), client_ip=_client_ip(websocket),
) )
register_active_session(user_id, state)
sender = asyncio.create_task(_sender(state, websocket)) sender = asyncio.create_task(_sender(state, websocket))
await state.send_queue.put({"type": "session", "id": str(contact_session.id)}) await state.send_queue.put({"type": "session", "id": str(contact_session.id)})
@@ -511,6 +542,7 @@ async def session_socket(websocket: WebSocket) -> None:
except WebSocketDisconnect: except WebSocketDisconnect:
pass pass
finally: finally:
unregister_active_session(user_id, state)
if state.ambient_task is not None: if state.ambient_task is not None:
state.ambient_task.cancel() state.ambient_task.cancel()
sender.cancel() sender.cancel()

View File

@@ -0,0 +1,255 @@
import uuid
import pytest
import app.ws as ws_module
from app.device_anomaly import (
BOOLEAN_ANOMALY_MAGNITUDE,
detect_boolean_transition,
detect_sensor_spike,
frequency_for_sensor_type,
is_boolean_sensor,
process_device_reading_for_summon,
reset_state,
)
from app.models.contact_session import ContactSession
from app.models.user import User
from app.ws import SeanceState, get_active_session, register_active_session
from .conftest import TestSessionLocal
@pytest.fixture(autouse=True)
def _clear_device_anomaly_state():
reset_state()
yield
reset_state()
# ---------------------------------------------------------------------------
# detect_sensor_spike — same three-guard shape as detect_wire_spike
# ---------------------------------------------------------------------------
def test_numeric_spike_needs_history():
assert detect_sensor_spike([], 99.0, "temperature") is None
assert detect_sensor_spike([21.0, 21.1, 21.0], 30.0, "temperature") is None
def test_numeric_spike_fires_on_real_surge():
history = [21.0, 21.1, 20.9, 21.2, 21.0, 21.1, 20.8]
deviation = detect_sensor_spike(history, 30.0, "temperature")
assert deviation is not None
assert deviation > 8.0
def test_numeric_spike_ignores_normal_fluctuation():
history = [21.0, 21.1, 20.9, 21.2, 21.0, 21.1, 20.8]
assert detect_sensor_spike(history, 21.3, "temperature") is None
def test_numeric_spike_has_absolute_floor_for_quiet_sensors():
# A rock-steady baseline with a tiny wobble must never cry ghost, even
# if that wobble technically clears a 3-sigma bar against near-zero std.
history = [20.00, 20.01, 19.99, 20.00, 20.02, 19.98, 20.01]
assert detect_sensor_spike(history, 20.30, "temperature") is None
def test_numeric_spike_adapts_to_loud_baseline():
# Once the sensor is already swinging widely, the same absolute jump is
# no longer anomalous relative to its own noisy baseline.
history = [15.0, 22.0, 14.0, 23.0, 16.0, 21.0, 15.5]
assert detect_sensor_spike(history, 24.0, "temperature") is None
def test_numeric_spike_uses_known_floor_per_sensor_type():
# Humidity's floor (3.0) is looser than temperature's (0.8) — a
# deviation that would fire for temperature must not fire for humidity.
history = [45.0, 46.0, 44.5, 45.5, 45.0, 44.8, 45.2]
assert detect_sensor_spike(history, 47.6, "humidity") is None
assert detect_sensor_spike(history, 47.6, "temperature") is not None
def test_numeric_spike_unknown_sensor_type_uses_relative_fallback_floor():
# No hand-picked floor for "voltage" — falls back to a fraction of the
# rolling mean, but a genuine multi-fold surge still fires.
history = [5.0, 5.1, 4.9, 5.0, 5.05, 4.95, 5.02]
assert detect_sensor_spike(history, 5.1, "voltage") is None
assert detect_sensor_spike(history, 12.0, "voltage") is not None
# ---------------------------------------------------------------------------
# detect_boolean_transition
# ---------------------------------------------------------------------------
def test_boolean_transition_false_to_true_is_anomalous():
assert detect_boolean_transition(False, True) is True
def test_boolean_transition_no_prior_state_is_not_anomalous():
assert detect_boolean_transition(None, True) is False
assert detect_boolean_transition(None, False) is False
def test_boolean_transition_true_to_true_is_not_anomalous():
assert detect_boolean_transition(True, True) is False
def test_boolean_transition_true_to_false_is_not_anomalous():
assert detect_boolean_transition(True, False) is False
# ---------------------------------------------------------------------------
# is_boolean_sensor classification rule
# ---------------------------------------------------------------------------
def test_is_boolean_sensor_by_unit():
assert is_boolean_sensor("bool", 1) is True
assert is_boolean_sensor("BOOL", 0) is True
def test_is_boolean_sensor_by_python_bool_value():
assert is_boolean_sensor("pct", True) is True
def test_is_boolean_sensor_numeric_zero_one_with_non_bool_unit_stays_numeric():
# A duty-cycle percentage reading exactly 0 or 1 must not be
# misclassified as boolean just because its value looks bool-like.
assert is_boolean_sensor("pct", 1) is False
assert is_boolean_sensor("c", 0) is False
# ---------------------------------------------------------------------------
# frequency_for_sensor_type
# ---------------------------------------------------------------------------
def test_frequency_stable_for_known_sensor_types():
assert frequency_for_sensor_type("temperature") == frequency_for_sensor_type("temperature")
assert frequency_for_sensor_type("temperature") != frequency_for_sensor_type("humidity")
def test_frequency_stable_and_deterministic_for_unknown_sensor_type():
freq1 = frequency_for_sensor_type("cosmic_ray_flux")
freq2 = frequency_for_sensor_type("cosmic_ray_flux")
assert freq1 == freq2
assert freq1 != frequency_for_sensor_type("other_unknown_sensor")
# ---------------------------------------------------------------------------
# process_device_reading_for_summon — integration
# ---------------------------------------------------------------------------
async def _make_user_and_session(db_session) -> tuple[uuid.UUID, uuid.UUID]:
user = User(username=f"devowner-{uuid.uuid4().hex[:8]}", password_hash="x")
db_session.add(user)
await db_session.flush()
session = ContactSession(user_id=user.id)
db_session.add(session)
await db_session.commit()
await db_session.refresh(session)
return user.id, session.id
@pytest.mark.asyncio
async def test_process_reading_pushes_anomaly_into_active_session(db_session, monkeypatch):
monkeypatch.setattr(ws_module, "session_maker", TestSessionLocal)
user_id, session_id = await _make_user_and_session(db_session)
state = SeanceState(user_id=user_id, session_id=session_id, client_ip="127.0.0.1")
register_active_session(user_id, state)
device_id = uuid.uuid4()
try:
# Warm up the baseline with unremarkable readings — none of these
# should produce an anomaly.
for value in [21.0, 21.1, 20.9, 21.2, 21.0, 21.1]:
await process_device_reading_for_summon(
user_id, device_id, "temperature", value, "c"
)
assert state.anomalies == []
# A genuine surge fires and gets pushed into state.anomalies in the
# same shape _handle_anomaly uses for the browser-based modes.
await process_device_reading_for_summon(
user_id, device_id, "temperature", 30.0, "c"
)
assert len(state.anomalies) == 1
anomaly = state.anomalies[0]
assert anomaly["source"] == "temperature"
assert anomaly["frequency"] == frequency_for_sensor_type("temperature")
assert anomaly["magnitude"] > 0
finally:
from app.ws import unregister_active_session
unregister_active_session(user_id, state)
@pytest.mark.asyncio
async def test_process_reading_boolean_presence_transition(db_session, monkeypatch):
monkeypatch.setattr(ws_module, "session_maker", TestSessionLocal)
user_id, session_id = await _make_user_and_session(db_session)
state = SeanceState(user_id=user_id, session_id=session_id, client_ip="127.0.0.1")
register_active_session(user_id, state)
device_id = uuid.uuid4()
try:
# First reading (False) establishes state, no prior value to
# transition from either way.
await process_device_reading_for_summon(user_id, device_id, "presence", 0, "bool")
assert state.anomalies == []
# false -> true is the anomaly.
await process_device_reading_for_summon(user_id, device_id, "presence", 1, "bool")
assert len(state.anomalies) == 1
anomaly = state.anomalies[0]
assert anomaly["source"] == "presence"
assert anomaly["magnitude"] == BOOLEAN_ANOMALY_MAGNITUDE
# true -> true is not.
await process_device_reading_for_summon(user_id, device_id, "presence", 1, "bool")
assert len(state.anomalies) == 1
finally:
from app.ws import unregister_active_session
unregister_active_session(user_id, state)
@pytest.mark.asyncio
async def test_process_reading_no_active_session_is_a_noop(db_session, monkeypatch):
monkeypatch.setattr(ws_module, "session_maker", TestSessionLocal)
user_id, _session_id = await _make_user_and_session(db_session)
device_id = uuid.uuid4()
assert get_active_session(user_id) is None
# Should not raise even though there's no active session and no
# baseline yet — feeding a single reading can't be anomalous anyway.
await process_device_reading_for_summon(user_id, device_id, "temperature", 21.0, "c")
assert get_active_session(user_id) is None
@pytest.mark.asyncio
async def test_process_reading_non_anomalous_value_does_not_touch_session(db_session, monkeypatch):
monkeypatch.setattr(ws_module, "session_maker", TestSessionLocal)
user_id, session_id = await _make_user_and_session(db_session)
state = SeanceState(user_id=user_id, session_id=session_id, client_ip="127.0.0.1")
register_active_session(user_id, state)
device_id = uuid.uuid4()
try:
for value in [21.0, 21.1, 20.9, 21.2, 21.0, 21.1, 21.0]:
await process_device_reading_for_summon(
user_id, device_id, "temperature", value, "c"
)
assert state.anomalies == []
finally:
from app.ws import unregister_active_session
unregister_active_session(user_id, state)