diff --git a/CHANGELOG.md b/CHANGELOG.md index ef0b191..d8955c9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,11 @@ # Changelog +## 0.7.10 - 2026-06-16 +- WebSocket-State-Changes aktualisieren einen internen Home-Assistant-State- + Cache und werten Aktoren direkt gegen diesen frischen Event-Zustand aus. +- Event-Auswertungen lösen keine REST-Statusabfrage mehr aus, bevor sie + aktive Aktoren schalten. + ## 0.7.9 - 2026-06-15 - Event-basierte Vorhersagen verwenden den frischen Sensorzustand direkt aus dem Home-Assistant-WebSocket-Event, damit Kontextwechsel ohne REST-Race sofort diff --git a/addon/config.yaml b/addon/config.yaml index 43084cc..c680163 100644 --- a/addon/config.yaml +++ b/addon/config.yaml @@ -1,5 +1,5 @@ name: SillyHome Next -version: "0.7.9" +version: "0.7.10" slug: sillyhome_next description: Lernt automatisch aus deinem Verhalten und steuert freigegebene Aktoren url: http://192.168.6.31:3000/pino/sillyhome-next diff --git a/app/behavior/engine.py b/app/behavior/engine.py index e7641dc..f4f15ce 100644 --- a/app/behavior/engine.py +++ b/app/behavior/engine.py @@ -1,6 +1,7 @@ from __future__ import annotations import logging +from collections.abc import Sequence from datetime import datetime, timedelta, timezone from zoneinfo import ZoneInfo @@ -18,6 +19,7 @@ from app.actuators.store import ActuatorStore from app.config import Settings from app.ha.exceptions import HaClientError from app.ha.history import LogbookEntry, StateHistoryPoint, StateHistorySeries +from app.ha.models import HaEntitySummary from app.ha.reader import HaReader _MAX_PATTERNS = 500 @@ -193,23 +195,26 @@ class BehaviorEngine: *, context_state_overrides: dict[str, str | None] | None = None, context_changed_at_overrides: dict[str, datetime | None] | None = None, + current_entities: Sequence[HaEntitySummary] | None = None, ) -> ActuatorRecord: record = self._store.get(actuator_entity_id) now = datetime.now(timezone.utc) - try: - entities = {entity.entity_id: entity for entity in self._ha_reader.read_entities()} - except HaClientError as exc: - logger.warning("Current HA state unavailable for %s: %s", actuator_entity_id, exc) - return self._save_behavior( - record, - record.behavior.model_copy( - update={ - "last_evaluated_at": now, - "prediction": None, - "reason": f"Aktueller Home-Assistant-Zustand ist nicht verfügbar: {exc}", - } - ), - ) + if current_entities is None: + try: + current_entities = self._ha_reader.read_entities() + except HaClientError as exc: + logger.warning("Current HA state unavailable for %s: %s", actuator_entity_id, exc) + return self._save_behavior( + record, + record.behavior.model_copy( + update={ + "last_evaluated_at": now, + "prediction": None, + "reason": f"Aktueller Home-Assistant-Zustand ist nicht verfügbar: {exc}", + } + ), + ) + entities = {entity.entity_id: entity for entity in current_entities} actuator = entities.get(actuator_entity_id) if actuator is None: return self._save_behavior( @@ -630,17 +635,25 @@ class BehaviorEngine: ) return self._store.upsert(updated) - def handle_state_change(self, entity_id: str, new_state: dict[str, object] | None) -> None: + def handle_state_change( + self, + entity_id: str, + new_state: dict[str, object] | None, + *, + current_entities: Sequence[HaEntitySummary] | None = None, + ) -> None: """Wird bei jedem HA-State-Change aufgerufen und löst sofortige Vorhersage aus. - Wenn entity_id ein Aktor ist: evaluate() direkt. - Wenn entity_id ein Kontext-Entity ist: alle betroffenen Aktoren evaluieren. + - Wenn current_entities gesetzt ist, kommt die Auswertung direkt aus dem + WebSocket-State-Cache statt aus einer frischen REST-Abfrage. """ # Aktor direkt evaluieren for record in self._store.list(): if record.actuator_entity_id == entity_id: try: - self.evaluate(record.actuator_entity_id) + self.evaluate(record.actuator_entity_id, current_entities=current_entities) except Exception: logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", record.actuator_entity_id) return @@ -661,6 +674,7 @@ class BehaviorEngine: actuator_entity_id, context_state_overrides={entity_id: event_state}, context_changed_at_overrides={entity_id: event_changed_at}, + current_entities=current_entities, ) except Exception: logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", actuator_entity_id) diff --git a/app/main.py b/app/main.py index 5e63f39..0169d60 100644 --- a/app/main.py +++ b/app/main.py @@ -3,6 +3,7 @@ import json import logging from contextlib import asynccontextmanager, suppress from collections.abc import AsyncIterator +from datetime import datetime, timezone from pathlib import Path from typing import cast @@ -19,6 +20,7 @@ from app.behavior.engine import BehaviorEngine from app.config import load_settings from app.core.exception_handlers import register_exception_handlers from app.ha.client import HaClient, HaClientSettings +from app.ha.models import HaEntitySummary from app.ha.reader import HaReader from app.ml.registry.model_registry import ModelRegistry from backend.routes.ml import init_ml_routes @@ -100,7 +102,7 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: app = FastAPI( title="SillyHome Next API", description="Lokales Smart-Home-Intelligenzsystem für Home Assistant.", - version="0.7.9", + version="0.7.10", lifespan=lifespan, ) app.state.settings = load_settings() @@ -155,14 +157,20 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None: """Hört auf Home-Assistant-Websocket-Events und löst sofortige Vorhersagen aus.""" settings = app.state.settings engine = app.state.behavior_engine + ha_reader = getattr(app.state, "ha_reader", None) store = app.state.actuator_store - if not isinstance(engine, BehaviorEngine) or not isinstance(store, ActuatorStore): + if ( + not isinstance(engine, BehaviorEngine) + or not isinstance(store, ActuatorStore) + or not isinstance(ha_reader, HaReader) + ): logger.error("BehaviorEngine oder ActuatorStore nicht initialisiert") ws_status = getattr(app.state, "ws_status", None) if ws_status is not None: ws_status.status = "error" ws_status.error = "BehaviorEngine oder ActuatorStore nicht initialisiert" return + state_cache: dict[str, HaEntitySummary] = {} ha_url = str(settings.ha_url).rstrip("/") ws_url = ha_url.replace("http://", "ws://").replace("https://", "wss://") + "/api/websocket" auth_token = cast(str, settings.ha_token) @@ -194,6 +202,7 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None: continue logger.info("WebSocket-Verbindung zu Home Assistant hergestellt") + state_cache = await asyncio.to_thread(_load_ha_state_cache, ha_reader) if ws_status is not None: ws_status.status = "connected" ws_status.error = None @@ -220,12 +229,15 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None: entity_id = event_data.get("entity_id") if not entity_id: continue + new_state = event_data.get("new_state") + _update_ha_state_cache(state_cache, entity_id, new_state) # Prüfe, ob Entity ein Aktor oder relevanter Kontext ist # Sofortige Vorhersage für betroffene Aktoren auslösen await asyncio.to_thread( engine.handle_state_change, entity_id, - event_data.get("new_state"), + new_state, + current_entities=list(state_cache.values()), ) except json.JSONDecodeError: logger.warning("Ungültige JSON-Nachricht von HA-WebSocket") @@ -269,3 +281,67 @@ async def _fallback_prediction(app: FastAPI) -> None: ws_status.status if ws_status else "unavailable", ) await asyncio.to_thread(engine.evaluate_all) + + +def _load_ha_state_cache(reader: HaReader) -> dict[str, HaEntitySummary]: + return {entity.entity_id: entity for entity in reader.read_entities()} + + +def _update_ha_state_cache( + state_cache: dict[str, HaEntitySummary], + entity_id: str, + new_state: object, +) -> None: + if not isinstance(new_state, dict): + state_cache.pop(entity_id, None) + return + state_cache[entity_id] = _ha_entity_from_event( + entity_id, + new_state, + state_cache.get(entity_id), + ) + + +def _ha_entity_from_event( + entity_id: str, + new_state: dict[str, object], + previous: HaEntitySummary | None, +) -> HaEntitySummary: + attributes = new_state.get("attributes") + attr = attributes if isinstance(attributes, dict) else {} + state_class = _optional_event_string(attr.get("state_class")) + device_class = _optional_event_string(attr.get("device_class")) + unit_of_measurement = _optional_event_string(attr.get("unit_of_measurement")) + friendly_name = _optional_event_string(attr.get("friendly_name")) + return HaEntitySummary( + entity_id=entity_id, + domain=entity_id.split(".", 1)[0], + state=_optional_event_string(new_state.get("state")), + last_changed=_event_datetime(new_state.get("last_changed")) + or _event_datetime(new_state.get("last_updated")), + state_class=state_class or (previous.state_class if previous else None), + device_class=device_class or (previous.device_class if previous else None), + unit_of_measurement=unit_of_measurement + or (previous.unit_of_measurement if previous else None), + friendly_name=friendly_name or (previous.friendly_name if previous else None), + area_id=previous.area_id if previous else None, + area_name=previous.area_name if previous else None, + device_id=previous.device_id if previous else None, + device_name=previous.device_name if previous else None, + ) + + +def _optional_event_string(value: object) -> str | None: + return value if isinstance(value, str) else None + + +def _event_datetime(value: object) -> datetime | None: + if not isinstance(value, str): + return None + try: + parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) + except ValueError: + return None + if parsed.tzinfo is None: + return parsed.replace(tzinfo=timezone.utc) + return parsed diff --git a/pyproject.toml b/pyproject.toml index 5c059ab..cba0e28 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "sillyhome-next" -version = "0.7.9" +version = "0.7.10" description = "Lokales Smart-Home-Intelligenzsystem für Home Assistant" requires-python = ">=3.11" dependencies = [ diff --git a/tests/behavior/test_engine.py b/tests/behavior/test_engine.py index f0dae46..7d809ef 100644 --- a/tests/behavior/test_engine.py +++ b/tests/behavior/test_engine.py @@ -587,3 +587,73 @@ def test_state_change_uses_websocket_context_state_for_immediate_action( assert reader.service_calls == [ ("light", "turn_on", {"entity_id": "light.storage"}) ] + + +def test_state_change_uses_event_cache_without_rest_state_query( + tmp_path: Path, +) -> None: + now = datetime.now(timezone.utc).replace(microsecond=0) + settings = _settings(tmp_path) + store = ActuatorStore(settings.actuator_store) + record = store.configure("light.storage") + record = record.model_copy( + update={ + "assignment": record.assignment.model_copy( + update={ + "selected_context_entity_ids": ["binary_sensor.storage_door"], + } + ), + "behavior": record.behavior.model_copy( + update={ + "mode": BehaviorMode.ACTIVE, + "status": BehaviorStatus.TRAINED, + "activation_ready": True, + "patterns": [ + BehaviorPattern( + target_state="on", + minute_of_day=60, + weekday=0, + context_states={"binary_sensor.storage_door": "on"}, + trigger_entity_id="binary_sensor.storage_door", + trigger_from_state="off", + trigger_to_state="on", + source="automation", + weight=1.0, + observed_at=now - timedelta(days=days_ago), + ) + for days_ago in (3, 2, 1) + ], + } + ), + } + ) + store.upsert(record) + reader = FakeBehaviorReader( + entities=[], + history=[], + logbook=[], + ) + + def fail_read_entities() -> list[HaEntitySummary]: + raise AssertionError("Event-Auswertung darf keinen REST-State lesen.") + + reader.read_entities = fail_read_entities # type: ignore[method-assign] + engine = BehaviorEngine(ha_reader=reader, store=store, settings=settings) + + engine.handle_state_change( + "binary_sensor.storage_door", + {"state": "on", "last_changed": now.isoformat()}, + current_entities=[ + HaEntitySummary(entity_id="light.storage", domain="light", state="off"), + HaEntitySummary( + entity_id="binary_sensor.storage_door", + domain="binary_sensor", + state="on", + last_changed=now, + ), + ], + ) + + assert reader.service_calls == [ + ("light", "turn_on", {"entity_id": "light.storage"}) + ] diff --git a/tests/test_main.py b/tests/test_main.py index b46bab4..c1c4ff0 100644 --- a/tests/test_main.py +++ b/tests/test_main.py @@ -1,4 +1,5 @@ import asyncio +from collections.abc import Sequence from pathlib import Path from unittest.mock import MagicMock, patch @@ -8,6 +9,8 @@ from fastapi.testclient import TestClient from app.actuators.store import ActuatorStore from app.behavior.engine import BehaviorEngine +from app.ha.models import HaEntitySummary +from app.ha.reader import HaReader from app.main import _ha_event_listener, app as fastapi_app, lifespan @@ -41,10 +44,32 @@ class _RecordingBehaviorEngine(BehaviorEngine): store=ActuatorStore(tmp_path / "actuators"), settings=MagicMock(), ) - self.state_changes: list[tuple[str, dict[str, object] | None]] = [] + self.state_changes: list[ + tuple[str, dict[str, object] | None, Sequence[HaEntitySummary] | None] + ] = [] - def handle_state_change(self, entity_id: str, new_state: dict[str, object] | None) -> None: - self.state_changes.append((entity_id, new_state)) + def handle_state_change( + self, + entity_id: str, + new_state: dict[str, object] | None, + *, + current_entities: Sequence[HaEntitySummary] | None = None, + ) -> None: + self.state_changes.append((entity_id, new_state, current_entities)) + + +class _FakeHaReader(HaReader): + def __init__(self) -> None: + pass + + def read_entities(self) -> list[HaEntitySummary]: + return [ + HaEntitySummary( + entity_id="light.test", + domain="light", + state="off", + ) + ] def test_ha_event_listener_processes_state_change(tmp_path: Path) -> None: @@ -83,12 +108,19 @@ def test_ha_event_listener_processes_state_change(tmp_path: Path) -> None: mock_app.state.ws_status = MagicMock() mock_engine = _RecordingBehaviorEngine(tmp_path) mock_app.state.behavior_engine = mock_engine + mock_app.state.ha_reader = _FakeHaReader() mock_store = ActuatorStore(tmp_path / "store") mock_app.state.actuator_store = mock_store mock_client = MagicMock() anyio.run(run_test) - assert mock_engine.state_changes == [("light.test", {"state": "on"})] + assert len(mock_engine.state_changes) == 1 + entity_id, new_state, current_entities = mock_engine.state_changes[0] + assert entity_id == "light.test" + assert new_state == {"state": "on"} + assert current_entities == [ + HaEntitySummary(entity_id="light.test", domain="light", state="on") + ] assert mock_app.state.ws_status.status == "connected" assert mock_app.state.ws_status.error is None