Compare commits

...

2 Commits

Author SHA1 Message Date
c5f42a39a9 Fix realtime HA state-change execution 2026-06-16 10:50:28 +02:00
309b33b812 Use fresh HA event state for behavior triggers 2026-06-15 19:37:45 +02:00
7 changed files with 346 additions and 27 deletions

View File

@@ -1,5 +1,18 @@
# 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
- Event-basierte Vorhersagen verwenden den frischen Sensorzustand direkt aus
dem Home-Assistant-WebSocket-Event, damit Kontextwechsel ohne REST-Race sofort
bewertet und geschaltet werden können
- Regressionstest stellt sicher, dass ein Türsensor-Event trotz veraltetem
HA-Snapshot direkt `light.turn_on` auslöst
## 0.7.8 - 2026-06-15 ## 0.7.8 - 2026-06-15
- Home-Assistant-WebSocket-Listener deaktiviert den clientseitigen Keepalive- - Home-Assistant-WebSocket-Listener deaktiviert den clientseitigen Keepalive-
Ping, damit stabile HA-Verbindungen nicht durch Ping-Timeouts ständig neu Ping, damit stabile HA-Verbindungen nicht durch Ping-Timeouts ständig neu

View File

@@ -1,5 +1,5 @@
name: SillyHome Next name: SillyHome Next
version: "0.7.8" 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
@@ -187,11 +189,19 @@ class BehaviorEngine:
results.append(record) results.append(record)
return results return results
def evaluate(self, actuator_entity_id: str) -> ActuatorRecord: def evaluate(
self,
actuator_entity_id: str,
*,
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) record = self._store.get(actuator_entity_id)
now = datetime.now(timezone.utc) now = datetime.now(timezone.utc)
if current_entities is None:
try: try:
entities = {entity.entity_id: entity for entity in self._ha_reader.read_entities()} current_entities = self._ha_reader.read_entities()
except HaClientError as exc: except HaClientError as exc:
logger.warning("Current HA state unavailable for %s: %s", actuator_entity_id, exc) logger.warning("Current HA state unavailable for %s: %s", actuator_entity_id, exc)
return self._save_behavior( return self._save_behavior(
@@ -204,6 +214,7 @@ class BehaviorEngine:
} }
), ),
) )
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(
@@ -230,6 +241,22 @@ class BehaviorEngine:
entity_id: entities[entity_id].last_changed entity_id: entities[entity_id].last_changed
for entity_id in current_context for entity_id in current_context
} }
selected_context_ids = {
entity_id
for entity_id in (
[
record.assignment.selected_numeric_entity_id,
*record.assignment.selected_context_entity_ids,
]
)
if entity_id
}
for entity_id, state in (context_state_overrides or {}).items():
if entity_id in selected_context_ids and state is not None:
current_context[entity_id] = state
for entity_id, changed_at in (context_changed_at_overrides or {}).items():
if entity_id in current_context:
current_context_changed_at[entity_id] = changed_at or now
prediction = predict_behavior( prediction = predict_behavior(
record.behavior.patterns, record.behavior.patterns,
current_context=current_context, current_context=current_context,
@@ -608,20 +635,30 @@ 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
event_state = _event_state(new_state)
event_changed_at = _event_changed_at(new_state) or datetime.now(timezone.utc)
# Kontext-Entity: alle Aktoren finden, die diesen Kontext nutzen # Kontext-Entity: alle Aktoren finden, die diesen Kontext nutzen
affected_actuators = [ affected_actuators = [
record.actuator_entity_id record.actuator_entity_id
@@ -633,11 +670,38 @@ class BehaviorEngine:
] ]
for actuator_entity_id in affected_actuators: for actuator_entity_id in affected_actuators:
try: try:
self.evaluate(actuator_entity_id) self.evaluate(
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: 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)
def _event_state(new_state: dict[str, object] | None) -> str | None:
if not isinstance(new_state, dict):
return None
state = new_state.get("state")
return state if isinstance(state, str) else None
def _event_changed_at(new_state: dict[str, object] | None) -> datetime | None:
if not isinstance(new_state, dict):
return None
value = new_state.get("last_changed") or new_state.get("last_updated")
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
def predict_behavior( def predict_behavior(
patterns: list[BehaviorPattern], patterns: list[BehaviorPattern],
*, *,

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.8", 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.8" 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

@@ -523,3 +523,137 @@ def test_prediction_ignores_stale_causal_context_state() -> None:
min_support=1, min_support=1,
window_minutes=30, window_minutes=30,
) is None ) is None
def test_state_change_uses_websocket_context_state_for_immediate_action(
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=[
HaEntitySummary(entity_id="light.storage", domain="light", state="off"),
HaEntitySummary(
entity_id="binary_sensor.storage_door",
domain="binary_sensor",
state="off",
last_changed=now - timedelta(minutes=5),
),
],
history=[],
logbook=[],
)
engine = BehaviorEngine(ha_reader=reader, store=store, settings=settings)
engine.handle_state_change(
"binary_sensor.storage_door",
{"state": "on", "last_changed": now.isoformat()},
)
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"})
]

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