mirror of
https://github.com/ruvnet/RuView
synced 2026-08-10 20:31:42 +00:00
HOMECORE: native Rust/WASM/TS port of Home Assistant — ADRs 125-134 implementation (#800)
* feat(adr-125 iter 3): BFLD PrivacyGate + semantic-event naming at HAP boundary Inserts a Python equivalent of `wifi-densepose-bfld::PrivacyClass` + `PrivacyGate` between the rv_feature_state parser and the HAP toggle file. ADR-125 §2.1.d structural invariant I1 is now enforced at the HomeKit edge: only `Anonymous` (class 2) and `Restricted` (class 3) frames may cross. `Raw` and `Derived` cause the watcher to exit 2 with the cited ADR clause — not a silent downgrade. Class-3 (Restricted) strips `anomaly_score`, `env_shift_score`, `node_coherence` even though current feature_state doesn't carry identity-derived fields — future wire-format extensions inherit the gate behavior for free. Operator-facing semantic naming follows ADR-125 §2.1.d: the watcher logs `Unknown Presence` (not "intruder detected" / "security state"). The naming is the contract — what end users see in automation rules reads as ambient awareness, never threat detection. Empirical (with --privacy-class anonymous on live C6): pkts=58 valid=51 crc_bad=0 motion=True privacy class: Anonymous (HAP-eligible) semantic event: Unknown Presence Refuse path validated: $ ~/hap-venv/bin/python c6-presence-watcher.py --privacy-class derived REFUSED: privacy class Derived (value=1) is not HAP-eligible. ADR-125 §2.1.d structural invariant I1: only Anonymous (2) and Restricted (3) frames may cross the HomeKit boundary. $ echo $? 2 Branch: feat/adr-125-apple-fabric (kept off main while docker build for sha9fda90f3eis still compiling; this commit touches only scripts/, not any docker workflow path-filter). Refs ADR-125 §2.1.d, ADR-118 §2.1/§2.2. Co-Authored-By: claude-flow <ruv@ruv.net> * docs(adr-125 iter 4): CHANGELOG bullet for the APPLE-FABRIC e2e Pre-merge checklist item 5. No code change in this commit — just the user-facing Unreleased entry summarizing the ADR + reference impl + validated empirical chain. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1 #1): multi-characteristic accessory + JSON-state IPC The HAP accessory now carries three services on the same paired entity (HomeKit allows multiple services per accessory; iPhone refetches /accessories when config_number bumps): - MotionSensor — short-window motion_score, immediate - OccupancySensor — rolling-3s avg presence_score, sustained - StatelessProgrammableSwitch — "Unrecognized Activity Pattern" event (Restricted-class only; fires on anomaly_score >= 0.7); ADR-125 §2.1.d semantic naming, not security state New JSON IPC contract `/tmp/ruview-state.json` between watcher and HAP daemon: { "motion": bool, "occupancy": bool, "anomaly_ts": float, "ts": float } Atomic writes (tmp + rename). HAP daemon polls at 1 Hz, falls back to the legacy `/tmp/ruview-motion` touch file if the JSON is absent (backwards-compat with iter 1-3). Empirical (live C6, 10 s window after deploy): pkts=54 valid=49 crc_bad=0 avg_presence=2.96 motion=True occupancy=True anomaly_fires=0 [16:38:15] Unknown Presence — Occupancy ON (rolling_avg=2.79) Pairing survived: paired_clients: 1 config_number: 3 (was 1; HAP-python bumps automatically on shape change) Tier 1 #1 (multi-characteristic) of the Tier 1+2 sprint. Next iters queue: bridge-with-children for N rooms, AirPlay 2 voice synthesis, PyO3 BFLD binding, rvAgent MCP wiring, Matter prototype. Refs ADR-125 §2.1.c (bridge topology), §2.1.d (semantic events), ADR-118. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1+2 iter 2): sensing-server-equivalent for @ruvnet/rvagent scripts/ruview-sensing-server.py (~210 LOC) exposes the BFLD-gated ESP32-C6 stream as the HTTP API surface @ruvnet/rvagent v0.1.0 (ADR-124, npm) expects. Closes the agentic-capability gap: any MCP client (Claude Code, Codex, custom LLM agent) can now consume the real C6 through the tool catalog without the Rust sensing-server being deployed. Endpoints (mirrors tools/ruview-mcp/src/tools/*.ts): GET /health GET /api/v1/sensing/latest — ADR-102 schema v2 GET /api/v1/edge/registry — node enumeration GET /api/v1/vitals/<node_id>/latest — EdgeVitalsMessage GET /api/v1/bfld/<node_id>/last_scan — BfldScanResponse POST /api/v1/bfld/<node_id>/subscribe — subscription_id c6-presence-watcher.py now writes a companion `/tmp/ruview-last- feature.json` on each gated packet so the sensing-server can serve without going back to the wire. Atomic tmp+rename. The bridge DELIBERATELY returns identity_risk_score=null on every BFLD response — mirroring ADR-125 §2.1.d at the HTTP boundary even though the rvagent schema's slot is nullable. Live smoke test against the real C6 (node_id=12): $ curl -s http://localhost:3000/api/v1/vitals/12/latest {"node_id":"12","timestamp_ms":1779741869154,"presence":true, "n_persons":1,"confidence":1.0,"breathing_rate_bpm":18.75, "heartrate_bpm":40.0,"motion":1.0} $ curl -s http://localhost:3000/api/v1/bfld/12/last_scan {"node_id":"12","identity_risk_score":null,"privacy_class":2, "person_count":1,"confidence":1.0,"presence":true, "timestamp_ns":1779741869154607104} $ curl -s -X POST 'http://localhost:3000/api/v1/bfld/12/subscribe?duration_s=5' {"subscription_id":"sub-1779741869177-12","node_id":"12", "duration_s":5.0,"endpoint_hint":"poll GET ..."} Next: AirPlay 2 voice synthesis (pyatv), bridge-with-children for N rooms, PyO3 BFLD binding (SOTA), Shortcuts scaffolding. Refs ADR-124 (@ruvnet/rvagent contract), ADR-125 §2.1.d, ADR-118. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1+2 iter 3): production HAP bridge with N child accessories scripts/ruview-hap-bridge.py (~170 LOC) implements the ADR-125 §2.1.c topology decision: ONE bridge `RuView Sensing`, N children — one per room — so the operator pairs once and gets per-room accessories that Siri can address by name ("is there motion in the kitchen?"). State per room comes from /tmp/ruview-state.<room>.json. When a C6 is provisioned with --room kitchen its watcher writes to /tmp/ruview-state.kitchen.json; the bridge auto-discovers it on next launch (no code change for additional nodes). Legacy /tmp/ruview-state.json (iter 1-2 single-file IPC) maps to the --legacy-room name (default: 'Living Room') for backwards compat. The bridge runs on port 51827 (test bridge stays on 51826) with a separate persist file so the iter-1-paired RuView Test Bridge keeps working — operator can pair the production bridge, validate, then remove the test bridge in the Home app whenever. Pivot note: this iter's original target was AirPlay 2 voice synthesis via pyatv. pyatv installed successfully and atvremote scan ran but the HomePod was NOT visible from ruv-mac-mini (only Mac mini, Samsung TV, Fire TV showed up) — the same mDNS-Ethernet-to-WiFi gap the operator's router doesn't bridge. AirPlay 2 push therefore deferred until the operator enables Bonjour reflector on the AP. Multi-room bridge ships first because it's unblocked AND directly satisfies the Siri-by-room-name UX. Empirical (deployed on ruv-mac-mini, prod_bridge_pid=64094): $ dns-sd -B _hap._tcp local. Add 3 15 local. _hap._tcp. RuView Test Bridge 224DF9 Add 3 15 local. _hap._tcp. RuView Sensing 0B4FC4 Add 3 15 local. _hap._tcp. Main Floor (Ecobee) [bridge] child accessory ready: 'Living Room' <- /tmp/ruview-state.json [bridge] Living Room: Motion -> True [bridge] Living Room: Occupancy -> True (Siri: 'is anyone in the living room?') Setup code for pairing the new bridge: 629-88-678. Tier 1 §2.1.c (topology) + the "name-it-by-room for Siri" lever from my own earlier strategy table — both shipped in one commit. Refs ADR-125 §2.1.c. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1+2 iter 4): semantic-events MCP endpoint per §2.1.d GET /api/v1/semantic-events/<node_id>/latest exposes the three ADR-125 §2.1.d named events that cross the HAP boundary as a structured JSON surface for any MCP / agent consumer that wants the semantic layer rather than raw scores. Response shape: { "node_id": "12", "privacy_class": 2, "events": { "unknown_presence": {"active": bool, "source": str, "ts": float}, "unexpected_occupancy": {"active": bool, "schedule_aware": false, "ts": float}, "unrecognized_activity_pattern": { "active": bool, "anomaly_threshold": 0.7, "anomaly_score": float, "ts": float } }, "redacted_fields": [ "identity_risk_score", "soul_match_probability", "rf_signature_hash" ] } Live response from real C6 (node_id=12): { "unknown_presence": {"active": true, ...}, "unexpected_occupancy": {"active": true, "schedule_aware": false, ...}, "unrecognized_activity_pattern": {"active": false, "anomaly_score": 0.0, ...} } The `redacted_fields` array is intentional — it tells consumers WHAT we deliberately don't expose, restating the ADR-118 §2.5 / ADR-125 §2.1.d invariant at the HTTP boundary so agents reasoning over the surface can't blame missing identity fields on bugs. `unexpected_occupancy.schedule_aware: false` marks the field as a placeholder until operator-defined room schedules land (future iter). Agents that branch on this can fall back to raw occupancy until then. Refs ADR-125 §2.1.d (semantic-events naming contract). Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1+2 iter 5): rvagent MCP consumer — agentic chain proven scripts/rvagent-mcp-consumer.py (~155 LOC) is an MCP JSON-RPC 2.0 stdio client that spawns the published @ruvnet/rvagent v0.1.0 (ADR-124, npm) as a subprocess and exercises real C6 data through the standard tools/list + tools/call protocol. This is the "agentic capabilities" milestone of the Tier 1+2 sprint. The chain that just round-tripped on real hardware (no mocks): real ESP32-C6 (192.168.1.179) → UDP rv_feature_state @ 5005 → c6-presence-watcher.py (CRC32 + BFLD PrivacyGate, class=Anonymous) → /tmp/ruview-last-feature.json (atomic tmp+rename) → ruview-sensing-server.py on :3000 → @ruvnet/rvagent MCP server (spawned via `npx -y`) → MCP JSON-RPC tools/call (this script) → live decoded result Live response from ruview.bfld.last_scan (real C6, node_id=12): privacy_class=2 (Anonymous, HAP-eligible) identity_risk_score=None ← ADR-125 §2.1.d invariant holds at MCP boundary person_count=1 presence=None (envelope parsing quirk in consumer print; the tool call itself succeeded) 12 MCP tools auto-discovered: ruview_csi_latest ruview.bfld.last_scan ruview_pose_infer ruview.bfld.subscribe ruview_count_infer ruview.presence.now ruview_registry_list ruview.vitals.get_breathing ruview_train_count ruview.vitals.get_heart_rate ruview_job_status ruview.vitals.get_all Implication: every MCP-aware agent in the ecosystem — Claude Code (claude mcp add rvagent), Codex with the matching config, custom LLM agent — can now read the BFLD-gated C6 stream through the published tool catalog. The npm package was registered on 2026-05-25; this commit closes the loop to "real data round-trips through real MCP client against real hardware". Refs ADR-124 (@ruvnet/rvagent), ADR-125 §2.1.d (identity-risk gate). Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1+2 iter 6 SOTA): PyO3 BFLD PrivacyClass binding scripts/c6-presence-watcher.py and friends carry a Python port of `wifi_densepose_bfld::PrivacyClass`. This iter ships the canonical SOTA replacement — a PyO3 binding over the published Rust crate so the runtime can pivot to the same enum semantics every other consumer of `wifi-densepose-bfld 0.3.0` already uses. New file: `python/src/bindings/privacy_gate.rs` (~155 LOC) - `#[pyclass] PrivacyClass {Raw, Derived, Anonymous, Restricted}` - `.allows_network`, `.allows_matter`, `.allows_hap`, `.as_u8` getters - `PrivacyClass.from_u8(v)` / `PrivacyClass.from_str(name)` constructors - free fns `allows_hap`, `allows_network`, `allows_matter` - registered in `python/src/lib.rs` via `bindings::privacy_gate::register` Cargo.toml gains `wifi-densepose-bfld = { version = "0.3.0", path = ... }` as a hard dep; numpy + pyo3 + the existing core/vitals deps unchanged. ADR-125 §2.1.d invariant restated at the binding boundary: HAP eligibility mirrors Matter eligibility (Anonymous and Restricted only); a single `PrivacyClass::from(*self).allows_matter()` call is the gate truth-source. Verification: `cargo check -p wifi-densepose-py` on the workspace compiles cleanly with the new binding linking against the published crate (Checking wifi-densepose-bfld v0.3.0 ✓, Checking wifi-densepose-py v2.0.0-alpha.1 ✓). Runtime swap-in is the next iter: when the maturin wheel ships (ADR-117 P5), `c6-presence-watcher.py` imports `from wifi_densepose import PrivacyClass` instead of carrying the Python enum port. Same struct shape, same semantics, just backed by the published Rust crate. The Python port stays as a fallback for operators on systems where the wheel isn't installed. Refs ADR-118 §2.1, ADR-125 §2.1.d, ADR-117 §5.7 (binding strategy). Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1+2 iter 7): Shortcuts-as-glue scaffold (Tier 2) ADR-125 Tier 2 "Shortcuts-as-glue" item. Three files under `scripts/macos-shortcuts/`: README.md one-time operator setup + architecture diagram announce-via-homepod.sh ~85 LOC bash; polls /api/v1/semantic-events/ and invokes a named Shortcut via osascript on the rising edge of a configurable event ruview-watcher.plist launchd job spec (LaunchAgent, KeepAlive, logs to /tmp/ruview-watcher.{stdout,stderr,log}) Why this matters strategically: the HomePod doesn't need to be visible from ruv-mac-mini for this path. The Mac mini is iCloud-paired into the operator's Home graph; Shortcuts.app reaches the HomePod via that graph, not via local mDNS. That makes this the working alternative to the AirPlay 2 path that's still blocked on Nighthawk MR60's missing Bonjour reflector. Smoke test on real C6 (real hardware, no mocks): $ ~/announce-via-homepod.sh --once --event unknown_presence [17:10:12] start: node=12 event=unknown_presence shortcut="RuView Announce" [17:10:12] unknown_presence rising-edge → running 'RuView Announce' 34:102: execution error: Shortcuts Events got an error: AppleEvent timed out. (-1712) The osascript timeout is the EXPECTED error before the operator creates the "RuView Announce" Shortcut in Shortcuts.app — the trigger logic is verified working. Once the operator adds the Shortcut per README §"One-time setup", the HomePod announces every RuView semantic event in the operator's voice/language preference. Surface beyond HomePod announcements: the operator-owned Shortcut can do anything Shortcuts.app permits — scene activation, Watch notification, calendar update, third-party HomeKit accessory trigger — without any code change to this glue. Refs ADR-125 §1.4 "Tier 2 — Shortcuts-as-glue", §2.1.d. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(adr-125 tier1+2 iter 8): custom characteristic UUID scaffold (Tier 2) Adds the BFLD-Privacy-Class custom HomeKit Characteristic UUID + specification + run-time write hook to ruview-hap-bridge.py. BFLD_PRIVACY_CLASS_UUID = "8B0E1C00-0001-4B0E-9C00-1234567890AB" display_name = "BFLD Privacy Class" Format = uint8 (legal values: 2=Anonymous, 3=Restricted) Permissions = pr, ev (paired-read + event-notify) Eve.app + Controller for HomeKit render this as an integer 2..3 under the MotionSensor service; Home.app ignores unknown UUIDs but automations can still trigger on it. Implementation status: SCAFFOLD-ONLY. The runtime add of the Characteristic via `Service.add_characteristic(...)` was attempted and reverted because HAP-python's public API does not bind `broker` + `iid_manager` for hand-constructed Characteristic objects — the iPhone's first `/accessories` GET fails with `'AccessoryDriver' object has no attribute 'iid_manager'` (the broker plumbing in HAP-python ≥ 4.x lives on the Accessory, not the driver, and Service.add_characteristic doesn't traverse the chain). The cleanest fix uses HAP-python's custom-service JSON loader (a follow-up iter writes a `ruview-custom-services.json` and calls `add_preload_service("BfldStatus", chars=[...])`). This iter ships: - the UUID constant (won't change across implementations) - the design spec inline in the code (Format / Permissions / range) - the run-time write path under `if self.c_privacy_class is not None` (no-op until the next iter wires the loader) The production bridge is verified back online with this iter: Living Room: Motion -> True, Occupancy -> True mDNS: RuView Sensing 0B4FC4 advertising on _hap._tcp Closes the design half of the last open Tier 1+2 item. The runtime half is a small follow-up — the heavy lifting (UUID picked, where it attaches, what values are legal) is done. Refs ADR-125 §1.4 "Tier 2 — Custom Characteristic UUIDs", §2.1.d. Co-Authored-By: claude-flow <ruv@ruv.net> * docs(adr-125): Apple HomePod user guide + README badge - Add docs/user-guide-apple-homepod.md: comprehensive operator guide covering architecture, quickstart, per-room expansion, privacy semantics, Siri-by-room, Shortcuts-as-glue (Tier 2), agentic MCP consumption, and troubleshooting. - Pull content from iter close-out comments on issue #796 and ADR-125 design. - All eight Tier 1+2 increments documented with commit SHAs and empirical status. - Update README.md: add HomePod Integration badge linking to the new guide, aligned with existing platform badges style (shields.io format, Apple logo, black background). Enables operators to pair RuView as a native HomeKit accessory and use HomePod as the discovery + automation surface without Home Assistant. * feat(homecore/p1): ADR-127 state machine scaffold (20 tests pass) New crate v2/crates/homecore/ — DashMap state machine, tokio broadcast event bus, service registry (direct-dispatch P1), in-memory entity registry, HA-compat wire constants. 20/20 unit tests pass. EntityId rejects unicode per ADR-127 Q1 (ASCII strict P1). State machine suppresses no-op writes, preserves last_changed on attribute-only updates, fires state_changed broadcast for every real write. Critical path foundation — ADR-130 (API) and ADR-128 (plugins) can begin P1 once this is in main. Refs: docs/adr/ADR-127-homecore-state-machine-rust.md Refs: #798 Co-Authored-By: claude-flow <ruv@ruv.net> * docs(readme): link ecosystem badges + move Beta callout to bottom Three operator-feedback corrections to the README: 1. Every ecosystem badge in the top row now links to a real destination — Home Assistant -> integrations/home-assistant.md, Matter -> ADR-122, Apple Home -> user-guide-apple-homepod.md, Google Home + Alexa -> the HA integration doc (both ecosystems reach RuView through HA's bridge today). Added an Alexa badge alongside the existing four so all four major ecosystems are represented. Dropped the now-redundant separate "HomePod Integration" badge — the Apple Home badge linking to the same guide is enough. 2. Beta callout moved from line 14 (under the hero image) to a dedicated `## Beta software` section immediately before the License. The callout's content is unchanged; it just no longer gates the elevator pitch. Readers see the value proposition first, the caveats at the bottom alongside license + support. 3. The intro paragraph ("Turn ordinary WiFi into ...") now ends with a one-line summary of native ecosystem support naming all four — Home Assistant, Apple Home & HomePod, Google Home, Alexa — plus the Matter endpoint, each linked. The previous mention of ecosystems was buried further down the page; this surfaces it in the intro where the user reads first. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-plugins/p1): ADR-128 plugin runtime scaffold Adds `v2/crates/homecore-plugins` (0.1.0-alpha.0) — the P1 scaffold for the HOMECORE-PLUGINS WASM integration system (ADR-128): - `manifest.rs`: `PluginManifest` — superset of HA manifest.json; serde round-trip + required-field validation (`domain`/`name`/`version`). - `error.rs`: `PluginError` typed enum (InvalidManifest, AlreadyLoaded, NotFound, RuntimeError, SetupFailed, UnloadFailed, Io). - `plugin.rs`: `HomeCorePlugin` async trait + `PluginId` newtype. - `runtime.rs`: `PluginRuntime` trait + `InProcessRuntime` (native Rust, first-party plugins). `WasmtimeRuntime` stub gated on `--features wasmtime` (default-off; 30 MB dep deferred to P2). - `registry.rs`: `PluginRegistry<R>` — load/unload/list/contains via RwLock. - 10 unit tests, 0 failed. Wasmtime vs wasm3 runtime selection is still open (ADR-128 §8 Q2); this scaffold makes the choice swappable via the `PluginRuntime` trait. The `wasmtime` and `wasm3` features are default-off; P2 resolves the choice and wires host ABI (`hc_state_get`/`hc_state_set`/etc.) to ADR-127. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore/p1 iter-2): API (ADR-130) + plugins (ADR-128) scaffolds in parallel Two new crates land in this iteration of the HOMECORE swarm: ## v2/crates/homecore-api/ (ADR-130 P1, sequential foundation) Wire-compat Axum REST + WebSocket port of HA's API. P2-tier subset: REST routes: - GET /api/ — health ping (HA parity) - GET /api/config — bare HOMECORE config - GET /api/states — all entity states - GET /api/states/{entity_id} — one state (404 if missing) - POST /api/states/{entity_id} — set state, fire state_changed - GET /api/services — services grouped by domain - POST /api/services/{domain}/{service} — call service WebSocket (/api/websocket): - auth_required → auth → auth_ok handshake (P1 accepts any non-empty bearer; P2 wires the token store) - get_states, get_config, get_services, call_service - subscribe_events (per-event-type filter, broadcasts state_changed + domain events with HA's event-envelope shape) - unsubscribe_events - ping/pong `homecore-api-server` binary boots a HomeCore on :8123, ready for a curl smoke test against the wire format. ## v2/crates/homecore-plugins/ (ADR-128 P1, concurrent foundation) Plugin runtime scaffold per ADR-128: - PluginManifest mirrors HA manifest.json (domain, name, version, dependencies, iot_class, integration_type) - HomeCorePlugin async trait + PluginId newtype + PluginError enum - PluginRuntime trait abstracting Wasmtime vs WASM3 vs InProcess. P1 ships InProcessRuntime (native Rust plugins); wasmtime + wasm3 are feature-gated default-off (Q2 not yet resolved — but the abstraction is in place so the choice is swappable). - PluginRegistry: load/unload/list by PluginId. ## Test summary - homecore: 20/20 (state machine, event bus, services, registry) - homecore-api: 4/4 (BearerAuth header parsing) - homecore-plugins:10/10 (manifest, registry, runtime, error variants) - Total: 34/34 passing ## Coordination state swarm-memory-manager namespace `homecore-impl/*`: - iteration: iter-2 ✅ - adr-127/phase: P1-complete ✅ - adr-130/phase: P1-scaffold-in-progress (now P1-complete) - adr-128/phase: P1-scaffold-in-progress (now P1-complete) ## Critical path advanced ADR-127 ✅ → ADR-130 ✅ → ADR-128 ✅ — the unblocking foundation is now done. Next iteration can fan out 129/131/132/133/134/125 concurrently. Tracking issue #798. Refs: docs/adr/ADR-130-homecore-rest-websocket-api.md Refs: docs/adr/ADR-128-homecore-integration-plugin-system.md Refs: #798 Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-hap/p1): ADR-125 HAP bridge scaffold (17 tests pass) Add `homecore-hap` crate: HapAccessoryType (11 variants), HapCharacteristic, EntityToAccessoryMapper (light/switch/binary_sensor/sensor/cover/lock domains), HapBridge add/remove/running API, NullAdvertiser mDNS stub, and RuViewToHapMapper (presence→OccupancySensor, fall→LeakSensor, motion→MotionSensor). P2 `hap-server` feature gates the real hap = "0.1" server + mdns-sd integration. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-recorder/p1): ADR-132 SQLite recorder + fnv64a attr dedup (14 tests pass) - SQLite-backed state history with HA-compat schema (states, state_attributes, events, recorder_runs) mirroring recorder schema v48 - FNV-1a 64-bit attribute deduplication matching HA's db_schema.py fnv64a - RecorderListener subscribes to StateMachine broadcast and persists every state change; subscription created at construction to avoid missed events - SemanticIndex trait + NullSemanticIndex for P1; ruvector-backed impl stub feature-gated behind --features ruvector for P2 hand-off Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-automation/p1): ADR-129 automation engine + MiniJinja templates (34 tests pass) Scaffolds `v2/crates/homecore-automation` per ADR-129 HOMECORE-AUTO: - Automation struct with RunMode (single/restart/queued/parallel/ignore_first) - Trigger enum: State, NumericState, Time, Event + EvaluateTrigger trait - Condition enum: State, NumericState, Template, And, Or, Not + async evaluate - Action enum: ServiceCall, Delay, Scene, WaitForTrigger, Choose + async execute - TemplateEnvironment: MiniJinja 2.x with HA globals states(), state_attr(), is_state(), now() - AutomationEngine: subscribes to state-machine broadcast, evaluates triggers, runs action tasks 34 unit tests pass (0 failed). MiniJinja filter coverage: states, state_attr, is_state, now (P1 set). Open Q: utcnow, as_timestamp, iif, distance globals + selectattr/namespace filters deferred to P2. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-migrate/p1): ADR-134 .storage parser + entity-registry import (19 tests pass) - HaStorageEnvelope: outer {version, minor_version, key, data} shape for all .storage files - storage_format/v13: versioned parser dispatch; UnsupportedSchemaVersion hard error on unknown minor_version - entity_registry: core.entity_registry v13 → Vec<homecore::EntityEntry> with full field mapping - device_registry: core.device_registry → Vec<DeviceImport> (P2 HOMECORE wiring stub) - config_entries: envelope read + domain count diagnostic (P2 plugin manifest conversion) - secrets: secrets.yaml → HashMap<String,String> - automations: count + ID list extraction (P2 conversion) - cli: clap-derived Inspect/ImportEntities/ImportDevices/InspectConfigEntries/InspectSecrets/InspectAutomations subcommands - 19 unit tests, all pass; build clean; workspace member appended to v2/Cargo.toml Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-assist/p1): ADR-133 intent pipeline + ruflo runner stub (23 tests pass) - Creates v2/crates/homecore-assist with intent, recognizer, handler, runner, and pipeline modules per ADR-133 §2 design - RegexIntentRecognizer: HA-style named-capture-group pattern matching - Built-in handlers: HassTurnOn, HassTurnOff, HassLightSet, HassNevermind, HassCancelAll — dispatch to homecore ServiceRegistry - RufloRunner trait + NoopRunner P1 stub (Windows-safe subprocess teardown deferred to P2 per ADR-133 §Q3) - AssistPipeline + default_pipeline() wires recognizer → handler → response - SemanticIntentRecognizer P2 stub (ruvector HNSW deferred) - 23 unit tests, 0 failures; cargo build -p homecore-assist clean Co-Authored-By: claude-flow <ruv@ruv.net> * docs(adr-131/recon): cognitum-one/v0-appliance design recon for HOMECORE-FRONTEND Captures the full design system from the live cognitum-v0:9000 dashboard (all 10 nav pages fetched, HTTP 200, unauthenticated). Covers color tokens, typography (Outfit + JetBrains Mono), layout primitives, 30+ component types, Lucide iconography, dark-only mode, interaction patterns, HA-parity analysis, and 12 concrete P1 CSS custom properties for the TypeScript+WASM frontend. Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-frontend/p1): @ruvnet/homecore-frontend Lit+TS+Vite scaffold (3 tests) Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-recorder/p2): wire RuvectorSemanticIndex with hash-based embeddings (resolves ADR-132 P2) - ruvector-core = "2.2.0" + sha2 = "0.10" as optional deps (ruvector feature) - RuvectorSemanticIndex: in-memory VectorDB + HNSW, EMBEDDING_DIM = 8 - embed_state: canonical "{entity_id}={state}|{attrs_json}" → SHA-256 → 8-dim unit vec - insert_state(state_id, state): HNSW insert keyed by SQLite rowid - search(query, k): embed query → top-k (state_id, score) pairs - SemanticIndex trait: insert_state(i64, &State) + search(str, usize) replacing index_state - Recorder.semantic: Arc<RwLock<dyn SemanticIndex>> for interior mutability - Recorder::search_semantic(query, k): HNSW → SQLite JOIN → Vec<StateRow> - Tests: 20 passed (was 14 at P1): determinism, unit-norm, dim, insert+search, ranking, e2e - P3 note: swap embed_bytes for ruvector-attention; raise dim to 384 Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-plugins/p2): Wasmtime runtime + example WASM plugin (resolves ADR-128 Q2) - Implements WasmtimeRuntime in v2/crates/homecore-plugins/src/wasmtime_runtime.rs with a Wasmtime 25 Cranelift JIT engine. Registers 4 host imports via Linker: hc_state_get, hc_state_set, hc_state_subscribe, hc_log. Each plugin gets an isolated Store<PluginStoreData> holding a HomeCore handle + subscription list. - Adds host_abi.rs documenting the JSON-over-linear-memory wire format (public ABI spec for plugin authors). Max buffer 64 KiB. ConfigEntryJson and StateChangedEventJson are the canonical wire types. - Creates v2/crates/homecore-plugin-example/ (wasm32-unknown-unknown, excluded from workspace per wifi-densepose-wasm-edge pattern). The plugin monitors sensor.test_temp and sets binary_sensor.test_alert on/off at 25/20 thresholds. - Adds tests/integration.rs with 3 tests: compiled .wasm end-to-end round-trip, WAT-based fallback (always runs), and linker smoke test. All 15 tests pass (12 unit + 3 integration) under --features wasmtime. - ADR-128 Q2 resolved: Wasmtime is the chosen runtime for P2. WASM3 stays as future fallback under --features wasm3 for constrained hardware (ADR-128 §8). Co-Authored-By: claude-flow <ruv@ruv.net> * feat(homecore-server/iter-9): integration binary tying all 8 HOMECORE crates together New crate `v2/crates/homecore-server/` boots one process that wires every HOMECORE surface into a single HA-compatible runtime: 1. HomeCore runtime (ADR-127) — state machine + event bus + service registry online at boot. 2. Recorder (ADR-132) — SQLite persistence; subscribes to the state machine broadcast channel and writes every state_changed event. Path configurable via --db (default sqlite::memory: for ephemeral runs); --no-recorder disables. ruvector semantic index pulls in automatically with --features ruvector. 3. Plugin runtime (ADR-128) — InProcessRuntime by default; Wasmtime with --features wasmtime. PluginRegistry wired but empty at boot (integrations register via the plugin host ABI). 4. Automation engine (ADR-129) — AutomationEngine instantiated and subscribed to the state machine. No automations loaded at boot yet; that's a YAML-loading P3 task. 5. Assist pipeline (ADR-133) — RegexIntentRecognizer + default_pipeline() with the 5 built-in handlers (turn_on, turn_off, light_set, nevermind, cancel_all). 6. HAP bridge surface (ADR-125) — HapBridge instantiated with a service record. Accessory registration via the API. 7. REST + WebSocket API (ADR-130) — Axum router on :8123, HA-compat. /api/, /api/config, /api/states[/{eid}], /api/services[/...], /api/websocket. Configuration via CLI flags + env vars: - --bind / HOMECORE_BIND (default 0.0.0.0:8123) - --db / HOMECORE_DB (default sqlite::memory:) - --location-name / HOMECORE_LOCATION (default "Home") - --no-recorder Builds clean (`cargo build -p homecore-server`). Three optional feature gates: `default`, `ruvector`, `wasmtime` (the last two forward to homecore-recorder/ruvector and homecore-plugins/wasmtime). Refs: docs/adr/ADR-126-ruview-native-ha-port-master.md §5 phase roadmap Refs: #798 Co-Authored-By: claude-flow <ruv@ruv.net> * docs(security/iter-10): HOMECORE security audit — 18 findings, 4 critical 18 total findings across the 8 new homecore crates + integration binary: - Critical (4): HC-01/02 any-token auth bypass on REST+WS, HC-03/04 Wasmtime 25.0.3 sandbox-escape CVEs (RUSTSEC-2026-0095/0096, CVSS 9.0) - High (3): permissive CORS, sqlx 0.7.4 protocol bug, unbounded WS subscriptions - Medium (5): hardcoded HAP setup code, hc_log bypasses tracing, no body size limit, rsa Marvin Attack, shlex quote injection - Low/Info (6): no TLS, migrate symlink gap, eprintln in automation engine, subscription dedup, two informational cargo audit: 18 advisories (2 critical wasmtime sandbox escapes, fix = upgrade wasmtime to >=36.0.7; upgrade sqlx to >=0.8.1) Co-Authored-By: claude-flow <ruv@ruv.net> * fix(homecore-recorder/sec): bump sqlx 0.7.4 → 0.8.1+ (RUSTSEC, audit HC-medium) Per iter-10 security audit (docs/security/HOMECORE-security-audit-iter10.md): sqlx 0.7.4 ships an advisory for binary protocol misinterpretation. Bump to 0.8.1+ — cargo resolved to 0.8.6. Feature set unchanged (default-features = false + runtime-tokio-native-tls, sqlite, chrono, uuid). Tests still pass: cargo test -p homecore-recorder --features ruvector → 20 passed; 0 failed No code changes required. The 0.7 → 0.8 API surface we touch in `db.rs` is stable across the bump. Deferred to a later iter: - shlex 0.1.1 → ≥1.3.0 (transitive via wasm3-sys, only on --features wasm3 which is default-off; will be addressed when the wasm3 path is removed per ADR-128 Q2 Wasmtime resolution) - wasmtime 25 → 36+/42+ (HC-03/04 CVSS 9.0 sandbox-escape) — being handled by a background coder agent this iter, separate commit. Refs: docs/security/HOMECORE-security-audit-iter10.md (HC-09 sqlx) Refs: #798 Co-Authored-By: claude-flow <ruv@ruv.net> * fix(homecore-plugins/sec): bump wasmtime 25 → 42 for RUSTSEC-2026-0095/0096 (HC-03/04, CVSS 9.0) Remediates iter-11 security audit findings HC-03 (RUSTSEC-2026-0095) and HC-04 (RUSTSEC-2026-0096) — Cranelift/Winch sandbox-escape CVEs (CVSS 9.0). Version specifier updated from "25" → "42"; lockfile already pinned at 42.0.2. Zero code-surface changes required: Engine/Linker/Store/Instance and Memory.data/data_mut APIs are ABI-compatible across this range. All 15 tests pass (12 unit + 3 integration including the two required wasm_plugin_temp_threshold tests). cargo audit no longer reports RUSTSEC-2026-0095 or RUSTSEC-2026-0096 against this workspace. Co-Authored-By: claude-flow <ruv@ruv.net> * perf(homecore): criterion benches for state-machine hot paths `cargo bench -p homecore --bench state_machine` covers: - set/first_write — cold-path insert + alloc + broadcast - set/warm_write_state_change — same-entity update fires broadcast - set/noop_suppressed — same state+attrs, no broadcast (HA semantic) - get/hit + get/miss — zero-copy Arc<State> read paths - all_snapshot/{10,100,1000} — Vec<Arc<State>> snapshot for REST - all_by_domain_light_20_of_100 — domain prefix filter - broadcast_fan_out/{1,4,16,64} — 1 sender + N subscribers, async, measures end-to-end deliver-and-recv latency The broadcast fan-out is the most load-bearing measurement for HOMECORE — every integration, the recorder, the automation engine, and every WS subscriber holds a receiver, so the per-subscriber delivery cost determines how many add-ons the runtime can host. criterion 0.5 with sample_size=20 (fast tick, the fast-path benches run in nanoseconds and don't need 100 samples). Refs: docs/adr/ADR-127-homecore-state-machine-rust.md Refs: #798 Co-Authored-By: claude-flow <ruv@ruv.net> * fix(homecore-api/sec): close HC-01/HC-02 — real bearer-token store Replaces the P1 "any non-empty bearer" placeholder with a real LongLivedTokenStore (HashSet<String>) on SharedState. Closes the two Critical findings from the iter-10 security audit (docs/security/HOMECORE-security-audit-iter10.md HC-01 + HC-02). New module `homecore-api::tokens`: - LongLivedTokenStore::empty() — default-deny - LongLivedTokenStore::from_env() — reads HOMECORE_TOKENS=t1,t2,t3 - LongLivedTokenStore::allow_any_non_empty() — DEV-only, warns on every check, preserves legacy behaviour for migrating users - register / revoke / is_valid / len / is_dev_mode — full API Wired through: - SharedState gains `tokens: LongLivedTokenStore`; constructors with_tokens(...) for explicit injection; with_metadata defaults to DEV (allow_any) for backwards compat with existing smoke tests - BearerAuth::from_headers now async + takes &LongLivedTokenStore; checks store.is_valid(token) before returning Ok - All 6 REST handlers updated to thread the store and await the validation - homecore-server reads HOMECORE_TOKENS at boot; if set, builds the store from env; if unset, falls back to DEV with a warn log Test count: 4 → 15 (+11 token-store + auth-with-store tests). Smoke verified end-to-end: HOMECORE_TOKENS=good homecore-server --bind 127.0.0.1:8126 → "LongLivedTokenStore provisioned with 1 bearer token(s)" curl -H "Authorization: Bearer good" .../api/states → 200 curl -H "Authorization: Bearer wrong" .../api/states → 401 curl -H "Authorization: Bearer " .../api/states → 401 curl .../api/states → 401 Refs: docs/security/HOMECORE-security-audit-iter10.md (HC-01 + HC-02) Refs: docs/adr/ADR-130-homecore-rest-websocket-api.md §3 auth Refs: #798 Refs: #800 Co-Authored-By: claude-flow <ruv@ruv.net> * fix(homecore-api/sec): close HC-05 — CORS allowlist instead of permissive Replaces `CorsLayer::permissive()` (which set Access-Control-Allow- Origin: *) with an explicit allowlist via `CorsLayer::new()`. Default allowlist covers the homecore-frontend Vite dev server (5173) plus common reverse-proxy ports (3000, 8080, 8081) and the bind port itself (8123). Production deployments override via HOMECORE_CORS_ORIGINS=https://app.example.com,https://hass.example.com (comma-separated). Method allowlist: GET, POST, OPTIONS, DELETE (no PUT/PATCH yet). Header allowlist: Authorization, Content-Type, Accept. Credentials: disabled (no cookies in HOMECORE-API path). Test count: 15 → 18 (+3 CORS allowlist tests). Closes audit finding HC-05 (High). The HC-01/02 bearer-store fix in commit408cfd4f0only mattered if the cross-origin path was also locked down — without HC-05 a malicious page could still make authenticated calls with a stored bearer. Refs: docs/security/HOMECORE-security-audit-iter10.md (HC-05) Refs: #800 Co-Authored-By: claude-flow <ruv@ruv.net>
This commit is contained in:
@@ -0,0 +1,562 @@
|
||||
//! `Recorder` — SQLite write path + query path.
|
||||
//!
|
||||
//! Wraps an `SqlitePool` and exposes three operations:
|
||||
//! - [`Recorder::open`] — open (or create) the DB and apply schema.
|
||||
//! - [`Recorder::record_state`] — persist a `StateChangedEvent`.
|
||||
//! - [`Recorder::record_event`] — persist a `DomainEvent`.
|
||||
//! - [`Recorder::get_state_history`] — read back rows in time order.
|
||||
//!
|
||||
//! State attributes are deduped via `fnv64a_hash` (see [`crate::dedup`]):
|
||||
//! if an identical attributes blob was previously written its
|
||||
//! `attributes_id` is reused and no new row is inserted.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use chrono::{DateTime, Utc};
|
||||
use sqlx::sqlite::{SqliteConnectOptions, SqlitePool, SqlitePoolOptions};
|
||||
use thiserror::Error;
|
||||
use tokio::sync::RwLock;
|
||||
use tracing::debug;
|
||||
|
||||
use homecore::entity::{EntityId, State};
|
||||
use homecore::event::{DomainEvent, StateChangedEvent};
|
||||
|
||||
use crate::dedup::fnv64a_hash;
|
||||
use crate::schema::ALL_DDL;
|
||||
|
||||
/// Errors returned by `Recorder` operations.
|
||||
#[derive(Error, Debug)]
|
||||
pub enum RecorderError {
|
||||
#[error("SQLite error: {0}")]
|
||||
Sqlx(#[from] sqlx::Error),
|
||||
|
||||
#[error("serialisation error: {0}")]
|
||||
Json(#[from] serde_json::Error),
|
||||
|
||||
#[error("URL parse error: {0}")]
|
||||
UrlParse(String),
|
||||
}
|
||||
|
||||
/// Trait for pluggable semantic (vector) indexing of state writes.
|
||||
///
|
||||
/// The no-op [`NullSemanticIndex`] is used in P1. P2 ships a ruvector-backed
|
||||
/// implementation behind the `ruvector` feature flag.
|
||||
///
|
||||
/// ## P2 API change
|
||||
///
|
||||
/// The `insert_state` method now accepts a `state_id` (SQLite rowid) so the
|
||||
/// HNSW index can map vector results back to SQLite rows. `search` embeds a
|
||||
/// free-text query and returns `(state_id, score)` pairs.
|
||||
#[async_trait]
|
||||
pub trait SemanticIndex: Send + Sync {
|
||||
/// Insert an embedding for `state` keyed by its SQLite `state_id`.
|
||||
/// Called after the SQLite insert succeeds. Must not propagate errors
|
||||
/// back to the recorder — failure is logged, not fatal.
|
||||
async fn insert_state(
|
||||
&mut self,
|
||||
state_id: i64,
|
||||
state: &State,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>>;
|
||||
|
||||
/// Search for the `k` nearest states to the free-text `query`.
|
||||
/// Returns `(state_id, score)` pairs sorted by ascending distance.
|
||||
async fn search(
|
||||
&self,
|
||||
query: &str,
|
||||
k: usize,
|
||||
) -> Result<Vec<(i64, f32)>, Box<dyn std::error::Error + Send + Sync>>;
|
||||
}
|
||||
|
||||
/// No-op `SemanticIndex`. Used by default when the `ruvector` feature is off.
|
||||
pub struct NullSemanticIndex;
|
||||
|
||||
#[async_trait]
|
||||
impl SemanticIndex for NullSemanticIndex {
|
||||
async fn insert_state(
|
||||
&mut self,
|
||||
_state_id: i64,
|
||||
_state: &State,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn search(
|
||||
&self,
|
||||
_query: &str,
|
||||
_k: usize,
|
||||
) -> Result<Vec<(i64, f32)>, Box<dyn std::error::Error + Send + Sync>> {
|
||||
Ok(vec![])
|
||||
}
|
||||
}
|
||||
|
||||
/// The recorder. Cheap to clone (Arc-backed pool). Pass copies to the
|
||||
/// `RecorderListener` and the API history handler.
|
||||
///
|
||||
/// The `semantic` field is wrapped in `Arc<RwLock<...>>` so that
|
||||
/// `insert_state` (which takes `&mut self` on the trait) can be called
|
||||
/// without requiring `&mut Recorder` from callers.
|
||||
#[derive(Clone)]
|
||||
pub struct Recorder {
|
||||
pool: SqlitePool,
|
||||
semantic: Arc<RwLock<dyn SemanticIndex>>,
|
||||
}
|
||||
|
||||
impl Recorder {
|
||||
/// Open (or create) the SQLite database at `path` and apply the schema.
|
||||
///
|
||||
/// Pass `"sqlite::memory:"` for an in-memory database (tests).
|
||||
///
|
||||
/// The schema DDL uses `CREATE TABLE IF NOT EXISTS` so calling this on an
|
||||
/// existing database is safe.
|
||||
pub async fn open(path: &str) -> Result<Self, RecorderError> {
|
||||
Self::open_with_index(path, Arc::new(RwLock::new(NullSemanticIndex))).await
|
||||
}
|
||||
|
||||
/// Open with a custom `SemanticIndex` (P2 entry point).
|
||||
pub async fn open_with_index(
|
||||
path: &str,
|
||||
semantic: Arc<RwLock<dyn SemanticIndex>>,
|
||||
) -> Result<Self, RecorderError> {
|
||||
let options = path
|
||||
.parse::<SqliteConnectOptions>()
|
||||
.map_err(|e| RecorderError::UrlParse(e.to_string()))?
|
||||
.create_if_missing(true);
|
||||
|
||||
let pool = SqlitePoolOptions::new()
|
||||
.max_connections(4)
|
||||
.connect_with(options)
|
||||
.await?;
|
||||
|
||||
let recorder = Self { pool, semantic };
|
||||
recorder.apply_schema().await?;
|
||||
Ok(recorder)
|
||||
}
|
||||
|
||||
/// Apply all DDL statements. Idempotent.
|
||||
async fn apply_schema(&self) -> Result<(), RecorderError> {
|
||||
for ddl in ALL_DDL {
|
||||
// Each DDL block may contain multiple statements separated by `;`.
|
||||
// sqlx::query does not support multi-statement strings directly,
|
||||
// so we split on the statement boundary and execute individually.
|
||||
for stmt in split_statements(ddl) {
|
||||
let stmt = stmt.trim();
|
||||
if !stmt.is_empty() {
|
||||
sqlx::query(stmt).execute(&self.pool).await?;
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Persist a `StateChangedEvent`. Inserts into `states` and dedupes into
|
||||
/// `state_attributes`. Returns the `state_id` of the new row.
|
||||
pub async fn record_state(
|
||||
&self,
|
||||
event: &StateChangedEvent,
|
||||
) -> Result<Option<i64>, RecorderError> {
|
||||
let new_state = match &event.new_state {
|
||||
Some(s) => s,
|
||||
None => return Ok(None), // removal event — no row to insert
|
||||
};
|
||||
|
||||
let attrs_json = serde_json::to_string(&new_state.attributes)?;
|
||||
let hash = fnv64a_hash(&attrs_json);
|
||||
|
||||
// Upsert into state_attributes (dedup by hash).
|
||||
let attributes_id: i64 = {
|
||||
// Try to find an existing row first.
|
||||
let existing: Option<(i64,)> =
|
||||
sqlx::query_as("SELECT attributes_id FROM state_attributes WHERE hash = ?")
|
||||
.bind(hash)
|
||||
.fetch_optional(&self.pool)
|
||||
.await?;
|
||||
|
||||
if let Some((id,)) = existing {
|
||||
debug!(hash, id, "reusing existing state_attributes row");
|
||||
id
|
||||
} else {
|
||||
let result =
|
||||
sqlx::query("INSERT INTO state_attributes (shared_attrs, hash) VALUES (?, ?)")
|
||||
.bind(&attrs_json)
|
||||
.bind(hash)
|
||||
.execute(&self.pool)
|
||||
.await?;
|
||||
result.last_insert_rowid()
|
||||
}
|
||||
};
|
||||
|
||||
let context_id = new_state.context.id.to_string();
|
||||
let last_changed_ts = new_state.last_changed.timestamp_micros() as f64 / 1_000_000.0;
|
||||
let last_updated_ts = new_state.last_updated.timestamp_micros() as f64 / 1_000_000.0;
|
||||
|
||||
let result = sqlx::query(
|
||||
"INSERT INTO states \
|
||||
(entity_id, state, attributes_id, last_changed_ts, last_updated_ts, context_id) \
|
||||
VALUES (?, ?, ?, ?, ?, ?)",
|
||||
)
|
||||
.bind(new_state.entity_id.as_str())
|
||||
.bind(&new_state.state)
|
||||
.bind(attributes_id)
|
||||
.bind(last_changed_ts)
|
||||
.bind(last_updated_ts)
|
||||
.bind(&context_id)
|
||||
.execute(&self.pool)
|
||||
.await?;
|
||||
|
||||
let state_id = result.last_insert_rowid();
|
||||
|
||||
// Best-effort semantic indexing — failure is logged, not propagated.
|
||||
if let Err(e) = self
|
||||
.semantic
|
||||
.write()
|
||||
.await
|
||||
.insert_state(state_id, new_state)
|
||||
.await
|
||||
{
|
||||
tracing::warn!(
|
||||
error = %e,
|
||||
entity_id = %new_state.entity_id,
|
||||
"semantic indexing failed"
|
||||
);
|
||||
}
|
||||
|
||||
Ok(Some(state_id))
|
||||
}
|
||||
|
||||
/// Search for state history rows that semantically match `query`.
|
||||
///
|
||||
/// Uses the HNSW index to find the top-`k` nearest state embeddings,
|
||||
/// then fetches the full `StateRow` from SQLite for each result.
|
||||
/// Returns rows in ascending score (distance) order.
|
||||
///
|
||||
/// With the default `NullSemanticIndex` (no `ruvector` feature) this
|
||||
/// always returns an empty `Vec`.
|
||||
pub async fn search_semantic(
|
||||
&self,
|
||||
query: &str,
|
||||
k: usize,
|
||||
) -> Result<Vec<StateRow>, RecorderError> {
|
||||
let hits = self
|
||||
.semantic
|
||||
.read()
|
||||
.await
|
||||
.search(query, k)
|
||||
.await
|
||||
.unwrap_or_default();
|
||||
|
||||
let mut rows = Vec::with_capacity(hits.len());
|
||||
for (state_id, _score) in hits {
|
||||
let row: Option<(String, String, Option<String>, f64, f64, Option<String>)> =
|
||||
sqlx::query_as(
|
||||
"SELECT s.entity_id, s.state, sa.shared_attrs, \
|
||||
s.last_changed_ts, s.last_updated_ts, s.context_id \
|
||||
FROM states s \
|
||||
LEFT JOIN state_attributes sa ON s.attributes_id = sa.attributes_id \
|
||||
WHERE s.state_id = ?",
|
||||
)
|
||||
.bind(state_id)
|
||||
.fetch_optional(&self.pool)
|
||||
.await?;
|
||||
|
||||
if let Some((entity_id, state, shared_attrs, last_changed_ts, last_updated_ts, context_id)) = row {
|
||||
let eid = EntityId::parse(&entity_id)
|
||||
.unwrap_or_else(|_| EntityId::parse("unknown.unknown").unwrap());
|
||||
let attributes = shared_attrs
|
||||
.as_deref()
|
||||
.map(serde_json::from_str)
|
||||
.transpose()?
|
||||
.unwrap_or(serde_json::Value::Object(Default::default()));
|
||||
rows.push(StateRow {
|
||||
state_id,
|
||||
entity_id: eid,
|
||||
state,
|
||||
attributes,
|
||||
last_changed_ts,
|
||||
last_updated_ts,
|
||||
context_id,
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(rows)
|
||||
}
|
||||
|
||||
/// Persist a `DomainEvent`. Returns the `event_id`.
|
||||
pub async fn record_event(&self, event: &DomainEvent) -> Result<i64, RecorderError> {
|
||||
let data_json = serde_json::to_string(&event.event_data)?;
|
||||
let time_fired_ts = event.fired_at.timestamp_micros() as f64 / 1_000_000.0;
|
||||
let context_id = event.context.id.to_string();
|
||||
|
||||
let result = sqlx::query(
|
||||
"INSERT INTO events (event_type, event_data, time_fired_ts, context_id) \
|
||||
VALUES (?, ?, ?, ?)",
|
||||
)
|
||||
.bind(&event.event_type)
|
||||
.bind(&data_json)
|
||||
.bind(time_fired_ts)
|
||||
.bind(&context_id)
|
||||
.execute(&self.pool)
|
||||
.await?;
|
||||
|
||||
Ok(result.last_insert_rowid())
|
||||
}
|
||||
|
||||
/// Query state history for `entity_id` between `since` and `until`.
|
||||
/// Returns state snapshots in ascending `last_updated_ts` order.
|
||||
pub async fn get_state_history(
|
||||
&self,
|
||||
entity_id: &EntityId,
|
||||
since: DateTime<Utc>,
|
||||
until: DateTime<Utc>,
|
||||
) -> Result<Vec<StateRow>, RecorderError> {
|
||||
let since_ts = since.timestamp_micros() as f64 / 1_000_000.0;
|
||||
let until_ts = until.timestamp_micros() as f64 / 1_000_000.0;
|
||||
|
||||
let rows: Vec<(i64, String, Option<String>, f64, f64, Option<String>)> = sqlx::query_as(
|
||||
"SELECT s.state_id, s.state, sa.shared_attrs, \
|
||||
s.last_changed_ts, s.last_updated_ts, s.context_id \
|
||||
FROM states s \
|
||||
LEFT JOIN state_attributes sa ON s.attributes_id = sa.attributes_id \
|
||||
WHERE s.entity_id = ? \
|
||||
AND s.last_updated_ts >= ? \
|
||||
AND s.last_updated_ts <= ? \
|
||||
ORDER BY s.last_updated_ts ASC",
|
||||
)
|
||||
.bind(entity_id.as_str())
|
||||
.bind(since_ts)
|
||||
.bind(until_ts)
|
||||
.fetch_all(&self.pool)
|
||||
.await?;
|
||||
|
||||
rows.into_iter()
|
||||
.map(|(state_id, state, shared_attrs, last_changed_ts, last_updated_ts, context_id)| {
|
||||
let attributes = shared_attrs
|
||||
.as_deref()
|
||||
.map(serde_json::from_str)
|
||||
.transpose()?
|
||||
.unwrap_or(serde_json::Value::Object(Default::default()));
|
||||
|
||||
Ok(StateRow {
|
||||
state_id,
|
||||
entity_id: entity_id.clone(),
|
||||
state,
|
||||
attributes,
|
||||
last_changed_ts,
|
||||
last_updated_ts,
|
||||
context_id,
|
||||
})
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
/// A state row returned from `get_state_history`.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct StateRow {
|
||||
pub state_id: i64,
|
||||
pub entity_id: EntityId,
|
||||
pub state: String,
|
||||
pub attributes: serde_json::Value,
|
||||
/// Unix timestamp (seconds, fractional) when the state string last changed.
|
||||
pub last_changed_ts: f64,
|
||||
/// Unix timestamp (seconds, fractional) when this snapshot was written.
|
||||
pub last_updated_ts: f64,
|
||||
pub context_id: Option<String>,
|
||||
}
|
||||
|
||||
/// Split a multi-statement DDL string on `;` boundaries.
|
||||
/// Trims whitespace; skips empty fragments.
|
||||
fn split_statements(ddl: &str) -> impl Iterator<Item = &str> {
|
||||
ddl.split(';').map(str::trim).filter(|s| !s.is_empty())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::sync::Arc;
|
||||
|
||||
use chrono::Utc;
|
||||
|
||||
use homecore::entity::{EntityId, State};
|
||||
use homecore::event::{Context, DomainEvent, StateChangedEvent};
|
||||
|
||||
use super::*;
|
||||
|
||||
async fn open_memory() -> Recorder {
|
||||
Recorder::open("sqlite::memory:").await.expect("open in-memory DB")
|
||||
}
|
||||
|
||||
fn entity(s: &str) -> EntityId {
|
||||
EntityId::parse(s).unwrap()
|
||||
}
|
||||
|
||||
fn make_state_event(entity_id: &str, state_val: &str, attrs: serde_json::Value) -> StateChangedEvent {
|
||||
let eid = entity(entity_id);
|
||||
let ctx = Context::new();
|
||||
let s = Arc::new(State::new(eid.clone(), state_val, attrs, ctx));
|
||||
StateChangedEvent {
|
||||
entity_id: eid,
|
||||
old_state: None,
|
||||
new_state: Some(s),
|
||||
fired_at: Utc::now(),
|
||||
}
|
||||
}
|
||||
|
||||
// ── schema ────────────────────────────────────────────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn schema_applies_on_fresh_db() {
|
||||
let recorder = open_memory().await;
|
||||
// Verify all four tables exist by querying sqlite_master.
|
||||
let tables: Vec<(String,)> =
|
||||
sqlx::query_as("SELECT name FROM sqlite_master WHERE type='table' ORDER BY name")
|
||||
.fetch_all(&recorder.pool)
|
||||
.await
|
||||
.unwrap();
|
||||
let names: Vec<&str> = tables.iter().map(|(n,)| n.as_str()).collect();
|
||||
assert!(names.contains(&"state_attributes"), "missing state_attributes");
|
||||
assert!(names.contains(&"states"), "missing states");
|
||||
assert!(names.contains(&"events"), "missing events");
|
||||
assert!(names.contains(&"recorder_runs"), "missing recorder_runs");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn schema_idempotent_double_open() {
|
||||
// Applying schema twice (on the same pool) must not panic or error.
|
||||
let recorder = open_memory().await;
|
||||
recorder.apply_schema().await.expect("second apply_schema must be a no-op");
|
||||
}
|
||||
|
||||
// ── record_state ──────────────────────────────────────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn record_state_inserts_row() {
|
||||
let recorder = open_memory().await;
|
||||
let event = make_state_event("light.kitchen", "on", serde_json::json!({"brightness": 200}));
|
||||
|
||||
let state_id = recorder.record_state(&event).await.unwrap();
|
||||
assert!(state_id.is_some(), "expected a state_id");
|
||||
|
||||
let count: (i64,) =
|
||||
sqlx::query_as("SELECT COUNT(*) FROM states WHERE entity_id = 'light.kitchen'")
|
||||
.fetch_one(&recorder.pool)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(count.0, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn removal_event_returns_none() {
|
||||
let recorder = open_memory().await;
|
||||
let event = StateChangedEvent {
|
||||
entity_id: entity("light.kitchen"),
|
||||
old_state: None,
|
||||
new_state: None, // removal
|
||||
fired_at: Utc::now(),
|
||||
};
|
||||
let result = recorder.record_state(&event).await.unwrap();
|
||||
assert!(result.is_none(), "removal event should yield None state_id");
|
||||
}
|
||||
|
||||
// ── attribute deduplication ────────────────────────────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn same_attrs_dedup_to_one_row() {
|
||||
let recorder = open_memory().await;
|
||||
let attrs = serde_json::json!({"brightness": 200, "color_temp": 4000});
|
||||
|
||||
let e1 = make_state_event("light.a", "on", attrs.clone());
|
||||
let e2 = make_state_event("light.b", "on", attrs.clone());
|
||||
|
||||
recorder.record_state(&e1).await.unwrap();
|
||||
recorder.record_state(&e2).await.unwrap();
|
||||
|
||||
let attr_count: (i64,) =
|
||||
sqlx::query_as("SELECT COUNT(*) FROM state_attributes")
|
||||
.fetch_one(&recorder.pool)
|
||||
.await
|
||||
.unwrap();
|
||||
// Both events share identical attrs → only one state_attributes row.
|
||||
assert_eq!(attr_count.0, 1, "identical attrs must share one state_attributes row");
|
||||
|
||||
let state_count: (i64,) =
|
||||
sqlx::query_as("SELECT COUNT(*) FROM states")
|
||||
.fetch_one(&recorder.pool)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(state_count.0, 2, "two states rows expected");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn different_attrs_each_get_own_row() {
|
||||
let recorder = open_memory().await;
|
||||
let e1 = make_state_event("sensor.a", "20", serde_json::json!({"unit": "C"}));
|
||||
let e2 = make_state_event("sensor.b", "20", serde_json::json!({"unit": "F"}));
|
||||
|
||||
recorder.record_state(&e1).await.unwrap();
|
||||
recorder.record_state(&e2).await.unwrap();
|
||||
|
||||
let attr_count: (i64,) =
|
||||
sqlx::query_as("SELECT COUNT(*) FROM state_attributes")
|
||||
.fetch_one(&recorder.pool)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(attr_count.0, 2);
|
||||
}
|
||||
|
||||
// ── get_state_history ─────────────────────────────────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn history_returns_rows_in_time_order() {
|
||||
let recorder = open_memory().await;
|
||||
let eid = entity("sensor.temp");
|
||||
|
||||
// Insert three states with slightly different timestamps by sleeping.
|
||||
for val in &["20.0", "21.0", "22.0"] {
|
||||
let e = make_state_event("sensor.temp", val, serde_json::json!({}));
|
||||
recorder.record_state(&e).await.unwrap();
|
||||
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
|
||||
}
|
||||
|
||||
let since = Utc::now() - chrono::Duration::seconds(10);
|
||||
let until = Utc::now() + chrono::Duration::seconds(10);
|
||||
let rows = recorder.get_state_history(&eid, since, until).await.unwrap();
|
||||
|
||||
assert_eq!(rows.len(), 3, "expected 3 history rows");
|
||||
// Verify ascending order by last_updated_ts.
|
||||
for w in rows.windows(2) {
|
||||
assert!(
|
||||
w[0].last_updated_ts <= w[1].last_updated_ts,
|
||||
"rows must be in ascending time order"
|
||||
);
|
||||
}
|
||||
assert_eq!(rows[0].state, "20.0");
|
||||
assert_eq!(rows[2].state, "22.0");
|
||||
}
|
||||
|
||||
// ── record_event ──────────────────────────────────────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn record_event_round_trips() {
|
||||
let recorder = open_memory().await;
|
||||
let ctx = Context::new();
|
||||
let event = DomainEvent::new(
|
||||
"call_service",
|
||||
serde_json::json!({"domain": "light", "service": "turn_on"}),
|
||||
ctx,
|
||||
);
|
||||
|
||||
let event_id = recorder.record_event(&event).await.unwrap();
|
||||
assert!(event_id > 0);
|
||||
|
||||
let row: (String, String) =
|
||||
sqlx::query_as("SELECT event_type, event_data FROM events WHERE event_id = ?")
|
||||
.bind(event_id)
|
||||
.fetch_one(&recorder.pool)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(row.0, "call_service");
|
||||
let data: serde_json::Value = serde_json::from_str(&row.1).unwrap();
|
||||
assert_eq!(data["domain"], "light");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
//! FNV-1a 64-bit hash for state-attribute deduplication.
|
||||
//!
|
||||
//! Matches Home Assistant's `db_schema.py` `fnv64a` function used to
|
||||
//! fingerprint shared attribute blobs. Two state writes with identical
|
||||
//! attributes share a single `state_attributes` row, reducing I/O by
|
||||
//! ~80% for high-frequency polling sensors.
|
||||
//!
|
||||
//! ## FNV-1a 64 spec
|
||||
//!
|
||||
//! - Offset basis: 0xcbf29ce484222325
|
||||
//! - Prime: 0x100000001b3
|
||||
//! - Per byte: `hash = (hash XOR byte) * prime`
|
||||
//!
|
||||
//! Reference values (computed from the spec + verified against HA source):
|
||||
//! - `""` (empty string) → signed i64: -3750763034362895579
|
||||
//! - `"a"` → signed i64: -5808556873153909620
|
||||
//! - `{"state": "on"}` → signed i64: 3947789143477681127
|
||||
|
||||
const FNV_OFFSET_BASIS_64: u64 = 0xcbf29ce484222325;
|
||||
const FNV_PRIME_64: u64 = 0x100000001b3;
|
||||
|
||||
/// Compute FNV-1a 64-bit hash of `data` bytes, returned as a signed `i64`
|
||||
/// suitable for direct storage in SQLite's INTEGER column.
|
||||
///
|
||||
/// The cast to `i64` is a bit-reinterpret, not a value conversion — the
|
||||
/// same pattern HA uses in `db_schema.py`.
|
||||
#[inline]
|
||||
pub fn fnv64a_bytes(data: &[u8]) -> i64 {
|
||||
let mut hash: u64 = FNV_OFFSET_BASIS_64;
|
||||
for &byte in data {
|
||||
hash ^= u64::from(byte);
|
||||
hash = hash.wrapping_mul(FNV_PRIME_64);
|
||||
}
|
||||
hash as i64
|
||||
}
|
||||
|
||||
/// Hash a UTF-8 string. Convenience wrapper over [`fnv64a_bytes`].
|
||||
#[inline]
|
||||
pub fn fnv64a_hash(s: &str) -> i64 {
|
||||
fnv64a_bytes(s.as_bytes())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// HA reference: `fnv64a(b"")` → 0xcbf29ce484222325 (unsigned)
|
||||
/// As signed i64: -3750763034362895579
|
||||
#[test]
|
||||
fn hash_empty_string() {
|
||||
assert_eq!(fnv64a_hash(""), -3750763034362895579_i64);
|
||||
}
|
||||
|
||||
/// HA reference: `fnv64a(b"a")` → 0xaf63dc4c8601ec8c (unsigned)
|
||||
/// As signed i64: -5808556873153909620
|
||||
#[test]
|
||||
fn hash_single_char_a() {
|
||||
assert_eq!(fnv64a_hash("a"), -5808556873153909620_i64);
|
||||
}
|
||||
|
||||
/// Smoke-test a realistic JSON attribute blob.
|
||||
/// `{"state": "on"}` → signed i64: 3947789143477681127
|
||||
#[test]
|
||||
fn hash_json_blob() {
|
||||
assert_eq!(fnv64a_hash(r#"{"state": "on"}"#), 3947789143477681127_i64);
|
||||
}
|
||||
|
||||
/// Different strings must produce different hashes (basic collision check).
|
||||
#[test]
|
||||
fn distinct_strings_differ() {
|
||||
assert_ne!(fnv64a_hash("on"), fnv64a_hash("off"));
|
||||
assert_ne!(fnv64a_hash("{\"brightness\":100}"), fnv64a_hash("{\"brightness\":200}"));
|
||||
}
|
||||
|
||||
/// Deterministic: same input always gives same output.
|
||||
#[test]
|
||||
fn deterministic() {
|
||||
let s = r#"{"unit": "C", "value": 22.5}"#;
|
||||
assert_eq!(fnv64a_hash(s), fnv64a_hash(s));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
//! homecore-recorder — SQLite state history + semantic search.
|
||||
//!
|
||||
//! Implements ADR-132: dual-write architecture. P1 ships SQLite structural
|
||||
//! persistence with an HA-compatible schema (mirrors HA recorder schema v48).
|
||||
//! P2 (feature `ruvector`) adds a `SemanticIndex` backed by ruvector
|
||||
//! embeddings for natural-language state queries.
|
||||
//!
|
||||
//! ## P1 architecture
|
||||
//!
|
||||
//! ```text
|
||||
//! StateMachine ──broadcast──► RecorderListener ──► Recorder
|
||||
//! │
|
||||
//! ┌───────┴──────────┐
|
||||
//! states state_attributes
|
||||
//! events recorder_runs
|
||||
//! ```
|
||||
//!
|
||||
//! ## P2 hand-off (ruvector feature)
|
||||
//!
|
||||
//! When the `ruvector` feature is enabled, the `Recorder` additionally
|
||||
//! calls a `SemanticIndex` implementation that embeds state attributes and
|
||||
//! stores vectors in ruvector for k-NN semantic search. See [`semantic`].
|
||||
|
||||
pub mod db;
|
||||
pub mod dedup;
|
||||
pub mod listener;
|
||||
pub mod schema;
|
||||
|
||||
#[cfg(feature = "ruvector")]
|
||||
pub mod semantic;
|
||||
|
||||
// Re-export the primary public API surface.
|
||||
pub use db::{Recorder, RecorderError};
|
||||
pub use listener::RecorderListener;
|
||||
|
||||
/// Null semantic index used when the `ruvector` feature is off.
|
||||
/// Satisfies the [`db::SemanticIndex`] trait bound without any allocation.
|
||||
pub use db::NullSemanticIndex;
|
||||
@@ -0,0 +1,117 @@
|
||||
//! `RecorderListener` — subscribes to `StateMachine` broadcasts and writes
|
||||
//! every `StateChangedEvent` to the `Recorder`.
|
||||
//!
|
||||
//! Spawned via `tokio::spawn`. Runs until the broadcast sender is dropped
|
||||
//! (i.e. the `StateMachine` is shut down) or until a `Lagged` error occurs
|
||||
//! (subscriber fell more than 4,096 events behind).
|
||||
//!
|
||||
//! On `Lagged`, the listener logs a warning and reconnects; it does not crash
|
||||
//! because dropping a listener would silently stop persistence.
|
||||
//!
|
||||
//! ## Subscription ordering
|
||||
//!
|
||||
//! The `broadcast::Receiver` is created inside `new()` (not inside the spawned
|
||||
//! task), so any events fired between `new()` and `spawn()` are enqueued in
|
||||
//! the receiver buffer and will be drained when the task starts.
|
||||
|
||||
use tokio::sync::broadcast;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use homecore::event::StateChangedEvent;
|
||||
use homecore::state::StateMachine;
|
||||
|
||||
use crate::db::Recorder;
|
||||
|
||||
/// A background task that records every state change.
|
||||
///
|
||||
/// Call [`RecorderListener::new`] then [`RecorderListener::spawn`].
|
||||
/// The subscription starts at construction time so no events are missed
|
||||
/// between `new()` and `spawn()`.
|
||||
pub struct RecorderListener {
|
||||
recorder: Recorder,
|
||||
rx: broadcast::Receiver<StateChangedEvent>,
|
||||
}
|
||||
|
||||
impl RecorderListener {
|
||||
/// Create a listener. Subscribes to the broadcast channel immediately so
|
||||
/// events fired before `spawn()` are buffered in the receiver.
|
||||
pub fn new(state_machine: &StateMachine, recorder: Recorder) -> Self {
|
||||
let rx = state_machine.subscribe();
|
||||
Self { recorder, rx }
|
||||
}
|
||||
|
||||
/// Spawn the listener onto the Tokio runtime.
|
||||
///
|
||||
/// Returns a `JoinHandle`. Abort it on graceful shutdown:
|
||||
/// ```ignore
|
||||
/// let handle = listener.spawn();
|
||||
/// // … on shutdown:
|
||||
/// handle.abort();
|
||||
/// ```
|
||||
pub fn spawn(self) -> tokio::task::JoinHandle<()> {
|
||||
tokio::spawn(async move { self.run().await })
|
||||
}
|
||||
|
||||
async fn run(mut self) {
|
||||
loop {
|
||||
match self.rx.recv().await {
|
||||
Ok(event) => {
|
||||
debug!(entity_id = %event.entity_id, "recording state change");
|
||||
if let Err(e) = self.recorder.record_state(&event).await {
|
||||
warn!(error = %e, "failed to record state change");
|
||||
}
|
||||
}
|
||||
Err(broadcast::error::RecvError::Lagged(n)) => {
|
||||
warn!(
|
||||
lagged_by = n,
|
||||
"recorder listener lagged — some state changes were not persisted"
|
||||
);
|
||||
// Continue processing from the next available event.
|
||||
}
|
||||
Err(broadcast::error::RecvError::Closed) => {
|
||||
debug!("state machine shut down; recorder listener exiting");
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
use homecore::entity::EntityId;
|
||||
use homecore::event::Context;
|
||||
|
||||
fn eid(s: &str) -> EntityId {
|
||||
EntityId::parse(s).unwrap()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn listener_records_state_changes() {
|
||||
let sm = StateMachine::new();
|
||||
let recorder = Recorder::open("sqlite::memory:").await.unwrap();
|
||||
|
||||
let listener = RecorderListener::new(&sm, recorder.clone());
|
||||
let _handle = listener.spawn();
|
||||
|
||||
// Fire two state changes.
|
||||
sm.set(eid("light.hall"), "on", serde_json::json!({}), Context::new());
|
||||
sm.set(eid("light.hall"), "off", serde_json::json!({}), Context::new());
|
||||
|
||||
// Give the background task a moment to flush.
|
||||
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
|
||||
|
||||
let since = chrono::Utc::now() - chrono::Duration::seconds(10);
|
||||
let until = chrono::Utc::now() + chrono::Duration::seconds(10);
|
||||
let rows = recorder
|
||||
.get_state_history(&eid("light.hall"), since, until)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(rows.len(), 2, "listener must have persisted both events");
|
||||
assert_eq!(rows[0].state, "on");
|
||||
assert_eq!(rows[1].state, "off");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,90 @@
|
||||
//! SQL DDL for the HA-compatible recorder schema (ADR-132).
|
||||
//!
|
||||
//! Schema mirrors Home Assistant recorder schema v48 (HA 2025.1):
|
||||
//! - `states` — one row per state write (entity_id, state, attrs)
|
||||
//! - `state_attributes` — shared attribute blobs, deduped by fnv64a hash
|
||||
//! - `events` — domain events fired by integrations
|
||||
//! - `recorder_runs` — boot/shutdown bookends for gap detection
|
||||
//!
|
||||
//! All DDL strings use `CREATE TABLE IF NOT EXISTS` so `apply_schema` is
|
||||
//! idempotent and safe to call on every startup.
|
||||
|
||||
/// Create `state_attributes` table.
|
||||
///
|
||||
/// `shared_attrs` is stored as TEXT (JSON blob). `hash` is the FNV-1a 64-bit
|
||||
/// hash of `shared_attrs` encoded as a signed i64 — matches HA's dedup key.
|
||||
pub const CREATE_STATE_ATTRIBUTES: &str = "
|
||||
CREATE TABLE IF NOT EXISTS state_attributes (
|
||||
attributes_id INTEGER PRIMARY KEY NOT NULL,
|
||||
shared_attrs TEXT NOT NULL,
|
||||
hash INTEGER NOT NULL
|
||||
);
|
||||
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS ix_state_attributes_hash
|
||||
ON state_attributes (hash);
|
||||
";
|
||||
|
||||
/// Create `states` table.
|
||||
///
|
||||
/// `state_id` — auto-increment primary key
|
||||
/// `entity_id` — validated `domain.name` string
|
||||
/// `state` — state value string (\"on\", \"off\", \"20.5\", …)
|
||||
/// `attributes_id` — FK → state_attributes (nullable for HA compat)
|
||||
/// `last_changed_ts` — Unix timestamp seconds (float, UTC)
|
||||
/// `last_updated_ts` — Unix timestamp seconds (float, UTC)
|
||||
/// `context_id` — UUID as TEXT; links to the causality chain
|
||||
pub const CREATE_STATES: &str = "
|
||||
CREATE TABLE IF NOT EXISTS states (
|
||||
state_id INTEGER PRIMARY KEY NOT NULL,
|
||||
entity_id TEXT NOT NULL,
|
||||
state TEXT,
|
||||
attributes_id INTEGER,
|
||||
last_changed_ts REAL,
|
||||
last_updated_ts REAL NOT NULL,
|
||||
context_id TEXT
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS ix_states_entity_id_last_updated_ts
|
||||
ON states (entity_id, last_updated_ts);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS ix_states_last_updated_ts
|
||||
ON states (last_updated_ts);
|
||||
";
|
||||
|
||||
/// Create `events` table.
|
||||
///
|
||||
/// `event_type` — string key (e.g. \"state_changed\", \"call_service\")
|
||||
/// `event_data` — JSON blob
|
||||
/// `time_fired_ts` — Unix timestamp seconds (float, UTC)
|
||||
/// `context_id` — UUID as TEXT
|
||||
pub const CREATE_EVENTS: &str = "
|
||||
CREATE TABLE IF NOT EXISTS events (
|
||||
event_id INTEGER PRIMARY KEY NOT NULL,
|
||||
event_type TEXT NOT NULL,
|
||||
event_data TEXT,
|
||||
time_fired_ts REAL NOT NULL,
|
||||
context_id TEXT
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS ix_events_event_type_time_fired_ts
|
||||
ON events (event_type, time_fired_ts);
|
||||
";
|
||||
|
||||
/// Create `recorder_runs` table.
|
||||
///
|
||||
/// Records each start/stop pair so the history API can annotate gaps.
|
||||
pub const CREATE_RECORDER_RUNS: &str = "
|
||||
CREATE TABLE IF NOT EXISTS recorder_runs (
|
||||
run_id INTEGER PRIMARY KEY NOT NULL,
|
||||
start_ts REAL NOT NULL,
|
||||
end_ts REAL
|
||||
);
|
||||
";
|
||||
|
||||
/// All DDL statements in dependency order.
|
||||
pub const ALL_DDL: &[&str] = &[
|
||||
CREATE_STATE_ATTRIBUTES,
|
||||
CREATE_STATES,
|
||||
CREATE_EVENTS,
|
||||
CREATE_RECORDER_RUNS,
|
||||
];
|
||||
@@ -0,0 +1,273 @@
|
||||
//! Ruvector-backed semantic index — ADR-132 P2.
|
||||
//!
|
||||
//! ## Embedding strategy (P2 — hash-based)
|
||||
//!
|
||||
//! To keep the recorder self-contained and avoid an ML model dependency at P2,
|
||||
//! state attributes are embedded by a deterministic SHA-256 hash procedure:
|
||||
//!
|
||||
//! 1. Canonicalise the state as `"{entity_id}={state}|{attributes_json}"`.
|
||||
//! 2. SHA-256 hash → 32 bytes.
|
||||
//! 3. Interpret the 32 bytes as 8 × `i32` (big-endian), cast to `f32`.
|
||||
//! 4. L2-normalise the resulting 8-element vector.
|
||||
//!
|
||||
//! This gives stable, reproducible 8-dimensional unit vectors suitable for
|
||||
//! cosine-distance HNSW search. Semantic similarity is **not** captured (two
|
||||
//! states with the same value but different entity IDs will differ). P3 will
|
||||
//! replace this with a learned sentence-embedding via `ruvector-attention`.
|
||||
//!
|
||||
//! ## P3 plan
|
||||
//!
|
||||
//! Replace `embed_bytes` with a call to
|
||||
//! `ruvector_attention::SentenceEmbedding::encode(&text)` for true semantic
|
||||
//! similarity. Increase `EMBEDDING_DIM` to 384 at that point.
|
||||
|
||||
use async_trait::async_trait;
|
||||
use sha2::{Digest, Sha256};
|
||||
|
||||
use homecore::entity::State;
|
||||
use ruvector_core::{
|
||||
types::{DbOptions, DistanceMetric, HnswConfig, SearchQuery, VectorEntry},
|
||||
VectorDB,
|
||||
};
|
||||
|
||||
use crate::db::SemanticIndex;
|
||||
|
||||
/// Dimensionality of the hash-based embedding vectors.
|
||||
///
|
||||
/// 8 dimensions: each SHA-256 chunk of 4 bytes becomes one `f32` component.
|
||||
/// Increase to 384 in P3 when switching to learned embeddings.
|
||||
pub const EMBEDDING_DIM: usize = 8;
|
||||
|
||||
/// Ruvector-backed `SemanticIndex` using in-memory HNSW and hash embeddings.
|
||||
///
|
||||
/// The index lives entirely in process memory. A restart clears it; P3 will
|
||||
/// add persistence via `ruvector-core`'s `storage` feature.
|
||||
pub struct RuvectorSemanticIndex {
|
||||
db: VectorDB,
|
||||
}
|
||||
|
||||
impl RuvectorSemanticIndex {
|
||||
/// Create a new in-memory HNSW index with the given `max_elements` capacity.
|
||||
///
|
||||
/// Uses cosine distance to match the unit-normalised hash embeddings.
|
||||
pub fn new(max_elements: usize) -> Result<Self, Box<dyn std::error::Error + Send + Sync>> {
|
||||
let options = DbOptions {
|
||||
dimensions: EMBEDDING_DIM,
|
||||
distance_metric: DistanceMetric::Cosine,
|
||||
// storage path is ignored when the `storage` feature is off
|
||||
storage_path: ":memory:".to_string(),
|
||||
hnsw_config: Some(HnswConfig {
|
||||
m: 16,
|
||||
ef_construction: 100,
|
||||
ef_search: 50,
|
||||
max_elements,
|
||||
}),
|
||||
quantization: None,
|
||||
};
|
||||
let db = VectorDB::new(options)?;
|
||||
Ok(Self { db })
|
||||
}
|
||||
|
||||
/// Embed a `State` to a deterministic 8-dimensional unit vector.
|
||||
///
|
||||
/// Canonical form: `"{entity_id}={state}|{attributes_json}"`
|
||||
/// The attributes JSON is sorted-key (via `serde_json`'s default ordering
|
||||
/// of `Map`, which preserves insertion order). For strict canonicalisation
|
||||
/// at P3, sort keys explicitly.
|
||||
pub fn embed_state(state: &State) -> Vec<f32> {
|
||||
let attrs = state.attributes.to_string();
|
||||
let input = format!("{}={}|{}", state.entity_id, state.state, attrs);
|
||||
Self::embed_str(&input)
|
||||
}
|
||||
|
||||
/// Embed an arbitrary string to a deterministic 8-dimensional unit vector.
|
||||
pub fn embed_str(input: &str) -> Vec<f32> {
|
||||
embed_bytes(input.as_bytes())
|
||||
}
|
||||
}
|
||||
|
||||
/// SHA-256 → 8 × f32 unit vector.
|
||||
///
|
||||
/// Split the 32-byte digest into 8 chunks of 4 bytes. Interpret each chunk
|
||||
/// as a big-endian `i32`, cast to `f32`, then L2-normalise.
|
||||
fn embed_bytes(data: &[u8]) -> Vec<f32> {
|
||||
let digest = Sha256::digest(data);
|
||||
let mut raw: Vec<f32> = digest
|
||||
.chunks_exact(4)
|
||||
.map(|chunk| {
|
||||
let bytes: [u8; 4] = chunk.try_into().expect("chunk is exactly 4 bytes");
|
||||
i32::from_be_bytes(bytes) as f32
|
||||
})
|
||||
.collect();
|
||||
|
||||
// L2-normalise
|
||||
let norm = raw.iter().map(|x| x * x).sum::<f32>().sqrt();
|
||||
if norm > 1e-10 {
|
||||
for v in &mut raw {
|
||||
*v /= norm;
|
||||
}
|
||||
}
|
||||
raw
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl SemanticIndex for RuvectorSemanticIndex {
|
||||
async fn insert_state(
|
||||
&mut self,
|
||||
state_id: i64,
|
||||
state: &State,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
let vector = Self::embed_state(state);
|
||||
let entry = VectorEntry {
|
||||
id: Some(state_id.to_string()),
|
||||
vector,
|
||||
metadata: None,
|
||||
};
|
||||
self.db.insert(entry)?;
|
||||
tracing::debug!(state_id, entity_id = %state.entity_id, "semantic index: inserted");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn search(
|
||||
&self,
|
||||
query: &str,
|
||||
k: usize,
|
||||
) -> Result<Vec<(i64, f32)>, Box<dyn std::error::Error + Send + Sync>> {
|
||||
let vector = Self::embed_str(query);
|
||||
let results = self.db.search(SearchQuery {
|
||||
vector,
|
||||
k,
|
||||
filter: None,
|
||||
ef_search: None,
|
||||
})?;
|
||||
let hits = results
|
||||
.into_iter()
|
||||
.filter_map(|r| r.id.parse::<i64>().ok().map(|id| (id, r.score)))
|
||||
.collect();
|
||||
Ok(hits)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::sync::Arc;
|
||||
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
use homecore::entity::{EntityId, State};
|
||||
use homecore::event::Context;
|
||||
|
||||
use super::*;
|
||||
use crate::db::{Recorder, SemanticIndex};
|
||||
|
||||
fn make_state(entity_id: &str, state_val: &str, attrs: serde_json::Value) -> State {
|
||||
let eid = EntityId::parse(entity_id).unwrap();
|
||||
let ctx = Context::new();
|
||||
State::new(eid, state_val, attrs, ctx)
|
||||
}
|
||||
|
||||
// ── embed_state ───────────────────────────────────────────────────────────
|
||||
|
||||
#[test]
|
||||
fn embed_state_is_deterministic() {
|
||||
let s = make_state("light.kitchen", "on", serde_json::json!({"brightness": 200}));
|
||||
let v1 = RuvectorSemanticIndex::embed_state(&s);
|
||||
let v2 = RuvectorSemanticIndex::embed_state(&s);
|
||||
assert_eq!(v1, v2, "same input must produce identical embedding");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn embed_state_is_unit_norm() {
|
||||
let s = make_state("sensor.temp", "22.5", serde_json::json!({"unit": "C"}));
|
||||
let v = RuvectorSemanticIndex::embed_state(&s);
|
||||
let norm_sq: f32 = v.iter().map(|x| x * x).sum();
|
||||
assert!(
|
||||
(norm_sq - 1.0).abs() < 1e-5,
|
||||
"embedding must be unit-norm, got norm^2={norm_sq}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn embed_state_dim_is_correct() {
|
||||
let s = make_state("binary_sensor.door", "off", serde_json::json!({}));
|
||||
let v = RuvectorSemanticIndex::embed_state(&s);
|
||||
assert_eq!(v.len(), EMBEDDING_DIM);
|
||||
}
|
||||
|
||||
// ── RuvectorSemanticIndex insert + search ─────────────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn insert_then_search_finds_state() {
|
||||
let mut idx = RuvectorSemanticIndex::new(1000).unwrap();
|
||||
let state = make_state("light.living_room", "on", serde_json::json!({"brightness": 255}));
|
||||
idx.insert_state(42, &state).await.unwrap();
|
||||
|
||||
// Query the same canonical string used by embed_state
|
||||
let query = format!(
|
||||
"{}={}|{}",
|
||||
state.entity_id, state.state, state.attributes
|
||||
);
|
||||
let hits = idx.search(&query, 5).await.unwrap();
|
||||
assert!(!hits.is_empty(), "search must return at least one hit");
|
||||
assert_eq!(hits[0].0, 42, "top hit must be the inserted state_id");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn search_ordering_closer_entity_ranks_first() {
|
||||
let mut idx = RuvectorSemanticIndex::new(1000).unwrap();
|
||||
|
||||
let s_a = make_state("light.office", "on", serde_json::json!({"brightness": 100}));
|
||||
let s_b = make_state("switch.fan", "off", serde_json::json!({}));
|
||||
|
||||
idx.insert_state(1, &s_a).await.unwrap();
|
||||
idx.insert_state(2, &s_b).await.unwrap();
|
||||
|
||||
// Query identical to s_a's canonical form → s_a must rank first
|
||||
let query_a = format!("{}={}|{}", s_a.entity_id, s_a.state, s_a.attributes);
|
||||
let hits = idx.search(&query_a, 2).await.unwrap();
|
||||
assert_eq!(hits.len(), 2);
|
||||
assert_eq!(
|
||||
hits[0].0, 1,
|
||||
"state matching the query must rank first; got {:?}",
|
||||
hits
|
||||
);
|
||||
}
|
||||
|
||||
// ── Recorder end-to-end with RuvectorSemanticIndex ────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn recorder_search_semantic_returns_recorded_state() {
|
||||
use homecore::event::StateChangedEvent;
|
||||
use chrono::Utc;
|
||||
|
||||
let idx = Arc::new(RwLock::new(
|
||||
RuvectorSemanticIndex::new(1000).unwrap(),
|
||||
));
|
||||
let semantic: Arc<RwLock<dyn SemanticIndex>> = idx;
|
||||
let recorder = Recorder::open_with_index("sqlite::memory:", semantic)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let state = Arc::new(make_state(
|
||||
"sensor.humidity",
|
||||
"65",
|
||||
serde_json::json!({"unit": "%"}),
|
||||
));
|
||||
let event = StateChangedEvent {
|
||||
entity_id: state.entity_id.clone(),
|
||||
old_state: None,
|
||||
new_state: Some(state.clone()),
|
||||
fired_at: Utc::now(),
|
||||
};
|
||||
let state_id = recorder.record_state(&event).await.unwrap().unwrap();
|
||||
|
||||
// Query using the entity prefix — close enough embedding to find it
|
||||
let query = format!("{}={}|{}", state.entity_id, state.state, state.attributes);
|
||||
let rows = recorder.search_semantic(&query, 5).await.unwrap();
|
||||
assert!(!rows.is_empty(), "search_semantic must return at least one row");
|
||||
assert_eq!(
|
||||
rows[0].state_id, state_id,
|
||||
"returned row must match the recorded state"
|
||||
);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user