From 0715eb08293dfa3ad7bed62ffbd8b814e3f9f295 Mon Sep 17 00:00:00 2001 From: Tim Schindler Date: Tue, 21 Jul 2026 20:59:10 +0200 Subject: [PATCH 1/2] test: remove timing race from coordinator observe-mode tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The observe tests raced a background thread (5 ms wall-clock sleep) against the coordinator reaching its resubscribe. try_enter_observe_mode clears _notified before subscribing, so a notify only counts if it lands between that clear and the end of the grace sleep. When the intervening update cycle (reconnect, cache sweep, executor hops) took longer than 5 ms, every notify was cleared and the mode stayed 'poll' — test_reconnect_from_observe_mode_resubscribes_immediately failed roughly 1 in 7 runs. Replace the threads with FakeObserveSession.notify_on_subscribe, which delivers the rep synchronously from subscribe() the way a real device answers a subscription. That puts the notify inside the grace window by construction rather than by timing. Tests that require the resubscribe to fail now set it to None explicitly instead of relying on the absence of a racing thread. No production code changed; no assertion weakened, and no sleeps, retries or timeouts added. Claude-Session: https://claude.ai/code/session_01PUSU6tDjHjtbExtyXDPT3N --- tests/localthings/conftest.py | 13 ++++ tests/localthings/test_coordinator.py | 108 ++++++-------------------- 2 files changed, 35 insertions(+), 86 deletions(-) diff --git a/tests/localthings/conftest.py b/tests/localthings/conftest.py index 4674f3c..f4fcc56 100644 --- a/tests/localthings/conftest.py +++ b/tests/localthings/conftest.py @@ -5,6 +5,7 @@ import json from pathlib import Path from unittest.mock import patch +import cbor2 import pytest from pytest_homeassistant_custom_component.common import MockConfigEntry @@ -150,12 +151,24 @@ class FakeObserveSession: self.subscribed: list[str] = [] self.fail_hrefs: set[str] = set() self.closed = False + # When set to a rep dict, subscribe() immediately delivers that rep + # as an OBSERVE notification for the href — what a real device does + # when it answers a subscription with the current representation. + # `try_enter_observe_mode` clears its notified set before it + # subscribes, so this is the only way to deliver a notify that + # reliably lands inside the grace period: a notify raced in from + # another thread can be wiped by that clear (or arrive after the + # grace period ends) depending on scheduling. Set to None to model + # a device that answers subscriptions but never notifies. + self.notify_on_subscribe: dict | None = None def subscribe(self, path_segs): href = '/' + '/'.join(path_segs) if href in self.fail_hrefs: raise ConnectionError("subscribe failed") self.subscribed.append(href) + if self.notify_on_subscribe is not None and self.on_notification is not None: + self.on_notification(href, cbor2.dumps(self.notify_on_subscribe)) return b'\x01' def refresh_observes(self, paths): diff --git a/tests/localthings/test_coordinator.py b/tests/localthings/test_coordinator.py index 20c7c54..6b1572f 100644 --- a/tests/localthings/test_coordinator.py +++ b/tests/localthings/test_coordinator.py @@ -1,8 +1,6 @@ """Tests for the LocalThingsCoordinator.""" from __future__ import annotations -import threading -import time from datetime import timedelta from unittest.mock import AsyncMock, patch @@ -261,20 +259,15 @@ async def test_enters_observe_mode_when_hot_warm_hrefs_notify( hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs # try_enter_observe_mode clears prior notifications as soon as it - # starts, so notifies must land *during* its grace-period sleep (same - # pattern test_observe.py uses), not before the call. - def _notify_during_grace_period(): - time.sleep(0.005) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - - notifier = threading.Thread(target=_notify_during_grace_period, daemon=True) - notifier.start() + # starts, so notifies must land *during* its grace-period sleep. The + # fake session delivers one synchronously from subscribe() (see + # FakeObserveSession.notify_on_subscribe in conftest.py), which puts + # them inside that window by construction rather than by timing. + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator._observe.mode == MODE_OBSERVE @@ -300,18 +293,11 @@ async def test_reconnect_while_observe_mode_downgrades_to_poll( # Get the coordinator into observe mode the same way # test_enters_observe_mode_when_hot_warm_hrefs_notify does. - def _notify_during_grace_period(): - time.sleep(0.005) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - - notifier = threading.Thread(target=_notify_during_grace_period, daemon=True) - notifier.start() + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator.observe_mode == MODE_OBSERVE @@ -324,12 +310,12 @@ async def test_reconnect_while_observe_mode_downgrades_to_poll( # Simulate the existing "poll failed, reconnecting" branch: _poll_once # fails once (triggering the reconnect/backoff path), then succeeds. - # No fresh notifies are supplied for the immediate resubscribe attempt - # this now triggers, so its grace period (shortened for tests by the - # `_fast_coordinator_timers` autouse fixture — see conftest.py) times - # out and mode stays 'poll' — this test only asserts the tear-down half of - # the fix; test_reconnect_from_observe_mode_resubscribes_immediately - # covers the successful-immediate-resubscribe half. + # The device stops notifying on subscribe, so the immediate resubscribe + # attempt this now triggers gets nothing back and mode stays 'poll' — + # this test only asserts the tear-down half of the fix; + # test_reconnect_from_observe_mode_resubscribes_immediately covers the + # successful-immediate-resubscribe half. + fake.notify_on_subscribe = None with ( patch( 'custom_components.localthings.coordinator.LocalThingsCoordinator._poll_once', @@ -362,18 +348,11 @@ async def test_poll_timeout_skips_reconnect_when_push_is_healthy( coordinator: LocalThingsCoordinator = hass.data[DOMAIN][mock_entry.entry_id] hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs - def _notify_during_grace_period(): - time.sleep(0.005) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - - notifier = threading.Thread(target=_notify_during_grace_period, daemon=True) - notifier.start() + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator.observe_mode == MODE_OBSERVE @@ -408,18 +387,11 @@ async def test_poll_timeout_reconnects_after_consecutive_limit_even_with_push( coordinator: LocalThingsCoordinator = hass.data[DOMAIN][mock_entry.entry_id] hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs - def _notify_during_grace_period(): - time.sleep(0.005) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - - notifier = threading.Thread(target=_notify_during_grace_period, daemon=True) - notifier.start() + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator.observe_mode == MODE_OBSERVE @@ -441,7 +413,10 @@ async def test_poll_timeout_reconnects_after_consecutive_limit_even_with_push( assert coordinator.observe_mode == MODE_OBSERVE # The final consecutive timeout crosses the limit and triggers the - # existing reconnect path, which then succeeds and downgrades. + # existing reconnect path, which then succeeds and downgrades. The + # device stops notifying on subscribe, so the resubscribe attempt that + # follows the reconnect finds nothing and mode stays 'poll'. + fake.notify_on_subscribe = None with ( patch( 'custom_components.localthings.coordinator.LocalThingsCoordinator._poll_once', @@ -473,18 +448,11 @@ async def test_poll_timeout_counter_resets_when_push_is_healthy_again( coordinator: LocalThingsCoordinator = hass.data[DOMAIN][mock_entry.entry_id] hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs - def _notify_during_grace_period(): - time.sleep(0.005) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - - notifier = threading.Thread(target=_notify_during_grace_period, daemon=True) - notifier.start() + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator.observe_mode == MODE_OBSERVE @@ -526,21 +494,11 @@ async def test_reconnect_from_observe_mode_resubscribes_immediately( coordinator: LocalThingsCoordinator = hass.data[DOMAIN][mock_entry.entry_id] hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs - def _notify(delay: float = 0.005) -> threading.Thread: - def _run(): - time.sleep(delay) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - t = threading.Thread(target=_run, daemon=True) - t.start() - return t - - notifier = _notify() + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator.observe_mode == MODE_OBSERVE @@ -563,10 +521,8 @@ async def test_reconnect_from_observe_mode_resubscribes_immediately( new=AsyncMock(), ), ): - notifier = _notify() await coordinator.async_request_refresh() await hass.async_block_till_done() - notifier.join() assert coordinator.observe_mode == MODE_OBSERVE @@ -587,21 +543,11 @@ async def test_sweep_mismatch_never_downgrades_a_live_observe_session( hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs assert len(hrefs) >= 2 - def _notify(delay: float = 0.005) -> threading.Thread: - def _run(): - time.sleep(delay) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - t = threading.Thread(target=_run, daemon=True) - t.start() - return t - - notifier = _notify() + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator.observe_mode == MODE_OBSERVE @@ -637,21 +583,11 @@ async def test_sweep_mismatch_forces_subpolls_on_a_live_observe_session( hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs assert len(hrefs) >= 2 - def _notify(delay: float = 0.005) -> threading.Thread: - def _run(): - time.sleep(delay) - for href in hrefs: - fake.on_notification(href, cbor2.dumps({'notified': True})) - t = threading.Thread(target=_run, daemon=True) - t.start() - return t - - notifier = _notify() + fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode, fake, hrefs, 0.02, 0.8, ) - notifier.join() assert entered is True assert coordinator.observe_mode == MODE_OBSERVE From 8dadea6796d16bc37463cff07b9b05973f69e120 Mon Sep 17 00:00:00 2001 From: Tim Schindler Date: Wed, 22 Jul 2026 07:48:43 +0200 Subject: [PATCH 2/2] test: address Copilot review on observe fixture typing and comment - Parameterize FakeObserveSession.notify_on_subscribe as dict[str, Any] | None to match the surrounding fully-typed attributes (the dict is delivered as the OBSERVE rep to on_notification); add the typing.Any import. - Reword the test_coordinator comment: subscribe() delivers a notification for every href it subscribes to (when notify_on_subscribe is set), not one. --- tests/localthings/conftest.py | 13 ++++++++----- tests/localthings/test_coordinator.py | 12 +++++++----- 2 files changed, 15 insertions(+), 10 deletions(-) diff --git a/tests/localthings/conftest.py b/tests/localthings/conftest.py index f4fcc56..a5a5e62 100644 --- a/tests/localthings/conftest.py +++ b/tests/localthings/conftest.py @@ -3,6 +3,7 @@ from __future__ import annotations import json from pathlib import Path +from typing import Any from unittest.mock import patch import cbor2 @@ -156,11 +157,13 @@ class FakeObserveSession: # when it answers a subscription with the current representation. # `try_enter_observe_mode` clears its notified set before it # subscribes, so this is the only way to deliver a notify that - # reliably lands inside the grace period: a notify raced in from - # another thread can be wiped by that clear (or arrive after the - # grace period ends) depending on scheduling. Set to None to model - # a device that answers subscriptions but never notifies. - self.notify_on_subscribe: dict | None = None + # reliably counts: it lands synchronously via on_notification, + # after that clear and before the post-sleep fraction check — a + # notify raced in from another thread can be wiped by the clear + # (or arrive after the check) depending on scheduling. Set to + # None to model a device that answers subscriptions but never + # notifies. + self.notify_on_subscribe: dict[str, Any] | None = None def subscribe(self, path_segs): href = '/' + '/'.join(path_segs) diff --git a/tests/localthings/test_coordinator.py b/tests/localthings/test_coordinator.py index 6b1572f..7739da5 100644 --- a/tests/localthings/test_coordinator.py +++ b/tests/localthings/test_coordinator.py @@ -258,11 +258,13 @@ async def test_enters_observe_mode_when_hot_warm_hrefs_notify( coordinator: LocalThingsCoordinator = hass.data[DOMAIN][mock_entry.entry_id] hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs - # try_enter_observe_mode clears prior notifications as soon as it - # starts, so notifies must land *during* its grace-period sleep. The - # fake session delivers one synchronously from subscribe() (see - # FakeObserveSession.notify_on_subscribe in conftest.py), which puts - # them inside that window by construction rather than by timing. + # try_enter_observe_mode clears prior notifications before it + # subscribes, then checks the notified fraction after the grace-period + # sleep — a notify counts as long as it arrives in that window, not + # specifically during the sleep itself. The fake session delivers one + # synchronously from subscribe() for every href it subscribes to (see + # FakeObserveSession.notify_on_subscribe in conftest.py), which lands + # it in that window by construction rather than by timing. fake.notify_on_subscribe = {'notified': True} entered = await hass.async_add_executor_job( coordinator._observe.try_enter_observe_mode,