"""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, }, )