Pace writes and retry a timed-out PUT in place
Issue #384: an AC reported `climate.set_temperature`/`set_hvac_mode`/ `set_swing_mode` failing with "command failed for <href> after reconnect: session operation timed out" across several hrefs. smartthings-local retransmits every block of a GET (3 attempts) but sends a POST exactly once, so a single dropped datagram is a write failure. Two things made that likely and then unrecoverable: - The write path never paced. RT-OCF drops requests that arrive inside the rate-limit interval, and the reconnect retry is the worst case -- `_connect_session()` reads /oic/p, /oic/d and /oic/res immediately before it, so the retry PUT lands milliseconds behind three GETs. - Any exception, timeout included, was read as session death: close the session, downgrade OBSERVE to poll, pause, re-handshake, one more single-shot PUT. The poll path already reads a TimeoutError the other way (`_defer_reconnect_for`). Pace before both POST sites, and retry a timed-out write once on the session in hand before spending a reconnect on it. The reconnect ladder is unchanged past that point; it just moves to a helper.
This commit is contained in:
@@ -9,6 +9,7 @@ import logging
|
||||
import threading
|
||||
import time
|
||||
import zlib
|
||||
from collections.abc import Callable
|
||||
from dataclasses import asdict
|
||||
from datetime import timedelta
|
||||
from typing import Any, cast
|
||||
@@ -1953,48 +1954,78 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
sess = self._session
|
||||
if sess is None:
|
||||
raise RuntimeError("no session")
|
||||
# Pace like every other send site (issue #384): RT-OCF drops a
|
||||
# request that lands inside the rate-limit interval, and nothing
|
||||
# retransmits a POST, so an unpaced write can only surface as a
|
||||
# timeout. Worst on a reconnect retry, where _connect_session()
|
||||
# has just read three identity resources.
|
||||
sess.pace()
|
||||
code, _ = sess.post(path_segs, cbor2.dumps(body), timeout=self._POST_TIMEOUT_S)
|
||||
self._log.info("PUT %s → code %#04x", write_href, code)
|
||||
|
||||
# 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)
|
||||
except TimeoutError as e:
|
||||
# A timeout is a late ACK, not a dead session -- the same
|
||||
# reading _defer_reconnect_for takes on the poll path. A
|
||||
# POST is sent exactly once (smartthings-local retransmits
|
||||
# every block of a GET, never a write), so one dropped
|
||||
# datagram lands here on a healthy session, and reconnecting
|
||||
# over it costs the pause, a fresh handshake, and this
|
||||
# device's OBSERVE subscriptions (issue #384).
|
||||
self._log.info("command timed out for %s, retrying: %s", write_href, e)
|
||||
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._reconnect_and_retry_write(_do_put, write_href, e2)
|
||||
else:
|
||||
self._rearm_write_settle(write_href)
|
||||
except Exception as e:
|
||||
await self._reconnect_and_retry_write(_do_put, write_href, e)
|
||||
await self.async_request_refresh()
|
||||
|
||||
def _rearm_write_settle(self, write_href: str) -> None:
|
||||
"""Re-arm the settle guard once a retried write finally lands.
|
||||
|
||||
The retries eat into the window armed before the first PUT, leaving
|
||||
too little of it for the confirming poll and reviving the
|
||||
revert-then-reapply symptom settle_s exists to prevent (issue #9)."""
|
||||
self._observe.mark_write_pending(
|
||||
write_href, settle_s=self._POST_TIMEOUT_S + self._POLL_TIMEOUT_S
|
||||
)
|
||||
|
||||
async def _reconnect_and_retry_write(
|
||||
self, do_put: Callable[[], None], write_href: str, e: Exception
|
||||
) -> None:
|
||||
"""Reconnect and retry a failed write once, raising if it fails again.
|
||||
|
||||
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."""
|
||||
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
|
||||
self._rearm_write_settle(write_href)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Debug raw write/read (issue #54, extended for issue #300): a
|
||||
# power-user escape hatch shared by the options-flow debug panel and
|
||||
@@ -2012,6 +2043,7 @@ class LocalThingsCoordinator(DataUpdateCoordinator[dict[str, Any]]):
|
||||
sess = self._session
|
||||
if sess is None:
|
||||
raise RuntimeError("no session")
|
||||
sess.pace()
|
||||
code, _ = sess.post(path_segs, cbor2.dumps(body), timeout=self._POST_TIMEOUT_S)
|
||||
self._log.warning("DEBUG raw write POST %s %r → code %#04x", href, body, code)
|
||||
new_rep: dict = {}
|
||||
|
||||
@@ -216,6 +216,10 @@ class FakeObserveSession:
|
||||
def refresh_observes(self, paths):
|
||||
return None
|
||||
|
||||
def pace(self):
|
||||
"""No-op stand-in for the real session's rate limiter, which the
|
||||
write path calls before every POST (issue #384)."""
|
||||
|
||||
def close(self):
|
||||
self.closed = True
|
||||
|
||||
|
||||
@@ -1349,6 +1349,151 @@ async def test_send_command_raises_after_reconnect_retry_also_fails(
|
||||
assert coordinator.observe_mode == MODE_POLL
|
||||
|
||||
|
||||
async def test_send_command_timeout_retries_in_place_without_reconnecting(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
mock_coordinator_observe_session,
|
||||
) -> None:
|
||||
"""Issue #384: a write timeout is a dropped datagram, not a dead session.
|
||||
|
||||
smartthings-local retransmits every block of a GET but sends a POST
|
||||
exactly once, so a single loss surfaces here as a timeout on a session
|
||||
that is still fine. Retrying in place has to be tried before spending a
|
||||
reconnect -- which costs the pause, a fresh handshake, and every OBSERVE
|
||||
subscription this device has."""
|
||||
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
|
||||
|
||||
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 SessionTimeoutError()
|
||||
return (0x44, b"")
|
||||
|
||||
with (
|
||||
patch.object(coordinator, "_close_session") as mock_close,
|
||||
patch.object(coordinator, "_connect_session") as mock_connect,
|
||||
patch(
|
||||
"custom_components.localthings.coordinator.asyncio.sleep",
|
||||
new=AsyncMock(),
|
||||
) as mock_sleep,
|
||||
):
|
||||
fake.post = _post
|
||||
await coordinator.async_send_command(bound, 5)
|
||||
|
||||
assert calls["n"] == 2
|
||||
assert not mock_close.called
|
||||
assert not mock_connect.called
|
||||
# No _RECONNECT_PAUSE_S burned on a session that never needed replacing.
|
||||
assert not mock_sleep.called
|
||||
# And the push subscriptions the reconnect would have thrown away survive.
|
||||
assert coordinator.observe_mode == MODE_OBSERVE
|
||||
assert coordinator._cache.get("/test/vs/0") == {"value": 5}
|
||||
|
||||
|
||||
async def test_send_command_timeout_falls_back_to_reconnect(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
mock_coordinator_observe_session,
|
||||
) -> None:
|
||||
"""The in-place retry (issue #384) is an extra rung on the ladder, not a
|
||||
replacement: a href that keeps timing out must still reach the reconnect
|
||||
retry and, failing that, the user (issue #294)."""
|
||||
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]
|
||||
|
||||
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
|
||||
raise SessionTimeoutError()
|
||||
|
||||
with (
|
||||
patch.object(coordinator, "_close_session") as mock_close,
|
||||
patch(
|
||||
"custom_components.localthings.coordinator.asyncio.sleep",
|
||||
new=AsyncMock(),
|
||||
),
|
||||
pytest.raises(HomeAssistantError),
|
||||
):
|
||||
fake.post = _post
|
||||
await coordinator.async_send_command(bound, 5)
|
||||
|
||||
# First attempt, in-place retry, then the reconnect's own retry.
|
||||
assert calls["n"] == 3
|
||||
assert mock_close.called
|
||||
|
||||
|
||||
async def test_write_paces_before_posting(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
mock_coordinator_observe_session,
|
||||
) -> None:
|
||||
"""Issue #384: RT-OCF silently drops a request that lands inside the
|
||||
rate-limit interval, and nothing retransmits a POST -- so an unpaced
|
||||
write can only ever surface as a timeout. Every other send site paces;
|
||||
the write path has to as well."""
|
||||
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)
|
||||
|
||||
order: list[str] = []
|
||||
|
||||
fake.pace = lambda: order.append("pace")
|
||||
|
||||
def _post(*args, **kwargs):
|
||||
order.append("post")
|
||||
return (0x44, b"")
|
||||
|
||||
fake.post = _post
|
||||
await coordinator.async_send_command(bound, 5)
|
||||
|
||||
assert order == ["pace", "post"]
|
||||
|
||||
|
||||
async def test_send_command_reconnect_downgrades_observe_mode(
|
||||
hass: HomeAssistant,
|
||||
mock_entry,
|
||||
|
||||
Reference in New Issue
Block a user