observe: keep silent hrefs on the sub-poll cadence
Issue #92: try_enter_observe_mode treated every subscribed href as push-covered once SUCCESS_FRACTION cleared, including ones that never notified. Those then skipped _run_subpolls for the rest of the session. Record subscribed-but-silent hrefs on fallback_hrefs (idle in observe mode) and keep just that set on the hot/warm cadence.
This commit is contained in:
@@ -975,16 +975,23 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
|||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
async def _run_subpolls(self, force: bool = False) -> None:
|
async def _run_subpolls(self, force: bool = False) -> None:
|
||||||
"""Poll hot/warm hrefs in the gaps between summary polls. No-op in
|
"""Poll hot/warm hrefs in the gaps between summary polls.
|
||||||
observe-primary mode (those hrefs are already covered by push)
|
|
||||||
unless `force` is set -- set when this cycle's sweep found the
|
In observe-primary mode this is a no-op for hrefs the device is
|
||||||
cache disagreeing with a still-live observe session (see
|
actually pushing, unless `force` is set (sweep disagreed with the
|
||||||
log_sweep_discrepancies): a bounded fallback for a channel gone
|
cache -- see log_sweep_discrepancies). Hrefs that were subscribed
|
||||||
silent without a reconnect."""
|
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:
|
if self._observe.mode == MODE_OBSERVE and not force:
|
||||||
return
|
silent = self._observe.fallback_hrefs
|
||||||
hot = self._hot_hrefs
|
if not silent:
|
||||||
warm = self._warm_hrefs
|
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:
|
if not hot and not warm:
|
||||||
return
|
return
|
||||||
step = self._SUBPOLL_STEP_S
|
step = self._SUBPOLL_STEP_S
|
||||||
|
|||||||
@@ -97,6 +97,11 @@ class ObserveManager:
|
|||||||
# Wakes try_enter_observe_mode's grace wait early once enough hrefs
|
# Wakes try_enter_observe_mode's grace wait early once enough hrefs
|
||||||
# have notified. Guards only `_notified` mutations + the `wait_for`.
|
# have notified. Guards only `_notified` mutations + the `wait_for`.
|
||||||
self._notify_cond = threading.Condition()
|
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.fallback_hrefs: set[str] = set()
|
||||||
self._on_applied: Callable[[str, dict, str], None] | None = None
|
self._on_applied: Callable[[str, dict, str], None] | None = None
|
||||||
self._refresh_task: ObserveRefreshTask | None = None
|
self._refresh_task: ObserveRefreshTask | None = None
|
||||||
@@ -270,6 +275,12 @@ class ObserveManager:
|
|||||||
#294) -- committing against a session a reconnect already replaced
|
#294) -- committing against a session a reconnect already replaced
|
||||||
would claim observe mode with nothing left to notice it's dead."""
|
would claim observe mode with nothing left to notice it's dead."""
|
||||||
self.subscribed_hrefs = set(subscribed)
|
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._set_mode(MODE_OBSERVE)
|
||||||
self.start_refresh_task(session)
|
self.start_refresh_task(session)
|
||||||
|
|
||||||
|
|||||||
@@ -1050,6 +1050,62 @@ async def test_sweep_mismatch_forces_subpolls_on_a_live_observe_session(
|
|||||||
mock_subpolls.assert_called_once_with(force=True)
|
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(
|
async def test_write_marks_href_pending_before_post(
|
||||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||||
) -> None:
|
) -> None:
|
||||||
|
|||||||
@@ -190,6 +190,7 @@ def test_try_enter_observe_mode_succeeds_when_all_hrefs_notify():
|
|||||||
assert entered is True
|
assert entered is True
|
||||||
assert mgr.mode == "observe"
|
assert mgr.mode == "observe"
|
||||||
assert mgr.subscribed_hrefs == set(hrefs)
|
assert mgr.subscribed_hrefs == set(hrefs)
|
||||||
|
assert mgr.fallback_hrefs == set()
|
||||||
finally:
|
finally:
|
||||||
mgr.close()
|
mgr.close()
|
||||||
|
|
||||||
@@ -218,6 +219,23 @@ def test_try_enter_observe_mode_falls_back_when_subscribe_fails_for_all():
|
|||||||
assert mgr.mode == "poll"
|
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():
|
def test_try_enter_observe_mode_meets_success_fraction_with_partial_notifies():
|
||||||
mgr = _manager()
|
mgr = _manager()
|
||||||
session = _FakeSession()
|
session = _FakeSession()
|
||||||
@@ -243,6 +261,7 @@ def test_try_enter_observe_mode_meets_success_fraction_with_partial_notifies():
|
|||||||
|
|
||||||
assert entered is True
|
assert entered is True
|
||||||
assert mgr.mode == "observe"
|
assert mgr.mode == "observe"
|
||||||
|
assert mgr.fallback_hrefs == {hrefs[3]}
|
||||||
finally:
|
finally:
|
||||||
mgr.close()
|
mgr.close()
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user