Add anomaly and performance monitoring
This commit is contained in:
@@ -9,7 +9,7 @@ from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from app.actuators.lifecycle import ActuatorReconciliationService
|
||||
from app.actuators.models import ActuatorRecord, ReconciliationState, SensorWeightGroup
|
||||
from app.actuators.models import ActuatorRecord, AnomalyEvent, ReconciliationState, SensorWeightGroup
|
||||
from app.actuators.models import JobQueueItem, JobQueueState, JobStatus, SafetyProfile
|
||||
from app.actuators.store import ActuatorStore
|
||||
from app.behavior.engine import BehaviorEngine
|
||||
@@ -89,6 +89,8 @@ class ActuatorSummary(BaseModel):
|
||||
activation_ready: bool
|
||||
activation_reason: str
|
||||
sample_count: int
|
||||
anomaly_count: int = 0
|
||||
critical_anomaly_count: int = 0
|
||||
prediction_target_state: str | None = None
|
||||
prediction_confidence: float | None = None
|
||||
updated_at: str
|
||||
@@ -108,6 +110,12 @@ class DashboardSystemStatus(BaseModel):
|
||||
configured_actuators: int = 0
|
||||
trained_models: int = 0
|
||||
review_required: int = 0
|
||||
performance_budget_ms: int = 3000
|
||||
job_p95_duration_ms: int | None = None
|
||||
slow_job_count: int = 0
|
||||
performance_status: str = "unknown"
|
||||
anomaly_count: int = 0
|
||||
critical_anomaly_count: int = 0
|
||||
|
||||
|
||||
class DashboardDiscoveryGroup(BaseModel):
|
||||
@@ -124,6 +132,12 @@ class DashboardOverview(BaseModel):
|
||||
jobs: JobQueueState = Field(default_factory=JobQueueState)
|
||||
|
||||
|
||||
class AnomalyOverview(BaseModel):
|
||||
actuator_entity_id: str
|
||||
friendly_name: str | None = None
|
||||
anomalies: list[AnomalyEvent] = Field(default_factory=list)
|
||||
|
||||
|
||||
@router.get("/discovery", response_model=list[HaEntitySummary])
|
||||
def discover_actuators(
|
||||
request: Request,
|
||||
@@ -266,6 +280,14 @@ def list_configured_summary(request: Request) -> list[ActuatorSummary]:
|
||||
activation_ready=record.behavior.activation_ready,
|
||||
activation_reason=record.behavior.activation_reason,
|
||||
sample_count=record.behavior.sample_count,
|
||||
anomaly_count=len([item for item in record.behavior.anomalies if not item.resolved]),
|
||||
critical_anomaly_count=len(
|
||||
[
|
||||
item
|
||||
for item in record.behavior.anomalies
|
||||
if not item.resolved and item.severity == "critical"
|
||||
]
|
||||
),
|
||||
prediction_target_state=(
|
||||
record.behavior.prediction.target_state
|
||||
if record.behavior.prediction is not None
|
||||
@@ -305,6 +327,9 @@ def dashboard_overview(request: Request) -> DashboardOverview:
|
||||
if isinstance(store, ActuatorStore)
|
||||
else JobQueueState()
|
||||
)
|
||||
job_p95_duration_ms, slow_job_count, performance_status = _performance_status(jobs)
|
||||
anomaly_count = sum(record.anomaly_count for record in actuators)
|
||||
critical_anomaly_count = sum(record.critical_anomaly_count for record in actuators)
|
||||
return DashboardOverview(
|
||||
system=DashboardSystemStatus(
|
||||
websocket_status=getattr(ws_status, "status", "unavailable"),
|
||||
@@ -317,6 +342,11 @@ def dashboard_overview(request: Request) -> DashboardOverview:
|
||||
configured_actuators=len(actuators),
|
||||
trained_models=reconciliation.trained_models,
|
||||
review_required=reconciliation.review_required,
|
||||
job_p95_duration_ms=job_p95_duration_ms,
|
||||
slow_job_count=slow_job_count,
|
||||
performance_status=performance_status,
|
||||
anomaly_count=anomaly_count,
|
||||
critical_anomaly_count=critical_anomaly_count,
|
||||
),
|
||||
cache=EntityCacheStatus(
|
||||
available=bool(raw_entities),
|
||||
@@ -329,6 +359,29 @@ def dashboard_overview(request: Request) -> DashboardOverview:
|
||||
)
|
||||
|
||||
|
||||
@router.get("/anomalies", response_model=list[AnomalyOverview])
|
||||
def list_anomalies(request: Request) -> list[AnomalyOverview]:
|
||||
records = _service(request).list_configured()
|
||||
entity_map = _load_cached_entity_map(
|
||||
request,
|
||||
{record.actuator_entity_id for record in records},
|
||||
)
|
||||
overview: list[AnomalyOverview] = []
|
||||
for record in records:
|
||||
active = [item for item in record.behavior.anomalies if not item.resolved]
|
||||
if not active:
|
||||
continue
|
||||
entity = entity_map.get(record.actuator_entity_id)
|
||||
overview.append(
|
||||
AnomalyOverview(
|
||||
actuator_entity_id=record.actuator_entity_id,
|
||||
friendly_name=entity.friendly_name if entity is not None else None,
|
||||
anomalies=active,
|
||||
)
|
||||
)
|
||||
return overview
|
||||
|
||||
|
||||
@router.get("", response_model=list[ActuatorRecord])
|
||||
def list_configured(request: Request) -> list[ActuatorRecord]:
|
||||
return _service(request).list_configured()
|
||||
@@ -635,6 +688,26 @@ def _finish_job(
|
||||
store.finish_job(job.job_id, status=status, summary=summary, error=error)
|
||||
|
||||
|
||||
def _performance_status(jobs: JobQueueState) -> tuple[int | None, int, str]:
|
||||
budget_ms = 3000
|
||||
durations = sorted(
|
||||
job.duration_ms
|
||||
for job in jobs.jobs
|
||||
if job.status is JobStatus.COMPLETED and job.duration_ms is not None
|
||||
)
|
||||
slow_count = sum(1 for duration in durations if duration >= budget_ms)
|
||||
if durations:
|
||||
index = min(len(durations) - 1, int(round((len(durations) - 1) * 0.95)))
|
||||
p95: int | None = durations[index]
|
||||
status_value = "slow" if slow_count else "ok"
|
||||
else:
|
||||
p95 = None
|
||||
status_value = "unknown"
|
||||
if any(job.status is JobStatus.RUNNING for job in jobs.jobs):
|
||||
status_value = "running" if status_value == "unknown" else status_value
|
||||
return p95, slow_count, status_value
|
||||
|
||||
|
||||
def _service(request: Request) -> ActuatorReconciliationService:
|
||||
service = getattr(request.app.state, "actuator_service", None)
|
||||
if not isinstance(service, ActuatorReconciliationService):
|
||||
|
||||
Reference in New Issue
Block a user