diff --git a/custom_components/localthings/coordinator.py b/custom_components/localthings/coordinator.py index a7da859..eabde62 100644 --- a/custom_components/localthings/coordinator.py +++ b/custom_components/localthings/coordinator.py @@ -975,16 +975,23 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]): # ------------------------------------------------------------------ async def _run_subpolls(self, force: bool = False) -> None: - """Poll hot/warm hrefs in the gaps between summary polls. No-op in - observe-primary mode (those hrefs are already covered by push) - unless `force` is set -- set when this cycle's sweep found the - cache disagreeing with a still-live observe session (see - log_sweep_discrepancies): a bounded fallback for a channel gone - silent without a reconnect.""" + """Poll hot/warm hrefs in the gaps between summary polls. + + In observe-primary mode this is a no-op for hrefs the device is + actually pushing, unless `force` is set (sweep disagreed with the + cache -- see log_sweep_discrepancies). Hrefs that were subscribed + but stayed silent through the grace period (issue #92) stay on + the hot/warm cadence via `fallback_hrefs`. + """ if self._observe.mode == MODE_OBSERVE and not force: - return - hot = self._hot_hrefs - warm = self._warm_hrefs + silent = self._observe.fallback_hrefs + if not silent: + return + hot = [h for h in self._hot_hrefs if h in silent] + warm = [h for h in self._warm_hrefs if h in silent] + else: + hot = self._hot_hrefs + warm = self._warm_hrefs if not hot and not warm: return step = self._SUBPOLL_STEP_S diff --git a/custom_components/localthings/observe.py b/custom_components/localthings/observe.py index 1215766..0f547b5 100644 --- a/custom_components/localthings/observe.py +++ b/custom_components/localthings/observe.py @@ -97,6 +97,11 @@ class ObserveManager: # Wakes try_enter_observe_mode's grace wait early once enough hrefs # have notified. Guards only `_notified` mutations + the `wait_for`. self._notify_cond = threading.Condition() + # Idle while polling, except after downgrade_to_poll (every href + # that was subscribed). While in observe mode this is the set of + # hrefs that were subscribed but never notified during the grace + # period (issue #92) -- they stay on the hot/warm sub-poll cadence + # instead of waiting for the 30s /device/0 sweep. self.fallback_hrefs: set[str] = set() self._on_applied: Callable[[str, dict, str], None] | None = None self._refresh_task: ObserveRefreshTask | None = None @@ -270,6 +275,12 @@ class ObserveManager: #294) -- committing against a session a reconnect already replaced would claim observe mode with nothing left to notice it's dead.""" self.subscribed_hrefs = set(subscribed) + with self._notify_cond: + notified = set(self._notified) + # Issue #92: subscribed-but-silent hrefs are counted as covered by + # push if we drop this, but they never emit a notify. Keep them on + # the poll cadence via fallback_hrefs (otherwise idle in observe). + self.fallback_hrefs = set(subscribed) - notified self._set_mode(MODE_OBSERVE) self.start_refresh_task(session) diff --git a/tests/localthings/test_coordinator.py b/tests/localthings/test_coordinator.py index ed1844e..16ec09e 100644 --- a/tests/localthings/test_coordinator.py +++ b/tests/localthings/test_coordinator.py @@ -1050,6 +1050,62 @@ async def test_sweep_mismatch_forces_subpolls_on_a_live_observe_session( mock_subpolls.assert_called_once_with(force=True) +async def test_observe_mode_subpolls_only_silent_hrefs( + hass: HomeAssistant, mock_entry, mock_coordinator_observe_session +) -> None: + """Issue #92: subscribed-but-silent hrefs keep the hot/warm cadence + instead of waiting for the 30s sweep.""" + await hass.config_entries.async_setup(mock_entry.entry_id) + await hass.async_block_till_done() + coordinator: LocalThingsCoordinator = hass.data[DOMAIN][mock_entry.entry_id] + assert coordinator._hot_hrefs + silent = coordinator._hot_hrefs[0] + coordinator._observe.mode = MODE_OBSERVE + coordinator._observe.fallback_hrefs = {silent} + + polled: list[list[str]] = [] + + def _capture(hrefs): + polled.append(list(hrefs)) + + with ( + patch.object(coordinator, "_poll_hrefs_blocking", side_effect=_capture), + patch( + "custom_components.localthings.coordinator.asyncio.sleep", + new_callable=AsyncMock, + ), + ): + await coordinator._run_subpolls() + + assert polled + for batch in polled: + assert set(batch) == {silent} + + +async def test_observe_mode_skips_subpolls_when_nothing_is_silent( + hass: HomeAssistant, mock_entry, mock_coordinator_observe_session +) -> None: + """The observe-mode no-op stays in place when every subscribed href + actually notified -- issue #92 only keeps the silent ones on poll.""" + await hass.config_entries.async_setup(mock_entry.entry_id) + await hass.async_block_till_done() + coordinator: LocalThingsCoordinator = hass.data[DOMAIN][mock_entry.entry_id] + coordinator._observe.mode = MODE_OBSERVE + coordinator._observe.fallback_hrefs = set() + + with ( + patch.object(coordinator, "_poll_hrefs_blocking") as mock_poll, + patch( + "custom_components.localthings.coordinator.asyncio.sleep", + new_callable=AsyncMock, + ) as mock_sleep, + ): + await coordinator._run_subpolls() + + mock_poll.assert_not_called() + mock_sleep.assert_not_called() + + async def test_write_marks_href_pending_before_post( hass: HomeAssistant, mock_entry, mock_coordinator_observe_session ) -> None: diff --git a/tests/localthings/test_observe.py b/tests/localthings/test_observe.py index 87af3e4..ec8e20f 100644 --- a/tests/localthings/test_observe.py +++ b/tests/localthings/test_observe.py @@ -190,6 +190,7 @@ def test_try_enter_observe_mode_succeeds_when_all_hrefs_notify(): assert entered is True assert mgr.mode == "observe" assert mgr.subscribed_hrefs == set(hrefs) + assert mgr.fallback_hrefs == set() finally: mgr.close() @@ -218,6 +219,23 @@ def test_try_enter_observe_mode_falls_back_when_subscribe_fails_for_all(): assert mgr.mode == "poll" +def test_enter_observe_mode_keeps_silent_hrefs_on_fallback(): + """Issue #92: subscribed hrefs that never notified stay on the poll + cadence via fallback_hrefs, rather than being treated as push-covered.""" + mgr = _manager() + session = _FakeSession() + subscribed = {"/a/vs/0", "/b/vs/0", "/c/vs/0"} + mgr.on_notification("/a/vs/0", cbor2.dumps({"x": 1})) + mgr.on_notification("/b/vs/0", cbor2.dumps({"x": 1})) + try: + mgr.enter_observe_mode(session, subscribed) + assert mgr.mode == "observe" + assert mgr.subscribed_hrefs == subscribed + assert mgr.fallback_hrefs == {"/c/vs/0"} + finally: + mgr.close() + + def test_try_enter_observe_mode_meets_success_fraction_with_partial_notifies(): mgr = _manager() session = _FakeSession() @@ -243,6 +261,7 @@ def test_try_enter_observe_mode_meets_success_fraction_with_partial_notifies(): assert entered is True assert mgr.mode == "observe" + assert mgr.fallback_hrefs == {hrefs[3]} finally: mgr.close()