278 files

This commit is contained in:
Home Assistant Version Control
2026-09-13 20:38:58 +00:00
parent d7a0a08372
commit 1b4f5f68c6
278 changed files with 36892 additions and 7150 deletions
@@ -5,13 +5,13 @@ from __future__ import annotations
import logging
import math
from abc import ABC, abstractmethod
from collections.abc import Coroutine
from datetime import datetime
from typing import TYPE_CHECKING, Any
from homeassistant.core import CALLBACK_TYPE, Event, HomeAssistant, State, callback
from homeassistant.helpers.event import (
EventStateChangedData,
async_call_later,
async_track_state_change_event,
)
@@ -24,11 +24,21 @@ from ...const import (
EVENT_TRIGGER_DEACTIVATED,
UNAVAILABLE_STATES,
)
from ...helpers.managed_timer import ManagedTimer
_LOGGER = logging.getLogger(__name__)
# The initial-evaluation retry: an entity that is unknown/unavailable at
# setup is re-checked this often, this many times, before the state-change
# listener alone is trusted to catch its recovery (issue #1 family).
RETRY_DELAY_SECONDS = 30.0
RETRY_MAX_ATTEMPTS = 10
class BaseTrigger(ABC):
# Class-level default: subclasses built without __init__ in tests still
# carry the flag (see the instance comment in __init__).
_recovered_since_reset: bool = True
"""Base class for all maintenance triggers."""
def __init__(
@@ -45,9 +55,18 @@ class BaseTrigger(ABC):
self.attribute = trigger_config.get("attribute")
self._triggered = False
# Bug audit 2026-09-12: after a completion resets the trigger, the
# sensor may still read beyond the threshold (the user tapped Complete
# before refilling). Such a re-activation is NOT a new edge and must
# not lift the coordinator's post-completion cooldown - only an
# activation after the value was seen on the other side is.
self._recovered_since_reset = True
self._current_value: float | None = None
self._unsub_listener: CALLBACK_TYPE | None = None
self._unsub_retry: CALLBACK_TYPE | None = None
# One timer per purpose (a retry and a subclass's for/hold window can
# be pending at the same time); the tasks a trigger spawns (history
# entry, refresh, persist) are tracked here and cancelled at teardown.
self._retry_timer = ManagedTimer(hass, f"{type(self).__name__}:{self.entity_id}:retry")
self._logged_unavailable = False # Log-once pattern for unavailable
@property
@@ -108,42 +127,47 @@ class BaseTrigger(ABC):
)
def _schedule_retry(self) -> None:
"""Schedule a retry of the initial evaluation after 30 seconds."""
self._cancel_retry()
"""Schedule a retry of the initial evaluation (30 s, capped).
@callback
def _retry_initial_evaluation(_now: datetime) -> None:
"""Re-check entity state after a delay."""
self._unsub_retry = None
state = self.hass.states.get(self.entity_id)
if state is None or state.state in UNAVAILABLE_STATES:
_LOGGER.debug(
"Trigger entity %s still %s after retry",
self.entity_id,
state.state if state else "missing",
)
return
value = self._get_numeric_value(state)
if value is not None:
self._current_value = value
self._evaluate_and_update(value)
_LOGGER.info(
"Trigger entity %s recovered after retry (value=%s)",
self.entity_id,
value,
)
Before the ManagedTimer this was a single hard-coded shot: an entity
still unavailable 30 s after setup was never re-checked by the timer
again (the state-change listener catches a LATER recovery, but not
one that happened while we were not looking). Now it re-arms until
the entity reports or the budget is spent.
"""
self._retry_timer.retry(self._retry_initial_evaluation, delay=RETRY_DELAY_SECONDS, max_attempts=RETRY_MAX_ATTEMPTS)
self._unsub_retry = async_call_later(self.hass, 30, _retry_initial_evaluation)
@callback
def _retry_initial_evaluation(self, _now: datetime) -> None:
"""Re-check entity state after a delay."""
state = self.hass.states.get(self.entity_id)
if state is None or state.state in UNAVAILABLE_STATES:
_LOGGER.debug(
"Trigger entity %s still %s after retry",
self.entity_id,
state.state if state else "missing",
)
self._schedule_retry()
return
self._retry_timer.reset_retries()
value = self._get_numeric_value(state)
if value is not None:
self._current_value = value
self._evaluate_and_update(value)
_LOGGER.info(
"Trigger entity %s recovered after retry (value=%s)",
self.entity_id,
value,
)
def _cancel_retry(self) -> None:
"""Cancel pending retry timer."""
if self._unsub_retry is not None:
self._unsub_retry()
self._unsub_retry = None
def _track(self, coro: Coroutine[Any, Any, Any], *, cancel_on_close: bool = True) -> None:
"""Spawn a fire-and-forget task the trigger owns (cancelled at
teardown unless it must land regardless — see the auto-complete)."""
self._retry_timer.track_task(coro, cancel_on_close=cancel_on_close)
async def async_teardown(self) -> None:
"""Remove the trigger listener."""
self._cancel_retry()
"""Remove the trigger listener, the retry timer and owned tasks."""
self._retry_timer.close()
if self._unsub_listener is not None:
self._unsub_listener()
self._unsub_listener = None
@@ -206,6 +230,8 @@ class BaseTrigger(ABC):
was_triggered = self._triggered
is_triggered = self.evaluate(value)
self._triggered = is_triggered
if not is_triggered:
self._recovered_since_reset = True
if is_triggered and not was_triggered:
self._on_trigger_activated(value)
@@ -229,7 +255,7 @@ class BaseTrigger(ABC):
full update interval — 30 s once, 4 min the next time. Debounced on
purpose (HA's ten-second window): a noisy sensor must not recompute the
object on every state change; user actions use async_refresh_now."""
self.hass.async_create_task(self._coordinator.async_request_refresh())
self._track(self._coordinator.async_request_refresh())
def _on_trigger_activated(self, value: float) -> None:
"""Handle trigger activation."""
@@ -248,8 +274,8 @@ class BaseTrigger(ABC):
)
# Add history entry for the trigger activation
self.hass.async_create_task(self._coordinator.async_add_trigger_history_entry(self._task_id, trigger_value=value))
self._coordinator.note_trigger_edge(self._task_id)
self._track(self._coordinator.async_add_trigger_history_entry(self._task_id, trigger_value=value))
self._coordinator.note_trigger_edge(self._task_id, recovered=self._recovered_since_reset)
self._request_coordinator_refresh()
# Fire event
@@ -301,7 +327,9 @@ class BaseTrigger(ABC):
# resets the trigger via reset(), which never lands here, so the
# manual flow cannot double-record.
if self.config.get("auto_complete_on_recovery"):
self.hass.async_create_task(self._coordinator.async_auto_complete_on_recovery(self._task_id, value))
# cancel_on_close=False: a completion in flight must land even
# when the recovery coincides with a reload of the entry.
self._track(self._coordinator.async_auto_complete_on_recovery(self._task_id, value), cancel_on_close=False)
def _get_numeric_value(self, state: State) -> float | None:
"""Extract numeric value from state or attribute."""
@@ -323,3 +351,4 @@ class BaseTrigger(ABC):
def reset(self) -> None:
"""Reset the trigger (called after maintenance completion)."""
self._triggered = False
self._recovered_since_reset = False
@@ -247,8 +247,8 @@ class CompoundTrigger(BaseTrigger):
current_value=None,
trigger_entity_id=None,
)
self.hass.async_create_task(self._coordinator.async_add_trigger_history_entry(self._task_id, trigger_value=None))
self._coordinator.note_trigger_edge(self._task_id)
self._track(self._coordinator.async_add_trigger_history_entry(self._task_id, trigger_value=None))
self._coordinator.note_trigger_edge(self._task_id, recovered=self._recovered_since_reset)
self._request_coordinator_refresh()
self.hass.bus.async_fire(
EVENT_TRIGGER_ACTIVATED,
@@ -17,14 +17,14 @@ import logging
from datetime import datetime, timedelta
from typing import TYPE_CHECKING, Any
from homeassistant.core import CALLBACK_TYPE, Event, HomeAssistant, State, callback
from homeassistant.core import Event, HomeAssistant, State, callback
from homeassistant.helpers.event import (
EventStateChangedData,
async_track_state_change_event,
async_track_time_interval,
)
from ...const import UNAVAILABLE_STATES
from ...helpers.managed_timer import ManagedTimer
if TYPE_CHECKING:
from ...sensor import MaintenanceSensor
@@ -88,7 +88,7 @@ class RuntimeTrigger(BaseTrigger):
else:
self._on_states = _DEFAULT_ON_STATES
self._unsub_periodic: CALLBACK_TYPE | None = None
self._periodic_timer = ManagedTimer(hass, f"RuntimeTrigger:{self.entity_id}:persist")
async def async_setup(self) -> None:
"""Set up runtime trigger with state restoration."""
@@ -160,17 +160,11 @@ class RuntimeTrigger(BaseTrigger):
def _start_periodic_timer(self) -> None:
"""Start the periodic persistence timer."""
self._unsub_periodic = async_track_time_interval(
self.hass,
self._periodic_callback,
_PERSIST_INTERVAL,
)
self._periodic_timer.schedule_interval(_PERSIST_INTERVAL, self._periodic_callback)
async def async_teardown(self) -> None:
"""Remove listeners and periodic timer."""
if self._unsub_periodic is not None:
self._unsub_periodic()
self._unsub_periodic = None
self._periodic_timer.close()
await super().async_teardown()
@callback
@@ -199,7 +193,7 @@ class RuntimeTrigger(BaseTrigger):
self._on_since_dt = now
self._on_since = now.isoformat()
self._session_booked = 0.0
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
elif not self._is_on(new_val) and self._on_since_dt is not None:
# Restored anchor but the device APPEARS off (deferred setup
# kept the anchor, then the first real state is OFF). Without
@@ -211,7 +205,7 @@ class RuntimeTrigger(BaseTrigger):
self._on_since_dt = None
self._on_since = None
self._session_booked = 0.0
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
self._update_evaluation()
return
@@ -224,7 +218,7 @@ class RuntimeTrigger(BaseTrigger):
self._on_since_dt = None
self._on_since = None
self._session_booked = 0.0
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
if not self._logged_unavailable:
_LOGGER.warning(
"Runtime trigger entity %s became %s (runtime paused)",
@@ -252,7 +246,7 @@ class RuntimeTrigger(BaseTrigger):
self._on_since_dt = None
self._on_since = None
self._session_booked = 0.0
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
_LOGGER.debug(
"Runtime trigger: %s turned OFF (accumulated=%.2fh)",
self.entity_id,
@@ -264,7 +258,7 @@ class RuntimeTrigger(BaseTrigger):
self._on_since_dt = now
self._on_since = now.isoformat()
self._session_booked = 0.0
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
_LOGGER.debug(
"Runtime trigger: %s turned ON (tracking started)",
self.entity_id,
@@ -278,7 +272,7 @@ class RuntimeTrigger(BaseTrigger):
self._on_since_dt = None
self._on_since = None
self._session_booked = 0.0
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
self._update_evaluation()
@@ -343,7 +337,7 @@ class RuntimeTrigger(BaseTrigger):
self._accumulate_elapsed(now)
self._on_since_dt = now
self._on_since = now.isoformat()
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
# Re-evaluate (runtime may have crossed threshold)
self._update_evaluation()
@@ -354,6 +348,11 @@ class RuntimeTrigger(BaseTrigger):
self._accumulated_seconds / 3600.0,
)
def _persist_runtime_soon(self) -> None:
"""Fire-and-forget persist, owned by the trigger (cancelled at teardown;
the Store write itself happens synchronously at task start)."""
self._track(self._persist_runtime())
async def _persist_runtime(self) -> None:
"""Persist accumulated runtime and on_since to the Store."""
data: dict[str, Any] = {
@@ -382,7 +381,7 @@ class RuntimeTrigger(BaseTrigger):
now = dt_util.utcnow()
self._on_since_dt = now
self._on_since = now.isoformat()
self.hass.async_create_task(self._persist_runtime())
self._persist_runtime_soon()
_LOGGER.debug(
"Runtime trigger reset: %s (accumulated hours cleared)",
self.entity_id,
@@ -6,15 +6,15 @@ import logging
from datetime import datetime
from typing import TYPE_CHECKING, Any
from homeassistant.core import CALLBACK_TYPE, Event, HomeAssistant, callback
from homeassistant.core import Event, HomeAssistant, callback
from homeassistant.helpers.event import (
EventStateChangedData,
async_call_later,
async_track_state_change_event,
)
from homeassistant.util import dt as dt_util
from ...const import UNAVAILABLE_STATES
from ...helpers.managed_timer import ManagedTimer
if TYPE_CHECKING:
from ...sensor import MaintenanceSensor
@@ -44,7 +44,6 @@ class StateChangeTrigger(BaseTrigger):
_pending_state: str | None = None
_pending_since: str | None = None
_interrupted_pending: str | None = None
_timer_cancel: CALLBACK_TYPE | None = None
def __init__(
self,
@@ -77,7 +76,7 @@ class StateChangeTrigger(BaseTrigger):
# sensors glitching for seconds at night) and the cycle counter
# (a flicker is not a wash cycle).
self._for_minutes: int = int(trigger_config.get("trigger_for_minutes", 0) or 0)
self._timer_cancel: CALLBACK_TYPE | None = None
self._hold_timer = ManagedTimer(hass, f"StateChangeTrigger:{self.entity_id}:hold")
# A window cut short by an unavailability blip — the only case that
# may re-open on recovery (see _handle_state_transition).
self._interrupted_pending: str | None = None
@@ -225,53 +224,45 @@ class StateChangeTrigger(BaseTrigger):
def _start_pending(self, new_val: str) -> None:
"""(Re)open the hold window for *new_val* — commits when the timer fires."""
self._cancel_timer()
self._hold_timer.cancel()
self._pending_state = new_val
self._pending_since = dt_util.utcnow().isoformat()
self._persist_runtime_soon()
self._start_hold_timer()
def _start_hold_timer(self, remaining_seconds: float | None = None) -> None:
self._cancel_timer()
duration = remaining_seconds if remaining_seconds is not None else self._for_minutes * 60
self._hold_timer.schedule(duration, self._hold_timer_fired)
@callback
def _timer_fired(_now: datetime) -> None:
pending = self._pending_state
self._pending_state = None
self._pending_since = None
self._timer_cancel = None
if pending is None:
return
# Safety net: only commit while the state still holds.
live = self.hass.states.get(self.entity_id)
if live is None or _norm_state(live.state) != _norm_state(pending):
self._persist_runtime_soon()
return
_LOGGER.debug(
"State hold timer fired: %s held %r for %d min",
self.entity_id,
pending,
self._for_minutes,
)
self._commit_transition(pending, None)
self._timer_cancel = async_call_later(self.hass, duration, _timer_fired)
@callback
def _hold_timer_fired(self, _now: datetime) -> None:
pending = self._pending_state
self._pending_state = None
self._pending_since = None
if pending is None:
return
# Safety net: only commit while the state still holds.
live = self.hass.states.get(self.entity_id)
if live is None or _norm_state(live.state) != _norm_state(pending):
self._persist_runtime_soon()
return
_LOGGER.debug(
"State hold timer fired: %s held %r for %d min",
self.entity_id,
pending,
self._for_minutes,
)
self._commit_transition(pending, None)
def _clear_pending(self) -> None:
"""Abandon the hold window (state moved on before it elapsed)."""
if self._pending_state is None and self._timer_cancel is None:
if self._pending_state is None and not self._hold_timer.pending:
return
self._cancel_timer()
self._hold_timer.cancel()
self._pending_state = None
self._pending_since = None
self._persist_runtime_soon()
def _cancel_timer(self) -> None:
if self._timer_cancel is not None:
self._timer_cancel()
self._timer_cancel = None
def _commit_transition(self, new_val: str, old_val: str | None) -> None:
"""Count one matching transition (immediately, or after its hold)."""
self._pending_state = None
@@ -442,7 +433,7 @@ class StateChangeTrigger(BaseTrigger):
def _persist_runtime_soon(self) -> None:
if self.hass.is_running:
self.hass.async_create_task(self._persist_runtime())
self._track(self._persist_runtime())
async def _persist_runtime(self) -> None:
"""Persist the full runtime dict (count + hold window) to the Store.
@@ -462,13 +453,13 @@ class StateChangeTrigger(BaseTrigger):
async def async_teardown(self) -> None:
"""Clean up the hold timer on teardown."""
self._cancel_timer()
self._hold_timer.close()
await super().async_teardown()
def reset(self) -> None:
"""Reset trigger, counter and any running hold window."""
super().reset()
self._cancel_timer()
self._hold_timer.cancel()
self._pending_state = None
self._pending_since = None
self._interrupted_pending = None
@@ -6,13 +6,13 @@ import logging
from datetime import datetime
from typing import TYPE_CHECKING, Any
from homeassistant.core import CALLBACK_TYPE, HomeAssistant, callback
from homeassistant.helpers.event import async_call_later
from homeassistant.core import HomeAssistant, callback
from homeassistant.util import dt as dt_util
if TYPE_CHECKING:
from ...sensor import MaintenanceSensor
from ...helpers.managed_timer import ManagedTimer
from ...helpers.trigger_fallback import threshold_exceeds
from .base_trigger import BaseTrigger
@@ -45,7 +45,7 @@ class ThresholdTrigger(BaseTrigger):
self._for_minutes: int = trigger_config.get("trigger_for_minutes", 0)
self._threshold_exceeded = False
self._timer_cancel: CALLBACK_TYPE | None = None
self._for_timer = ManagedTimer(hass, f"ThresholdTrigger:{self.entity_id}:for")
# Restore persisted exceeded-since timestamp (survives HA restarts)
from ...helpers.dates import parse_persisted_utc
@@ -99,7 +99,7 @@ class ThresholdTrigger(BaseTrigger):
self._exceeded_since = dt_util.utcnow().isoformat()
self._exceeded_since_dt = None
if self.hass.is_running:
self.hass.async_create_task(self._persist_exceeded_since())
self._track(self._persist_exceeded_since())
self._start_for_timer()
return False
# Timer running or already triggered
@@ -111,9 +111,9 @@ class ThresholdTrigger(BaseTrigger):
self._exceeded_since = None
self._exceeded_since_dt = None
if self.hass.is_running:
self.hass.async_create_task(self._persist_exceeded_since())
self._track(self._persist_exceeded_since())
self._threshold_exceeded = False
self._cancel_timer()
self._for_timer.cancel()
return False
def _start_for_timer(self, remaining_seconds: float | None = None) -> None:
@@ -124,45 +124,37 @@ class ThresholdTrigger(BaseTrigger):
The exceeded-since timestamp is persisted so the timer survives HA
restarts.
"""
self._cancel_timer()
duration = remaining_seconds if remaining_seconds is not None else self._for_minutes * 60
self._for_timer.schedule(duration, self._for_timer_fired)
@callback
def _timer_fired(_now: datetime) -> None:
"""Handle timer completion."""
# Safety net (mirrors the state_change hold timer): only commit
# while the premise still HOLDS. _threshold_exceeded is cleared
# only by a numeric in-range reading, so a sensor that went
# unavailable right after crossing kept it True and the timer
# activated on a value nobody had observed for the whole window
# (bug audit 2026-08-22). Discard the window entirely — a bare
# return would leave the latch set and evaluate() would swallow
# every future exceeding reading; the next one re-arms fresh.
state = self.hass.states.get(self.entity_id)
live = self._get_numeric_value(state) if state is not None else None
if live is None or not self._value_exceeds_threshold(live):
self._threshold_exceeded = False
self._exceeded_since = None
self._exceeded_since_dt = None
if self.hass.is_running:
self.hass.async_create_task(self._persist_exceeded_since())
return
if self._threshold_exceeded:
_LOGGER.debug(
"Threshold for-timer fired: %s (%d min)",
self.entity_id,
self._for_minutes,
)
self._triggered = True
self._on_trigger_activated(self._current_value or 0.0)
self._timer_cancel = async_call_later(self.hass, duration, _timer_fired)
def _cancel_timer(self) -> None:
"""Cancel the for-duration timer."""
if self._timer_cancel is not None:
self._timer_cancel()
self._timer_cancel = None
@callback
def _for_timer_fired(self, _now: datetime) -> None:
"""Handle for-timer completion."""
# Safety net (mirrors the state_change hold timer): only commit
# while the premise still HOLDS. _threshold_exceeded is cleared
# only by a numeric in-range reading, so a sensor that went
# unavailable right after crossing kept it True and the timer
# activated on a value nobody had observed for the whole window
# (bug audit 2026-08-22). Discard the window entirely — a bare
# return would leave the latch set and evaluate() would swallow
# every future exceeding reading; the next one re-arms fresh.
state = self.hass.states.get(self.entity_id)
live = self._get_numeric_value(state) if state is not None else None
if live is None or not self._value_exceeds_threshold(live):
self._threshold_exceeded = False
self._exceeded_since = None
self._exceeded_since_dt = None
if self.hass.is_running:
self._track(self._persist_exceeded_since())
return
if self._threshold_exceeded:
_LOGGER.debug(
"Threshold for-timer fired: %s (%d min)",
self.entity_id,
self._for_minutes,
)
self._triggered = True
self._on_trigger_activated(self._current_value or 0.0)
async def _persist_exceeded_since(self) -> None:
"""Persist exceeded-since timestamp for survival across restarts."""
@@ -174,7 +166,7 @@ class ThresholdTrigger(BaseTrigger):
async def async_teardown(self) -> None:
"""Clean up timer on teardown."""
self._cancel_timer()
self._for_timer.close()
await super().async_teardown()
def reset(self) -> None:
@@ -184,5 +176,5 @@ class ThresholdTrigger(BaseTrigger):
self._exceeded_since = None
self._exceeded_since_dt = None
if self.hass.is_running:
self.hass.async_create_task(self._persist_exceeded_since())
self._cancel_timer()
self._track(self._persist_exceeded_since())
self._for_timer.cancel()