Compare commits

...

5 Commits

9 changed files with 394 additions and 33 deletions

View File

@@ -1,5 +1,34 @@
# 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
- Home-Assistant-WebSocket-Listener deaktiviert den clientseitigen Keepalive-
Ping, damit stabile HA-Verbindungen nicht durch Ping-Timeouts ständig neu
aufgebaut werden
- Fallback-Auswertung läuft bei getrenntem WebSocket kurzfristig alle 5 Sekunden,
damit übernommene Aktoren nicht ohne Steuerung bleiben
## 0.7.7 - 2026-06-15
- WebSocket-State-Changes lesen jetzt das echte Home-Assistant-Eventformat
(`event.data.entity_id`), damit Kontextwechsel wie Türsensoren sofort
Vorhersagen und Schaltungen auslösen statt erst beim nächsten Statusabruf
## 0.7.6 - 2026-06-14
- Kontextvorschläge blenden zusätzlich Batterie-, Status-, Node-, Last-Seen-
und Basic-Entities aus, sofern sie nicht bewusst manuell ausgewählt wurden
## 0.7.5 - 2026-06-14
- Kontextvorschläge weiter geschärft: Standardliste zeigt nur gleiche Räume,
gemeinsame Geräte/Tokens oder echte globale Außenwerte

View File

@@ -1,5 +1,5 @@
name: SillyHome Next
version: "0.7.5"
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

View File

@@ -81,20 +81,28 @@ _MANUAL_CONTEXT_DOMAINS = frozenset({
_CONTEXT_SUGGESTION_LIMIT = 120
_OUTDOOR_TOKENS = frozenset({"aussen", "außen", "outdoor", "garten", "terrasse", "balkon"})
_DIAGNOSTIC_TOKENS = frozenset({
"basic",
"battery",
"connect",
"count",
"diagnostic",
"firmware",
"gesehen",
"last",
"linkquality",
"knoten",
"knotens",
"mqtt",
"node",
"reason",
"restart",
"rssi",
"signal",
"ssid",
"status",
"uptime",
"wifi",
"zuletzt",
})

View File

@@ -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
@@ -187,23 +189,32 @@ class BehaviorEngine:
results.append(record)
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)
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(
@@ -230,6 +241,22 @@ class BehaviorEngine:
entity_id: entities[entity_id].last_changed
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(
record.behavior.patterns,
current_context=current_context,
@@ -608,20 +635,30 @@ 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
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
affected_actuators = [
record.actuator_entity_id
@@ -633,11 +670,38 @@ class BehaviorEngine:
]
for actuator_entity_id in affected_actuators:
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:
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(
patterns: list[BehaviorPattern],
*,

View File

@@ -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.5",
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)
@@ -171,7 +179,7 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None:
if ws_status is not None:
ws_status.status = "connecting"
try:
async with websockets.connect(ws_url) as websocket:
async with websockets.connect(ws_url, ping_interval=None) as websocket:
auth_required_msg = await websocket.recv()
auth_required_data = json.loads(auth_required_msg)
if auth_required_data.get("type") != "auth_required":
@@ -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
@@ -213,12 +222,23 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None:
event = data.get("event", {})
if event.get("event_type") != "state_changed":
continue
entity_id = event.get("entity_id")
event_data = event.get("data", {})
if not isinstance(event_data, dict):
logger.warning("State-Changed-Event ohne gültige Daten empfangen")
continue
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.get("new_state"))
await asyncio.to_thread(
engine.handle_state_change,
entity_id,
new_state,
current_entities=list(state_cache.values()),
)
except json.JSONDecodeError:
logger.warning("Ungültige JSON-Nachricht von HA-WebSocket")
except Exception as exc:
@@ -244,7 +264,13 @@ async def _fallback_prediction(app: FastAPI) -> None:
Dies verhindert kompletten Ausfall der Vorhersagen bei Netzwerkproblemen.
"""
while True:
await asyncio.sleep(app.state.settings.prediction_interval_seconds)
ws_status = getattr(app.state, "ws_status", None)
websocket_connected = ws_status is not None and ws_status.status == "connected"
await asyncio.sleep(
app.state.settings.prediction_interval_seconds
if websocket_connected
else min(5, app.state.settings.prediction_interval_seconds)
)
# Nur ausführen, wenn WebSocket nicht verbunden ist
ws_status = getattr(app.state, "ws_status", None)
if ws_status is None or ws_status.status != "connected":
@@ -255,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

View File

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

View File

@@ -78,7 +78,7 @@ def _service(
model_store=str(tmp_path / "models"),
automation_store=str(tmp_path / "automations"),
actuator_store=str(tmp_path / "actuators"),
history_days=14,
history_days=31,
min_training_points=5,
retrain_stale_hours=24,
reconcile_interval_seconds=900,

View File

@@ -523,3 +523,137 @@ def test_prediction_ignores_stale_causal_context_state() -> None:
min_support=1,
window_minutes=30,
) 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
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:
@@ -55,18 +80,22 @@ def test_ha_event_listener_processes_state_change(tmp_path: Path) -> None:
'{"type":"auth_ok"}',
(
'{"type":"event","event":{"event_type":"state_changed",'
'"entity_id":"light.test","new_state":{"state":"on"}}}'
'"data":{"entity_id":"light.test","new_state":{"state":"on"}}}}'
),
asyncio.CancelledError(),
]
)
with patch("websockets.connect", return_value=fake_ws):
with patch("websockets.connect", return_value=fake_ws) as connect:
try:
await _ha_event_listener(mock_app, mock_client)
except asyncio.CancelledError:
pass
connect.assert_called_once_with(
"ws://homeassistant:8123/api/websocket",
ping_interval=None,
)
assert fake_ws.sent == [
{"type": "auth", "access_token": "test-token"},
{"id": 1, "type": "subscribe_events", "event_type": "state_changed"},
@@ -79,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