feat(observe): add ObserveManager with write-settle guard

This commit is contained in:
Marc Billow
2026-07-09 18:28:46 -05:00
parent 3ab0d5553e
commit eba73118e9
2 changed files with 111 additions and 0 deletions
+61
View File
@@ -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)
+50
View File
@@ -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}