Compare commits

..

3 Commits

Author SHA1 Message Date
9ddb065f62 Speed up HA event processing
Some checks failed
quality / test (3.11) (push) Has been cancelled
quality / test (3.13) (push) Has been cancelled
2026-06-16 13:58:58 +02:00
8222f24ebe Group configured actuator overview 2026-06-16 13:51:02 +02:00
a7a2f8c78a Make SillyHome startup resilient 2026-06-16 13:43:41 +02:00
6 changed files with 112 additions and 19 deletions

View File

@@ -1,5 +1,23 @@
# Changelog
## 0.7.17 - 2026-06-16
- WebSocket-Eventpfad ist schneller: irrelevante HA-State-Changes werden vor
dem teuren State-Cache-Listenbau verworfen.
- WebSocket nutzt Keepalive und reconnectet nach Abbrüchen nach 1s statt 5s.
## 0.7.16 - 2026-06-16
- Beobachtete Aktoren werden in der Übersicht nach Raum oder Typ gruppiert und
mit Friendly Name angezeigt.
## 0.7.15 - 2026-06-16
- Add-on-Start ist robust gegen Home-Assistant-Core-502 beim Systemboot:
API und WebSocket-Listener starten trotzdem, Reconciliation/Training werden
im Hintergrund mit Retry nachgeholt.
- Periodische Reconciliation und Fallback-Auswertung beenden den Dienst nicht
mehr bei temporären HA-Fehlern.
- Add-on-Watchdog prüft `/health`, damit Supervisor den Dienst nach Absturz
wieder starten kann.
## 0.7.14 - 2026-06-16
- Onboarding-Vorschläge laden im Dashboard nachgelagert, damit Status,
Aktor-Auswahl und bestehende Geräte nicht auf Automation-Discovery warten.

View File

@@ -1,5 +1,5 @@
name: SillyHome Next
version: "0.7.14"
version: "0.7.17"
slug: sillyhome_next
description: Lernt automatisch aus deinem Verhalten und steuert freigegebene Aktoren
url: http://192.168.6.31:3000/pino/sillyhome-next
@@ -7,6 +7,7 @@ arch:
- amd64
startup: application
boot: auto
watchdog: http://[HOST]:[PORT:8000]/health
init: false
ingress: true
ingress_port: 8000

View File

@@ -43,6 +43,7 @@ class _WsStatus:
async def lifespan(app: FastAPI) -> AsyncIterator[None]:
settings = app.state.settings
client: HaClient | None = None
startup_task: asyncio.Task[None] | None = None
reconcile_task: asyncio.Task[None] | None = None
event_listener_task: asyncio.Task[None] | None = None
fallback_task: asyncio.Task[None] | None = None
@@ -74,15 +75,17 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]:
settings=settings,
)
app.state.ws_status = _WsStatus()
await asyncio.to_thread(app.state.actuator_service.reconcile_all, "startup")
await asyncio.to_thread(app.state.behavior_engine.train_all)
await asyncio.to_thread(app.state.behavior_engine.evaluate_all)
startup_task = asyncio.create_task(_startup_reconciliation(app))
reconcile_task = asyncio.create_task(_periodic_reconciliation(app))
event_listener_task = asyncio.create_task(_ha_event_listener(app, client))
fallback_task = asyncio.create_task(_fallback_prediction(app))
try:
yield
finally:
if startup_task is not None:
startup_task.cancel()
with suppress(asyncio.CancelledError):
await startup_task
if reconcile_task is not None:
reconcile_task.cancel()
with suppress(asyncio.CancelledError):
@@ -102,7 +105,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.14",
version="0.7.17",
lifespan=lifespan,
)
app.state.settings = load_settings()
@@ -147,10 +150,39 @@ async def _periodic_reconciliation(app: FastAPI) -> None:
service = getattr(app.state, "actuator_service", None)
if not isinstance(service, ActuatorReconciliationService):
continue
await asyncio.to_thread(service.reconcile_all, "scheduled")
try:
await asyncio.to_thread(service.reconcile_all, "scheduled")
engine = getattr(app.state, "behavior_engine", None)
if isinstance(engine, BehaviorEngine):
await asyncio.to_thread(engine.train_all)
except Exception:
logger.exception("Geplante Reconciliation fehlgeschlagen; nächster Lauf versucht es erneut.")
async def _startup_reconciliation(app: FastAPI) -> None:
delay_seconds = 5
while True:
service = getattr(app.state, "actuator_service", None)
engine = getattr(app.state, "behavior_engine", None)
if isinstance(engine, BehaviorEngine):
if not isinstance(service, ActuatorReconciliationService) or not isinstance(
engine,
BehaviorEngine,
):
return
try:
await asyncio.to_thread(service.reconcile_all, "startup")
await asyncio.to_thread(engine.train_all)
await asyncio.to_thread(engine.evaluate_all)
logger.info("Startup-Reconciliation erfolgreich abgeschlossen.")
return
except Exception as exc:
logger.warning(
"Startup-Reconciliation verschoben: %s. Neuer Versuch in %ss.",
exc,
delay_seconds,
)
await asyncio.sleep(delay_seconds)
delay_seconds = min(delay_seconds * 2, 60)
async def _ha_event_listener(app: FastAPI, client: HaClient) -> None:
@@ -179,7 +211,11 @@ 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, ping_interval=None) as websocket:
async with websockets.connect(
ws_url,
ping_interval=20,
ping_timeout=10,
) as websocket:
auth_required_msg = await websocket.recv()
auth_required_data = json.loads(auth_required_msg)
if auth_required_data.get("type") != "auth_required":
@@ -231,6 +267,8 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None:
continue
new_state = event_data.get("new_state")
_update_ha_state_cache(state_cache, entity_id, new_state)
if not _is_relevant_state_change(store, str(entity_id)):
continue
# Prüfe, ob Entity ein Aktor oder relevanter Kontext ist
# Sofortige Vorhersage für betroffene Aktoren auslösen
await asyncio.to_thread(
@@ -243,18 +281,22 @@ async def _ha_event_listener(app: FastAPI, client: HaClient) -> None:
logger.warning("Ungültige JSON-Nachricht von HA-WebSocket")
except Exception as exc:
logger.exception("Fehler bei Event-Verarbeitung: %s", exc)
except (websockets.exceptions.ConnectionClosed, OSError) as exc:
logger.warning("WebSocket-Verbindung unterbrochen: %s. Wiederholung in 5s...", exc)
except (
websockets.exceptions.ConnectionClosed,
websockets.exceptions.InvalidStatus,
OSError,
) as exc:
logger.warning("WebSocket-Verbindung unterbrochen: %s. Wiederholung in 1s...", exc)
if ws_status is not None:
ws_status.status = "reconnecting"
ws_status.error = str(exc)
await asyncio.sleep(5)
await asyncio.sleep(1)
except Exception as exc:
logger.exception("Unerwarteter Fehler im Event-Listener: %s", exc)
if ws_status is not None:
ws_status.status = "error"
ws_status.error = str(exc)
await asyncio.sleep(5)
await asyncio.sleep(1)
# Fallback: periodische Vorhersage falls Event-Stream ausfällt
@@ -280,7 +322,10 @@ async def _fallback_prediction(app: FastAPI) -> None:
"Fallback-Vorhersage aktiv (WebSocket-Status: %s)",
ws_status.status if ws_status else "unavailable",
)
await asyncio.to_thread(engine.evaluate_all)
try:
await asyncio.to_thread(engine.evaluate_all)
except Exception:
logger.exception("Fallback-Vorhersage fehlgeschlagen.")
def _load_ha_state_cache(reader: HaReader) -> dict[str, HaEntitySummary]:
@@ -302,6 +347,17 @@ def _update_ha_state_cache(
)
def _is_relevant_state_change(store: ActuatorStore, entity_id: str) -> bool:
for record in store.list():
if record.actuator_entity_id == entity_id:
return True
if record.assignment.selected_numeric_entity_id == entity_id:
return True
if entity_id in record.assignment.selected_context_entity_ids:
return True
return False
def _ha_entity_from_event(
entity_id: str,
new_state: dict[str, object],

View File

@@ -465,13 +465,28 @@ async function configureActuator() {
async function loadConfiguredActuators() {
const box = document.getElementById("configured-actuators");
try {
const rows = await api("v1/actuators");
const [rows, entities] = await Promise.all([
api("v1/actuators"),
api("v1/entities"),
]);
const entityMap = new Map(entities.map(entity => [entity.entity_id, entity]));
const groups = new Map();
for (const record of rows) {
const entity = entityMap.get(record.actuator_entity_id) || {};
const group = entity.area_name || actuatorGroupLabel(record.actuator_entity_id.split(".", 1)[0]);
if (!groups.has(group)) groups.set(group, []);
groups.get(group).push({record, entity});
}
const groupedRows = [...groups.entries()].sort(([left], [right]) => left.localeCompare(right));
box.innerHTML = rows.length ? `
<div class="card-list">
${rows.map(record => `
${groupedRows.map(([group, items]) => `
<h3>${escapeHtml(group)}</h3>
<div class="card-list">
${items.map(({record, entity}) => `
<article class="actuator-card ${currentActuatorId === record.actuator_entity_id ? "selected" : ""}">
<div class="card-title">
<div>
<div><strong>${escapeHtml(entity.friendly_name || record.actuator_entity_id)}</strong></div>
<div class="entity-id">${escapeHtml(record.actuator_entity_id)}</div>
<div class="${record.behavior.status === "trained" ? "ok" : "warn"}">${escapeHtml(behaviorLabel(record))}</div>
</div>
@@ -493,7 +508,8 @@ async function loadConfiguredActuators() {
</div>
</article>
`).join("")}
</div>` : "<p>Noch keine Aktoren ausgewählt.</p>";
</div>
`).join("")}` : "<p>Noch keine Aktoren ausgewählt.</p>";
} catch (error) {
box.textContent = error.message;
}

View File

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

View File

@@ -94,7 +94,8 @@ def test_ha_event_listener_processes_state_change(tmp_path: Path) -> None:
connect.assert_called_once_with(
"ws://homeassistant:8123/api/websocket",
ping_interval=None,
ping_interval=20,
ping_timeout=10,
)
assert fake_ws.sent == [
{"type": "auth", "access_token": "test-token"},
@@ -110,6 +111,7 @@ def test_ha_event_listener_processes_state_change(tmp_path: Path) -> None:
mock_app.state.behavior_engine = mock_engine
mock_app.state.ha_reader = _FakeHaReader()
mock_store = ActuatorStore(tmp_path / "store")
mock_store.configure("light.test")
mock_app.state.actuator_store = mock_store
mock_client = MagicMock()