Fix realtime HA state-change execution
Some checks failed
quality / test (3.11) (push) Has been cancelled
quality / test (3.13) (push) Has been cancelled

This commit is contained in:
2026-06-16 10:50:28 +02:00
parent 309b33b812
commit c5f42a39a9
7 changed files with 223 additions and 25 deletions

View File

@@ -1,5 +1,11 @@
# Changelog # 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 ## 0.7.9 - 2026-06-15
- Event-basierte Vorhersagen verwenden den frischen Sensorzustand direkt aus - Event-basierte Vorhersagen verwenden den frischen Sensorzustand direkt aus
dem Home-Assistant-WebSocket-Event, damit Kontextwechsel ohne REST-Race sofort dem Home-Assistant-WebSocket-Event, damit Kontextwechsel ohne REST-Race sofort

View File

@@ -1,5 +1,5 @@
name: SillyHome Next name: SillyHome Next
version: "0.7.9" version: "0.7.10"
slug: sillyhome_next slug: sillyhome_next
description: Lernt automatisch aus deinem Verhalten und steuert freigegebene Aktoren description: Lernt automatisch aus deinem Verhalten und steuert freigegebene Aktoren
url: http://192.168.6.31:3000/pino/sillyhome-next url: http://192.168.6.31:3000/pino/sillyhome-next

View File

@@ -1,6 +1,7 @@
from __future__ import annotations from __future__ import annotations
import logging import logging
from collections.abc import Sequence
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta, timezone
from zoneinfo import ZoneInfo from zoneinfo import ZoneInfo
@@ -18,6 +19,7 @@ from app.actuators.store import ActuatorStore
from app.config import Settings from app.config import Settings
from app.ha.exceptions import HaClientError from app.ha.exceptions import HaClientError
from app.ha.history import LogbookEntry, StateHistoryPoint, StateHistorySeries from app.ha.history import LogbookEntry, StateHistoryPoint, StateHistorySeries
from app.ha.models import HaEntitySummary
from app.ha.reader import HaReader from app.ha.reader import HaReader
_MAX_PATTERNS = 500 _MAX_PATTERNS = 500
@@ -193,23 +195,26 @@ class BehaviorEngine:
*, *,
context_state_overrides: dict[str, str | None] | None = None, context_state_overrides: dict[str, str | None] | None = None,
context_changed_at_overrides: dict[str, datetime | None] | None = None, context_changed_at_overrides: dict[str, datetime | None] | None = None,
current_entities: Sequence[HaEntitySummary] | None = None,
) -> ActuatorRecord: ) -> ActuatorRecord:
record = self._store.get(actuator_entity_id) record = self._store.get(actuator_entity_id)
now = datetime.now(timezone.utc) now = datetime.now(timezone.utc)
try: if current_entities is None:
entities = {entity.entity_id: entity for entity in self._ha_reader.read_entities()} try:
except HaClientError as exc: current_entities = self._ha_reader.read_entities()
logger.warning("Current HA state unavailable for %s: %s", actuator_entity_id, exc) except HaClientError as exc:
return self._save_behavior( logger.warning("Current HA state unavailable for %s: %s", actuator_entity_id, exc)
record, return self._save_behavior(
record.behavior.model_copy( record,
update={ record.behavior.model_copy(
"last_evaluated_at": now, update={
"prediction": None, "last_evaluated_at": now,
"reason": f"Aktueller Home-Assistant-Zustand ist nicht verfügbar: {exc}", "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) actuator = entities.get(actuator_entity_id)
if actuator is None: if actuator is None:
return self._save_behavior( return self._save_behavior(
@@ -630,17 +635,25 @@ class BehaviorEngine:
) )
return self._store.upsert(updated) 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. """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 Aktor ist: evaluate() direkt.
- Wenn entity_id ein Kontext-Entity ist: alle betroffenen Aktoren evaluieren. - 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 # Aktor direkt evaluieren
for record in self._store.list(): for record in self._store.list():
if record.actuator_entity_id == entity_id: if record.actuator_entity_id == entity_id:
try: try:
self.evaluate(record.actuator_entity_id) self.evaluate(record.actuator_entity_id, current_entities=current_entities)
except Exception: except Exception:
logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", record.actuator_entity_id) logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", record.actuator_entity_id)
return return
@@ -661,6 +674,7 @@ class BehaviorEngine:
actuator_entity_id, actuator_entity_id,
context_state_overrides={entity_id: event_state}, context_state_overrides={entity_id: event_state},
context_changed_at_overrides={entity_id: event_changed_at}, context_changed_at_overrides={entity_id: event_changed_at},
current_entities=current_entities,
) )
except Exception: except Exception:
logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", actuator_entity_id) logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", actuator_entity_id)

View File

@@ -3,6 +3,7 @@ import json
import logging import logging
from contextlib import asynccontextmanager, suppress from contextlib import asynccontextmanager, suppress
from collections.abc import AsyncIterator from collections.abc import AsyncIterator
from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from typing import cast from typing import cast
@@ -19,6 +20,7 @@ from app.behavior.engine import BehaviorEngine
from app.config import load_settings from app.config import load_settings
from app.core.exception_handlers import register_exception_handlers from app.core.exception_handlers import register_exception_handlers
from app.ha.client import HaClient, HaClientSettings from app.ha.client import HaClient, HaClientSettings
from app.ha.models import HaEntitySummary
from app.ha.reader import HaReader from app.ha.reader import HaReader
from app.ml.registry.model_registry import ModelRegistry from app.ml.registry.model_registry import ModelRegistry
from backend.routes.ml import init_ml_routes from backend.routes.ml import init_ml_routes
@@ -100,7 +102,7 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]:
app = FastAPI( app = FastAPI(
title="SillyHome Next API", title="SillyHome Next API",
description="Lokales Smart-Home-Intelligenzsystem für Home Assistant.", description="Lokales Smart-Home-Intelligenzsystem für Home Assistant.",
version="0.7.9", version="0.7.10",
lifespan=lifespan, lifespan=lifespan,
) )
app.state.settings = load_settings() 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.""" """Hört auf Home-Assistant-Websocket-Events und löst sofortige Vorhersagen aus."""
settings = app.state.settings settings = app.state.settings
engine = app.state.behavior_engine engine = app.state.behavior_engine
ha_reader = getattr(app.state, "ha_reader", None)
store = app.state.actuator_store 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") logger.error("BehaviorEngine oder ActuatorStore nicht initialisiert")
ws_status = getattr(app.state, "ws_status", None) ws_status = getattr(app.state, "ws_status", None)
if ws_status is not None: if ws_status is not None:
ws_status.status = "error" ws_status.status = "error"
ws_status.error = "BehaviorEngine oder ActuatorStore nicht initialisiert" ws_status.error = "BehaviorEngine oder ActuatorStore nicht initialisiert"
return return
state_cache: dict[str, HaEntitySummary] = {}
ha_url = str(settings.ha_url).rstrip("/") ha_url = str(settings.ha_url).rstrip("/")
ws_url = ha_url.replace("http://", "ws://").replace("https://", "wss://") + "/api/websocket" ws_url = ha_url.replace("http://", "ws://").replace("https://", "wss://") + "/api/websocket"
auth_token = cast(str, settings.ha_token) auth_token = cast(str, settings.ha_token)
@@ -194,6 +202,7 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None:
continue continue
logger.info("WebSocket-Verbindung zu Home Assistant hergestellt") 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: if ws_status is not None:
ws_status.status = "connected" ws_status.status = "connected"
ws_status.error = None 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") entity_id = event_data.get("entity_id")
if not entity_id: if not entity_id:
continue 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 # Prüfe, ob Entity ein Aktor oder relevanter Kontext ist
# Sofortige Vorhersage für betroffene Aktoren auslösen # Sofortige Vorhersage für betroffene Aktoren auslösen
await asyncio.to_thread( await asyncio.to_thread(
engine.handle_state_change, engine.handle_state_change,
entity_id, entity_id,
event_data.get("new_state"), new_state,
current_entities=list(state_cache.values()),
) )
except json.JSONDecodeError: except json.JSONDecodeError:
logger.warning("Ungültige JSON-Nachricht von HA-WebSocket") 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", ws_status.status if ws_status else "unavailable",
) )
await asyncio.to_thread(engine.evaluate_all) 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

View File

@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "sillyhome-next" name = "sillyhome-next"
version = "0.7.9" version = "0.7.10"
description = "Lokales Smart-Home-Intelligenzsystem für Home Assistant" description = "Lokales Smart-Home-Intelligenzsystem für Home Assistant"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = [ dependencies = [

View File

@@ -587,3 +587,73 @@ def test_state_change_uses_websocket_context_state_for_immediate_action(
assert reader.service_calls == [ assert reader.service_calls == [
("light", "turn_on", {"entity_id": "light.storage"}) ("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"})
]

View File

@@ -1,4 +1,5 @@
import asyncio import asyncio
from collections.abc import Sequence
from pathlib import Path from pathlib import Path
from unittest.mock import MagicMock, patch from unittest.mock import MagicMock, patch
@@ -8,6 +9,8 @@ from fastapi.testclient import TestClient
from app.actuators.store import ActuatorStore from app.actuators.store import ActuatorStore
from app.behavior.engine import BehaviorEngine 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 from app.main import _ha_event_listener, app as fastapi_app, lifespan
@@ -41,10 +44,32 @@ class _RecordingBehaviorEngine(BehaviorEngine):
store=ActuatorStore(tmp_path / "actuators"), store=ActuatorStore(tmp_path / "actuators"),
settings=MagicMock(), 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: def handle_state_change(
self.state_changes.append((entity_id, new_state)) 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: 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_app.state.ws_status = MagicMock()
mock_engine = _RecordingBehaviorEngine(tmp_path) mock_engine = _RecordingBehaviorEngine(tmp_path)
mock_app.state.behavior_engine = mock_engine mock_app.state.behavior_engine = mock_engine
mock_app.state.ha_reader = _FakeHaReader()
mock_store = ActuatorStore(tmp_path / "store") mock_store = ActuatorStore(tmp_path / "store")
mock_app.state.actuator_store = mock_store mock_app.state.actuator_store = mock_store
mock_client = MagicMock() mock_client = MagicMock()
anyio.run(run_test) 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.status == "connected"
assert mock_app.state.ws_status.error is None assert mock_app.state.ws_status.error is None