Files
2026-07-08 10:43:39 -04:00

458 lines
18 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Home Assistant event bridge - translates HA events to orchestrator calls.
This infrastructure component listens to HA entity state changes and delegates
to the orchestrator.
"""
from __future__ import annotations
import asyncio
import logging
from collections.abc import Callable
from typing import TYPE_CHECKING
from homeassistant.core import Event, EventStateChangedData, HomeAssistant, callback
from homeassistant.helpers.event import async_track_state_change_event
from homeassistant.util import dt as dt_util
from .vtherm_compat import get_vtherm_attribute
if TYPE_CHECKING:
from datetime import datetime
from ..application import HeatingOrchestrator
_LOGGER = logging.getLogger(__name__)
# Minimum change in monitored sensor value (humidity %, cloud coverage %) to
# trigger a recalculation. Small fluctuations below this threshold are ignored.
_MONITORED_ENTITY_CHANGE_THRESHOLD = 3.0
# Tolerance in seconds for anticipated_start_time comparison. Changes smaller
# than this are not considered meaningful enough to re-publish the event.
_ANTICIPATION_TIME_TOLERANCE_SECONDS = 60
class HAEventBridge:
"""Bridges Home Assistant events to application service.
This infrastructure component:
- Listens to relevant HA entity state changes
- Translates events to application service calls
- Manages state change listeners lifecycle
NO business logic - pure event routing.
"""
def __init__(
self,
hass: HomeAssistant,
orchestrator: HeatingOrchestrator,
vtherm_entity_id: str,
scheduler_entity_ids: list[str],
monitored_entity_ids: list[str] | None = None,
entry_id: str | None = None,
get_ihp_enabled_func: Callable[[], bool] | None = None,
) -> None:
"""Initialize the event bridge.
Args:
hass: Home Assistant instance
orchestrator: Orchestrator to delegate to
vtherm_entity_id: VTherm entity to monitor for slopes
scheduler_entity_ids: Scheduler entities to monitor
monitored_entity_ids: Additional entities to monitor (humidity, etc.)
entry_id: Config entry ID for event filtering
get_ihp_enabled_func: Callback function to get current IHP enabled state
"""
self._hass = hass
self._orchestrator = orchestrator
self._vtherm_entity_id = vtherm_entity_id
self._scheduler_entity_ids = scheduler_entity_ids
self._monitored_entity_ids = monitored_entity_ids or []
self._get_ihp_enabled = get_ihp_enabled_func or (lambda: True)
self._entry_id = entry_id
# Track all entities that should trigger updates
self._tracked_entities = (
[vtherm_entity_id] + scheduler_entity_ids + self._monitored_entity_ids
)
# Listener cleanup callbacks
self._listeners: list = []
# Debouncing state
self._ignore_vtherm_until: datetime | None = None
# Deduplication: remember the last event data that was published so we
# can skip firing when nothing meaningful has changed.
self._last_published_data: dict | None = None
# Track active tasks for proper shutdown
self._active_tasks: set = set()
self._is_shutting_down = False
self._recalculate_task: asyncio.Task | None = None
self._recalculate_pending = False
def setup_listeners(self) -> None:
"""Setup all event listeners."""
@callback
def _on_entity_changed(event: Event[EventStateChangedData]) -> None:
"""Handle entity state change events.
All listened entities trigger _recalculate_and_publish(), which routes
to appropriate orchestrator method based on event source.
"""
entity_id = event.data.get("entity_id")
if entity_id not in self._tracked_entities:
return
# VTherm-specific handling for slope learning
if entity_id == self._vtherm_entity_id:
self._handle_vtherm_change(event)
elif entity_id in self._scheduler_entity_ids:
# Only trigger for meaningful scheduler state changes
self._trigger_recalculate_if_meaningful(
self._has_meaningful_scheduler_change(event), entity_id
)
else:
# Monitored entity (humidity, cloud cover) apply threshold filter
self._trigger_recalculate_if_meaningful(
self._has_meaningful_monitored_change(event), entity_id
)
# Register state change listener
unsub = async_track_state_change_event(
self._hass, self._tracked_entities, _on_entity_changed
)
self._listeners.append(unsub)
_LOGGER.debug("Event bridge tracking %d entities", len(self._tracked_entities))
# ------------------------------------------------------------------
# Smart filtering helpers
# ------------------------------------------------------------------
def _trigger_recalculate_if_meaningful(self, meaningful: bool, entity_id: str) -> None:
"""Trigger recalculation when the entity change is meaningful.
Logs the outcome at DEBUG level to aid diagnostics without noise.
"""
if meaningful and not self._is_shutting_down:
_LOGGER.debug("Entity %s changed meaningfully, triggering update", entity_id)
self._request_recalculate()
elif not meaningful:
_LOGGER.debug("Entity %s change not actionable, skipping recalculation", entity_id)
def _request_recalculate(self) -> None:
"""Request a recalculation with coalescing.
If a recalculation is already running, mark one pending pass instead of
spawning additional concurrent tasks.
"""
if self._is_shutting_down:
return
if self._recalculate_task is not None and not self._recalculate_task.done():
self._recalculate_pending = True
_LOGGER.debug("Recalculation already running, coalescing trigger")
return
task = self._hass.async_create_task(self._run_recalculate_loop())
self._recalculate_task = task
self._active_tasks.add(task)
task.add_done_callback(self._on_recalculate_task_done)
def _on_recalculate_task_done(self, task: asyncio.Task) -> None:
"""Track recalculation task completion and cleanup references."""
self._active_tasks.discard(task)
if self._recalculate_task is task:
self._recalculate_task = None
async def _run_recalculate_loop(self) -> None:
"""Run recalculation and process one coalesced pending request if needed."""
while not self._is_shutting_down:
try:
await self._recalculate_and_publish()
except Exception as exc: # noqa: BLE001
_LOGGER.error("Error during recalculation workflow: %s", exc, exc_info=True)
if self._recalculate_pending:
self._recalculate_pending = False
_LOGGER.debug("Processing coalesced recalculation trigger")
continue
break
def _has_meaningful_scheduler_change(self, event: Event[EventStateChangedData]) -> bool:
"""Return True only when a scheduler state change is actionable for IHP.
Ignores attribute-only updates that do not affect the IHP schedule
(e.g., internal counters, last_triggered, etc.). Only the following
changes are considered meaningful:
- The enabled / disabled state (on / off) toggled.
- The ``next_trigger`` timestamp changed (different upcoming slot).
- The ``actions`` attribute changed (different target temperature / preset).
"""
old_state = event.data.get("old_state")
new_state = event.data.get("new_state")
if not old_state or not new_state:
return True
# Enabled / disabled toggle
if old_state.state != new_state.state:
return True
old_attrs = old_state.attributes
new_attrs = new_state.attributes
# Next occurrence time changed
if old_attrs.get("next_trigger") != new_attrs.get("next_trigger"):
return True
# Target temperature / actions changed
return bool(old_attrs.get("actions") != new_attrs.get("actions"))
def _has_meaningful_monitored_change(self, event: Event[EventStateChangedData]) -> bool:
"""Return True only when a monitored sensor value changed significantly.
Tiny fluctuations in humidity or cloud coverage sensors are ignored.
Availability transitions are always considered meaningful.
"""
old_state = event.data.get("old_state")
new_state = event.data.get("new_state")
if not old_state or not new_state:
return True
_unavailable = {"unavailable", "unknown"}
old_unavailable = old_state.state in _unavailable
new_unavailable = new_state.state in _unavailable
# Availability transition is always actionable
if old_unavailable != new_unavailable:
return True
# Both unavailable nothing changed
if old_unavailable and new_unavailable:
return False
try:
old_val = float(str(old_state.state))
new_val = float(str(new_state.state))
return abs(new_val - old_val) >= _MONITORED_ENTITY_CHANGE_THRESHOLD
except (TypeError, ValueError):
# Non-numeric state: trigger on any state string change
return bool(old_state.state != new_state.state)
def _is_meaningful_change_from_last(self, new_data: dict) -> bool:
"""Return True when *new_data* differs meaningfully from the last published event.
Prevents flooding sensors with identical or near-identical events.
"""
last = self._last_published_data
if last is None:
return True
# Core schedule fields any difference is significant
for key in ("next_schedule_time", "next_target_temperature", "scheduler_entity"):
if new_data.get(key) != last.get(key):
return True
# Current indoor temperature (0.1 °C tolerance)
new_temp = new_data.get("current_temp")
last_temp = last.get("current_temp")
# Transition between "no data" and a real temperature is always meaningful
if (new_temp is None) != (last_temp is None):
return True
if new_temp is not None and last_temp is not None and abs(new_temp - last_temp) >= 0.1:
return True
# Learned heating slope (0.05 °C/h tolerance)
new_lhs = new_data.get("learned_heating_slope")
last_lhs = last.get("learned_heating_slope")
# Transition between "no data" and a real slope is always meaningful
if (new_lhs is None) != (last_lhs is None):
return True
if new_lhs is not None and last_lhs is not None and abs(new_lhs - last_lhs) >= 0.05:
return True
# Anticipated start time (1-minute tolerance)
new_start = new_data.get("anticipated_start_time")
last_start = last.get("anticipated_start_time")
# If both are equal (including both None or identical strings/datetimes), no meaningful change
if new_start == last_start:
return False
# If one is None and the other is not, this is a meaningful change
if new_start is None or last_start is None:
return True
try:
new_dt = dt_util.parse_datetime(new_start) if isinstance(new_start, str) else new_start
last_dt = (
dt_util.parse_datetime(last_start) if isinstance(last_start, str) else last_start
)
if new_dt is None or last_dt is None:
return True
if abs((new_dt - last_dt).total_seconds()) >= _ANTICIPATION_TIME_TOLERANCE_SECONDS:
return True
except (TypeError, ValueError, AttributeError):
return True
return False
def _handle_vtherm_change(self, event: Event[EventStateChangedData]) -> None:
"""Handle VTherm state changes (temperature filter + recalculation trigger).
Args:
event: State change event
"""
if self._is_shutting_down:
return
old_state = event.data.get("old_state")
new_state = event.data.get("new_state")
if not old_state or not new_state:
return
# Check if we should ignore (self-induced change)
if self._ignore_vtherm_until and dt_util.now() < self._ignore_vtherm_until:
_LOGGER.debug("Ignoring self-induced VTherm change")
return
# Extract temperature changes (v8.0.0+ compatible)
old_temp = get_vtherm_attribute(old_state, "current_temperature")
new_temp = get_vtherm_attribute(new_state, "current_temperature")
if old_temp == new_temp:
_LOGGER.debug("VTherm change but temperature unchanged, skipping")
return
_LOGGER.debug("VTherm temperature changed: %s -> %s", old_temp, new_temp)
self._request_recalculate()
async def _recalculate_and_publish(self) -> None:
"""Recalculate anticipation and publish event for sensors if data changed.
Always fires the same unified event structure. Scheduling/clearing is
expressed via None values (never a ``clear_values`` flag) so that all
sensors share a single, consistent event shape.
"""
anticipation_data = await self._orchestrator.calculate_and_schedule_anticipation(
ihp_enabled=self._get_ihp_enabled()
)
if not anticipation_data:
anticipation_data = {}
# Build unified event structure (same keys always, None for missing values)
has_complete_data = (
anticipation_data.get("anticipated_start_time") is not None
and anticipation_data.get("next_schedule_time") is not None
)
if has_complete_data:
event_data: dict = {
"entry_id": self._entry_id,
"anticipated_start_time": anticipation_data["anticipated_start_time"].isoformat(),
"next_schedule_time": anticipation_data["next_schedule_time"].isoformat(),
"next_target_temperature": anticipation_data.get("next_target_temperature"),
"anticipation_minutes": anticipation_data.get("anticipation_minutes"),
"current_temp": anticipation_data.get("current_temp"),
"learned_heating_slope": anticipation_data.get("learned_heating_slope"),
"confidence_level": anticipation_data.get("confidence_level"),
"scheduler_entity": anticipation_data.get("scheduler_entity"),
}
else:
event_data = {
"entry_id": self._entry_id,
"anticipated_start_time": None,
"next_schedule_time": None,
"next_target_temperature": None,
"anticipation_minutes": None,
"current_temp": anticipation_data.get("current_temp"),
"learned_heating_slope": anticipation_data.get("learned_heating_slope"),
"confidence_level": None,
"scheduler_entity": None,
}
if self._is_meaningful_change_from_last(event_data):
self._hass.bus.async_fire(
"intelligent_heating_pilot_anticipation_calculated",
event_data,
)
self._last_published_data = event_data
_LOGGER.debug("Published anticipation event for sensors (data changed)")
else:
_LOGGER.debug(
"Skipped anticipation event: no meaningful change from last published data"
)
def ignore_vtherm_changes_for(self, seconds: int = 10) -> None:
"""Temporarily ignore VTherm changes (used after self-induced changes).
Args:
seconds: How long to ignore changes
"""
from datetime import timedelta
self._ignore_vtherm_until = dt_util.now() + timedelta(seconds=seconds)
async def async_cleanup(self) -> None:
"""Cleanup all event listeners and cancel pending tasks.
This is called during integration unload to ensure:
1. No more events trigger new calculations
2. All pending tasks are cancelled with timeout
"""
_LOGGER.debug("Starting event bridge cleanup with %d active tasks", len(self._active_tasks))
self._is_shutting_down = True
# Unsubscribe from all listeners to prevent new tasks
for unsub in self._listeners:
unsub()
self._listeners.clear()
# Cancel all active tasks with timeout
if self._active_tasks:
_LOGGER.debug("Cancelling %d active recalculation tasks", len(self._active_tasks))
for task in self._active_tasks:
if not task.done():
task.cancel()
# Wait for tasks to complete (with timeout to prevent hangs)
try:
await asyncio.wait_for(
asyncio.gather(*self._active_tasks, return_exceptions=True),
timeout=5.0, # 5-second timeout for all tasks
)
except asyncio.TimeoutError:
_LOGGER.warning(
"Timeout waiting for %d tasks to shutdown; forcing cleanup",
len(self._active_tasks),
)
finally:
self._active_tasks.clear()
_LOGGER.debug("Event bridge cleanup completed")
def cleanup(self) -> None:
"""Synchronous cleanup for backwards compatibility.
Note: This is called by code that doesn't support async, but doesn't
cancel pending tasks. Use async_cleanup() during normal unload.
"""
if self._is_shutting_down:
return
self._is_shutting_down = True
for unsub in self._listeners:
unsub()
self._listeners.clear()
_LOGGER.debug("Event bridge cleaned up (sync mode)")