Merge pull request #306 from mbillow/claude/issue-triage-backoff-xkeddl
Fix reconnect/retry gaps found in issue triage (#291, #287, #294)
This commit is contained in:
@@ -163,6 +163,17 @@ def _fetch_samsung_uuid() -> str:
|
||||
raise RuntimeError(f"UUID not found in {_SAMSUNG_CLOUD_HOST} certificate subject")
|
||||
|
||||
|
||||
def _normalize_pem(text: str) -> str:
|
||||
"""Strip a pasted PEM's BOM, CRLF endings, and blank lines before
|
||||
`cryptography` sees it -- a text editor's copy carries all three and
|
||||
fails with an opaque InvalidHeader, while the same file dumped via
|
||||
`type` doesn't (issue #291)."""
|
||||
text = text.lstrip("\ufeff")
|
||||
text = text.replace("\r\n", "\n").replace("\r", "\n")
|
||||
lines = [line for line in text.split("\n") if line.strip()]
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def _mint_leaf_cert(ca_cert_pem: str, ca_key_pem: str, uuid: str) -> tuple[str, str]:
|
||||
"""Mint a fresh RSA-2048 leaf cert signed by the CA.
|
||||
|
||||
@@ -722,8 +733,10 @@ class LocalThingsConfigFlow(config_entries.ConfigFlow, domain=DOMAIN):
|
||||
if leaf_cert and leaf_key:
|
||||
existing_leaf = (leaf_cert, leaf_key)
|
||||
else:
|
||||
self._ca_cert_pem = user_input[CONF_CA_CERT_PEM].strip()
|
||||
self._ca_key_pem = user_input[CONF_CA_KEY_PEM].strip()
|
||||
# Normalized here, not just before minting: this is also
|
||||
# what gets stored and reused to re-mint the leaf later.
|
||||
self._ca_cert_pem = _normalize_pem(user_input[CONF_CA_CERT_PEM])
|
||||
self._ca_key_pem = _normalize_pem(user_input[CONF_CA_KEY_PEM])
|
||||
|
||||
try:
|
||||
info = await self.hass.async_add_executor_job(
|
||||
|
||||
@@ -15,7 +15,7 @@ from typing import Any
|
||||
import cbor2
|
||||
from homeassistant.config_entries import ConfigEntry
|
||||
from homeassistant.core import HomeAssistant, callback
|
||||
from homeassistant.exceptions import ServiceValidationError
|
||||
from homeassistant.exceptions import HomeAssistantError, ServiceValidationError
|
||||
from homeassistant.helpers import issue_registry as ir
|
||||
from homeassistant.helpers.device_registry import DeviceInfo
|
||||
from homeassistant.helpers.update_coordinator import DataUpdateCoordinator, UpdateFailed
|
||||
@@ -204,6 +204,14 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
self._consecutive_poll_timeouts = 0
|
||||
self._unbound_hrefs: list[str] = []
|
||||
self._reconnect_times: list[float] = []
|
||||
# See _maybe_retry_observe_mode: last_mode_change_ts alone doesn't
|
||||
# move on a failed attempt, so this tracks attempts too.
|
||||
self._last_observe_attempt_ts = 0.0
|
||||
# Set by both reconnect paths (poll and command) that hand back a
|
||||
# session with zero OBSERVE registrations while mode was still
|
||||
# observe; consumed once to trigger an immediate resubscribe
|
||||
# instead of waiting out _RECOVERY_RETRY_S.
|
||||
self._resubscribe_due = False
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Session management (all blocking — must run in executor)
|
||||
@@ -731,29 +739,72 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
ir.async_delete_issue(self.hass, DOMAIN, issue_id)
|
||||
|
||||
async def _attempt_observe_mode(self) -> None:
|
||||
"""Called once, right after first discovery. Blocking (sleeps for
|
||||
the whole grace period) — must run in an executor."""
|
||||
"""Called once, right after first discovery, after a reconnect
|
||||
downgrades from observe, and periodically while polling. Two
|
||||
phases: subscribing holds `_session_lock` (each send is fire-and-
|
||||
forget, not a network round trip); the grace wait that follows
|
||||
does not, so a concurrent command write isn't blocked for the
|
||||
whole ~15s wait (issue #294) -- only the brief subscribe burst.
|
||||
|
||||
The wait can outlast a reconnect elsewhere (the poll path's own
|
||||
recovery, or a command retry), which would otherwise let a stale
|
||||
success commit observe mode against a session that's already been
|
||||
replaced -- claiming "Push" with nothing left to notice it's dead.
|
||||
`self._session is sess` re-checked under the lock right before
|
||||
committing closes that: `sess` keeps the old object alive, so
|
||||
identity can't be recycled onto a new one.
|
||||
"""
|
||||
hrefs = self._hot_hrefs + self._warm_hrefs
|
||||
if not hrefs:
|
||||
return
|
||||
if self._session is None:
|
||||
# _poll_once already connects on a real poll; only fires if the
|
||||
# session was closed out from under us concurrently.
|
||||
await self.hass.async_add_executor_job(self._connect_session)
|
||||
sess = self._session
|
||||
if sess is None:
|
||||
self._last_observe_attempt_ts = time.monotonic()
|
||||
async with self._session_lock:
|
||||
if self._session is None:
|
||||
# _poll_once already connects on a real poll; only fires if
|
||||
# the session was closed out from under us concurrently.
|
||||
await self.hass.async_add_executor_job(self._connect_session)
|
||||
sess = self._session
|
||||
if sess is None:
|
||||
return
|
||||
subscribed = await self.hass.async_add_executor_job(
|
||||
self._observe.subscribe_hrefs, sess, hrefs
|
||||
)
|
||||
if not subscribed:
|
||||
self._observe.abandon_observe_attempt()
|
||||
return
|
||||
await self.hass.async_add_executor_job(
|
||||
self._observe.try_enter_observe_mode,
|
||||
sess,
|
||||
hrefs,
|
||||
self._OBSERVE_GRACE_PERIOD_S,
|
||||
reached = await self.hass.async_add_executor_job(
|
||||
self._observe.await_observe_notifies, subscribed, self._OBSERVE_GRACE_PERIOD_S
|
||||
)
|
||||
async with self._session_lock:
|
||||
stale_session = self._session is not sess
|
||||
if not reached or stale_session:
|
||||
self._observe.abandon_observe_attempt()
|
||||
if stale_session:
|
||||
# A reconnect elsewhere replaced the session while this
|
||||
# attempt waited -- that session has never been tried,
|
||||
# so retry it next cycle instead of leaving it
|
||||
# unsubscribed for up to _RECOVERY_RETRY_S, which
|
||||
# _last_observe_attempt_ts (already stamped above, for
|
||||
# the now-abandoned session) would otherwise throttle
|
||||
# for (issue #294).
|
||||
self._resubscribe_due = True
|
||||
return
|
||||
self._observe.enter_observe_mode(sess, subscribed)
|
||||
|
||||
async def _maybe_retry_observe_mode(self) -> None:
|
||||
"""While in poll-only mode, periodically re-attempt observe mode
|
||||
so a device that gains internet access recovers push automatically."""
|
||||
if time.monotonic() - self._observe.last_mode_change_ts < _RECOVERY_RETRY_S:
|
||||
so a device that gains internet access recovers push automatically.
|
||||
|
||||
Gated on the more recent of the two timestamps, not just
|
||||
`last_mode_change_ts`: `_set_mode` only stamps that on an actual
|
||||
transition, so a device that never successfully enters observe
|
||||
mode would otherwise leave this throttle open forever after the
|
||||
first `_RECOVERY_RETRY_S` window -- re-attempting (and paying the
|
||||
subscribe-burst lock) on every single poll cycle instead of every
|
||||
`_RECOVERY_RETRY_S`.
|
||||
"""
|
||||
last_attempt = max(self._observe.last_mode_change_ts, self._last_observe_attempt_ts)
|
||||
if time.monotonic() - last_attempt < _RECOVERY_RETRY_S:
|
||||
return
|
||||
await self._attempt_observe_mode()
|
||||
|
||||
@@ -805,7 +856,6 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
self._subpoll_task.cancel()
|
||||
self._subpoll_task = None
|
||||
|
||||
just_downgraded_from_observe = False
|
||||
async with self._session_lock:
|
||||
try:
|
||||
resources = await self.hass.async_add_executor_job(self._poll_once)
|
||||
@@ -834,6 +884,13 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
resources = await self.hass.async_add_executor_job(self._poll_once)
|
||||
except Exception as e2:
|
||||
self._log.error("poll failed after reconnect: %s", e2)
|
||||
# Without this, a fully unreachable device left the
|
||||
# connection-mode sensor stuck on "Push" forever -- only
|
||||
# the success branch below ever downgraded it (issue
|
||||
# #287). No just_downgraded_from_observe here: there's no
|
||||
# live session this cycle to resubscribe on.
|
||||
if self._observe.mode == MODE_OBSERVE:
|
||||
self._observe.downgrade_to_poll()
|
||||
snapshot = self._cache.snapshot()
|
||||
# Same precondition as _defer_reconnect_for (issue #254):
|
||||
# degraded-but-successful data only makes sense once
|
||||
@@ -855,7 +912,7 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
"poll and resubscribing on the new session"
|
||||
)
|
||||
self._observe.downgrade_to_poll()
|
||||
just_downgraded_from_observe = True
|
||||
self._resubscribe_due = True
|
||||
|
||||
if not self._discovered:
|
||||
# One-time (issue #177): find sibling subdevices before the
|
||||
@@ -888,7 +945,8 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
for href, rep in resources.items():
|
||||
self._observe.apply(href, rep, source=source)
|
||||
|
||||
if first_cycle or just_downgraded_from_observe:
|
||||
if first_cycle or self._resubscribe_due:
|
||||
self._resubscribe_due = False
|
||||
await self._attempt_observe_mode()
|
||||
elif self._observe.mode == MODE_POLL:
|
||||
await self._maybe_retry_observe_mode()
|
||||
@@ -910,7 +968,8 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def async_send_command(self, bound_entity: BoundEntity, payload: Any) -> None:
|
||||
"""Write a value to the device. Fire-and-forget.
|
||||
"""Write a value to the device. Retries once on a dead session
|
||||
(issue #294); raises HomeAssistantError if that retry fails too.
|
||||
|
||||
A description-level validate_fn (SwitchDesc only, currently) rejects
|
||||
a write with a user-facing message ahead of write_fn's silent
|
||||
@@ -1012,18 +1071,52 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
)
|
||||
|
||||
def _do_put():
|
||||
if self._session is None:
|
||||
self._connect_session()
|
||||
sess = self._session
|
||||
if sess is None:
|
||||
raise RuntimeError("no session")
|
||||
code, _ = sess.post(path_segs, cbor2.dumps(body), timeout=self._POST_TIMEOUT_S)
|
||||
self._log.info("PUT %s → code %#04x", write_href, code)
|
||||
|
||||
try:
|
||||
await self.hass.async_add_executor_job(_do_put)
|
||||
except Exception as e:
|
||||
self._log.error("command failed for %s: %s", write_href, e)
|
||||
else:
|
||||
await self.async_request_refresh()
|
||||
# Mirrors the poll path's reconnect-and-retry (issue #294): a PUT
|
||||
# landing on a session Samsung's firmware closed between polls used
|
||||
# to be silently lost -- no retry, no user-facing error.
|
||||
async with self._session_lock:
|
||||
try:
|
||||
await self.hass.async_add_executor_job(_do_put)
|
||||
except Exception as e:
|
||||
self._log.warning("command failed for %s, reconnecting: %s", write_href, e)
|
||||
await self.hass.async_add_executor_job(self._close_session)
|
||||
# The session is dead the moment it's closed, so any OBSERVE
|
||||
# subscriptions on it are too -- downgrade here, before the
|
||||
# retry, so a retry that also fails doesn't leave mode
|
||||
# claiming "Push" on a session that no longer exists
|
||||
# (issue #294; the poll path handles the same fact for its
|
||||
# own reconnect the same way, unconditionally on close).
|
||||
if self._observe.mode == MODE_OBSERVE:
|
||||
self._observe.downgrade_to_poll()
|
||||
self._resubscribe_due = True
|
||||
await asyncio.sleep(self._RECONNECT_PAUSE_S)
|
||||
try:
|
||||
await self.hass.async_add_executor_job(_do_put)
|
||||
except Exception as e2:
|
||||
self._log.error("command failed for %s after reconnect: %s", write_href, e2)
|
||||
raise HomeAssistantError(
|
||||
translation_domain=DOMAIN,
|
||||
translation_key="command_failed",
|
||||
translation_placeholders={"href": write_href, "error": str(e2)},
|
||||
) from e2
|
||||
# The retry's own reconnect pause + second PUT can eat well
|
||||
# into the settle window armed above, leaving too little of
|
||||
# it for the confirming poll below and reviving the
|
||||
# revert-then-reapply symptom settle_s exists to prevent
|
||||
# (issue #9). Re-arm it fresh now that the write actually
|
||||
# landed.
|
||||
self._observe.mark_write_pending(
|
||||
write_href, settle_s=self._POST_TIMEOUT_S + self._POLL_TIMEOUT_S
|
||||
)
|
||||
await self.async_request_refresh()
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Debug raw write (issue #54): a power-user escape hatch for the
|
||||
|
||||
@@ -171,17 +171,16 @@ class ObserveManager:
|
||||
self._last_notify_ts is not None and time.monotonic() - self._last_notify_ts < window_s
|
||||
)
|
||||
|
||||
def try_enter_observe_mode(
|
||||
self,
|
||||
session,
|
||||
hrefs: list[str],
|
||||
grace_period_s: float = GRACE_PERIOD_S,
|
||||
success_fraction: float = SUCCESS_FRACTION,
|
||||
) -> bool:
|
||||
"""Blocking — subscribes to every href then waits up to
|
||||
`grace_period_s`, returning early once `success_fraction` of hrefs
|
||||
have notified. Caller must run this in an executor, never on the
|
||||
event loop."""
|
||||
def subscribe_hrefs(self, session, hrefs: list[str]) -> set[str]:
|
||||
"""Register OBSERVE on every href; returns the ones that took.
|
||||
Blocking — run in an executor.
|
||||
|
||||
Split from the grace wait below (issue #294) so the coordinator can
|
||||
hold its session lock for just these sends -- each is a fire-and-
|
||||
forget UDP datagram (DtlsCoapSession.subscribe doesn't wait for the
|
||||
device's ack), unlike the wait, which can block for the whole grace
|
||||
period and must not hold a lock a command write is also waiting on.
|
||||
"""
|
||||
with self._notify_cond:
|
||||
self._notified.clear()
|
||||
subscribed: set[str] = set()
|
||||
@@ -192,30 +191,65 @@ class ObserveManager:
|
||||
subscribed.add(href)
|
||||
except Exception as e:
|
||||
self.log.warning("subscribe %s failed: %s", href, e)
|
||||
return subscribed
|
||||
|
||||
def await_observe_notifies(
|
||||
self,
|
||||
subscribed: set[str],
|
||||
grace_period_s: float = GRACE_PERIOD_S,
|
||||
success_fraction: float = SUCCESS_FRACTION,
|
||||
) -> bool:
|
||||
"""Blocking — waits up to `grace_period_s`, returning early once
|
||||
`success_fraction` of `subscribed` have notified. Touches no
|
||||
session; safe to run without holding a session lock."""
|
||||
if not subscribed:
|
||||
self._stop_refresh_task()
|
||||
self._set_mode(MODE_POLL)
|
||||
self.subscribed_hrefs = set()
|
||||
return False
|
||||
|
||||
def _fraction_reached() -> bool:
|
||||
return len(set(self._notified) & subscribed) / len(subscribed) >= success_fraction
|
||||
|
||||
with self._notify_cond:
|
||||
reached = self._notify_cond.wait_for(
|
||||
_fraction_reached,
|
||||
timeout=grace_period_s,
|
||||
)
|
||||
return self._notify_cond.wait_for(_fraction_reached, timeout=grace_period_s)
|
||||
|
||||
if reached:
|
||||
self.subscribed_hrefs = subscribed
|
||||
self._set_mode(MODE_OBSERVE)
|
||||
self.start_refresh_task(session)
|
||||
return True
|
||||
def enter_observe_mode(self, session, subscribed: set[str]) -> None:
|
||||
"""Commit a successful attempt. Caller must have re-confirmed
|
||||
`session` is still the live one under its session lock (issue
|
||||
#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)
|
||||
self._set_mode(MODE_OBSERVE)
|
||||
self.start_refresh_task(session)
|
||||
|
||||
def abandon_observe_attempt(self) -> None:
|
||||
"""Drop a failed or stale attempt: no subscriptions worth keeping."""
|
||||
self._stop_refresh_task()
|
||||
self.subscribed_hrefs = set()
|
||||
self._set_mode(MODE_POLL)
|
||||
|
||||
def try_enter_observe_mode(
|
||||
self,
|
||||
session,
|
||||
hrefs: list[str],
|
||||
grace_period_s: float = GRACE_PERIOD_S,
|
||||
success_fraction: float = SUCCESS_FRACTION,
|
||||
) -> bool:
|
||||
"""Blocking — subscribes to every href then waits up to
|
||||
`grace_period_s`, returning early once `success_fraction` of hrefs
|
||||
have notified. Caller must run this in an executor, never on the
|
||||
event loop.
|
||||
|
||||
Single-threaded convenience wrapper around the phase split above
|
||||
(subscribe_hrefs / await_observe_notifies / enter_observe_mode /
|
||||
abandon_observe_attempt) for callers -- direct and most existing
|
||||
tests -- that don't need the lock-scoping those phases exist for."""
|
||||
subscribed = self.subscribe_hrefs(session, hrefs)
|
||||
if not subscribed:
|
||||
self.abandon_observe_attempt()
|
||||
return False
|
||||
if self.await_observe_notifies(subscribed, grace_period_s, success_fraction):
|
||||
self.enter_observe_mode(session, subscribed)
|
||||
return True
|
||||
self.abandon_observe_attempt()
|
||||
return False
|
||||
|
||||
def _set_mode(self, mode: str) -> None:
|
||||
|
||||
@@ -1356,6 +1356,9 @@
|
||||
},
|
||||
"intensive_unavailable_for_cycle": {
|
||||
"message": "Intenzivní není u vybraného cyklu k dispozici."
|
||||
},
|
||||
"command_failed": {
|
||||
"message": "Příkaz pro {href} selhal i po opětovném připojení: {error}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1356,6 +1356,9 @@
|
||||
},
|
||||
"intensive_unavailable_for_cycle": {
|
||||
"message": "Intensive isn't available on the selected cycle."
|
||||
},
|
||||
"command_failed": {
|
||||
"message": "The command to {href} failed even after reconnecting: {error}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -119,6 +119,9 @@
|
||||
},
|
||||
"intensive_unavailable_for_cycle": {
|
||||
"message": "El modo intensivo no está disponible en el ciclo seleccionado."
|
||||
},
|
||||
"command_failed": {
|
||||
"message": "El comando para {href} falló incluso después de reconectar: {error}"
|
||||
}
|
||||
},
|
||||
"entity": {
|
||||
|
||||
@@ -1356,6 +1356,9 @@
|
||||
},
|
||||
"intensive_unavailable_for_cycle": {
|
||||
"message": "Intensivo non è disponibile per il ciclo selezionato."
|
||||
},
|
||||
"command_failed": {
|
||||
"message": "Il comando per {href} è fallito anche dopo la riconnessione: {error}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1356,6 +1356,9 @@
|
||||
},
|
||||
"intensive_unavailable_for_cycle": {
|
||||
"message": "선택한 코스에서는 강력 세탁을 사용할 수 없습니다."
|
||||
},
|
||||
"command_failed": {
|
||||
"message": "재연결 후에도 {href} 명령이 실패했습니다: {error}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1356,6 +1356,9 @@
|
||||
},
|
||||
"intensive_unavailable_for_cycle": {
|
||||
"message": "Intensief is niet beschikbaar voor het geselecteerde programma."
|
||||
},
|
||||
"command_failed": {
|
||||
"message": "Het commando naar {href} is ook na opnieuw verbinden mislukt: {error}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -75,6 +75,41 @@ async def test_successful_setup(hass: HomeAssistant, mock_probe) -> None:
|
||||
assert result["data"][CONF_CA_CERT_PEM] == MOCK_CA_CERT_PEM
|
||||
|
||||
|
||||
async def test_setup_normalizes_messy_pasted_pem(hass: HomeAssistant, mock_probe) -> None:
|
||||
"""A PEM with a leading UTF-8 BOM, CRLF line endings, and a stray blank
|
||||
line -- the kind a Windows text editor's copy produces, as opposed to a
|
||||
`type` dump (issue #291) -- must still be accepted and stored in its
|
||||
normalized form, not rejected with an opaque InvalidHeader."""
|
||||
messy_cert = "\ufeff" + MOCK_CA_CERT_PEM.replace("\n", "\r\n") + "\r\n\r\n"
|
||||
messy_key = "\ufeff" + MOCK_CA_KEY_PEM.replace("\n", "\r\n")
|
||||
|
||||
result = await hass.config_entries.flow.async_init(DOMAIN, context={"source": "user"})
|
||||
result = await hass.config_entries.flow.async_configure(
|
||||
result["flow_id"],
|
||||
{
|
||||
CONF_HOST: MOCK_HOST,
|
||||
CONF_CA_CERT_PEM: messy_cert,
|
||||
CONF_CA_KEY_PEM: messy_key,
|
||||
},
|
||||
)
|
||||
assert result["type"] == FlowResultType.CREATE_ENTRY
|
||||
assert result["data"][CONF_CA_CERT_PEM] == MOCK_CA_CERT_PEM
|
||||
assert result["data"][CONF_CA_KEY_PEM] == MOCK_CA_KEY_PEM
|
||||
|
||||
|
||||
def test_normalize_pem_strips_bom_crlf_and_blank_lines() -> None:
|
||||
"""Unit-level check of the helper itself, isolated from the flow."""
|
||||
from custom_components.localthings.config_flow import _normalize_pem
|
||||
|
||||
messy = "\ufeff-----BEGIN CERTIFICATE-----\r\nTEST-CA\r\n\r\n-----END CERTIFICATE-----\r\n"
|
||||
assert _normalize_pem(messy) == (
|
||||
"-----BEGIN CERTIFICATE-----\nTEST-CA\n-----END CERTIFICATE-----"
|
||||
)
|
||||
# A clean PEM (the `type`-dump case) passes through unchanged.
|
||||
clean = "-----BEGIN CERTIFICATE-----\nTEST-CA\n-----END CERTIFICATE-----"
|
||||
assert _normalize_pem(clean) == clean
|
||||
|
||||
|
||||
def test_order_candidates_prefers_known_ports() -> None:
|
||||
"""Live ports are ordered with the historically known DTLS ports first,
|
||||
then the rest ascending."""
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from datetime import timedelta
|
||||
from unittest.mock import AsyncMock, patch
|
||||
|
||||
@@ -21,6 +22,7 @@ from custom_components.localthings.const import (
|
||||
SUMMARY_INTERVAL_S,
|
||||
)
|
||||
from custom_components.localthings.coordinator import (
|
||||
_RECOVERY_RETRY_S,
|
||||
LocalThingsCoordinator,
|
||||
_local_source_port,
|
||||
)
|
||||
@@ -30,7 +32,7 @@ from custom_components.localthings.registry.capabilities.common import (
|
||||
remote_control_required_for_write,
|
||||
)
|
||||
|
||||
from .conftest import ENTRY_DATA, MOCK_MODEL, MOCK_SERIAL
|
||||
from .conftest import ENTRY_DATA, MOCK_MODEL, MOCK_SERIAL, FakeObserveSession
|
||||
from .conftest import _load_fridge_resources as _load_fridge
|
||||
|
||||
|
||||
@@ -478,6 +480,57 @@ async def test_reconnect_while_observe_mode_downgrades_to_poll(
|
||||
assert coordinator._observe.mode == MODE_POLL
|
||||
|
||||
|
||||
async def test_total_poll_failure_downgrades_observe_mode_to_poll(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session, fridge_resources
|
||||
) -> None:
|
||||
"""A device that drops off the network entirely -- both the poll and its
|
||||
reconnect retry fail -- must not leave the connection-mode sensor
|
||||
reporting 'Push' forever (issue #287). Only the *successful* reconnect
|
||||
branch used to touch observe mode (see
|
||||
test_reconnect_while_observe_mode_downgrades_to_poll); this covers the
|
||||
branch where the device stays unreachable."""
|
||||
fake = mock_coordinator_observe_session
|
||||
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]
|
||||
hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs
|
||||
|
||||
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,
|
||||
)
|
||||
assert entered is True
|
||||
assert coordinator.observe_mode == MODE_OBSERVE
|
||||
|
||||
last_notify_ts = coordinator._observe._last_notify_ts
|
||||
assert last_notify_ts is not None
|
||||
coordinator._observe._last_notify_ts = last_notify_ts - (PUSH_HEALTH_WINDOW_S + 1)
|
||||
|
||||
with (
|
||||
patch(
|
||||
"custom_components.localthings.coordinator.LocalThingsCoordinator._poll_once",
|
||||
side_effect=[RuntimeError("connection lost"), RuntimeError("still lost")],
|
||||
),
|
||||
patch(
|
||||
"custom_components.localthings.coordinator.asyncio.sleep",
|
||||
new=AsyncMock(),
|
||||
),
|
||||
):
|
||||
await coordinator.async_request_refresh()
|
||||
await hass.async_block_till_done()
|
||||
|
||||
# The update still "succeeds" with the last-known snapshot (issue #254's
|
||||
# degraded-data path) -- but the connection mode must reflect reality
|
||||
# now, not the stale OBSERVE state from before the outage.
|
||||
assert coordinator.observe_mode == MODE_POLL
|
||||
assert coordinator.last_update_success is True
|
||||
|
||||
|
||||
async def test_poll_timeout_skips_reconnect_when_push_is_healthy(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||
) -> None:
|
||||
@@ -691,6 +744,149 @@ async def test_reconnect_from_observe_mode_resubscribes_immediately(
|
||||
assert coordinator.observe_mode == MODE_OBSERVE
|
||||
|
||||
|
||||
async def test_attempt_observe_mode_discards_stale_commit_after_session_swap(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||
) -> None:
|
||||
"""A reconnect (the poll path's own, or a command's retry) can swap
|
||||
self._session while this attempt's grace wait is in flight -- it runs
|
||||
without holding _session_lock precisely so a write isn't blocked behind
|
||||
it (issue #294). Committing observe mode against the now-stale local
|
||||
`sess` reference would claim "Push" on a session that's already gone,
|
||||
with nothing left to notice -- the identity re-check under the lock
|
||||
right before committing must catch this and abandon instead.
|
||||
|
||||
The new session is never-tried, though, not just abandoned: it must
|
||||
flag an immediate resubscribe rather than let _last_observe_attempt_ts
|
||||
(stamped for the now-abandoned attempt) throttle it for up to
|
||||
_RECOVERY_RETRY_S.
|
||||
|
||||
Simulates the swap from inside await_observe_notifies itself rather
|
||||
than via real concurrency: subscribe_hrefs (and its lock) has already
|
||||
returned by the time that call runs, so this lands exactly in the
|
||||
window the identity check exists to cover, deterministically."""
|
||||
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._resubscribe_due is False
|
||||
|
||||
other = FakeObserveSession()
|
||||
|
||||
def _swap_session_mid_wait(subscribed, grace_period_s, success_fraction=None):
|
||||
coordinator._session = other # ty: ignore[invalid-assignment]
|
||||
return True
|
||||
|
||||
with patch.object(
|
||||
coordinator._observe, "await_observe_notifies", side_effect=_swap_session_mid_wait
|
||||
):
|
||||
await coordinator._attempt_observe_mode()
|
||||
|
||||
assert coordinator.observe_mode == MODE_POLL
|
||||
assert coordinator._observe.subscribed_hrefs == set()
|
||||
assert coordinator._observe._refresh_thread is None
|
||||
assert coordinator._resubscribe_due is True
|
||||
|
||||
|
||||
async def test_maybe_retry_observe_mode_uses_most_recent_attempt_not_just_mode_change(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||
) -> None:
|
||||
"""_set_mode only stamps last_mode_change_ts on an actual transition,
|
||||
so a device that never successfully enters observe mode leaves that
|
||||
timestamp stuck at construction time forever -- a failed attempt keeps
|
||||
calling _set_mode(MODE_POLL) while already in MODE_POLL, a no-op.
|
||||
Gating solely on that timestamp would make the 600s throttle open once
|
||||
and then never close again, re-attempting (and paying the subscribe
|
||||
burst) on every single poll cycle instead of every _RECOVERY_RETRY_S."""
|
||||
fake = mock_coordinator_observe_session
|
||||
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.observe_mode == MODE_POLL # never notified during setup
|
||||
|
||||
# Simulate exactly the scenario above: mode_change_ts is old (as it
|
||||
# would be forever, for a device that never gets push), but an attempt
|
||||
# really did just run.
|
||||
coordinator._observe.last_mode_change_ts = time.monotonic() - _RECOVERY_RETRY_S - 1
|
||||
coordinator._last_observe_attempt_ts = time.monotonic()
|
||||
|
||||
with patch.object(fake, "subscribe") as mock_subscribe:
|
||||
await coordinator._maybe_retry_observe_mode()
|
||||
|
||||
mock_subscribe.assert_not_called()
|
||||
|
||||
|
||||
async def test_maybe_retry_observe_mode_also_respects_a_mode_change_outside_an_attempt(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||
) -> None:
|
||||
"""The mirror of the case above: last_mode_change_ts can be the more
|
||||
recent of the two as well, e.g. right after the poll or command path's
|
||||
own downgrade (neither goes through _attempt_observe_mode, so neither
|
||||
stamps _last_observe_attempt_ts). Dropping last_mode_change_ts from the
|
||||
max() would let a device that was *just* downgraded get re-attempted
|
||||
immediately instead of respecting _RECOVERY_RETRY_S."""
|
||||
fake = mock_coordinator_observe_session
|
||||
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.observe_mode == MODE_POLL
|
||||
|
||||
coordinator._last_observe_attempt_ts = time.monotonic() - _RECOVERY_RETRY_S - 1
|
||||
coordinator._observe.last_mode_change_ts = time.monotonic()
|
||||
|
||||
with patch.object(fake, "subscribe") as mock_subscribe:
|
||||
await coordinator._maybe_retry_observe_mode()
|
||||
|
||||
mock_subscribe.assert_not_called()
|
||||
|
||||
|
||||
async def test_attempt_observe_mode_releases_lock_before_the_grace_wait(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||
) -> None:
|
||||
"""The subscribe burst holds _session_lock (it touches the session);
|
||||
the grace wait after it must not, or a command write could stall
|
||||
behind up to _OBSERVE_GRACE_PERIOD_S of an unrelated observe-mode-entry
|
||||
attempt (issue #294). By construction, subscribe_hrefs's own
|
||||
`async with self._session_lock:` has already exited by the time
|
||||
await_observe_notifies is even called -- checked here rather than
|
||||
inferred from timing."""
|
||||
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]
|
||||
|
||||
locked_during_wait = {"value": None}
|
||||
|
||||
def _check_lock(subscribed, grace_period_s, success_fraction=None):
|
||||
locked_during_wait["value"] = coordinator._session_lock.locked()
|
||||
return True
|
||||
|
||||
with patch.object(coordinator._observe, "await_observe_notifies", side_effect=_check_lock):
|
||||
await coordinator._attempt_observe_mode()
|
||||
|
||||
assert locked_during_wait["value"] is False
|
||||
|
||||
|
||||
async def test_attempt_observe_mode_holds_lock_during_the_subscribe_burst(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||
) -> None:
|
||||
"""The other half of the split above: the subscribe burst does touch
|
||||
the session, so it must hold _session_lock -- that's what actually
|
||||
stops a concurrent close from landing mid-subscribe (issue #294)."""
|
||||
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]
|
||||
|
||||
locked_during_subscribe = {"value": None}
|
||||
real_subscribe_hrefs = coordinator._observe.subscribe_hrefs
|
||||
|
||||
def _check_lock(session, hrefs):
|
||||
locked_during_subscribe["value"] = coordinator._session_lock.locked()
|
||||
return real_subscribe_hrefs(session, hrefs)
|
||||
|
||||
with patch.object(coordinator._observe, "subscribe_hrefs", side_effect=_check_lock):
|
||||
await coordinator._attempt_observe_mode()
|
||||
|
||||
assert locked_during_subscribe["value"] is True
|
||||
|
||||
|
||||
async def test_sweep_mismatch_never_downgrades_a_live_observe_session(
|
||||
hass: HomeAssistant, mock_entry, mock_coordinator_observe_session
|
||||
) -> None:
|
||||
@@ -956,6 +1152,192 @@ async def test_send_command_survives_stale_confirm_poll(
|
||||
assert coordinator._cache.get("/test/vs/0") == {"value": 5}
|
||||
|
||||
|
||||
async def test_send_command_reconnects_and_retries_after_socket_closed(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
mock_coordinator_observe_session,
|
||||
) -> None:
|
||||
"""A command lost to a session Samsung's firmware closed between polls
|
||||
must not just vanish (issue #294): `_do_put` failing once is now
|
||||
followed by a reconnect and a single retry, mirroring the poll path's
|
||||
own recovery in `_async_update_data`.
|
||||
|
||||
`mock_coordinator_observe_session` patches `_close_session` to a no-op,
|
||||
which would leave `self._session` never actually going `None` -- and
|
||||
with it, `_do_put`'s own `if self._session is None: self._connect_session()`
|
||||
guard never exercised, so a broken reconnect could still pass. Overridden
|
||||
here to actually drop the session, so the retry only succeeds if that
|
||||
guard really rebuilds it."""
|
||||
from custom_components.localthings.registry.discovery import BoundEntity
|
||||
from custom_components.localthings.registry.entities import NumberDesc
|
||||
|
||||
fake = mock_coordinator_observe_session
|
||||
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]
|
||||
|
||||
def _write_fn(payload, rep, href=None):
|
||||
return (["test", "vs", "0"], {"value": payload})
|
||||
|
||||
desc = NumberDesc(key="test", field="value", write_fn=_write_fn)
|
||||
bound = BoundEntity(href="/test/vs/0", capability=coordinator.bound[0].capability, desc=desc)
|
||||
|
||||
calls = {"n": 0}
|
||||
|
||||
def _post(*args, **kwargs):
|
||||
calls["n"] += 1
|
||||
if calls["n"] == 1:
|
||||
raise ConnectionError("socket closed")
|
||||
return (0x44, b"")
|
||||
|
||||
def _drop_session():
|
||||
coordinator._session = None
|
||||
|
||||
reconnects = {"n": 0}
|
||||
|
||||
def _reconnect():
|
||||
reconnects["n"] += 1
|
||||
coordinator._session = fake
|
||||
|
||||
with (
|
||||
patch.object(fake, "subscribe"),
|
||||
patch.object(coordinator, "_close_session", side_effect=_drop_session),
|
||||
patch.object(coordinator, "_connect_session", side_effect=_reconnect),
|
||||
patch(
|
||||
"custom_components.localthings.coordinator.asyncio.sleep",
|
||||
new=AsyncMock(),
|
||||
),
|
||||
):
|
||||
fake.post = _post
|
||||
await coordinator.async_send_command(bound, 5)
|
||||
|
||||
assert reconnects["n"] == 1
|
||||
|
||||
assert calls["n"] == 2
|
||||
assert coordinator._cache.get("/test/vs/0") == {"value": 5}
|
||||
|
||||
|
||||
async def test_send_command_raises_after_reconnect_retry_also_fails(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
mock_coordinator_observe_session,
|
||||
) -> None:
|
||||
"""If the command still fails on the reconnected session, the user must
|
||||
see it -- previously this was swallowed into a log line with no
|
||||
feedback at all (issue #294).
|
||||
|
||||
Also covers a sibling bug the fix for that same issue introduced: the
|
||||
session is closed the moment the first attempt fails, so any OBSERVE
|
||||
subscriptions on it are already dead regardless of whether the retry
|
||||
that follows succeeds -- a failed retry must still downgrade mode, or
|
||||
it's left claiming "Push" on a session that no longer exists."""
|
||||
from homeassistant.exceptions import HomeAssistantError
|
||||
|
||||
from custom_components.localthings.registry.discovery import BoundEntity
|
||||
from custom_components.localthings.registry.entities import NumberDesc
|
||||
|
||||
fake = mock_coordinator_observe_session
|
||||
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]
|
||||
hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs
|
||||
|
||||
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,
|
||||
)
|
||||
assert entered is True
|
||||
assert coordinator.observe_mode == MODE_OBSERVE
|
||||
|
||||
def _write_fn(payload, rep, href=None):
|
||||
return (["test", "vs", "0"], {"value": payload})
|
||||
|
||||
desc = NumberDesc(key="test", field="value", write_fn=_write_fn)
|
||||
bound = BoundEntity(href="/test/vs/0", capability=coordinator.bound[0].capability, desc=desc)
|
||||
|
||||
def _post(*args, **kwargs):
|
||||
raise ConnectionError("socket closed")
|
||||
|
||||
with (
|
||||
patch(
|
||||
"custom_components.localthings.coordinator.asyncio.sleep",
|
||||
new=AsyncMock(),
|
||||
),
|
||||
pytest.raises(HomeAssistantError),
|
||||
):
|
||||
fake.post = _post
|
||||
await coordinator.async_send_command(bound, 5)
|
||||
|
||||
assert coordinator.observe_mode == MODE_POLL
|
||||
|
||||
|
||||
async def test_send_command_reconnect_downgrades_observe_mode(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
mock_coordinator_observe_session,
|
||||
) -> None:
|
||||
"""A command's own successful reconnect hands back a session with zero
|
||||
OBSERVE registrations too, same as the poll path's reconnect -- must
|
||||
downgrade the same way and flag a resubscribe, or observe mode stays
|
||||
claimed against a session the write just replaced underneath it
|
||||
(issue #294)."""
|
||||
from custom_components.localthings.registry.discovery import BoundEntity
|
||||
from custom_components.localthings.registry.entities import NumberDesc
|
||||
|
||||
fake = mock_coordinator_observe_session
|
||||
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]
|
||||
hrefs = coordinator._hot_hrefs + coordinator._warm_hrefs
|
||||
|
||||
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,
|
||||
)
|
||||
assert entered is True
|
||||
assert coordinator.observe_mode == MODE_OBSERVE
|
||||
|
||||
def _write_fn(payload, rep, href=None):
|
||||
return (["test", "vs", "0"], {"value": payload})
|
||||
|
||||
desc = NumberDesc(key="test", field="value", write_fn=_write_fn)
|
||||
bound = BoundEntity(href="/test/vs/0", capability=coordinator.bound[0].capability, desc=desc)
|
||||
|
||||
calls = {"n": 0}
|
||||
|
||||
def _post(*args, **kwargs):
|
||||
calls["n"] += 1
|
||||
if calls["n"] == 1:
|
||||
raise ConnectionError("socket closed")
|
||||
return (0x44, b"")
|
||||
|
||||
with (
|
||||
patch.object(fake, "subscribe") as mock_subscribe,
|
||||
patch(
|
||||
"custom_components.localthings.coordinator.asyncio.sleep",
|
||||
new=AsyncMock(),
|
||||
),
|
||||
):
|
||||
fake.post = _post
|
||||
await coordinator.async_send_command(bound, 5)
|
||||
|
||||
# _resubscribe_due is consumed by this same call's own trailing
|
||||
# refresh (async_request_refresh is awaited, not fire-and-forget),
|
||||
# so the visible effect is a resubscribe attempt, not a lingering
|
||||
# flag value to assert on afterward.
|
||||
assert mock_subscribe.called
|
||||
|
||||
assert coordinator.observe_mode == MODE_POLL
|
||||
|
||||
|
||||
async def test_second_write_to_same_href_lands_during_first_writes_settle_window(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
|
||||
Reference in New Issue
Block a user