diff --git a/custom_components/localthings/observe.py b/custom_components/localthings/observe.py index e82800f..94f1d20 100644 --- a/custom_components/localthings/observe.py +++ b/custom_components/localthings/observe.py @@ -87,6 +87,7 @@ class ObserveManager: """Blocking — subscribes to every href then sleeps for the whole grace period. Caller must run this in an executor, never on the event loop.""" + self._notified.clear() subscribed: set[str] = set() for href in hrefs: segs = [s for s in href.strip('/').split('/') if s] diff --git a/tests/localthings/test_observe.py b/tests/localthings/test_observe.py index f9cafb2..c6fc68f 100644 --- a/tests/localthings/test_observe.py +++ b/tests/localthings/test_observe.py @@ -1,6 +1,7 @@ """Tests for ObserveManager: write-settle guard and mode defaults.""" from __future__ import annotations +import threading import time import cbor2 @@ -78,15 +79,18 @@ def test_try_enter_observe_mode_succeeds_when_all_hrefs_notify(): session = _FakeSession() hrefs = ['/oven/vs/0', '/power/vs/0'] - def _notify_all(): + def _notify_during_grace_period(): + # Simulate notifications arriving during the grace period + time.sleep(0.005) # Let the sleep start, then notify partway through for href in hrefs: mgr.on_notification(href, cbor2.dumps({'x': 1})) - # Grace period is real time.sleep() in the manager; keep it tiny and - # notify synchronously beforehand so notifies are already recorded - # when the (short) sleep completes. - _notify_all() - entered = mgr.try_enter_observe_mode(session, hrefs, grace_period_s=0.01) + # Start background thread to deliver notifications during grace period + notifier = threading.Thread(target=_notify_during_grace_period, daemon=True) + notifier.start() + + entered = mgr.try_enter_observe_mode(session, hrefs, grace_period_s=0.02) + notifier.join() assert entered is True assert mgr.mode == 'observe' @@ -121,12 +125,20 @@ def test_try_enter_observe_mode_meets_success_fraction_with_partial_notifies(): mgr = _manager() session = _FakeSession() hrefs = ['/a/vs/0', '/b/vs/0', '/c/vs/0', '/d/vs/0'] - for href in hrefs[:3]: # 3/4 = 0.75, below default 0.8 threshold - mgr.on_notification(href, cbor2.dumps({'x': 1})) + + def _notify_partial(): + # Simulate partial notifications arriving during grace period + time.sleep(0.005) + for href in hrefs[:3]: # 3/4 = 0.75 + mgr.on_notification(href, cbor2.dumps({'x': 1})) + + notifier = threading.Thread(target=_notify_partial, daemon=True) + notifier.start() entered = mgr.try_enter_observe_mode( - session, hrefs, grace_period_s=0.01, success_fraction=0.7, + session, hrefs, grace_period_s=0.02, success_fraction=0.7, ) + notifier.join() assert entered is True assert mgr.mode == 'observe' @@ -136,3 +148,44 @@ def test_on_notification_ignores_malformed_cbor(): mgr = _manager() mgr.on_notification('/oven/vs/0', b'\xff\xff\xff not cbor') assert mgr.cache.get('/oven/vs/0') is None + + +def test_try_enter_observe_mode_clears_stale_notifications_on_retry(): + """Regression test: verify second call to try_enter_observe_mode() doesn't + get polluted by notifications from the first call. This happens in real + usage when poll-only mode periodically retries entering observe mode.""" + mgr = _manager() + session = _FakeSession() + hrefs = ['/a/vs/0', '/b/vs/0'] + + # First call: all hrefs notify -> succeeds and enters observe mode + def _notify_all_first(): + time.sleep(0.005) + for href in hrefs: + mgr.on_notification(href, cbor2.dumps({'x': 1})) + + notifier1 = threading.Thread(target=_notify_all_first, daemon=True) + notifier1.start() + entered = mgr.try_enter_observe_mode(session, hrefs, grace_period_s=0.02) + notifier1.join() + assert entered is True + assert mgr.mode == 'observe' + + # Second call: retry with same hrefs, but only one notifies. + # Without the fix (missing self._notified.clear()), the old notifications + # would leak in, making this appear successful (2/2 instead of 1/2). + # With the fix, it should fail because only 1/2 < 0.8. + session.subscribed = [] # reset for clean test + + def _notify_partial_second(): + time.sleep(0.005) + mgr.on_notification(hrefs[0], cbor2.dumps({'x': 2})) # only one notifies + + notifier2 = threading.Thread(target=_notify_partial_second, daemon=True) + notifier2.start() + entered = mgr.try_enter_observe_mode(session, hrefs, grace_period_s=0.02) + notifier2.join() + + assert entered is False + assert mgr.mode == 'poll' + assert mgr.subscribed_hrefs == set()