from __future__ import annotations import logging from itertools import product from collections.abc import Sequence from datetime import datetime, timedelta, timezone from time import perf_counter from zoneinfo import ZoneInfo from app.actuators.models import ( ActuatorRecord, AdaptiveWeightUpdate, AgentInsight, AnomalyEvent, ActuatorGroup, AutomationConflict, BehaviorMode, BehaviorPattern, BehaviorPrediction, BehaviorState, BehaviorStatus, DecisionFactor, DecisionTrace, ExecutionEvent, FeedbackKind, LatencyMeasurement, ManualOverride, ModelSnapshot, RelatedAutomation, SafetyProfile, SafetyStage, SceneSuggestion, SimulationOutcome, TimeProfile, ) from app.actuators.store import ActuatorStore from app.config import Settings from app.ha.exceptions import HaClientError from app.ha.history import LogbookEntry, StateHistoryPoint, StateHistorySeries from app.ha.models import HaEntitySummary from app.ha.reader import HaReader _MAX_PATTERNS = 500 _MAX_MODEL_SNAPSHOTS = 3 _MAX_SNAPSHOT_PATTERNS = 120 _MAX_EXECUTION_EVENTS = 100 _MAX_DECISION_TRACES = 30 _MAX_LATENCY_MEASUREMENTS = 50 _MAX_FEEDBACK_LOG = 50 _ACTION_LOGBOOK_TOLERANCE = timedelta(seconds=10) _CONTEXT_TRIGGER_TOLERANCE = timedelta(seconds=3) _OWN_ACTION_TOLERANCE = timedelta(seconds=20) _SAFE_ACTIVE_DOMAINS = frozenset({"cover", "fan", "humidifier", "light", "switch"}) _AUTOMATION_CONTEXT_DOMAINS = frozenset({"automation", "script"}) logger = logging.getLogger(__name__) class BehaviorEngine: def __init__( self, *, ha_reader: HaReader, store: ActuatorStore, settings: Settings, ) -> None: self._ha_reader = ha_reader self._store = store self._settings = settings def train_all(self) -> list[ActuatorRecord]: results: list[ActuatorRecord] = [] for record in self._store.list(): try: results.append(self.train(record.actuator_entity_id)) except Exception: logger.exception("Behavior training failed for %s", record.actuator_entity_id) results.append(record) return results def train(self, actuator_entity_id: str) -> ActuatorRecord: record = self._store.get(actuator_entity_id) now = datetime.now(timezone.utc) raw_context_ids = list( dict.fromkeys( [ record.assignment.selected_numeric_entity_id, *record.assignment.selected_context_entity_ids, ] ) ) context_ids = [ entity_id for entity_id in raw_context_ids if isinstance(entity_id, str) ] if not context_ids: return self._save_behavior( record, record.behavior.model_copy( update={ "status": BehaviorStatus.COLLECTING, "activation_ready": False, "activation_reason": ( "Freigabe gesperrt: Noch kein geeigneter Kontext erkannt." ), "last_trained_at": now, "reason": "Noch kein geeigneter Kontext für Verhaltenslernen vorhanden.", "anomalies": _detect_anomalies( record, now=now, min_behavior_actions=self._settings.min_behavior_actions, stale_hours=self._settings.retrain_stale_hours, sample_count=0, trusted_actions=0, prediction=None, safety_blockers=[], ), } ), ) start = now - timedelta(days=self._settings.history_days) history_ids = [actuator_entity_id, *context_ids] try: history = { series.entity_id: series for series in self._ha_reader.read_state_history(history_ids, start, now) } except (HaClientError, ValueError) as exc: logger.warning("Behavior history unavailable for %s: %s", actuator_entity_id, exc) return self._save_behavior( record, record.behavior.model_copy( update={ "status": BehaviorStatus.BLOCKED, "last_trained_at": now, "reason": f"Home-Assistant-Historie konnte nicht gelesen werden: {exc}", } ), ) actuator_history = history.get(actuator_entity_id) if actuator_history is None or len(actuator_history.points) < 2: return self._save_behavior( record, record.behavior.model_copy( update={ "status": BehaviorStatus.COLLECTING, "sample_count": 0, "high_confidence_sample_count": 0, "activation_ready": False, "activation_reason": ( "Freigabe gesperrt: Noch keine historischen " "Aktorhandlungen gefunden." ), "patterns": [], "last_trained_at": now, "reason": "Noch keine historischen Aktorhandlungen gefunden.", "anomalies": _detect_anomalies( record, now=now, min_behavior_actions=self._settings.min_behavior_actions, stale_hours=self._settings.retrain_stale_hours, sample_count=0, trusted_actions=0, prediction=None, safety_blockers=[], ), } ), ) try: logbook = list(self._ha_reader.read_logbook(actuator_entity_id, start, now)) except (HaClientError, ValueError) as exc: logger.warning("Logbook unavailable for %s: %s", actuator_entity_id, exc) logbook = [] patterns = self._build_patterns( actuator_history=actuator_history, context_history=history, context_ids=context_ids, logbook=logbook, own_executions=record.behavior.execution_events, ) trusted_actions = sum( 1 for pattern in patterns if pattern.source in {"user", "automation"} ) status = ( BehaviorStatus.TRAINED if len(patterns) >= self._settings.min_behavior_actions else BehaviorStatus.COLLECTING ) reason = ( f"{len(patterns)} Handlungen mit automatisch erfasstem Kontext gelernt." if status is BehaviorStatus.TRAINED else ( f"{len(patterns)} von mindestens {self._settings.min_behavior_actions} " "benötigten Handlungen gelernt." ) ) activation_ready = ( status is BehaviorStatus.TRAINED and trusted_actions >= self._settings.min_behavior_actions ) activation_reason = ( "Freigabe bereit: Genügend eindeutig zugeordnete Handlungen gelernt." if activation_ready else ( "Freigabe gesperrt: " f"{max(0, self._settings.min_behavior_actions - trusted_actions)} " "eindeutig zugeordnete Handlungen fehlen." ) ) model_version_id = f"model-{now.strftime('%Y%m%d%H%M%S')}" behavior = record.behavior.model_copy( update={ "status": status, "sample_count": len(patterns), "high_confidence_sample_count": trusted_actions, "activation_ready": activation_ready, "activation_reason": activation_reason, "patterns": patterns[-_MAX_PATTERNS:], "last_trained_at": now, "reason": reason, "sample_trend": [*record.behavior.sample_trend, len(patterns)][-30:], "knowledge": _knowledge_lines(record, len(patterns), trusted_actions), "assumptions": _assumption_lines(record), "uncertainties": _uncertainty_lines(record, len(patterns), trusted_actions), "time_profiles": _time_profiles(patterns), "model_snapshots": _next_model_snapshots( record.behavior.model_snapshots, model_version_id, patterns[-_MAX_PATTERNS:], len(patterns), trusted_actions, _average(record.behavior.confidence_trend), record.behavior.incorrect_feedback_count, reason, ), "active_model_version": model_version_id, "anomalies": _detect_anomalies( record, now=now, min_behavior_actions=self._settings.min_behavior_actions, stale_hours=self._settings.retrain_stale_hours, sample_count=len(patterns), trusted_actions=trusted_actions, prediction=record.behavior.prediction, safety_blockers=record.behavior.safety_blockers, ), } ) return self._save_behavior(record, behavior) def evaluate_all(self) -> list[ActuatorRecord]: results: list[ActuatorRecord] = [] for record in self._store.list(): try: results.append(self.evaluate(record.actuator_entity_id)) except Exception: logger.exception("Behavior evaluation failed for %s", record.actuator_entity_id) results.append(record) return results def evaluate( self, actuator_entity_id: str, *, context_state_overrides: dict[str, str | None] | None = None, context_changed_at_overrides: dict[str, datetime | None] | None = None, current_entities: Sequence[HaEntitySummary] | None = None, trigger_entity_id: str | None = None, trigger_state: str | None = None, event_received_at: datetime | None = None, ) -> ActuatorRecord: started_perf = perf_counter() record = self._store.get(actuator_entity_id) now = datetime.now(timezone.utc) if current_entities is None: try: current_entities = self._ha_reader.read_entities() except HaClientError as exc: logger.warning("Current HA state unavailable for %s: %s", actuator_entity_id, exc) return self._save_behavior( record, record.behavior.model_copy( update={ "last_evaluated_at": now, "prediction": None, "reason": f"Aktueller Home-Assistant-Zustand ist nicht verfügbar: {exc}", } ), ) entities = {entity.entity_id: entity for entity in current_entities} actuator = entities.get(actuator_entity_id) if actuator is None: return self._save_behavior( record, record.behavior.model_copy( update={ "last_evaluated_at": now, "prediction": None, "reason": "Aktor ist aktuell nicht in Home Assistant verfügbar.", } ), ) current_context = { entity_id: entities[entity_id].state for entity_id in ( [ record.assignment.selected_numeric_entity_id, *record.assignment.selected_context_entity_ids, ] ) if entity_id and entity_id in entities and entities[entity_id].state is not None } current_context_changed_at = { entity_id: entities[entity_id].last_changed for entity_id in current_context } selected_context_ids = { entity_id for entity_id in ( [ record.assignment.selected_numeric_entity_id, *record.assignment.selected_context_entity_ids, ] ) if entity_id } for entity_id, state in (context_state_overrides or {}).items(): if entity_id in selected_context_ids and state is not None: current_context[entity_id] = state for entity_id, changed_at in (context_changed_at_overrides or {}).items(): if entity_id in current_context: current_context_changed_at[entity_id] = changed_at or now prediction = predict_behavior( record.behavior.patterns, current_context=current_context, current_context_changed_at=current_context_changed_at, context_weights=_context_weights_for(record), now=now, min_support=self._settings.min_behavior_actions, window_minutes=self._settings.prediction_window_minutes, causal_window_seconds=self._settings.prediction_interval_seconds * 2, timezone_name=self._settings.timezone, ) if prediction is not None: safety_allowed, safety_blockers = self._assess_safety( record, actuator.state, prediction, now, ) prediction = prediction.model_copy( update={ "execution_reason": ( "Ausführung ist freigegeben." if safety_allowed else "Nicht ausgeführt: " + " ".join(safety_blockers) ) } ) else: safety_allowed = False safety_blockers = ["Keine fällige Vorhersage."] decision_to_service_ms: int | None = None decision_factors = _decision_factors_for(record, current_context, prediction) behavior = record.behavior.model_copy( update={ "last_evaluated_at": now, "prediction": prediction, "reason": ( prediction.reason if prediction is not None else "Aktuell ist kein gelerntes Handlungsmuster fällig." ), "decision_factors": decision_factors, "knowledge": _knowledge_lines(record, record.behavior.sample_count, record.behavior.high_confidence_sample_count), "assumptions": _assumption_lines(record), "uncertainties": _uncertainty_lines(record, record.behavior.sample_count, record.behavior.high_confidence_sample_count), "safety_blockers": safety_blockers if prediction is not None else [], "anomalies": _detect_anomalies( record, now=now, min_behavior_actions=self._settings.min_behavior_actions, stale_hours=self._settings.retrain_stale_hours, sample_count=record.behavior.sample_count, trusted_actions=record.behavior.high_confidence_sample_count, prediction=prediction, safety_blockers=safety_blockers if prediction is not None else [], ), "confidence_trend": ( [*record.behavior.confidence_trend, round(prediction.confidence, 4)][-30:] if prediction is not None else record.behavior.confidence_trend ), } ) if ( prediction is not None and safety_allowed ): domain = actuator_entity_id.split(".", 1)[0] service = service_for_state(domain, prediction.target_state) if service is not None: if record.behavior.dry_run_enabled: behavior = behavior.model_copy( update={ "prediction": prediction.model_copy( update={ "executed": False, "execution_reason": ( "Dry-run: Aktion wäre ausgeführt worden." ), } ), "dry_run_sample_count": record.behavior.dry_run_sample_count + 1, "reason": ( f"Dry-run hätte {prediction.target_state!r} mit " f"{prediction.confidence:.0%} Sicherheit ausgeführt." ), } ) return self._save_behavior( record, _append_decision_trace( behavior, trigger_entity_id=trigger_entity_id, trigger_state=trigger_state, prediction=prediction, safety_blockers=safety_blockers, duration_ms=_elapsed_ms(started_perf), event_received_at=event_received_at, decision_to_service_ms=None, executed=False, source="event" if event_received_at is not None else "manual", ), ) try: service_started_perf = perf_counter() self._ha_reader.call_service( domain, service, {"entity_id": actuator_entity_id}, ) decision_to_service_ms = _elapsed_ms(service_started_perf) except (HaClientError, ValueError) as exc: logger.error( "Predicted action failed for %s: %s", actuator_entity_id, exc, ) behavior = behavior.model_copy( update={ "reason": f"Vorhersage wurde aus Sicherheitsgründen nicht ausgeführt: {exc}" } ) return self._save_behavior( record, _append_decision_trace( behavior, trigger_entity_id=trigger_entity_id, trigger_state=trigger_state, prediction=prediction, safety_blockers=[str(exc)], duration_ms=_elapsed_ms(started_perf), event_received_at=event_received_at, decision_to_service_ms=None, executed=False, source="event" if event_received_at is not None else "manual", ), ) event = ExecutionEvent( target_state=prediction.target_state, executed_at=now, ) behavior = behavior.model_copy( update={ "prediction": prediction.model_copy( update={ "executed": True, "execution_reason": ( f"Ausgeführt mit {prediction.confidence:.0%} Sicherheit." ), } ), "last_executed_at": now, "execution_events": [ *behavior.execution_events, event, ][-_MAX_EXECUTION_EVENTS:], "reason": ( f"Vorhersage mit {prediction.confidence:.0%} Sicherheit ausgeführt." ), } ) else: behavior = behavior.model_copy( update={ "reason": ( f"Der vorhergesagte Zustand {prediction.target_state!r} " "ist für autonomes Schalten nicht freigegeben." ) } ) behavior = _append_decision_trace( behavior, trigger_entity_id=trigger_entity_id, trigger_state=trigger_state, prediction=prediction, safety_blockers=safety_blockers, duration_ms=_elapsed_ms(started_perf), event_received_at=event_received_at, decision_to_service_ms=( decision_to_service_ms ), executed=bool(prediction is not None and behavior.prediction is not None and behavior.prediction.executed), source="event" if event_received_at is not None else "manual", ) return self._save_behavior(record, behavior) def simulate( self, actuator_entity_id: str, *, sensor_states: dict[str, str], sensor_weights: dict[str, float], state_options: dict[str, list[str]], max_results: int, include_current: bool = True, ) -> list[SimulationOutcome]: record = self._store.get(actuator_entity_id) now = datetime.now(timezone.utc) current_entities = self._ha_reader.read_entities() entities = {entity.entity_id: entity for entity in current_entities} actuator = entities.get(actuator_entity_id) if actuator is None: raise KeyError("Aktor ist aktuell nicht in Home Assistant verfügbar.") selected_context_ids = [ entity_id for entity_id in [ record.assignment.selected_numeric_entity_id, *record.assignment.selected_context_entity_ids, ] if entity_id ] if not selected_context_ids: return [] base_context = { entity_id: entities[entity_id].state for entity_id in selected_context_ids if entity_id in entities and entities[entity_id].state is not None } base_changed_at = { entity_id: entities[entity_id].last_changed for entity_id in base_context } context_weights = _context_weights_for(record) for entity_id, weight in sensor_weights.items(): if entity_id in selected_context_ids: context_weights[entity_id] = max(0.0, min(1.0, weight)) scenarios = _simulation_contexts( base_context, sensor_states=sensor_states, state_options=state_options, selected_context_ids=selected_context_ids, include_current=include_current, ) outcomes: list[SimulationOutcome] = [] for index, context in enumerate(scenarios[:64], start=1): prediction_context: dict[str, str | None] = dict(context) changed_at = dict(base_changed_at) for entity_id, state in context.items(): if base_context.get(entity_id) != state: changed_at[entity_id] = now prediction = predict_behavior( record.behavior.patterns, current_context=prediction_context, current_context_changed_at=changed_at, context_weights=context_weights, now=now, min_support=self._settings.min_behavior_actions, window_minutes=self._settings.prediction_window_minutes, causal_window_seconds=self._settings.prediction_interval_seconds * 2, timezone_name=self._settings.timezone, ) if prediction is not None: would_execute, blockers = self._assess_safety(record, actuator.state, prediction, now) recommendation = ( f"Bestes Szenario: {prediction.target_state} mit {prediction.confidence:.0%}." if would_execute else ( f"Vorhersage {prediction.target_state} mit {prediction.confidence:.0%}, " "aber blockiert: " + " ".join(blockers) ) ) else: would_execute = False blockers = ["Keine fällige Vorhersage."] recommendation = "Dieses Szenario erzeugt keine fällige Vorhersage." outcomes.append( SimulationOutcome( scenario_id=f"scenario-{index}", actuator_entity_id=actuator_entity_id, sensor_states=context, sensor_weights={ entity_id: round(context_weights.get(entity_id, 1.0), 4) for entity_id in context }, prediction=prediction, decision_factors=_decision_factors_for( record, prediction_context, prediction, context_weights=context_weights, ), would_execute=would_execute, blockers=blockers, score=round(prediction.confidence if prediction is not None else 0.0, 4), recommendation=recommendation, ) ) return sorted( outcomes, key=lambda item: ( item.prediction is None, -item.score, item.scenario_id, ), )[:max_results] def record_feedback( self, actuator_entity_id: str, *, correct: bool, expected_state: str | None = None, kind: FeedbackKind | None = None, ) -> ActuatorRecord: record = self._store.get(actuator_entity_id) now = datetime.now(timezone.utc) entities = {entity.entity_id: entity for entity in self._ha_reader.read_entities()} actuator = entities.get(actuator_entity_id) if actuator is None: raise KeyError("Aktor ist aktuell nicht in Home Assistant verfügbar.") context_ids = [ entity_id for entity_id in [ record.assignment.selected_numeric_entity_id, *record.assignment.selected_context_entity_ids, ] if entity_id ] current_context = { entity_id: entities[entity_id].state for entity_id in context_ids if entity_id in entities and entities[entity_id].state is not None } prediction = record.behavior.prediction patterns = list(record.behavior.patterns) reason = "Nutzerfeedback gespeichert." if correct and prediction is not None: local = now.astimezone(ZoneInfo(self._settings.timezone)) patterns.append( BehaviorPattern( target_state=prediction.target_state, minute_of_day=local.hour * 60 + local.minute, weekday=local.weekday(), context_states={ entity_id: state for entity_id, state in current_context.items() if state is not None }, source="user_feedback", weight=1.0, observed_at=now, ) ) reason = "Vorhersage wurde vom Nutzer als korrekt bestätigt." correct_count = record.behavior.correct_feedback_count + 1 incorrect_count = record.behavior.incorrect_feedback_count feedback_kind = kind or FeedbackKind.CORRECT else: target = prediction.target_state if prediction is not None else None if target: patterns = [ pattern.model_copy(update={"weight": 0.1}) if pattern.target_state == target and _pattern_context_matches(pattern, current_context) else pattern for pattern in patterns ] if expected_state: local = now.astimezone(ZoneInfo(self._settings.timezone)) patterns.append( BehaviorPattern( target_state=expected_state, minute_of_day=local.hour * 60 + local.minute, weekday=local.weekday(), context_states={ entity_id: state for entity_id, state in current_context.items() if state is not None }, source="user_correction", weight=1.0, observed_at=now, ) ) reason = "Vorhersage wurde vom Nutzer als falsch markiert." correct_count = record.behavior.correct_feedback_count incorrect_count = record.behavior.incorrect_feedback_count + 1 feedback_kind = kind or FeedbackKind.WRONG if feedback_kind is FeedbackKind.NEVER_AUTOMATE: safety = record.behavior.safety.model_copy( update={ "manual_block": True, "updated_at": now, "note": "Durch Nutzerfeedback dauerhaft blockiert.", } ) else: safety = record.behavior.safety adaptive_updates, manual_override = _adapt_sensor_weights( record, current_context, correct=correct, ) if correct and prediction is not None: safety = record.behavior.safety behavior = record.behavior.model_copy( update={ "patterns": patterns[-_MAX_PATTERNS:], "prediction": ( prediction.model_copy(update={"execution_reason": reason}) if prediction is not None else None ), "reason": reason, "last_trained_at": now, "correct_feedback_count": correct_count, "incorrect_feedback_count": incorrect_count, "feedback_log": [ *record.behavior.feedback_log, feedback_kind, ][-_MAX_FEEDBACK_LOG:], "safety": safety, "adaptive_weight_updates": [ *record.behavior.adaptive_weight_updates, *adaptive_updates, ][-50:], "anomalies": _detect_anomalies( record, now=now, min_behavior_actions=self._settings.min_behavior_actions, stale_hours=self._settings.retrain_stale_hours, sample_count=len(patterns), trusted_actions=record.behavior.high_confidence_sample_count, prediction=prediction, safety_blockers=record.behavior.safety_blockers, correct_feedback_count=correct_count, incorrect_feedback_count=incorrect_count, ), } ) record_for_save = ( record.model_copy(update={"manual_override": manual_override}) if manual_override is not None else record ) return self._save_behavior(record_for_save, behavior) def set_dry_run(self, actuator_entity_id: str, *, enabled: bool) -> ActuatorRecord: record = self._store.get(actuator_entity_id) now = datetime.now(timezone.utc) behavior = record.behavior.model_copy( update={ "dry_run_enabled": enabled, "dry_run_started_at": now if enabled else record.behavior.dry_run_started_at, "reason": ( "Dry-run aktiv; freigegebene Aktionen werden protokolliert, aber nicht geschaltet." if enabled else "Dry-run beendet." ), } ) return self._save_behavior(record, behavior) def refresh_planning_insights(self) -> list[ActuatorRecord]: records = self._store.list() groups = _derive_actuator_groups(records) scenes = _derive_scene_suggestions(records) insights_by_actuator = _derive_agent_insights(records) updated: list[ActuatorRecord] = [] for record in records: behavior = record.behavior.model_copy( update={ "actuator_groups": [ group for group in groups if record.actuator_entity_id in group.member_entity_ids ], "scene_suggestions": [ scene for scene in scenes if record.actuator_entity_id in scene.member_entity_ids ], "agent_insights": insights_by_actuator.get(record.actuator_entity_id, []), } ) updated.append(self._save_behavior(record, behavior)) return updated def rollback_model( self, actuator_entity_id: str, *, version_id: str, ) -> ActuatorRecord: record = self._store.get(actuator_entity_id) snapshot = next( (item for item in record.behavior.model_snapshots if item.version_id == version_id), None, ) if snapshot is None: raise ValueError("Modell-Snapshot nicht gefunden.") behavior = record.behavior.model_copy( update={ "patterns": snapshot.patterns, "sample_count": snapshot.sample_count, "high_confidence_sample_count": snapshot.high_confidence_sample_count, "active_model_version": snapshot.version_id, "reason": f"Rollback auf Modell-Snapshot {snapshot.version_id}.", } ) return self._save_behavior(record, behavior) def set_safety_profile( self, actuator_entity_id: str, *, profile: SafetyProfile, ) -> ActuatorRecord: record = self._store.get(actuator_entity_id) behavior = record.behavior.model_copy( update={ "safety": profile.model_copy(update={"updated_at": datetime.now(timezone.utc)}), "reason": "Sicherheitsprofil wurde manuell aktualisiert.", } ) return self._save_behavior(record, behavior) def refresh_related_automations(self, actuator_entity_id: str) -> ActuatorRecord: record = self._store.get(actuator_entity_id) related = [ RelatedAutomation( entity_id=item.entity_id, config_id=item.config_id, friendly_name=item.friendly_name, enabled=item.enabled, ) for item in self._ha_reader.find_automations_for_entity( actuator_entity_id ) ] behavior = record.behavior.model_copy( update={ "related_automations": related, "automation_conflicts": _automation_conflicts(record, related), } ) behavior = behavior.model_copy( update={ "anomalies": _detect_anomalies( record.model_copy(update={"behavior": behavior}), now=datetime.now(timezone.utc), min_behavior_actions=self._settings.min_behavior_actions, stale_hours=self._settings.retrain_stale_hours, sample_count=behavior.sample_count, trusted_actions=behavior.high_confidence_sample_count, prediction=behavior.prediction, safety_blockers=behavior.safety_blockers, ) } ) return self._save_behavior(record, behavior) def set_automation_enabled( self, actuator_entity_id: str, automation_entity_id: str, *, enabled: bool, ) -> ActuatorRecord: record = self.refresh_related_automations(actuator_entity_id) if automation_entity_id not in { item.entity_id for item in record.behavior.related_automations }: raise ValueError( "Die Automation ist diesem Aktor nicht eindeutig zugeordnet." ) self._ha_reader.call_service( "automation", "turn_on" if enabled else "turn_off", {"entity_id": automation_entity_id}, ) related = [ item.model_copy(update={"enabled": enabled}) if item.entity_id == automation_entity_id else item for item in record.behavior.related_automations ] paused = [ entity_id for entity_id in record.behavior.paused_automation_entity_ids if entity_id != automation_entity_id ] behavior = record.behavior.model_copy( update={ "related_automations": related, "paused_automation_entity_ids": paused, } ) return self._save_behavior(record, behavior) def set_active( self, actuator_entity_id: str, *, active: bool, pause_matching_automations: bool = False, restore_paused_automations: bool = False, ) -> ActuatorRecord: record = self.refresh_related_automations(actuator_entity_id) now = datetime.now(timezone.utc) if active: domain = actuator_entity_id.split(".", 1)[0] if domain not in _SAFE_ACTIVE_DOMAINS: raise ValueError( f"Automatisches Schalten ist für die Domain {domain} nicht freigegeben." ) if record.behavior.status is not BehaviorStatus.TRAINED: raise ValueError("Das Verhaltensmodell hat noch nicht genügend Handlungen gelernt.") if not record.behavior.activation_ready: raise ValueError(record.behavior.activation_reason) mode = BehaviorMode.ACTIVE approved_at = now behavior = record.behavior.model_copy( update={ "mode": mode, "approved_at": approved_at, "safety": record.behavior.safety.model_copy( update={"stage": SafetyStage.ACTIVE, "updated_at": now} ), "reason": ( "Autonomes Lernen und Schalten wurde ausdrücklich freigegeben." ), } ) record = self._save_behavior(record, behavior) if pause_matching_automations: paused: list[str] = [] try: for automation in record.behavior.related_automations: if not automation.enabled: continue self._ha_reader.call_service( "automation", "turn_off", {"entity_id": automation.entity_id}, ) paused.append(automation.entity_id) except (HaClientError, ValueError): for entity_id in paused: try: self._ha_reader.call_service( "automation", "turn_on", {"entity_id": entity_id}, ) except (HaClientError, ValueError): logger.exception( "Failed to restore automation %s after handoff error", entity_id, ) rollback = record.behavior.model_copy( update={ "mode": BehaviorMode.SHADOW, "approved_at": None, "reason": ( "Übernahme fehlgeschlagen; SillyHome bleibt im " "Shadow-Modus." ), } ) self._save_behavior(record, rollback) raise related = [ automation.model_copy(update={"enabled": False}) if automation.entity_id in paused else automation for automation in record.behavior.related_automations ] behavior = record.behavior.model_copy( update={ "related_automations": related, "paused_automation_entity_ids": paused, "reason": ( "SillyHome steuert aktiv; passende HA-Automationen " "wurden pausiert." ), } ) return self._save_behavior(record, behavior) return record else: if restore_paused_automations: for entity_id in record.behavior.paused_automation_entity_ids: self._ha_reader.call_service( "automation", "turn_on", {"entity_id": entity_id}, ) mode = BehaviorMode.SHADOW approved_at = None reason = ( "Shadow-Modus aktiv; pausierte HA-Automationen wurden fortgesetzt." if restore_paused_automations else "Shadow-Modus aktiv; Vorhersagen werden nicht ausgeführt." ) behavior = record.behavior.model_copy( update={ "mode": mode, "approved_at": approved_at, "safety": record.behavior.safety.model_copy( update={"stage": SafetyStage.SHADOW, "updated_at": now} ), "related_automations": [ automation.model_copy(update={"enabled": True}) if ( restore_paused_automations and automation.entity_id in record.behavior.paused_automation_entity_ids ) else automation for automation in record.behavior.related_automations ], "paused_automation_entity_ids": ( [] if restore_paused_automations else record.behavior.paused_automation_entity_ids ), "reason": reason, } ) return self._save_behavior(record, behavior) def _prediction_execution_reason( self, record: ActuatorRecord, current_state: str | None, prediction: BehaviorPrediction, now: datetime, ) -> str: if record.behavior.mode is not BehaviorMode.ACTIVE: return "Nicht ausgeführt: SillyHome ist im Shadow-Modus." if prediction.confidence < self._settings.prediction_confidence: return ( "Nicht ausgeführt: Sicherheit liegt unter der " f"Schaltschwelle von {self._settings.prediction_confidence:.0%}." ) if current_state == prediction.target_state: return "Nicht ausgeführt: Zielzustand ist bereits erreicht." if not self._cooldown_elapsed( record.behavior, now, prediction.target_state, ): return "Nicht ausgeführt: Sicherheits-Cooldown ist noch aktiv." return "Ausführung ist freigegeben." def _assess_safety( self, record: ActuatorRecord, current_state: str | None, prediction: BehaviorPrediction, now: datetime, ) -> tuple[bool, list[str]]: profile = record.behavior.safety blockers: list[str] = [] domain = record.actuator_entity_id.split(".", 1)[0] if not record.enabled: blockers.append("Aktor ist in SillyHome deaktiviert.") if domain not in _SAFE_ACTIVE_DOMAINS: blockers.append(f"Domain {domain} ist nicht für autonomes Schalten freigegeben.") if profile.manual_block: blockers.append("Manuelle Sicherheitssperre ist aktiv.") stage = profile.stage if ( record.behavior.mode is BehaviorMode.ACTIVE and profile.updated_at is None and stage is SafetyStage.SHADOW ): stage = SafetyStage.ACTIVE if stage not in {SafetyStage.ACTIVE, SafetyStage.PARTIAL}: blockers.append(f"Safety-Stufe {stage.value} erlaubt noch kein Schalten.") if record.behavior.mode is not BehaviorMode.ACTIVE: blockers.append("SillyHome ist im Shadow-Modus.") if not record.behavior.activation_ready: blockers.append(record.behavior.activation_reason) threshold = _confidence_threshold_for(profile, prediction.target_state) if prediction.confidence < threshold: blockers.append( f"Sicherheit {prediction.confidence:.0%} liegt unter der Schwelle {threshold:.0%}." ) if current_state == prediction.target_state: blockers.append("Zielzustand ist bereits erreicht.") if not self._cooldown_elapsed( record.behavior, now, prediction.target_state, cooldown_seconds=profile.cooldown_seconds, ): blockers.append("Sicherheits-Cooldown ist noch aktiv.") return not blockers, blockers def _build_patterns( self, *, actuator_history: StateHistorySeries, context_history: dict[str, StateHistorySeries], context_ids: list[str], logbook: list[LogbookEntry], own_executions: list[ExecutionEvent], ) -> list[BehaviorPattern]: patterns: list[BehaviorPattern] = [] previous_state = actuator_history.points[0].state for point in actuator_history.points[1:]: if point.state == previous_state: continue previous_state = point.state if _matches_own_execution(point, own_executions): continue source, weight = _action_source(point, logbook) trigger = _recent_context_transition( context_history, context_ids, point.timestamp, ) contexts = { entity_id: state for entity_id in context_ids if (state := _state_at(context_history.get(entity_id), point.timestamp)) is not None } local = point.timestamp.astimezone(ZoneInfo(self._settings.timezone)) patterns.append( BehaviorPattern( target_state=point.state, minute_of_day=local.hour * 60 + local.minute, weekday=local.weekday(), context_states=contexts, trigger_entity_id=trigger[0] if trigger else None, trigger_from_state=trigger[1] if trigger else None, trigger_to_state=trigger[2] if trigger else None, source=source, weight=weight, observed_at=point.timestamp, ) ) return patterns def _cooldown_elapsed( self, behavior: BehaviorState, now: datetime, target_state: str, *, cooldown_seconds: int | None = None, ) -> bool: if behavior.last_executed_at is None: return True last_event = behavior.execution_events[-1] if behavior.execution_events else None if last_event is not None and last_event.target_state != target_state: return True return (now - behavior.last_executed_at) >= timedelta( seconds=cooldown_seconds if cooldown_seconds is not None else self._settings.execution_cooldown_seconds ) def _save_behavior( self, record: ActuatorRecord, behavior: BehaviorState, ) -> ActuatorRecord: behavior = behavior.model_copy( update={ "model_snapshots": _compact_model_snapshots(behavior.model_snapshots), } ) updated = record.model_copy( update={ "behavior": behavior, "updated_at": datetime.now(timezone.utc), } ) return self._store.upsert(updated) def handle_state_change( self, entity_id: str, new_state: dict[str, object] | None, *, current_entities: Sequence[HaEntitySummary] | None = None, ) -> None: """Wird bei jedem HA-State-Change aufgerufen und löst sofortige Vorhersage aus. - Wenn entity_id ein Aktor ist: evaluate() direkt. - Wenn entity_id ein Kontext-Entity ist: alle betroffenen Aktoren evaluieren. - Wenn current_entities gesetzt ist, kommt die Auswertung direkt aus dem WebSocket-State-Cache statt aus einer frischen REST-Abfrage. """ event_received_at = datetime.now(timezone.utc) records = self._store.list() # Aktor direkt evaluieren for record in records: if record.actuator_entity_id == entity_id: try: self.evaluate( record.actuator_entity_id, current_entities=current_entities, trigger_entity_id=entity_id, trigger_state=_event_state(new_state), event_received_at=event_received_at, ) except Exception: logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", record.actuator_entity_id) return event_state = _event_state(new_state) event_changed_at = _event_changed_at(new_state) or datetime.now(timezone.utc) # Kontext-Entity: alle Aktoren finden, die diesen Kontext nutzen affected_actuators = [ record.actuator_entity_id for record in records if ( record.assignment.selected_numeric_entity_id == entity_id or entity_id in record.assignment.selected_context_entity_ids ) ] for actuator_entity_id in affected_actuators: try: self.evaluate( actuator_entity_id, context_state_overrides={entity_id: event_state}, context_changed_at_overrides={entity_id: event_changed_at}, current_entities=current_entities, trigger_entity_id=entity_id, trigger_state=event_state, event_received_at=event_received_at, ) except Exception: logger.exception("Event-basierte Vorhersage fehlgeschlagen für %s", actuator_entity_id) def _event_state(new_state: dict[str, object] | None) -> str | None: if not isinstance(new_state, dict): return None state = new_state.get("state") return state if isinstance(state, str) else None def _event_changed_at(new_state: dict[str, object] | None) -> datetime | None: if not isinstance(new_state, dict): return None value = new_state.get("last_changed") or new_state.get("last_updated") if not isinstance(value, str): return None try: parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) except ValueError: return None if parsed.tzinfo is None: return parsed.replace(tzinfo=timezone.utc) return parsed def _elapsed_ms(started_perf: float) -> int: return max(0, int((perf_counter() - started_perf) * 1000)) def _append_decision_trace( behavior: BehaviorState, *, trigger_entity_id: str | None, trigger_state: str | None, prediction: BehaviorPrediction | None, safety_blockers: list[str], duration_ms: int, event_received_at: datetime | None, decision_to_service_ms: int | None, executed: bool, source: str, ) -> BehaviorState: now = datetime.now(timezone.utc) blocked = prediction is None or bool(safety_blockers) trace = DecisionTrace( trace_id=f"{now.strftime('%Y%m%d%H%M%S%f')}.{trigger_entity_id or 'manual'}", created_at=now, trigger_entity_id=trigger_entity_id, trigger_state=trigger_state, target_state=prediction.target_state if prediction is not None else None, confidence=prediction.confidence if prediction is not None else None, executed=executed, blocked=blocked, reason=( prediction.execution_reason if prediction is not None else behavior.reason ), blockers=safety_blockers if prediction is not None else ["Keine fällige Vorhersage."], duration_ms=duration_ms, ) updated = behavior.model_copy( update={ "decision_timeline": [ *behavior.decision_timeline, trace, ][-_MAX_DECISION_TRACES:], } ) if event_received_at is None: return updated return _append_latency_measurement( updated, trigger_entity_id=trigger_entity_id, event_received_at=event_received_at, event_to_decision_ms=duration_ms, decision_to_service_ms=decision_to_service_ms, executed=executed, source=source, ) def _append_latency_measurement( behavior: BehaviorState, *, trigger_entity_id: str | None, event_received_at: datetime | None, event_to_decision_ms: int | None, decision_to_service_ms: int | None, executed: bool, source: str, ) -> BehaviorState: if event_received_at is None: return behavior now = datetime.now(timezone.utc) event_to_done_ms = max(0, int((now - event_received_at).total_seconds() * 1000)) measurement = LatencyMeasurement( measured_at=now, trigger_entity_id=trigger_entity_id, event_to_decision_ms=event_to_decision_ms, decision_to_service_ms=decision_to_service_ms, event_to_done_ms=event_to_done_ms, executed=executed, source=source, ) return behavior.model_copy( update={ "latency_measurements": [ *behavior.latency_measurements, measurement, ][-_MAX_LATENCY_MEASUREMENTS:], } ) def _derive_actuator_groups(records: list[ActuatorRecord]) -> list[ActuatorGroup]: by_area: dict[str, list[str]] = {} for record in records: area = _area_hint(record) if area: by_area.setdefault(area, []).append(record.actuator_entity_id) return [ ActuatorGroup( group_id=_slug(f"area_{area}"), name=f"Raum {area}", area_name=area, member_entity_ids=sorted(entity_ids), reason="Aktor-Gruppe aus gemeinsamer Raum-/Kontextzuordnung abgeleitet.", ) for area, entity_ids in sorted(by_area.items()) if len(entity_ids) >= 2 ] def _derive_scene_suggestions(records: list[ActuatorRecord]) -> list[SceneSuggestion]: scenes: list[SceneSuggestion] = [] by_context: dict[tuple[str, str], list[str]] = {} for record in records: for pattern in record.behavior.patterns: for entity_id, state in pattern.context_states.items(): by_context.setdefault((entity_id, state), []).append(record.actuator_entity_id) for (entity_id, state), members in sorted(by_context.items()): unique_members = sorted(set(members)) if len(unique_members) < 2: continue scenes.append( SceneSuggestion( scene_id=_slug(f"{entity_id}_{state}"), label=f"{entity_id} ist {state}", member_entity_ids=unique_members, confidence=min(1.0, len(members) / max(3, len(unique_members) * 2)), reason="Mehrere Aktoren reagieren historisch auf denselben Kontext.", last_seen_at=max( ( pattern.observed_at for record in records for pattern in record.behavior.patterns if pattern.context_states.get(entity_id) == state ), default=None, ), ) ) return scenes[-20:] def _derive_agent_insights(records: list[ActuatorRecord]) -> dict[str, list[AgentInsight]]: result: dict[str, list[AgentInsight]] = {} for record in records: insights: list[AgentInsight] = [] if record.behavior.automation_conflicts: insights.append( AgentInsight( insight_id=f"{record.actuator_entity_id}.automation_conflict", severity="warning", title="Automation-Konflikt prüfen", detail="Eine passende HA-Automation kann parallel zu SillyHome schalten.", action="Automation pausieren oder SillyHome im Shadow-Modus lassen.", ) ) if record.behavior.latency_measurements: durations = [ item.event_to_done_ms for item in record.behavior.latency_measurements if item.event_to_done_ms is not None ] if durations and max(durations) > 1500: insights.append( AgentInsight( insight_id=f"{record.actuator_entity_id}.latency", severity="warning", title="Schalt-Latenz beobachten", detail=f"Letzte maximale Event-Latenz: {max(durations)} ms.", action="WebSocket-Status, HA-Servicezeit und Sensor-Routing pruefen.", ) ) if record.behavior.incorrect_feedback_count > record.behavior.correct_feedback_count: insights.append( AgentInsight( insight_id=f"{record.actuator_entity_id}.feedback", severity="warning", title="Viele negative Feedbacks", detail="Das Modell trifft aktuell mehr falsche als richtige Entscheidungen.", action="Kontextzuordnung, Gewichtung oder Modell-Rollback pruefen.", ) ) result[record.actuator_entity_id] = insights[:5] return result def _area_hint(record: ActuatorRecord) -> str | None: for candidate in [*record.context_candidates, *record.numeric_candidates]: if candidate.area_name: return candidate.area_name return None def _slug(value: str) -> str: result = [] for char in value.lower(): if char.isalnum(): result.append(char) elif char in {".", "_", "-", " "}: result.append("_") return "".join(result).strip("_")[:64] or "item" def _confidence_threshold_for(profile: SafetyProfile, target_state: str) -> float: if target_state == "on" and profile.min_confidence_on is not None: return profile.min_confidence_on if target_state in {"off", "closed"} and profile.min_confidence_off is not None: return profile.min_confidence_off return profile.min_confidence def _decision_factors_for( record: ActuatorRecord, current_context: dict[str, str | None], prediction: BehaviorPrediction | None, *, context_weights: dict[str, float] | None = None, ) -> list[DecisionFactor]: factors: list[DecisionFactor] = [] weights = context_weights or {} candidates = { candidate.entity_id: candidate for candidate in [*record.numeric_candidates, *record.context_candidates] } for entity_id, state in current_context.items(): candidate = candidates.get(entity_id) weight = weights.get( entity_id, candidate.effective_weight if candidate is not None else 1.0, ) relevance = candidate.confidence if candidate is not None else 0.5 contribution = round(min(1.0, weight * relevance), 4) factors.append( DecisionFactor( entity_id=entity_id, label=( candidate.friendly_name if candidate is not None and candidate.friendly_name else entity_id ), factor_type="context", state=state, weight=round(weight, 4), contribution=contribution, evidence=( candidate.evidence[:4] if candidate is not None else ["Aktuell ausgewähltes Kontextsignal."] ), ) ) if prediction is not None: factors.append( DecisionFactor( label=f"Vorhersage {prediction.target_state}", factor_type="prediction", state=prediction.target_state, weight=1.0, contribution=prediction.confidence, evidence=[prediction.reason], ) ) return sorted(factors, key=lambda item: (-item.contribution, item.label))[:12] def _context_weights_for(record: ActuatorRecord) -> dict[str, float]: weights = { candidate.entity_id: candidate.effective_weight for candidate in [*record.numeric_candidates, *record.context_candidates] } override = record.manual_override if override is not None: for entity_id, weight in override.sensor_weights.items(): weights[entity_id] = max(0.0, min(1.0, weight)) for group in override.sensor_weight_groups: for entity_id in group.entity_ids: weights[entity_id] = max(0.0, min(1.0, group.weight)) return weights def _simulation_contexts( base_context: dict[str, str | None], *, sensor_states: dict[str, str], state_options: dict[str, list[str]], selected_context_ids: list[str], include_current: bool, ) -> list[dict[str, str]]: selected = set(selected_context_ids) base = { entity_id: state for entity_id, state in base_context.items() if entity_id in selected and state is not None } for entity_id, state in sensor_states.items(): if entity_id in selected: base[entity_id] = state option_items = [ ( entity_id, list(dict.fromkeys(state for state in states if state))[:6], ) for entity_id, states in state_options.items() if entity_id in selected and states ][:6] contexts: list[dict[str, str]] = [] if include_current or not option_items: contexts.append(dict(base)) if option_items: keys = [item[0] for item in option_items] value_lists = [item[1] for item in option_items] for values in product(*value_lists): context = dict(base) context.update(dict(zip(keys, values, strict=True))) if context not in contexts: contexts.append(context) if len(contexts) >= 64: break return contexts def _knowledge_lines( record: ActuatorRecord, sample_count: int, trusted_actions: int, ) -> list[str]: lines = [ f"{sample_count} historische Aktorhandlungen sind ausgewertet.", f"{trusted_actions} Handlungen stammen eindeutig von Nutzer oder HA-Automationen.", ] if record.assignment.selected_numeric_entity_id: lines.append(f"Hauptsensor: {record.assignment.selected_numeric_entity_id}.") if record.assignment.selected_context_entity_ids: lines.append( f"{len(record.assignment.selected_context_entity_ids)} Kontextsignale sind verbunden." ) return lines def _assumption_lines(record: ActuatorRecord) -> list[str]: lines = [ "Ähnliche Zeitfenster und ähnliche Kontextzustände deuten auf ähnliche Nutzerabsicht hin." ] if record.manual_override is not None: lines.append("Manuelle Sensor-/Kontextkorrekturen werden höher gewichtet.") if record.behavior.related_automations: lines.append("Passende HA-Automationen gelten als starker Hinweis auf vorhandene Logik.") return lines def _uncertainty_lines( record: ActuatorRecord, sample_count: int, trusted_actions: int, ) -> list[str]: lines: list[str] = [] if sample_count < trusted_actions + 3: lines.append("Noch wenig Varianz in den gelernten Handlungen.") if trusted_actions < sample_count: lines.append("Ein Teil der Handlungen ist nicht eindeutig Nutzer oder Automation zugeordnet.") if record.assignment.review_required: lines.append("Die automatische Kontextzuordnung verlangt noch Prüfung.") if record.behavior.incorrect_feedback_count: lines.append( f"{record.behavior.incorrect_feedback_count} negative Feedbacks senken Vertrauen." ) return lines or ["Keine kritische Unsicherheit aus den lokalen Daten erkannt."] def _next_model_snapshots( existing: list[ModelSnapshot], version_id: str, patterns: list[BehaviorPattern], sample_count: int, trusted_actions: int, average_confidence: float, incorrect_feedback_count: int, reason: str, ) -> list[ModelSnapshot]: snapshot = ModelSnapshot( version_id=version_id, sample_count=sample_count, high_confidence_sample_count=trusted_actions, average_confidence=round(average_confidence, 4), incorrect_feedback_count=incorrect_feedback_count, patterns=patterns[-_MAX_SNAPSHOT_PATTERNS:], reason=reason, ) return _compact_model_snapshots([*existing, snapshot]) def _compact_model_snapshots(existing: list[ModelSnapshot]) -> list[ModelSnapshot]: return [ snapshot.model_copy( update={"patterns": snapshot.patterns[-_MAX_SNAPSHOT_PATTERNS:]} ) for snapshot in existing[-_MAX_MODEL_SNAPSHOTS:] ] def _average(values: list[float]) -> float: return sum(values) / len(values) if values else 0.0 def _time_profiles(patterns: list[BehaviorPattern]) -> list[TimeProfile]: buckets = { "night": ("Nacht", range(0, 360)), "morning": ("Morgen", range(360, 720)), "day": ("Tag", range(720, 1080)), "evening": ("Abend", range(1080, 1440)), } profiles: list[TimeProfile] = [] for profile_id, (label, minutes) in buckets.items(): selected = [pattern for pattern in patterns if pattern.minute_of_day in minutes] if not selected: profiles.append(TimeProfile(profile_id=profile_id, label=label)) continue by_state: dict[str, int] = {} for pattern in selected: by_state[pattern.target_state] = by_state.get(pattern.target_state, 0) + 1 dominant_state, count = max(by_state.items(), key=lambda item: (item[1], item[0])) profiles.append( TimeProfile( profile_id=profile_id, label=label, sample_count=len(selected), dominant_state=dominant_state, confidence=round(count / len(selected), 4), ) ) weekend = [pattern for pattern in patterns if pattern.weekday >= 5] profiles.append( TimeProfile( profile_id="weekend", label="Wochenende", sample_count=len(weekend), dominant_state=( max( {pattern.target_state: 0 for pattern in weekend}, key=lambda state: sum(pattern.target_state == state for pattern in weekend), ) if weekend else None ), confidence=round(len(weekend) / len(patterns), 4) if patterns else 0.0, ) ) return profiles def _adapt_sensor_weights( record: ActuatorRecord, current_context: dict[str, str | None], *, correct: bool, ) -> tuple[list[AdaptiveWeightUpdate], ManualOverride | None]: if not current_context: return [], record.manual_override candidates = { candidate.entity_id: candidate for candidate in [*record.numeric_candidates, *record.context_candidates] } previous = record.manual_override weights = dict(previous.sensor_weights if previous is not None else {}) updates: list[AdaptiveWeightUpdate] = [] delta = 0.03 if correct else -0.08 for entity_id in current_context: candidate = candidates.get(entity_id) base = weights.get( entity_id, candidate.effective_weight if candidate is not None else 1.0, ) new_weight = round(min(1.0, max(0.1, base + delta)), 4) if new_weight == base: continue weights[entity_id] = new_weight updates.append( AdaptiveWeightUpdate( entity_id=entity_id, previous_weight=round(base, 4), new_weight=new_weight, reason=( "Feedback korrekt: Kontextsignal leicht höher gewichtet." if correct else "Feedback falsch: Kontextsignal vorsichtig abgewertet." ), ) ) if not updates: return [], previous return updates, ManualOverride( numeric_entity_id=( previous.numeric_entity_id if previous is not None else record.assignment.selected_numeric_entity_id ), context_entity_ids=( previous.context_entity_ids if previous is not None else record.assignment.selected_context_entity_ids ), sensor_weights=weights, sensor_weight_groups=previous.sensor_weight_groups if previous is not None else [], note="Sensor-Gewichtungen automatisch aus Feedback angepasst.", ) def _automation_conflicts( record: ActuatorRecord, related: list[RelatedAutomation], ) -> list[AutomationConflict]: conflicts: list[AutomationConflict] = [] for automation in related: if record.behavior.mode is BehaviorMode.ACTIVE and automation.enabled: conflicts.append( AutomationConflict( automation_entity_id=automation.entity_id, severity="warning", status="open", reason=( "SillyHome ist aktiv, aber diese passende HA-Automation " "ist ebenfalls aktiv. Das kann zu konkurrierenden Schaltungen führen." ), ) ) elif automation.entity_id in record.behavior.paused_automation_entity_ids: conflicts.append( AutomationConflict( automation_entity_id=automation.entity_id, severity="info", status="controlled", reason="Automation ist durch SillyHome pausiert.", ) ) return conflicts def _detect_anomalies( record: ActuatorRecord, *, now: datetime, min_behavior_actions: int, stale_hours: int, sample_count: int, trusted_actions: int, prediction: BehaviorPrediction | None, safety_blockers: list[str], correct_feedback_count: int | None = None, incorrect_feedback_count: int | None = None, ) -> list[AnomalyEvent]: anomalies: list[AnomalyEvent] = [] def add(category: str, severity: str, title: str, detail: str) -> None: anomalies.append( AnomalyEvent( anomaly_id=f"{record.actuator_entity_id}.{category}", category=category, severity=severity, title=title, detail=detail, detected_at=now, ) ) if not record.assignment.selected_context_entity_ids and not record.assignment.selected_numeric_entity_id: add( "missing_context", "warning", "Kein Kontext verbunden", "Der Aktor hat keine Sensor-/Kontextbasis. Entscheidungen bleiben unsicher.", ) if sample_count < min_behavior_actions: add( "low_samples", "info", "Zu wenig Lernbeispiele", f"{sample_count} von {min_behavior_actions} benoetigten Handlungen gelernt.", ) if trusted_actions < sample_count: add( "unclear_sources", "info", "Unklare Aktorhandlungen", "Ein Teil der gelernten Handlungen stammt nicht eindeutig von Nutzer oder Automation.", ) if record.behavior.last_trained_at is not None: age = now - record.behavior.last_trained_at if age > timedelta(hours=stale_hours): add( "stale_training", "warning", "Training ist veraltet", f"Letztes Training liegt mehr als {stale_hours} Stunden zurueck.", ) if prediction is not None and prediction.matching_patterns and prediction.confidence < record.behavior.safety.min_confidence: add( "low_confidence_prediction", "warning", "Vorhersage unter Sicherheitsgrenze", ( f"Confidence {prediction.confidence:.0%} liegt unter " f"{record.behavior.safety.min_confidence:.0%}." ), ) if record.behavior.safety.manual_block: add( "manual_block", "info", "Manuelle Sicherheitssperre aktiv", "Der Aktor ist bewusst gegen automatisches Schalten gesperrt.", ) if safety_blockers: add( "safety_blockers", "info", "Safety blockiert aktuelle Aktion", " ".join(safety_blockers)[:500], ) if any(conflict.severity == "warning" for conflict in record.behavior.automation_conflicts): add( "automation_conflict", "critical", "Parallele Automation erkannt", "SillyHome und mindestens eine passende HA-Automation koennen parallel schalten.", ) correct = ( record.behavior.correct_feedback_count if correct_feedback_count is None else correct_feedback_count ) incorrect = ( record.behavior.incorrect_feedback_count if incorrect_feedback_count is None else incorrect_feedback_count ) total = correct + incorrect if total >= 3 and incorrect / total >= 0.35: add( "feedback_error_rate", "critical", "Viele falsche Vorhersagen", f"{incorrect} von {total} Feedbacks waren negativ. Modell pruefen oder Rollback nutzen.", ) return anomalies[-30:] def predict_behavior( patterns: list[BehaviorPattern], *, current_context: dict[str, str | None], now: datetime, min_support: int, window_minutes: int, current_context_changed_at: dict[str, datetime | None] | None = None, context_weights: dict[str, float] | None = None, causal_window_seconds: int = 120, timezone_name: str = "Europe/Berlin", ) -> BehaviorPrediction | None: if not patterns: return None local = now.astimezone(ZoneInfo(timezone_name)) minute_of_day = local.hour * 60 + local.minute changed_at = current_context_changed_at or {} by_state: dict[str, list[float]] = {} causal_support_by_state: dict[str, int] = {} for pattern in patterns: if pattern.trigger_entity_id and pattern.trigger_to_state: trigger_changed_at = changed_at.get(pattern.trigger_entity_id) trigger_age = ( (now - trigger_changed_at).total_seconds() if trigger_changed_at is not None else None ) if not ( current_context.get(pattern.trigger_entity_id) == pattern.trigger_to_state and trigger_age is not None and 0 <= trigger_age <= causal_window_seconds ): continue comparable = [ (entity_id, expected) for entity_id, expected in pattern.context_states.items() if entity_id in current_context ] context_score = _weighted_context_score( comparable, current_context, context_weights or {}, ) score = pattern.weight * (0.85 + 0.15 * context_score) by_state.setdefault(pattern.target_state, []).append(score) causal_support_by_state[pattern.target_state] = ( causal_support_by_state.get(pattern.target_state, 0) + 1 ) continue distance = _circular_minute_distance(minute_of_day, pattern.minute_of_day) if distance > window_minutes: continue time_score = 1.0 - (distance / max(window_minutes, 1)) weekday_score = ( 1.0 if local.weekday() == pattern.weekday else 0.5 if (local.weekday() >= 5) == (pattern.weekday >= 5) else 0.0 ) comparable = [ (entity_id, expected) for entity_id, expected in pattern.context_states.items() if entity_id in current_context ] context_score = _weighted_context_score( comparable, current_context, context_weights or {}, ) score = pattern.weight * ( 0.45 * time_score + 0.45 * context_score + 0.10 * weekday_score ) by_state.setdefault(pattern.target_state, []).append(score) if not by_state: return None target_state, scores = max( by_state.items(), key=lambda item: (sum(item[1]), len(item[1]), item[0]), ) support = len(scores) causal_support = causal_support_by_state.get(target_state, 0) confidence = min(1.0, (sum(scores) / support) * min(1.0, support / min_support)) if confidence <= 0: return None return BehaviorPrediction( target_state=target_state, confidence=round(confidence, 4), generated_at=now, matching_patterns=support, reason=( ( f"{causal_support} historische Handlungen folgten demselben " "frischen Sensorwechsel." ) if causal_support else f"{support} ähnliche Handlungsmuster passen zu Zeit und aktuellem Kontext." ), ) def _weighted_context_score( comparable: list[tuple[str, str]], current_context: dict[str, str | None], context_weights: dict[str, float], ) -> float: if not comparable: return 0.5 total_weight = 0.0 matched_weight = 0.0 for entity_id, expected in comparable: weight = max(0.0, min(1.0, context_weights.get(entity_id, 1.0))) total_weight += weight if current_context.get(entity_id) == expected: matched_weight += weight if total_weight <= 0: return 0.5 return matched_weight / total_weight def service_for_state(domain: str, target_state: str) -> str | None: if domain in {"fan", "humidifier", "light", "media_player", "remote", "switch"}: return {"on": "turn_on", "off": "turn_off"}.get(target_state) if domain == "scene": return "turn_on" if target_state == "on" else None if domain == "cover": return {"open": "open_cover", "closed": "close_cover"}.get(target_state) return None def _state_at(series: StateHistorySeries | None, timestamp: datetime) -> str | None: if series is None: return None state: str | None = None for point in series.points: if point.timestamp > timestamp: break state = point.state return state def _action_source( point: StateHistoryPoint, logbook: list[LogbookEntry], ) -> tuple[str, float]: nearest = min( logbook, key=lambda item: abs(item.timestamp - point.timestamp), default=None, ) if nearest is None or abs(nearest.timestamp - point.timestamp) > _ACTION_LOGBOOK_TOLERANCE: return "physical_or_unknown", 0.7 if nearest.context_user_id: return "user", 1.0 if nearest.context_domain in _AUTOMATION_CONTEXT_DOMAINS: return "automation", 1.0 return "physical_or_unknown", 0.7 def _matches_own_execution( point: StateHistoryPoint, own_executions: list[ExecutionEvent], ) -> bool: return any( event.target_state == point.state and abs(event.executed_at - point.timestamp) <= _OWN_ACTION_TOLERANCE for event in own_executions ) def _pattern_context_matches( pattern: BehaviorPattern, current_context: dict[str, str | None], ) -> bool: comparable = [ (entity_id, expected) for entity_id, expected in pattern.context_states.items() if entity_id in current_context ] if not comparable: return False return all(current_context[entity_id] == expected for entity_id, expected in comparable) def _recent_context_transition( history: dict[str, StateHistorySeries], context_ids: list[str], timestamp: datetime, ) -> tuple[str, str, str] | None: nearest: tuple[timedelta, str, str, str] | None = None for entity_id in context_ids: series = history.get(entity_id) if series is None: continue previous_state: str | None = None for point in series.points: if point.timestamp > timestamp: break if previous_state is not None and point.state != previous_state: age = timestamp - point.timestamp if age <= _CONTEXT_TRIGGER_TOLERANCE and ( nearest is None or age < nearest[0] ): nearest = (age, entity_id, previous_state, point.state) previous_state = point.state if nearest is None: return None return nearest[1], nearest[2], nearest[3] def _circular_minute_distance(left: int, right: int) -> int: direct = abs(left - right) return min(direct, 1440 - direct)