mirror of
https://github.com/home-assistant/core.git
synced 2026-08-24 10:13:52 -05:00
Improve condition history manager (#174069)
This commit is contained in:
@@ -491,7 +491,10 @@ class _HistoryPrimingManager:
|
||||
tracking its entities, or the read could miss a change still queued in the
|
||||
recorder and compute too generous an anchor. A condition therefore never
|
||||
rides a flush that was already running when it arrived (the lobby); it waits
|
||||
that one out and joins the next.
|
||||
that one out and joins the next, and re-attempts if the flush it rode was
|
||||
cancelled before completing. This mirrors `ReloadServiceHelper` minus its
|
||||
target de-duplication, which does not apply because each condition reads its
|
||||
own entities.
|
||||
"""
|
||||
|
||||
def __init__(self, hass: HomeAssistant) -> None:
|
||||
@@ -499,6 +502,7 @@ class _HistoryPrimingManager:
|
||||
self._hass = hass
|
||||
self._flush_condition = asyncio.Condition()
|
||||
self._flushing = False
|
||||
self._flush_ok = False
|
||||
self._query_lock = asyncio.Lock()
|
||||
|
||||
async def async_prime[_T](
|
||||
@@ -520,29 +524,30 @@ class _HistoryPrimingManager:
|
||||
if self._flushing:
|
||||
await self._flush_condition.wait()
|
||||
|
||||
do_flush = False
|
||||
while True:
|
||||
async with self._flush_condition:
|
||||
if not self._flushing:
|
||||
# First past the lobby this generation: we run the flush.
|
||||
self._flushing = True
|
||||
do_flush = True
|
||||
break
|
||||
# A peer began a fresh flush after we cleared the lobby; it
|
||||
# covers us too, so wait for it and ride it.
|
||||
# A peer began a fresh flush after we cleared the lobby; ride it.
|
||||
await self._flush_condition.wait()
|
||||
break
|
||||
|
||||
if not do_flush:
|
||||
return
|
||||
if self._flush_ok:
|
||||
return
|
||||
# The flush we waited for was cancelled before completing (its owner
|
||||
# timed out): loop and start or wait for a fresh one rather than read
|
||||
# against a queue that was never flushed.
|
||||
|
||||
instance = get_instance(self._hass)
|
||||
flushed = False
|
||||
try:
|
||||
if (commit_future := instance.async_get_commit_future()) is not None:
|
||||
await commit_future
|
||||
flushed = True
|
||||
finally:
|
||||
async with self._flush_condition:
|
||||
self._flushing = False
|
||||
self._flush_ok = flushed
|
||||
self._flush_condition.notify_all()
|
||||
|
||||
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
"""Test the condition helper."""
|
||||
|
||||
import asyncio
|
||||
from collections.abc import Mapping
|
||||
from collections.abc import Callable, Mapping
|
||||
from contextlib import AbstractContextManager, nullcontext as does_not_raise
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timedelta
|
||||
@@ -5908,6 +5908,19 @@ async def test_history_priming_manager_serializes_queries(
|
||||
assert max_running == 1
|
||||
|
||||
|
||||
async def _advance_until(predicate: Callable[[], bool]) -> None:
|
||||
"""Pump the event loop until predicate holds, failing if it never does.
|
||||
|
||||
Avoids coupling tests to an exact number of internal await hops while still
|
||||
failing cleanly rather than hanging on a regression.
|
||||
"""
|
||||
for _ in range(1000):
|
||||
if predicate():
|
||||
return
|
||||
await asyncio.sleep(0)
|
||||
pytest.fail("condition was not reached")
|
||||
|
||||
|
||||
async def test_history_priming_manager_does_not_ride_in_flight_flush(
|
||||
recorder_mock: Recorder, hass: HomeAssistant
|
||||
) -> None:
|
||||
@@ -5917,8 +5930,8 @@ async def test_history_priming_manager_does_not_ride_in_flight_flush(
|
||||
sees them. A condition that started tracking after an in-flight flush began
|
||||
could miss its own just-queued change if it rode that flush, so it waits the
|
||||
flush out and a fresh one is performed for it. Without the lobby step this
|
||||
test fails: the late arrivals would ride the first flush (one flush total)
|
||||
instead of sharing a second, fresh one.
|
||||
test fails: the late arrivals would ride the first flush (it would stay at
|
||||
one flush total) instead of sharing a second, fresh one.
|
||||
"""
|
||||
manager = _HistoryPrimingManager(hass)
|
||||
instance = get_instance(hass)
|
||||
@@ -5936,32 +5949,69 @@ async def test_history_priming_manager_does_not_ride_in_flight_flush(
|
||||
with patch.object(instance, "async_get_commit_future", _spy_commit_future):
|
||||
# C0 claims the flush and is mid-flush (its commit future is pending).
|
||||
c0 = asyncio.create_task(manager.async_prime(_job))
|
||||
for _ in range(10):
|
||||
await asyncio.sleep(0)
|
||||
if flush_futures:
|
||||
break
|
||||
assert len(flush_futures) == 1
|
||||
await _advance_until(lambda: len(flush_futures) == 1)
|
||||
|
||||
# Two conditions arrive while C0's flush runs; they must not ride it.
|
||||
c1 = asyncio.create_task(manager.async_prime(_job))
|
||||
c2 = asyncio.create_task(manager.async_prime(_job))
|
||||
for _ in range(5):
|
||||
await asyncio.sleep(0)
|
||||
# Parked in the lobby: no new flush yet, none finished.
|
||||
assert len(flush_futures) == 1
|
||||
assert not c1.done()
|
||||
assert not c2.done()
|
||||
|
||||
# C0's flush completes; C1 now performs a fresh flush and C2 rides it.
|
||||
# C0's flush completes; C1 then performs a fresh flush and C2 rides it.
|
||||
flush_futures[0].set_result(None)
|
||||
assert await c0 == "done"
|
||||
for _ in range(10):
|
||||
await asyncio.sleep(0)
|
||||
# Exactly one fresh flush is shared by C1 and C2, not one each: this is
|
||||
# the assertion that fails without the lobby (it would stay 1).
|
||||
assert len(flush_futures) == 2
|
||||
await _advance_until(lambda: len(flush_futures) == 2)
|
||||
|
||||
flush_futures[1].set_result(None)
|
||||
assert await asyncio.gather(c1, c2) == ["done", "done"]
|
||||
# One fresh flush shared by C1 and C2, not one each (and not C0's stale
|
||||
# one): C1 flushed, C2 rode it.
|
||||
assert len(flush_futures) == 2
|
||||
|
||||
|
||||
async def test_history_priming_manager_retries_after_cancelled_flush(
|
||||
recorder_mock: Recorder, hass: HomeAssistant
|
||||
) -> None:
|
||||
"""A rider re-flushes when the flush it rode was cancelled before completing.
|
||||
|
||||
If the condition performing a generation's shared flush is cancelled by its
|
||||
timeout while awaiting the commit, the riders must not read against the
|
||||
unflushed queue — they perform a fresh flush instead. Without that retry this
|
||||
test fails: the rider would proceed on the cancelled flush and never make a
|
||||
second one.
|
||||
"""
|
||||
manager = _HistoryPrimingManager(hass)
|
||||
instance = get_instance(hass)
|
||||
|
||||
flush_futures: list[asyncio.Future[None]] = []
|
||||
|
||||
def _spy_commit_future() -> asyncio.Future[None]:
|
||||
fut = hass.loop.create_future()
|
||||
flush_futures.append(fut)
|
||||
return fut
|
||||
|
||||
async def _job(_recorder: Recorder) -> str:
|
||||
return "done"
|
||||
|
||||
with patch.object(instance, "async_get_commit_future", _spy_commit_future):
|
||||
# C0 takes the lobby so c1 and c2 form one generation behind it.
|
||||
c0 = asyncio.create_task(manager.async_prime(_job))
|
||||
await _advance_until(lambda: len(flush_futures) == 1)
|
||||
c1 = asyncio.create_task(manager.async_prime(_job))
|
||||
c2 = asyncio.create_task(manager.async_prime(_job))
|
||||
flush_futures[0].set_result(None)
|
||||
assert await c0 == "done"
|
||||
|
||||
# c1 performs the generation's flush (the second one) and c2 rides it.
|
||||
await _advance_until(lambda: len(flush_futures) == 2)
|
||||
|
||||
# c1 is cancelled mid-flush, as its timeout would do. c2 must then run
|
||||
# its own fresh flush rather than ride c1's cancelled one.
|
||||
c1.cancel()
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await c1
|
||||
await _advance_until(lambda: len(flush_futures) == 3)
|
||||
|
||||
flush_futures[2].set_result(None)
|
||||
assert await c2 == "done"
|
||||
|
||||
|
||||
async def test_history_priming_manager_cancelled_lobby_waiter(
|
||||
@@ -5987,14 +6037,11 @@ async def test_history_priming_manager_cancelled_lobby_waiter(
|
||||
|
||||
with patch.object(instance, "async_get_commit_future", _spy_commit_future):
|
||||
c0 = asyncio.create_task(manager.async_prime(_job))
|
||||
for _ in range(10):
|
||||
await asyncio.sleep(0)
|
||||
if flush_futures:
|
||||
break
|
||||
# A second priming parks in the lobby, then its timeout cancels it.
|
||||
await _advance_until(lambda: len(flush_futures) == 1)
|
||||
# A second priming parks in the lobby (reached in one step, as its lock
|
||||
# acquire is uncontended), then its timeout cancels it.
|
||||
waiter = asyncio.create_task(manager.async_prime(_job))
|
||||
for _ in range(3):
|
||||
await asyncio.sleep(0)
|
||||
await asyncio.sleep(0)
|
||||
waiter.cancel()
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await waiter
|
||||
@@ -6003,9 +6050,7 @@ async def test_history_priming_manager_cancelled_lobby_waiter(
|
||||
flush_futures[0].set_result(None)
|
||||
assert await c0 == "done"
|
||||
later = asyncio.create_task(manager.async_prime(_job))
|
||||
for _ in range(10):
|
||||
await asyncio.sleep(0)
|
||||
assert len(flush_futures) == 2
|
||||
await _advance_until(lambda: len(flush_futures) == 2)
|
||||
flush_futures[1].set_result(None)
|
||||
assert await later == "done"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user