Files
sillyhome-future/app/main.py

583 lines
22 KiB
Python

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 pydantic import BaseModel, Field
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,
BackupBundle,
BehaviorPatternV2,
ControlProfile,
ControlState,
EntityState,
HandoffMode,
LearningProfile,
LearningState,
RuntimeState,
SafetyStage,
StateEvent,
)
from app.core.stores import FutureStores
stores = FutureStores(os.getenv("SILLYHOME_FUTURE_STORE", ".future_store"))
event_core = EventCore(stores)
handoff = HandoffMatrix()
ha_client: FutureHaClient | None = None
class ControlStageUpdate(BaseModel):
stage: SafetyStage
min_confidence: float = Field(default=0.82, ge=0.0, le=1.0)
manual_block: bool = False
cooldown_seconds: int = Field(default=900, ge=0)
class PatternCreateRequest(BaseModel):
trigger_entity_id: str
trigger_state: str | None = None
target_state: str
confidence: float = Field(default=0.9, ge=0.0, le=1.0)
support: int = Field(default=3, ge=1)
source: str = Field(default="dashboard", max_length=40)
@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.8",
lifespan=lifespan,
)
@app.get("/health")
def health() -> dict[str, str]:
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:]
actuator_entities = [
entity
for entity in runtime.entities.values()
if entity.domain in {"light", "switch", "fan", "cover", "humidifier"}
]
return {
"websocket_status": runtime.websocket_status,
"entity_count": len(runtime.entities),
"actuator_count": len(actuator_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, execute=_execute_ha_decision)
@app.get("/v2/runtime", response_model=RuntimeState)
def get_runtime() -> RuntimeState:
return stores.runtime()
@app.get("/v2/entities", response_model=list[EntityState])
def list_entities(domain: str | None = None, q: str | None = None) -> list[EntityState]:
entities = list(stores.runtime().entities.values())
if domain:
wanted = {item.strip() for item in domain.split(",") if item.strip()}
entities = [entity for entity in entities if entity.domain in wanted]
if q:
needle = q.casefold()
entities = [
entity
for entity in entities
if needle in entity.entity_id.casefold()
or (entity.friendly_name is not None and needle in entity.friendly_name.casefold())
]
return sorted(entities, key=lambda entity: entity.entity_id)[:500]
@app.get("/v2/learning", response_model=LearningState)
def get_learning() -> LearningState:
return stores.learning()
@app.put("/v2/learning", response_model=LearningState)
def put_learning(state: LearningState) -> LearningState:
return stores.save_learning(state)
@app.get("/v2/control", response_model=ControlState)
def get_control() -> ControlState:
return stores.control()
@app.put("/v2/control/{actuator_entity_id}", response_model=ControlProfile)
def put_control(actuator_entity_id: str, profile: ControlProfile) -> ControlProfile:
state = stores.control()
state.profiles[actuator_entity_id] = profile
stores.save_control(state)
return profile
@app.post("/v2/control/{actuator_entity_id}/stage", response_model=ControlProfile)
def update_control_stage(
actuator_entity_id: str,
update: ControlStageUpdate,
) -> ControlProfile:
state = stores.control()
profile = state.profiles.get(
actuator_entity_id,
ControlProfile(actuator_entity_id=actuator_entity_id),
)
updated = profile.model_copy(
update={
"stage": update.stage,
"min_confidence": update.min_confidence,
"manual_block": update.manual_block,
"cooldown_seconds": update.cooldown_seconds,
"handoff_mode": handoff.classify(profile),
}
)
state.profiles[actuator_entity_id] = updated
stores.save_control(state)
return updated
@app.post("/v2/learning/{actuator_entity_id}/patterns", response_model=LearningProfile)
def create_learning_pattern(
actuator_entity_id: str,
pattern: PatternCreateRequest,
) -> LearningProfile:
state = stores.learning()
profile = state.profiles.get(
actuator_entity_id,
LearningProfile(actuator_entity_id=actuator_entity_id),
)
updated = profile.model_copy(
update={
"patterns": [
*profile.patterns,
BehaviorPatternV2(
actuator_entity_id=actuator_entity_id,
target_state=pattern.target_state,
trigger_entity_id=pattern.trigger_entity_id,
trigger_state=pattern.trigger_state,
support=pattern.support,
confidence=pattern.confidence,
source=pattern.source,
),
],
"model_version": "dashboard-v1",
}
)
state.profiles[actuator_entity_id] = updated
stores.save_learning(state)
return updated
@app.post("/v2/handoff/{actuator_entity_id}/assume", response_model=ControlProfile)
def assume_control(actuator_entity_id: str) -> ControlProfile:
state = stores.control()
profile = state.profiles.get(
actuator_entity_id,
ControlProfile(actuator_entity_id=actuator_entity_id),
)
updated = handoff.assume_control(profile)
state.profiles[actuator_entity_id] = updated
stores.save_control(state)
return updated
@app.post("/v2/handoff/{actuator_entity_id}/rollback", response_model=ControlProfile)
def rollback_control(actuator_entity_id: str) -> ControlProfile:
state = stores.control()
profile = state.profiles.get(
actuator_entity_id,
ControlProfile(actuator_entity_id=actuator_entity_id),
)
updated = handoff.rollback(profile)
state.profiles[actuator_entity_id] = updated
stores.save_control(state)
return updated
@app.get("/v2/handoff/{actuator_entity_id}", response_model=HandoffMode)
def classify_handoff(actuator_entity_id: str) -> HandoffMode:
profile = stores.control().profiles.get(
actuator_entity_id,
ControlProfile(actuator_entity_id=actuator_entity_id),
)
return handoff.classify(profile)
@app.get("/v2/backup/export", response_model=BackupBundle)
def export_backup() -> BackupBundle:
return stores.export_backup()
@app.post("/v2/backup/restore")
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:
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>
:root { color-scheme: dark; --bg: #101418; --panel: #171d23; --line: #2c3640; --text: #eef3f6; --muted: #9fb0bd; --accent: #42d392; --warn: #f3c969; }
* { box-sizing: border-box; }
body { font-family: system-ui, sans-serif; margin: 0; background: var(--bg); color: var(--text); }
main { max-width: 1180px; margin: 0 auto; padding: 20px; }
h1 { font-size: 28px; margin: 0 0 18px; }
h2 { font-size: 18px; margin: 24px 0 10px; }
.grid { display: grid; grid-template-columns: repeat(auto-fit, minmax(170px, 1fr)); gap: 12px; }
.split { display: grid; grid-template-columns: minmax(260px, 1fr) minmax(320px, 1.2fr); gap: 14px; align-items: start; }
.card { border: 1px solid var(--line); border-radius: 8px; padding: 14px; background: var(--panel); }
.label { color: #9fb0bd; font-size: 13px; }
.value { font-size: 24px; margin-top: 4px; }
label { display: block; color: var(--muted); font-size: 13px; margin: 10px 0 5px; }
input, select, button { width: 100%; border: 1px solid var(--line); border-radius: 7px; background: #0c1014; color: var(--text); font: inherit; padding: 10px; }
button { cursor: pointer; background: #20303a; }
button:hover { border-color: var(--accent); }
.row { display: grid; grid-template-columns: repeat(2, minmax(0, 1fr)); gap: 10px; }
.stages { display: grid; grid-template-columns: repeat(4, 1fr); gap: 8px; margin-top: 8px; }
.stages button.active { border-color: var(--accent); color: var(--accent); }
.primary { background: #17402c; border-color: #256f4a; }
.status { min-height: 24px; color: var(--warn); margin-top: 10px; }
pre { max-height: 360px; overflow: auto; white-space: pre-wrap; word-break: break-word; background: #0c1014; padding: 14px; border-radius: 8px; }
@media (max-width: 800px) { .split, .row, .stages { grid-template-columns: 1fr; } }
</style>
</head>
<body>
<main>
<h1>SillyHome Future</h1>
<section class="grid" id="metrics"></section>
<section class="split">
<div class="card">
<h2>Control</h2>
<label for="entitySearch">Suche</label>
<input id="entitySearch" placeholder="light., switch., sensor..." autocomplete="off">
<label for="actuator">Aktor</label>
<select id="actuator"></select>
<div class="row">
<div>
<label for="minConfidence">Min. Confidence</label>
<input id="minConfidence" type="number" min="0" max="1" step="0.01" value="0.82">
</div>
<div>
<label for="cooldown">Cooldown Sekunden</label>
<input id="cooldown" type="number" min="0" step="30" value="900">
</div>
</div>
<label><input id="manualBlock" type="checkbox" style="width:auto;margin-right:6px"> Manuell blockieren</label>
<div class="stages" id="stages"></div>
<button class="primary" id="saveControl">Control speichern</button>
<div class="status" id="controlStatus"></div>
</div>
<div class="card">
<h2>Learning Pattern</h2>
<div class="row">
<div>
<label for="trigger">Trigger</label>
<select id="trigger"></select>
</div>
<div>
<label for="triggerState">Trigger State</label>
<input id="triggerState" placeholder="on, off, open...">
</div>
</div>
<div class="row">
<div>
<label for="targetState">Zielzustand</label>
<input id="targetState" value="on">
</div>
<div>
<label for="patternConfidence">Confidence</label>
<input id="patternConfidence" type="number" min="0" max="1" step="0.01" value="0.9">
</div>
</div>
<button class="primary" id="addPattern">Pattern anlegen</button>
<div class="status" id="patternStatus"></div>
<h2>Aktuelles Profil</h2>
<pre id="profile">Lade...</pre>
</div>
</section>
<h2>Audit</h2>
<pre id="audit">Lade...</pre>
</main>
<script>
const stages = ['observe', 'dry_run', 'active', 'blocked'];
const state = { entities: [], control: {}, learning: {}, selectedStage: 'observe' };
function optionText(entity) {
const name = entity.friendly_name ? ` - ${entity.friendly_name}` : '';
const current = entity.state == null ? '' : ` (${entity.state})`;
return `${entity.entity_id}${current}${name}`;
}
function setStatus(id, text) {
document.getElementById(id).textContent = text;
if (text) setTimeout(() => document.getElementById(id).textContent = '', 4000);
}
function renderStages() {
document.getElementById('stages').innerHTML = stages.map(stage =>
`<button type="button" data-stage="${stage}" class="${stage === state.selectedStage ? 'active' : ''}">${stage}</button>`
).join('');
document.querySelectorAll('[data-stage]').forEach(button => {
button.onclick = () => { state.selectedStage = button.dataset.stage; renderStages(); };
});
}
function renderEntities() {
const actuator = document.getElementById('actuator');
const trigger = document.getElementById('trigger');
const query = document.getElementById('entitySearch').value.toLowerCase();
const shown = state.entities.filter(entity =>
!query || optionText(entity).toLowerCase().includes(query)
);
const actuators = shown.filter(entity => ['light', 'switch', 'fan', 'cover', 'humidifier'].includes(entity.domain));
actuator.innerHTML = actuators.map(entity => `<option value="${entity.entity_id}">${optionText(entity)}</option>`).join('');
trigger.innerHTML = shown.map(entity => `<option value="${entity.entity_id}">${optionText(entity)}</option>`).join('');
renderProfile();
}
function renderProfile() {
const id = document.getElementById('actuator').value;
const profile = {
control: state.control.profiles?.[id] || null,
learning: state.learning.profiles?.[id] || null,
};
if (profile.control) {
state.selectedStage = profile.control.stage;
document.getElementById('minConfidence').value = profile.control.min_confidence;
document.getElementById('cooldown').value = profile.control.cooldown_seconds;
document.getElementById('manualBlock').checked = profile.control.manual_block;
renderStages();
}
document.getElementById('profile').textContent = JSON.stringify(profile, null, 2);
}
async function loadDashboard() {
const [dashResponse, entitiesResponse, controlResponse, learningResponse] = await Promise.all([
fetch('/v2/dashboard'),
fetch('/v2/entities?domain=light,switch,fan,cover,humidifier,binary_sensor,sensor'),
fetch('/v2/control'),
fetch('/v2/learning'),
]);
const data = await dashResponse.json();
state.entities = await entitiesResponse.json();
state.control = await controlResponse.json();
state.learning = await learningResponse.json();
const metrics = [
['WebSocket', data.websocket_status],
['Entities', data.entity_count],
['Actuators', data.actuator_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('');
renderEntities();
document.getElementById('audit').textContent = JSON.stringify(data.audit, null, 2);
}
document.getElementById('entitySearch').oninput = renderEntities;
document.getElementById('actuator').onchange = renderProfile;
document.getElementById('saveControl').onclick = async () => {
const id = document.getElementById('actuator').value;
const response = await fetch(`/v2/control/${encodeURIComponent(id)}/stage`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
stage: state.selectedStage,
min_confidence: Number(document.getElementById('minConfidence').value),
manual_block: document.getElementById('manualBlock').checked,
cooldown_seconds: Number(document.getElementById('cooldown').value),
}),
});
if (!response.ok) throw new Error(await response.text());
setStatus('controlStatus', 'Gespeichert');
await loadDashboard();
};
document.getElementById('addPattern').onclick = async () => {
const id = document.getElementById('actuator').value;
const response = await fetch(`/v2/learning/${encodeURIComponent(id)}/patterns`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
trigger_entity_id: document.getElementById('trigger').value,
trigger_state: document.getElementById('triggerState').value || null,
target_state: document.getElementById('targetState').value,
confidence: Number(document.getElementById('patternConfidence').value),
}),
});
if (!response.ok) throw new Error(await response.text());
setStatus('patternStatus', 'Pattern angelegt');
await loadDashboard();
};
renderStages();
loadDashboard();
setInterval(loadDashboard, 5000);
</script>
</body>
</html>"""