From eba73118e95f28809206553679c82b9dd34aa16b Mon Sep 17 00:00:00 2001 From: Marc Billow Date: Thu, 9 Jul 2026 18:28:46 -0500 Subject: [PATCH] feat(observe): add ObserveManager with write-settle guard --- custom_components/localthings/observe.py | 61 ++++++++++++++++++++++++ tests/localthings/test_observe.py | 50 +++++++++++++++++++ 2 files changed, 111 insertions(+) create mode 100644 custom_components/localthings/observe.py create mode 100644 tests/localthings/test_observe.py diff --git a/custom_components/localthings/observe.py b/custom_components/localthings/observe.py new file mode 100644 index 0000000..bed99b1 --- /dev/null +++ b/custom_components/localthings/observe.py @@ -0,0 +1,61 @@ +"""Observe-mode (CoAP OBSERVE) support layered on top of StateCache. + +Owns mode selection (push vs. poll), the write-settle guard (drops a +just-written href's incoming updates for a few seconds so a slow-to-settle +device doesn't revert an optimistic write), and missed-notification +detection. `smartthings_local`'s StateCache/DtlsCoapSession/ObserveRefreshTask +are an external pip dependency we don't own, so behavior that would +naturally live inside StateCache.apply_rep lives here instead, gating +whether apply_rep is called at all. +""" +from __future__ import annotations + +import logging +import threading +import time + +from smartthings_local.ocf.state_cache import StateCache + +_LOGGER = logging.getLogger(__name__) + +MODE_OBSERVE = 'observe' +MODE_POLL = 'poll' + +DEFAULT_SETTLE_S = 4.0 + + +class ObserveManager: + """Per-device observe-mode state: mode, write-settle guard, and (later) + subscription/staleness tracking. Pure sync logic — safe to call from + any thread; callers on the event loop must still marshal any HA state + push through `hass.add_job`, this class does not touch asyncio.""" + + def __init__(self, cache: StateCache, logger: logging.Logger | None = None): + self.cache = cache + self.log = logger or _LOGGER + self.mode = MODE_POLL + self.last_mode_change_ts = time.monotonic() + self.last_mode_change_wall = time.time() + self._settle_until: dict[str, float] = {} + self._settle_lock = threading.Lock() + + def mark_write_pending(self, href: str, settle_s: float = DEFAULT_SETTLE_S) -> None: + with self._settle_lock: + self._settle_until[href] = time.monotonic() + settle_s + + def _is_settling(self, href: str) -> bool: + with self._settle_lock: + until = self._settle_until.get(href) + if until is None: + return False + if time.monotonic() >= until: + del self._settle_until[href] + return False + return True + + def apply(self, href: str, rep: dict, source: str) -> bool: + """Gate a StateCache.apply_rep call through the write-settle guard.""" + if self._is_settling(href): + self.log.debug("dropping %s update for %s (settling)", source, href) + return False + return self.cache.apply_rep(href, rep, source=source) diff --git a/tests/localthings/test_observe.py b/tests/localthings/test_observe.py new file mode 100644 index 0000000..102020e --- /dev/null +++ b/tests/localthings/test_observe.py @@ -0,0 +1,50 @@ +"""Tests for ObserveManager: write-settle guard and mode defaults.""" +from __future__ import annotations + +import time + +from smartthings_local.ocf.state_cache import StateCache + +from custom_components.localthings.observe import ObserveManager, MODE_POLL + + +class _NullDescriptor: + def on_observation(self, state, href, rep): + return None + + +def _manager() -> ObserveManager: + return ObserveManager(StateCache(_NullDescriptor())) + + +def test_starts_in_poll_mode(): + mgr = _manager() + assert mgr.mode == MODE_POLL + + +def test_apply_writes_through_when_not_settling(): + mgr = _manager() + assert mgr.apply('/oven/vs/0', {'a': 1}, source='poll') is True + assert mgr.cache.get('/oven/vs/0') == {'a': 1} + + +def test_apply_drops_update_during_settle_window(): + mgr = _manager() + mgr.cache.apply_rep('/oven/vs/0', {'a': 1}, source='seed') + mgr.mark_write_pending('/oven/vs/0', settle_s=1.0) + + result = mgr.apply('/oven/vs/0', {'a': 2}, source='poll') + + assert result is False + assert mgr.cache.get('/oven/vs/0') == {'a': 1} + + +def test_apply_accepts_update_after_settle_window_elapses(): + mgr = _manager() + mgr.mark_write_pending('/oven/vs/0', settle_s=0.05) + time.sleep(0.1) + + result = mgr.apply('/oven/vs/0', {'a': 2}, source='poll') + + assert result is True + assert mgr.cache.get('/oven/vs/0') == {'a': 2}