diff --git a/addon/config.yaml b/addon/config.yaml index 1e6d0ab..07a2585 100644 --- a/addon/config.yaml +++ b/addon/config.yaml @@ -1,5 +1,5 @@ name: SillyHome Future -version: "2.0.0-alpha.3" +version: "2.0.0-alpha.4" slug: sillyhome_future description: Event-first SillyHome v2 test controller url: http://192.168.6.31:3000/Otto/sillyhome-future diff --git a/app/core/event_core.py b/app/core/event_core.py index 14e6801..324c767 100644 --- a/app/core/event_core.py +++ b/app/core/event_core.py @@ -1,12 +1,14 @@ from __future__ import annotations from datetime import datetime, timezone +from collections.abc import Callable from app.core.decision import DecisionEngineV2 from app.core.handoff import HandoffMatrix from app.core.models import ( AuditEvent, ControlProfile, + Decision, EntityState, LearningProfile, StateEvent, @@ -31,7 +33,12 @@ class EventCore: self._decision = DecisionEngineV2() self._handoff = HandoffMatrix() - def process_state_event(self, event: StateEvent) -> list[AuditEvent]: + def process_state_event( + self, + event: StateEvent, + *, + execute: Callable[[Decision], bool] | None = None, + ) -> list[AuditEvent]: runtime = self._stores.runtime() learning = self._stores.learning() control = self._stores.control() @@ -69,6 +76,8 @@ class EventCore: learning=learning_profile, control=control_profile, ) + if callable(execute) and decision.allowed and not decision.dry_run: + decision = _execute_decision(decision, execute) audit.append( AuditEvent( event_id=_event_id("decision"), @@ -87,3 +96,17 @@ class EventCore: def _event_id(prefix: str) -> str: return f"{prefix}-{datetime.now(timezone.utc).strftime('%Y%m%d%H%M%S%f')}" + +def _execute_decision(decision: Decision, execute: Callable[[Decision], bool]) -> Decision: + try: + result = execute(decision) + except Exception as exc: + return decision.model_copy( + update={ + "executed": False, + "allowed": False, + "reason": f"Ausfuehrung fehlgeschlagen: {exc}", + "blockers": [*decision.blockers, str(exc)], + } + ) + return decision.model_copy(update={"executed": bool(result)}) diff --git a/app/core/ha_client.py b/app/core/ha_client.py new file mode 100644 index 0000000..c3be304 --- /dev/null +++ b/app/core/ha_client.py @@ -0,0 +1,163 @@ +from __future__ import annotations + +import json +from collections.abc import AsyncIterator +from dataclasses import dataclass +from datetime import datetime, timezone + +import httpx +import websockets + +from app.core.models import EntityState, StateEvent + + +@dataclass(frozen=True) +class HaClientConfig: + core_url: str + token: str + websocket_url: str | None = None + timeout_seconds: float = 15.0 + + +class FutureHaClient: + def __init__(self, config: HaClientConfig) -> None: + self._config = config + self._core_url = config.core_url.rstrip("/") + + def read_states(self) -> list[EntityState]: + with httpx.Client(timeout=self._config.timeout_seconds) as client: + response = client.get( + f"{self._core_url}/api/states", + headers=self._headers, + ) + response.raise_for_status() + payload = response.json() + result: list[EntityState] = [] + for item in payload if isinstance(payload, list) else []: + if not isinstance(item, dict): + continue + entity_id = item.get("entity_id") + if not isinstance(entity_id, str) or "." not in entity_id: + continue + attributes = item.get("attributes") + attr = attributes if isinstance(attributes, dict) else {} + result.append( + EntityState( + entity_id=entity_id, + domain=entity_id.split(".", 1)[0], + state=item.get("state") if isinstance(item.get("state"), str) else None, + changed_at=_parse_datetime(item.get("last_changed")), + area_name=attr.get("area_name") if isinstance(attr.get("area_name"), str) else None, + device_id=attr.get("device_id") if isinstance(attr.get("device_id"), str) else None, + friendly_name=( + attr.get("friendly_name") + if isinstance(attr.get("friendly_name"), str) + else None + ), + ) + ) + return result + + def call_service( + self, + domain: str, + service: str, + service_data: dict[str, object], + ) -> None: + with httpx.Client(timeout=self._config.timeout_seconds) as client: + response = client.post( + f"{self._core_url}/api/services/{domain}/{service}", + headers=self._headers, + json=service_data, + ) + response.raise_for_status() + + async def listen_state_events(self) -> AsyncIterator[StateEvent]: + websocket_url = self._config.websocket_url or _default_websocket_url(self._core_url) + async with websockets.connect(websocket_url, ping_interval=None) as websocket: + auth_required = json.loads(await websocket.recv()) + if auth_required.get("type") != "auth_required": + raise RuntimeError("Home Assistant websocket did not request authentication") + await websocket.send(json.dumps({"type": "auth", "access_token": self._config.token})) + auth_result = json.loads(await websocket.recv()) + if auth_result.get("type") != "auth_ok": + raise RuntimeError("Home Assistant websocket authentication failed") + await websocket.send( + json.dumps({"id": 1, "type": "subscribe_events", "event_type": "state_changed"}) + ) + async for raw_message in websocket: + data = json.loads(raw_message) + if data.get("type") != "event": + continue + event = data.get("event") + if not isinstance(event, dict) or event.get("event_type") != "state_changed": + continue + event_data = event.get("data") + if not isinstance(event_data, dict): + continue + entity_id = event_data.get("entity_id") + new_state = event_data.get("new_state") + if not isinstance(entity_id, str) or not isinstance(new_state, dict): + continue + state = new_state.get("state") + attributes = new_state.get("attributes") + attr = attributes if isinstance(attributes, dict) else {} + yield StateEvent( + entity_id=entity_id, + new_state=state if isinstance(state, str) else None, + changed_at=_parse_datetime( + new_state.get("last_changed") or new_state.get("last_updated") + ), + attributes={ + key: value + for key, value in { + "area_name": attr.get("area_name"), + "device_id": attr.get("device_id"), + "friendly_name": attr.get("friendly_name"), + }.items() + if isinstance(value, str) + }, + ) + + @property + def _headers(self) -> dict[str, str]: + return { + "Authorization": f"Bearer {self._config.token}", + "Content-Type": "application/json", + } + + +def service_for_state(domain: str, target_state: str | None) -> tuple[str, str] | None: + if target_state is None: + return None + if domain in {"light", "switch", "fan", "humidifier"}: + if target_state == "on": + return domain, "turn_on" + if target_state == "off": + return domain, "turn_off" + if domain == "cover": + if target_state in {"open", "opening"}: + return domain, "open_cover" + if target_state in {"closed", "closing"}: + return domain, "close_cover" + return None + + +def _default_websocket_url(core_url: str) -> str: + stripped = core_url.rstrip("/") + if stripped.endswith("/core"): + return stripped.replace("http://", "ws://").replace("https://", "wss://") + "/websocket" + return stripped.replace("http://", "ws://").replace("https://", "wss://") + "/api/websocket" + + +def _parse_datetime(value: object) -> datetime: + if not isinstance(value, str): + return datetime.now(timezone.utc) + try: + parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) + except ValueError: + return datetime.now(timezone.utc) + if parsed.tzinfo is None: + return parsed.replace(tzinfo=timezone.utc) + return parsed + diff --git a/app/main.py b/app/main.py index 9097bae..20a45ea 100644 --- a/app/main.py +++ b/app/main.py @@ -1,10 +1,15 @@ from __future__ import annotations +import asyncio import os +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager, suppress from fastapi import FastAPI +from fastapi.responses import HTMLResponse from app.core.event_core import EventCore +from app.core.ha_client import FutureHaClient, HaClientConfig, service_for_state from app.core.handoff import HandoffMatrix from app.core.models import ( AuditEvent, @@ -18,24 +23,72 @@ from app.core.models import ( ) from app.core.stores import FutureStores -app = FastAPI( - title="SillyHome Future API", - description="SillyHome v2 event-core side project.", - version="2.0.0-alpha.3", -) stores = FutureStores(os.getenv("SILLYHOME_FUTURE_STORE", ".future_store")) event_core = EventCore(stores) handoff = HandoffMatrix() +ha_client: FutureHaClient | None = None + + +@asynccontextmanager +async def lifespan(app: FastAPI) -> AsyncIterator[None]: + global ha_client + listener_task: asyncio.Task[None] | None = None + ha_client = _ha_client_from_env() + if ha_client is not None: + _load_initial_ha_states(ha_client) + listener_task = asyncio.create_task(_ha_listener_loop(ha_client)) + try: + yield + finally: + if listener_task is not None: + listener_task.cancel() + with suppress(asyncio.CancelledError): + await listener_task + + +app = FastAPI( + title="SillyHome Future API", + description="SillyHome v2 event-core side project.", + version="2.0.0-alpha.4", + lifespan=lifespan, +) @app.get("/health") def health() -> dict[str, str]: - return {"status": "ok", "version": app.version} + runtime = stores.runtime() + return { + "status": "ok", + "version": app.version, + "ha": runtime.websocket_status, + } + + +@app.get("/", response_class=HTMLResponse) +def dashboard() -> str: + return _dashboard_html() + + +@app.get("/v2/dashboard") +def dashboard_data() -> dict[str, object]: + runtime = stores.runtime() + learning = stores.learning() + control = stores.control() + latest_audit = runtime.audit[-20:] + return { + "websocket_status": runtime.websocket_status, + "entity_count": len(runtime.entities), + "learning_profiles": len(learning.profiles), + "control_profiles": len(control.profiles), + "rooms": list(learning.rooms.values()), + "scenes": list(learning.scenes.values()), + "audit": latest_audit, + } @app.post("/v2/events/state", response_model=list[AuditEvent]) def ingest_state_event(event: StateEvent) -> list[AuditEvent]: - return event_core.process_state_event(event) + return event_core.process_state_event(event, execute=_execute_ha_decision) @app.get("/v2/runtime", response_model=RuntimeState) @@ -110,3 +163,107 @@ def export_backup() -> BackupBundle: def restore_backup(bundle: BackupBundle) -> dict[str, str]: stores.restore_backup(bundle) return {"status": "restored"} + + +def _ha_client_from_env() -> FutureHaClient | None: + token = os.getenv("SUPERVISOR_TOKEN") or os.getenv("SILLYHOME_FUTURE_HA_TOKEN") + if not token: + return None + core_url = os.getenv("SILLYHOME_FUTURE_HA_CORE_URL", "http://supervisor/core") + websocket_url = os.getenv("SILLYHOME_FUTURE_HA_WS_URL") + return FutureHaClient( + HaClientConfig(core_url=core_url, token=token, websocket_url=websocket_url) + ) + + +def _load_initial_ha_states(client: FutureHaClient) -> None: + runtime = stores.runtime() + try: + for entity in client.read_states(): + runtime.entities[entity.entity_id] = entity + runtime.websocket_status = "initial_state_loaded" + except Exception as exc: + runtime.websocket_status = f"initial_state_error: {exc}" + stores.save_runtime(runtime) + + +async def _ha_listener_loop(client: FutureHaClient) -> None: + while True: + runtime = stores.runtime() + runtime.websocket_status = "connecting" + stores.save_runtime(runtime) + try: + async for event in client.listen_state_events(): + runtime = stores.runtime() + runtime.websocket_status = "connected" + stores.save_runtime(runtime) + event_core.process_state_event(event, execute=_execute_ha_decision) + except asyncio.CancelledError: + raise + except Exception as exc: + runtime = stores.runtime() + runtime.websocket_status = f"reconnecting: {exc}" + stores.save_runtime(runtime) + await asyncio.sleep(2) + + +def _execute_ha_decision(decision: object) -> bool: + if ha_client is None or not hasattr(decision, "actuator_entity_id"): + return False + actuator_entity_id = str(decision.actuator_entity_id) + target_state = getattr(decision, "target_state", None) + service = service_for_state(actuator_entity_id.split(".", 1)[0], target_state) + if service is None: + return False + domain, service_name = service + ha_client.call_service(domain, service_name, {"entity_id": actuator_entity_id}) + return True + + +def _dashboard_html() -> str: + return """ + + + + + SillyHome Future + + + +
+

SillyHome Future

+
+

Audit

+
Lade...
+
+ + +""" diff --git a/pyproject.toml b/pyproject.toml index 452d716..d90d5fd 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,13 +4,15 @@ build-backend = "setuptools.build_meta" [project] name = "sillyhome-future" -version = "2.0.0-alpha.3" +version = "2.0.0-alpha.4" description = "SillyHome v2 event-core prototype" requires-python = ">=3.11" dependencies = [ "fastapi>=0.115", + "httpx>=0.27", "pydantic>=2.8", "uvicorn[standard]>=0.30", + "websockets>=12.0", "pyyaml>=6.0", ] diff --git a/tests/test_addon_config.py b/tests/test_addon_config.py index 1822c36..5842144 100644 --- a/tests/test_addon_config.py +++ b/tests/test_addon_config.py @@ -9,7 +9,7 @@ def test_addon_config_declares_future_addon() -> None: config = yaml.safe_load(Path("addon/config.yaml").read_text(encoding="utf-8")) assert config["slug"] == "sillyhome_future" - assert config["version"] == "2.0.0-alpha.3" + assert config["version"] == "2.0.0-alpha.4" assert config["ingress"] is True assert config["ingress_port"] == 8099 assert config["homeassistant_api"] is True diff --git a/tests/test_future_api.py b/tests/test_future_api.py index a08db30..3f1f729 100644 --- a/tests/test_future_api.py +++ b/tests/test_future_api.py @@ -15,7 +15,7 @@ def test_health_and_backup_roundtrip(tmp_path, monkeypatch) -> None: # type: ig assert stores is not None assert health.status_code == 200 - assert health.json()["version"] == "2.0.0-alpha.3" + assert health.json()["version"] == "2.0.0-alpha.4" assert backup.status_code == 200 assert restore.status_code == 200 assert restore.json() == {"status": "restored"} diff --git a/tests/test_future_core.py b/tests/test_future_core.py index afdffc0..ef22cf7 100644 --- a/tests/test_future_core.py +++ b/tests/test_future_core.py @@ -7,6 +7,7 @@ from app.core.handoff import HandoffMatrix from app.core.models import ( BehaviorPatternV2, ControlProfile, + Decision, EntityState, HandoffMode, LearningProfile, @@ -91,3 +92,64 @@ def test_handoff_matrix_detects_conflict_and_rollback() -> None: rolled_back = matrix.rollback(controlled) assert rolled_back.handoff_mode is HandoffMode.ROLLBACK assert rolled_back.paused_automation_ids == [] + + +def test_event_core_executes_allowed_active_decision(tmp_path: Path) -> None: + stores = FutureStores(tmp_path) + stores.save_runtime( + RuntimeState( + entities={ + "light.storage": EntityState( + entity_id="light.storage", + domain="light", + state="off", + ) + } + ) + ) + stores.save_learning( + LearningState( + profiles={ + "light.storage": LearningProfile( + actuator_entity_id="light.storage", + patterns=[ + BehaviorPatternV2( + actuator_entity_id="light.storage", + target_state="on", + trigger_entity_id="binary_sensor.storage_door", + trigger_state="on", + support=3, + confidence=0.95, + ) + ], + ) + } + ) + ) + stores.save_control( + stores.control().model_copy( + update={ + "profiles": { + "light.storage": ControlProfile( + actuator_entity_id="light.storage", + stage=SafetyStage.ACTIVE, + min_confidence=0.8, + ) + } + } + ) + ) + calls: list[str] = [] + + def execute(decision: Decision) -> bool: + calls.append(decision.actuator_entity_id) + return True + + audit = EventCore(stores).process_state_event( + StateEvent(entity_id="binary_sensor.storage_door", new_state="on"), + execute=execute, + ) + + decision = [item.decision for item in audit if item.decision is not None][0] + assert calls == ["light.storage"] + assert decision.executed is True diff --git a/tests/test_ha_client.py b/tests/test_ha_client.py new file mode 100644 index 0000000..38ec66b --- /dev/null +++ b/tests/test_ha_client.py @@ -0,0 +1,12 @@ +from __future__ import annotations + +from app.core.ha_client import service_for_state + + +def test_service_for_state_maps_safe_domains() -> None: + assert service_for_state("light", "on") == ("light", "turn_on") + assert service_for_state("switch", "off") == ("switch", "turn_off") + assert service_for_state("cover", "open") == ("cover", "open_cover") + assert service_for_state("cover", "closed") == ("cover", "close_cover") + assert service_for_state("lock", "unlocked") is None +