Add safety dashboard and decision transparency
Some checks failed
quality / test (3.11) (push) Has been cancelled
quality / test (3.13) (push) Has been cancelled

This commit is contained in:
2026-06-17 18:26:49 +02:00
parent ca253d1e6c
commit 0101596e93
12 changed files with 792 additions and 23 deletions

View File

@@ -10,6 +10,7 @@ from pydantic import BaseModel, Field
from app.actuators.lifecycle import ActuatorReconciliationService
from app.actuators.models import ActuatorRecord, ReconciliationState, SensorWeightGroup
from app.actuators.models import JobQueueItem, JobQueueState, JobStatus, SafetyProfile
from app.actuators.store import ActuatorStore
from app.behavior.engine import BehaviorEngine
from app.config import Settings
@@ -55,6 +56,10 @@ class FeedbackRequest(BaseModel):
expected_state: str | None = Field(default=None, max_length=100)
class SafetyProfileRequest(BaseModel):
safety: SafetyProfile
class ActuatorSuggestion(BaseModel):
entity_id: str
domain: str
@@ -112,6 +117,7 @@ class DashboardOverview(BaseModel):
cache: EntityCacheStatus
actuators: list[ActuatorSummary]
discovery_groups: list[DashboardDiscoveryGroup]
jobs: JobQueueState = Field(default_factory=JobQueueState)
@router.get("/discovery", response_model=list[HaEntitySummary])
@@ -124,8 +130,19 @@ def discover_actuators(
if cached_entities:
entities = {entity.entity_id: entity for entity in cached_entities}
else:
fresh_entities = list(ha_reader.read_entities())
_save_cached_entities(request, fresh_entities)
job = _start_job(
request,
kind="discovery",
trigger="manual" if refresh else "cache-miss",
summary="Home-Assistant-Entities werden gelesen und klassifiziert.",
)
try:
fresh_entities = list(ha_reader.read_entities())
_save_cached_entities(request, fresh_entities)
except Exception as exc:
_finish_job(job, request, status=JobStatus.FAILED, summary="Discovery fehlgeschlagen.", error=str(exc))
raise
_finish_job(job, request, status=JobStatus.COMPLETED, summary=f"{len(fresh_entities)} Entities klassifiziert.")
entities = {entity.entity_id: entity for entity in fresh_entities}
discovered = discover_entities(list(entities.values()))
actuator_ids = _deduplicate_actuator_ids(
@@ -278,6 +295,12 @@ def dashboard_overview(request: Request) -> DashboardOverview:
reconciliation = _reconciliation_state_or_default(request)
ws_status = getattr(request.app.state, "ws_status", None)
actuators = list_configured_summary(request)
store = getattr(request.app.state, "actuator_store", None)
jobs = (
store.load_job_queue()
if isinstance(store, ActuatorStore)
else JobQueueState()
)
return DashboardOverview(
system=DashboardSystemStatus(
websocket_status=getattr(ws_status, "status", "unavailable"),
@@ -298,6 +321,7 @@ def dashboard_overview(request: Request) -> DashboardOverview:
),
actuators=actuators,
discovery_groups=cached_groups,
jobs=jobs,
)
@@ -372,6 +396,18 @@ def record_feedback(
raise HTTPException(status_code=404, detail=str(exc)) from exc
@router.post("/{actuator_entity_id}/safety", response_model=ActuatorRecord)
def set_safety_profile(
actuator_entity_id: str,
payload: SafetyProfileRequest,
request: Request,
) -> ActuatorRecord:
try:
return _behavior(request).set_safety_profile(actuator_entity_id, profile=payload.safety)
except KeyError as exc:
raise HTTPException(status_code=404, detail=str(exc)) from exc
@router.post("/{actuator_entity_id}/activation", response_model=ActuatorRecord)
def set_activation(
actuator_entity_id: str,
@@ -441,11 +477,27 @@ def refresh_related_automations(
actuator_entity_id: str,
request: Request,
) -> ActuatorRecord:
job = _start_job(
request,
kind="automation_refresh",
trigger="manual",
target=actuator_entity_id,
summary="Passende HA-Automationen werden gesucht.",
)
try:
return _behavior(request).refresh_related_automations(actuator_entity_id)
record = _behavior(request).refresh_related_automations(actuator_entity_id)
_finish_job(
job,
request,
status=JobStatus.COMPLETED,
summary=f"{len(record.behavior.related_automations)} Automationen gefunden.",
)
return record
except KeyError as exc:
_finish_job(job, request, status=JobStatus.FAILED, summary="Automation-Refresh fehlgeschlagen.", error=str(exc))
raise HTTPException(status_code=404, detail=str(exc)) from exc
except (ValueError, HaClientError) as exc:
_finish_job(job, request, status=JobStatus.FAILED, summary="Automation-Refresh fehlgeschlagen.", error=str(exc))
raise HTTPException(status_code=409, detail=str(exc)) from exc
@@ -486,12 +538,85 @@ def run_reconciliation(
request: Request,
trigger: str = Query(default="manual", pattern=r"^[a-z0-9_-]{1,32}$"),
) -> ReconciliationState:
state = _service(request).reconcile_all(trigger=trigger)
_behavior(request).train_all()
_behavior(request).evaluate_all()
reconciliation_job = _start_job(
request,
kind="reconciliation",
trigger=trigger,
summary="Kontext, Zuordnung und Modelle werden abgeglichen.",
)
training_job: JobQueueItem | None = None
evaluation_job: JobQueueItem | None = None
try:
state = _service(request).reconcile_all(trigger=trigger)
_finish_job(reconciliation_job, request, status=JobStatus.COMPLETED, summary=state.last_summary)
reconciliation_job = None
training_job = _start_job(
request,
kind="training",
trigger=trigger,
summary="Gelernte Aktorhandlungen werden aktualisiert.",
)
_behavior(request).train_all()
_finish_job(training_job, request, status=JobStatus.COMPLETED, summary="Training abgeschlossen.")
training_job = None
evaluation_job = _start_job(
request,
kind="evaluation",
trigger=trigger,
summary="Aktuelle Vorhersagen werden neu berechnet.",
)
_behavior(request).evaluate_all()
_finish_job(evaluation_job, request, status=JobStatus.COMPLETED, summary="Evaluation abgeschlossen.")
evaluation_job = None
except Exception as exc:
for job in [reconciliation_job, training_job, evaluation_job]:
if isinstance(job, JobQueueItem) and job.status is JobStatus.RUNNING:
_finish_job(job, request, status=JobStatus.FAILED, summary="Job fehlgeschlagen.", error=str(exc))
raise
return state
@router.get("/job-queue/state", response_model=JobQueueState)
def get_job_queue(request: Request) -> JobQueueState:
store = getattr(request.app.state, "actuator_store", None)
if not isinstance(store, ActuatorStore):
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Actuator Store nicht initialisiert.",
)
return store.load_job_queue()
def _start_job(
request: Request,
*,
kind: str,
trigger: str,
target: str | None = None,
summary: str = "",
) -> JobQueueItem | None:
store = getattr(request.app.state, "actuator_store", None)
if not isinstance(store, ActuatorStore):
return None
return store.start_job(kind=kind, trigger=trigger, target=target, summary=summary)
def _finish_job(
job: JobQueueItem | None,
request: Request,
*,
status: JobStatus,
summary: str,
error: str | None = None,
) -> None:
if job is None:
return
store = getattr(request.app.state, "actuator_store", None)
if not isinstance(store, ActuatorStore):
return
store.finish_job(job.job_id, status=status, summary=summary, error=error)
def _service(request: Request) -> ActuatorReconciliationService:
service = getattr(request.app.state, "actuator_service", None)
if not isinstance(service, ActuatorReconciliationService):