Compare commits
6 Commits
v2.0.0-alp
...
v2.0.0-alp
| Author | SHA1 | Date | |
|---|---|---|---|
| 937c79a19b | |||
| 149f11a18e | |||
| 38a85fce74 | |||
| b77992606a | |||
| 9d01af98b2 | |||
| 1336bc36fa |
13
README.md
13
README.md
@@ -22,6 +22,18 @@ pip install -e ".[dev]"
|
|||||||
uvicorn app.main:app --reload
|
uvicorn app.main:app --reload
|
||||||
```
|
```
|
||||||
|
|
||||||
|
## Home Assistant Add-on
|
||||||
|
|
||||||
|
Add this repository in Home Assistant:
|
||||||
|
|
||||||
|
```text
|
||||||
|
http://192.168.6.31:3000/Otto/sillyhome-future
|
||||||
|
```
|
||||||
|
|
||||||
|
Then install **SillyHome Future**. The add-on is intentionally independent from
|
||||||
|
`sillyhome-next` and starts its own v2 event-core API on port `8099` behind
|
||||||
|
Ingress.
|
||||||
|
|
||||||
## Test
|
## Test
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
@@ -29,4 +41,3 @@ pytest
|
|||||||
ruff check .
|
ruff check .
|
||||||
mypy app tests
|
mypy app tests
|
||||||
```
|
```
|
||||||
|
|
||||||
|
|||||||
13
addon/Dockerfile
Normal file
13
addon/Dockerfile
Normal file
@@ -0,0 +1,13 @@
|
|||||||
|
FROM python:3.13-slim
|
||||||
|
|
||||||
|
WORKDIR /app
|
||||||
|
|
||||||
|
ARG SILLYHOME_FUTURE_REF=v2.0.0-alpha.7
|
||||||
|
RUN python -m pip install --no-cache-dir \
|
||||||
|
"http://192.168.6.31:3000/Otto/sillyhome-future/archive/${SILLYHOME_FUTURE_REF}.tar.gz"
|
||||||
|
|
||||||
|
COPY run.sh /run.sh
|
||||||
|
RUN chmod 0755 /run.sh \
|
||||||
|
&& useradd --create-home --uid 10001 sillyhome
|
||||||
|
|
||||||
|
ENTRYPOINT ["/run.sh"]
|
||||||
28
addon/config.yaml
Normal file
28
addon/config.yaml
Normal file
@@ -0,0 +1,28 @@
|
|||||||
|
name: SillyHome Future
|
||||||
|
version: "2.0.0-alpha.7"
|
||||||
|
slug: sillyhome_future
|
||||||
|
description: Event-first SillyHome v2 test controller
|
||||||
|
url: http://192.168.6.31:3000/Otto/sillyhome-future
|
||||||
|
arch:
|
||||||
|
- amd64
|
||||||
|
startup: application
|
||||||
|
boot: manual
|
||||||
|
watchdog: http://[HOST]:[PORT:8099]/health
|
||||||
|
init: false
|
||||||
|
ingress: true
|
||||||
|
ingress_port: 8099
|
||||||
|
panel_title: SillyHome Future
|
||||||
|
panel_icon: mdi:home-lightning-bolt
|
||||||
|
panel_admin: true
|
||||||
|
homeassistant_api: true
|
||||||
|
hassio_api: false
|
||||||
|
auth_api: false
|
||||||
|
map:
|
||||||
|
- type: addon_config
|
||||||
|
read_only: false
|
||||||
|
options:
|
||||||
|
store_path: /data/future_store
|
||||||
|
log_level: info
|
||||||
|
schema:
|
||||||
|
store_path: str
|
||||||
|
log_level: list(debug|info|warning|error)
|
||||||
31
addon/run.sh
Normal file
31
addon/run.sh
Normal file
@@ -0,0 +1,31 @@
|
|||||||
|
#!/bin/sh
|
||||||
|
set -eu
|
||||||
|
|
||||||
|
export PYTHONUNBUFFERED=1
|
||||||
|
|
||||||
|
STORE_PATH="/data/future_store"
|
||||||
|
LOG_LEVEL="info"
|
||||||
|
if [ -f /data/options.json ]; then
|
||||||
|
eval "$(python - <<'PY'
|
||||||
|
import json
|
||||||
|
import shlex
|
||||||
|
|
||||||
|
with open("/data/options.json", encoding="utf-8") as stream:
|
||||||
|
options = json.load(stream)
|
||||||
|
for name, env_name in {
|
||||||
|
"store_path": "STORE_PATH",
|
||||||
|
"log_level": "LOG_LEVEL",
|
||||||
|
}.items():
|
||||||
|
if name in options:
|
||||||
|
print(f"{env_name}={shlex.quote(str(options[name]))}")
|
||||||
|
PY
|
||||||
|
)"
|
||||||
|
fi
|
||||||
|
|
||||||
|
export SILLYHOME_FUTURE_STORE="$STORE_PATH"
|
||||||
|
mkdir -p "$SILLYHOME_FUTURE_STORE"
|
||||||
|
chown -R sillyhome:sillyhome "$SILLYHOME_FUTURE_STORE"
|
||||||
|
|
||||||
|
exec su -s /bin/sh sillyhome -c \
|
||||||
|
"uvicorn app.main:app --host 0.0.0.0 --port 8099 --log-level ${LOG_LEVEL}"
|
||||||
|
|
||||||
@@ -1,12 +1,14 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
|
from collections.abc import Callable
|
||||||
|
|
||||||
from app.core.decision import DecisionEngineV2
|
from app.core.decision import DecisionEngineV2
|
||||||
from app.core.handoff import HandoffMatrix
|
from app.core.handoff import HandoffMatrix
|
||||||
from app.core.models import (
|
from app.core.models import (
|
||||||
AuditEvent,
|
AuditEvent,
|
||||||
ControlProfile,
|
ControlProfile,
|
||||||
|
Decision,
|
||||||
EntityState,
|
EntityState,
|
||||||
LearningProfile,
|
LearningProfile,
|
||||||
StateEvent,
|
StateEvent,
|
||||||
@@ -31,7 +33,12 @@ class EventCore:
|
|||||||
self._decision = DecisionEngineV2()
|
self._decision = DecisionEngineV2()
|
||||||
self._handoff = HandoffMatrix()
|
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()
|
runtime = self._stores.runtime()
|
||||||
learning = self._stores.learning()
|
learning = self._stores.learning()
|
||||||
control = self._stores.control()
|
control = self._stores.control()
|
||||||
@@ -69,6 +76,8 @@ class EventCore:
|
|||||||
learning=learning_profile,
|
learning=learning_profile,
|
||||||
control=control_profile,
|
control=control_profile,
|
||||||
)
|
)
|
||||||
|
if callable(execute) and decision.allowed and not decision.dry_run:
|
||||||
|
decision = _execute_decision(decision, execute)
|
||||||
audit.append(
|
audit.append(
|
||||||
AuditEvent(
|
AuditEvent(
|
||||||
event_id=_event_id("decision"),
|
event_id=_event_id("decision"),
|
||||||
@@ -87,3 +96,17 @@ class EventCore:
|
|||||||
def _event_id(prefix: str) -> str:
|
def _event_id(prefix: str) -> str:
|
||||||
return f"{prefix}-{datetime.now(timezone.utc).strftime('%Y%m%d%H%M%S%f')}"
|
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)})
|
||||||
|
|||||||
163
app/core/ha_client.py
Normal file
163
app/core/ha_client.py
Normal file
@@ -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
|
||||||
|
|
||||||
222
app/main.py
222
app/main.py
@@ -1,16 +1,22 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
import os
|
import os
|
||||||
|
from collections.abc import AsyncIterator
|
||||||
|
from contextlib import asynccontextmanager, suppress
|
||||||
|
|
||||||
from fastapi import FastAPI
|
from fastapi import FastAPI
|
||||||
|
from fastapi.responses import HTMLResponse
|
||||||
|
|
||||||
from app.core.event_core import EventCore
|
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.handoff import HandoffMatrix
|
||||||
from app.core.models import (
|
from app.core.models import (
|
||||||
AuditEvent,
|
AuditEvent,
|
||||||
BackupBundle,
|
BackupBundle,
|
||||||
ControlProfile,
|
ControlProfile,
|
||||||
ControlState,
|
ControlState,
|
||||||
|
EntityState,
|
||||||
HandoffMode,
|
HandoffMode,
|
||||||
LearningState,
|
LearningState,
|
||||||
RuntimeState,
|
RuntimeState,
|
||||||
@@ -18,24 +24,72 @@ from app.core.models import (
|
|||||||
)
|
)
|
||||||
from app.core.stores import FutureStores
|
from app.core.stores import FutureStores
|
||||||
|
|
||||||
app = FastAPI(
|
|
||||||
title="SillyHome Future API",
|
|
||||||
description="SillyHome v2 event-core side project.",
|
|
||||||
version="2.0.0-alpha.1",
|
|
||||||
)
|
|
||||||
stores = FutureStores(os.getenv("SILLYHOME_FUTURE_STORE", ".future_store"))
|
stores = FutureStores(os.getenv("SILLYHOME_FUTURE_STORE", ".future_store"))
|
||||||
event_core = EventCore(stores)
|
event_core = EventCore(stores)
|
||||||
handoff = HandoffMatrix()
|
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.7",
|
||||||
|
lifespan=lifespan,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@app.get("/health")
|
@app.get("/health")
|
||||||
def health() -> dict[str, str]:
|
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])
|
@app.post("/v2/events/state", response_model=list[AuditEvent])
|
||||||
def ingest_state_event(event: StateEvent) -> 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)
|
@app.get("/v2/runtime", response_model=RuntimeState)
|
||||||
@@ -111,3 +165,157 @@ def restore_backup(bundle: BackupBundle) -> dict[str, str]:
|
|||||||
stores.restore_backup(bundle)
|
stores.restore_backup(bundle)
|
||||||
return {"status": "restored"}
|
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:
|
||||||
|
routed_triggers: set[str] = set()
|
||||||
|
next_route_refresh = 0.0
|
||||||
|
pending_unrouted: list[StateEvent] = []
|
||||||
|
last_unrouted_flush = 0.0
|
||||||
|
while True:
|
||||||
|
await asyncio.to_thread(_set_websocket_status, "connecting")
|
||||||
|
try:
|
||||||
|
async for event in client.listen_state_events():
|
||||||
|
loop_time = asyncio.get_running_loop().time()
|
||||||
|
if loop_time >= next_route_refresh:
|
||||||
|
routed_triggers = await asyncio.to_thread(_routed_trigger_ids)
|
||||||
|
next_route_refresh = loop_time + 5
|
||||||
|
await asyncio.to_thread(_set_websocket_status, "connected")
|
||||||
|
if event.entity_id in routed_triggers:
|
||||||
|
if pending_unrouted:
|
||||||
|
await asyncio.to_thread(_merge_unrouted_events, pending_unrouted)
|
||||||
|
pending_unrouted = []
|
||||||
|
await asyncio.to_thread(
|
||||||
|
event_core.process_state_event,
|
||||||
|
event,
|
||||||
|
execute=_execute_ha_decision,
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
pending_unrouted.append(event)
|
||||||
|
if len(pending_unrouted) >= 100 or loop_time - last_unrouted_flush >= 2:
|
||||||
|
await asyncio.to_thread(_merge_unrouted_events, pending_unrouted)
|
||||||
|
pending_unrouted = []
|
||||||
|
last_unrouted_flush = loop_time
|
||||||
|
await asyncio.sleep(0)
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise
|
||||||
|
except Exception as exc:
|
||||||
|
await asyncio.to_thread(_set_websocket_status, f"reconnecting: {exc}")
|
||||||
|
await asyncio.sleep(2)
|
||||||
|
|
||||||
|
|
||||||
|
def _set_websocket_status(status: str) -> None:
|
||||||
|
runtime = stores.runtime()
|
||||||
|
if runtime.websocket_status == status:
|
||||||
|
return
|
||||||
|
runtime.websocket_status = status
|
||||||
|
stores.save_runtime(runtime)
|
||||||
|
|
||||||
|
|
||||||
|
def _routed_trigger_ids() -> set[str]:
|
||||||
|
result: set[str] = set()
|
||||||
|
for profile in stores.learning().profiles.values():
|
||||||
|
for pattern in profile.patterns:
|
||||||
|
if pattern.trigger_entity_id is not None:
|
||||||
|
result.add(pattern.trigger_entity_id)
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
def _merge_unrouted_events(events: list[StateEvent]) -> None:
|
||||||
|
if not events:
|
||||||
|
return
|
||||||
|
runtime = stores.runtime()
|
||||||
|
for event in events:
|
||||||
|
runtime.entities[event.entity_id] = EntityState(
|
||||||
|
entity_id=event.entity_id,
|
||||||
|
domain=event.entity_id.split(".", 1)[0],
|
||||||
|
state=event.new_state,
|
||||||
|
changed_at=event.changed_at,
|
||||||
|
area_name=event.attributes.get("area_name"),
|
||||||
|
device_id=event.attributes.get("device_id"),
|
||||||
|
friendly_name=event.attributes.get("friendly_name"),
|
||||||
|
)
|
||||||
|
stores.save_runtime(runtime)
|
||||||
|
|
||||||
|
|
||||||
|
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 """<!doctype html>
|
||||||
|
<html lang="de">
|
||||||
|
<head>
|
||||||
|
<meta charset="utf-8">
|
||||||
|
<meta name="viewport" content="width=device-width, initial-scale=1">
|
||||||
|
<title>SillyHome Future</title>
|
||||||
|
<style>
|
||||||
|
body { font-family: system-ui, sans-serif; margin: 0; background: #101418; color: #eef3f6; }
|
||||||
|
main { max-width: 1080px; margin: 0 auto; padding: 24px; }
|
||||||
|
h1 { font-size: 28px; margin: 0 0 18px; }
|
||||||
|
.grid { display: grid; grid-template-columns: repeat(auto-fit, minmax(180px, 1fr)); gap: 12px; }
|
||||||
|
.card { border: 1px solid #2c3640; border-radius: 8px; padding: 14px; background: #171d23; }
|
||||||
|
.label { color: #9fb0bd; font-size: 13px; }
|
||||||
|
.value { font-size: 24px; margin-top: 4px; }
|
||||||
|
pre { white-space: pre-wrap; word-break: break-word; background: #0c1014; padding: 14px; border-radius: 8px; }
|
||||||
|
</style>
|
||||||
|
</head>
|
||||||
|
<body>
|
||||||
|
<main>
|
||||||
|
<h1>SillyHome Future</h1>
|
||||||
|
<section class="grid" id="metrics"></section>
|
||||||
|
<h2>Audit</h2>
|
||||||
|
<pre id="audit">Lade...</pre>
|
||||||
|
</main>
|
||||||
|
<script>
|
||||||
|
async function loadDashboard() {
|
||||||
|
const response = await fetch('/v2/dashboard');
|
||||||
|
const data = await response.json();
|
||||||
|
const metrics = [
|
||||||
|
['WebSocket', data.websocket_status],
|
||||||
|
['Entities', data.entity_count],
|
||||||
|
['Learning', data.learning_profiles],
|
||||||
|
['Control', data.control_profiles],
|
||||||
|
['Rooms', data.rooms.length],
|
||||||
|
['Scenes', data.scenes.length],
|
||||||
|
];
|
||||||
|
document.getElementById('metrics').innerHTML = metrics.map(([label, value]) =>
|
||||||
|
`<div class="card"><div class="label">${label}</div><div class="value">${value}</div></div>`
|
||||||
|
).join('');
|
||||||
|
document.getElementById('audit').textContent = JSON.stringify(data.audit, null, 2);
|
||||||
|
}
|
||||||
|
loadDashboard();
|
||||||
|
setInterval(loadDashboard, 5000);
|
||||||
|
</script>
|
||||||
|
</body>
|
||||||
|
</html>"""
|
||||||
|
|||||||
@@ -27,3 +27,9 @@ state, but switching decisions must remain inspectable and testable.
|
|||||||
6. Dashboard v2
|
6. Dashboard v2
|
||||||
7. Add-on hardening
|
7. Add-on hardening
|
||||||
|
|
||||||
|
## Add-on Contract
|
||||||
|
|
||||||
|
The Home Assistant add-on is owned by this repository and does not copy code
|
||||||
|
from `sillyhome-next`. It starts the v2 API with `SILLYHOME_FUTURE_STORE`
|
||||||
|
pointing at `/data/future_store`, so runtime, learning and control stores are
|
||||||
|
kept inside the add-on data volume.
|
||||||
|
|||||||
@@ -4,13 +4,16 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "sillyhome-future"
|
name = "sillyhome-future"
|
||||||
version = "2.0.0-alpha.1"
|
version = "2.0.0-alpha.7"
|
||||||
description = "SillyHome v2 event-core prototype"
|
description = "SillyHome v2 event-core prototype"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"fastapi>=0.115",
|
"fastapi>=0.115",
|
||||||
|
"httpx>=0.27",
|
||||||
"pydantic>=2.8",
|
"pydantic>=2.8",
|
||||||
"uvicorn[standard]>=0.30",
|
"uvicorn[standard]>=0.30",
|
||||||
|
"websockets>=12.0",
|
||||||
|
"pyyaml>=6.0",
|
||||||
]
|
]
|
||||||
|
|
||||||
[project.optional-dependencies]
|
[project.optional-dependencies]
|
||||||
@@ -32,4 +35,3 @@ target-version = "py311"
|
|||||||
python_version = "3.11"
|
python_version = "3.11"
|
||||||
strict = true
|
strict = true
|
||||||
plugins = []
|
plugins = []
|
||||||
|
|
||||||
|
|||||||
4
repository.yaml
Normal file
4
repository.yaml
Normal file
@@ -0,0 +1,4 @@
|
|||||||
|
name: SillyHome Future Add-ons
|
||||||
|
url: http://192.168.6.31:3000/Otto/sillyhome-future
|
||||||
|
maintainer: Otto
|
||||||
|
|
||||||
21
tests/test_addon_config.py
Normal file
21
tests/test_addon_config.py
Normal file
@@ -0,0 +1,21 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import yaml # type: ignore[import-untyped]
|
||||||
|
|
||||||
|
|
||||||
|
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.7"
|
||||||
|
assert config["ingress"] is True
|
||||||
|
assert config["ingress_port"] == 8099
|
||||||
|
assert config["homeassistant_api"] is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_repository_points_to_gitea_repo() -> None:
|
||||||
|
repository = yaml.safe_load(Path("repository.yaml").read_text(encoding="utf-8"))
|
||||||
|
|
||||||
|
assert repository["url"] == "http://192.168.6.31:3000/Otto/sillyhome-future"
|
||||||
@@ -15,8 +15,7 @@ def test_health_and_backup_roundtrip(tmp_path, monkeypatch) -> None: # type: ig
|
|||||||
|
|
||||||
assert stores is not None
|
assert stores is not None
|
||||||
assert health.status_code == 200
|
assert health.status_code == 200
|
||||||
assert health.json()["version"] == "2.0.0-alpha.1"
|
assert health.json()["version"] == "2.0.0-alpha.7"
|
||||||
assert backup.status_code == 200
|
assert backup.status_code == 200
|
||||||
assert restore.status_code == 200
|
assert restore.status_code == 200
|
||||||
assert restore.json() == {"status": "restored"}
|
assert restore.json() == {"status": "restored"}
|
||||||
|
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ from app.core.handoff import HandoffMatrix
|
|||||||
from app.core.models import (
|
from app.core.models import (
|
||||||
BehaviorPatternV2,
|
BehaviorPatternV2,
|
||||||
ControlProfile,
|
ControlProfile,
|
||||||
|
Decision,
|
||||||
EntityState,
|
EntityState,
|
||||||
HandoffMode,
|
HandoffMode,
|
||||||
LearningProfile,
|
LearningProfile,
|
||||||
@@ -91,3 +92,64 @@ def test_handoff_matrix_detects_conflict_and_rollback() -> None:
|
|||||||
rolled_back = matrix.rollback(controlled)
|
rolled_back = matrix.rollback(controlled)
|
||||||
assert rolled_back.handoff_mode is HandoffMode.ROLLBACK
|
assert rolled_back.handoff_mode is HandoffMode.ROLLBACK
|
||||||
assert rolled_back.paused_automation_ids == []
|
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
|
||||||
|
|||||||
12
tests/test_ha_client.py
Normal file
12
tests/test_ha_client.py
Normal file
@@ -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
|
||||||
|
|
||||||
Reference in New Issue
Block a user