From 46a930eb719ca03902671069052dd06a18139ea5 Mon Sep 17 00:00:00 2001 From: Jack Nagy Date: Tue, 30 Jun 2026 19:27:59 +0100 Subject: [PATCH] Production-hardening pass: half-open detection, cascade throttle, OBSERVE refresh, lamp/door coupling Four failure modes observed in-house since the polling-first refactor (b00c2fd): 1. Half-open DTLS sessions where the socket stays writable but the peer has gone silent. Ping sends succeed against a wedged peer because RT-OCF doesn't reliably emit a RST; only successful polls prove the session is live. - PollScheduler exposes last_success_ts (bumped on every 2.05). - KeepaliveTask takes liveness_fn(); ticks fail if no 2.05 in the last 60s, even when the ping send succeeded. - Bridge force-closes the session after 120s unreachable so run_forever() breaks out of sess.join() and reconnects. 2. RT-OCF cascade under load. One wedged path can eat 8s of timeout, the next tier tick fires immediately and stacks another attempt, and the device wedges harder. - PollTier.timeout_s per-tier override (hot=2s, warm=4s, sweep=15s). - On TimeoutError, the href goes into a 5-60s cooldown via the existing _defer_until mechanism. - take_window_stats() now reports successful-poll RTT separately from a timeout count, exposed as the "Poll Timeouts (window)" diagnostic entity in HA. - Active-window throttle: if the previous health window saw >=3 timeouts and is_active=True, drop back to idle cadence -- stops stacking polls on a stalled responder. 3. OBSERVE table aging across cloud-auth blips. The device stays DTLS-reachable but the on-device stack clears its observer table during the blip, so push delivery stays dead even after upstream recovers. - New ObserveRefreshTask per bridge; every 6h derregs all current observer tokens and re-subscribes on the existing session. 4. Oven lamp/door coupling. Oven hardware auto-drives the lamp from door state but /mode/vs/0 is warm-tier (30s) so HA showed stale lamp during a cook. - Track door + lamp value-change timestamps in descriptor_state. When the door transition is newer, derive lamp from door_open. When an HA optimistic write is newer, the cache value wins. Also: ANSI-coloured WARNING/ERROR lines (NO_COLOR=1 opt-out), jittered reconnect backoff so dryer + oven don't sync up after a router blip. In-house verification: running on dryer + oven since 2026-06-03. --- samsung_appliance/appliances/base.py | 3 + samsung_appliance/appliances/dryer.py | 6 ++ samsung_appliance/appliances/oven.py | 86 ++++++++++++++---- samsung_appliance/bridge.py | 87 ++++++++++++++++-- samsung_appliance/coap_dtls.py | 58 ++++++++---- samsung_appliance/keepalive.py | 38 ++++++-- samsung_appliance/logger.py | 50 +++++++++-- samsung_appliance/observe_refresh.py | 50 +++++++++++ samsung_appliance/poll_scheduler.py | 122 ++++++++++++++++++++++++-- 9 files changed, 438 insertions(+), 62 deletions(-) create mode 100644 samsung_appliance/observe_refresh.py diff --git a/samsung_appliance/appliances/base.py b/samsung_appliance/appliances/base.py index 6f891cb..fe3977b 100644 --- a/samsung_appliance/appliances/base.py +++ b/samsung_appliance/appliances/base.py @@ -216,6 +216,9 @@ def bridge_diagnostic_discovery(topic_prefix: str, sensor('slow_polls', 'Slow Polls (window)', "{{ value_json.poll_window_slow_count | default(0) }}", icon='mdi:timer-sand'), + sensor('poll_timeouts', 'Poll Timeouts (window)', + "{{ value_json.poll_window_timeout_count | default(0) }}", + icon='mdi:timer-off-outline'), sensor('poll_errors', 'Poll Errors (window)', "{{ value_json.poll_window_errors | default(0) }}", icon='mdi:alert-circle-outline'), diff --git a/samsung_appliance/appliances/dryer.py b/samsung_appliance/appliances/dryer.py index c182289..0186680 100644 --- a/samsung_appliance/appliances/dryer.py +++ b/samsung_appliance/appliances/dryer.py @@ -439,11 +439,15 @@ def command_handlers(): # Empirical ceiling on this firmware is ~14 req/s (probe_poll_rate_combined.py # 2026-06-03). The hot tier sits at 1s idle / 0.5s active — comfortably under # the ceiling and leaves headroom for the warm + sweep budgets. +# Per-tier timeouts are scaled to cadence: hot tier retries every 1s, so +# a tight 2s ceiling caps the cascade damage from one wedged poll. Warm +# tier has more headroom; sweep is multi-block Block2 and tolerates ~15s. DRYER_POLL_TIERS = [ PollTier( name='hot', interval_s=1.0, active_interval_s=0.5, + timeout_s=2.0, paths=( ('operational', 'state', 'vs', '0'), ), @@ -451,6 +455,7 @@ DRYER_POLL_TIERS = [ PollTier( name='warm', interval_s=15.0, + timeout_s=4.0, paths=( ('power', 'vs', '0'), ('kidslock', 'vs', '0'), @@ -467,6 +472,7 @@ DRYER_POLL_TIERS = [ PollTier( name='sweep', interval_s=300.0, + timeout_s=15.0, paths=(('device', '0'),), is_sweep=True, ), diff --git a/samsung_appliance/appliances/oven.py b/samsung_appliance/appliances/oven.py index e5eed9e..9da542c 100644 --- a/samsung_appliance/appliances/oven.py +++ b/samsung_appliance/appliances/oven.py @@ -1,6 +1,6 @@ """Oven descriptor (Samsung NV7000BS-class). -Resource map captured 2026-05-31 via DTLS-CoAP with the ab0b0ac4 cert. +Resource map captured 2026-05-31 via DTLS-CoAP with the client cert. See `local-tools/comparisons/oven-tree.md` for the full field reference. Write surfaces this descriptor exposes: @@ -188,6 +188,11 @@ def flatten(links): current_mode = modes[0] if modes else None options = g('/mode/vs/0', 'x.com.samsung.da.options') or [] lamp = _option_value(options, 'UpperLamp') # 'On' / 'Off' + # The door-coupling override (open → On, close → Off) lives in + # project() below — it needs descriptor_state to compare the + # latest door transition against the latest lamp value change so + # an HA-initiated optimistic write isn't clobbered by stale door + # state. sound = _option_value(options, 'Sound') # 'On' / 'Off' fastpreheat = _option_value(options, 'fastpreheat') # 'On' / 'Off' # NaturalSteam only appears in the options array after it's been @@ -273,30 +278,65 @@ def flatten(links): # extrapolate downward while machine_state == active. # --------------------------------------------------------------------- def on_observation(state, href, rep): - if href != '/operational/state/vs/0': + now = time.time() + if href == '/operational/state/vs/0': + rem = rep.get('x.com.samsung.da.remainingTime') + if isinstance(rem, str): + try: + h, m, s = rem.split(':') + state['remaining_anchor'] = (now, + int(h) * 3600 + int(m) * 60 + int(s)) + except (ValueError, AttributeError): + pass return - rem = rep.get('x.com.samsung.da.remainingTime') - if not isinstance(rem, str): + # Door + lamp tracking feeds the lamp/door coupling in project(). + # Both timestamps bump only on value CHANGES so the comparison + # tells us which event happened more recently. /doors is hot-tier + # (1s) and would otherwise dominate; /mode is warm-tier (30s) and + # picks up HA optimistic writes via apply_optimistic → apply_rep. + if href == '/doors/vs/0': + items = rep.get('x.com.samsung.da.items') or [] + door = items[0].get('x.com.samsung.da.openState') if items else None + if door != state.get('_door_last'): + state['_door_last'] = door + state['_door_change_ts'] = now return - try: - h, m, s = rem.split(':') - state['remaining_anchor'] = (time.time(), - int(h) * 3600 + int(m) * 60 + int(s)) - except (ValueError, AttributeError): - pass + if href == '/mode/vs/0': + options = rep.get('x.com.samsung.da.options') or [] + lamp = _option_value(options, 'UpperLamp') + if lamp != state.get('_lamp_last'): + state['_lamp_last'] = lamp + state['_lamp_change_ts'] = now def project(state, sensors): - anchor = state.get('remaining_anchor') - if sensors.get('machine_state') != 'active' or anchor is None: - return sensors - ts, total = anchor - remaining = max(0, int(total - (time.time() - ts))) - h, rest = divmod(remaining, 3600) - m, s = divmod(rest, 60) sensors = dict(sensors) - sensors['completion_time'] = f"{h}:{m:02d}:{s:02d}" - sensors['completion_minutes'] = h * 60 + m + (1 if s > 0 else 0) + # Remaining-time projection: the oven pushes /operational/state on + # state transitions but not on remainingTime ticks. Extrapolate + # from the most recent anchor while the machine is active. + anchor = state.get('remaining_anchor') + if sensors.get('machine_state') == 'active' and anchor is not None: + ts, total = anchor + remaining = max(0, int(total - (time.time() - ts))) + h, rest = divmod(remaining, 3600) + m, s = divmod(rest, 60) + sensors['completion_time'] = f"{h}:{m:02d}:{s:02d}" + sensors['completion_minutes'] = h * 60 + m + (1 if s > 0 else 0) + # Lamp / door coupling. The oven hardware auto-drives the lamp from + # the door state, but /mode/vs/0 only polls every 30s. When a door + # TRANSITION is more recent than the last lamp VALUE change, derive + # lamp from door for sub-second freshness. When a lamp toggle is + # more recent (HA optimistic write, or panel-driven /mode diff), + # the cache value wins — preserves HA toggle responsiveness even + # while the door is closed. + door_ts = state.get('_door_change_ts') + lamp_ts = state.get('_lamp_change_ts') + door_open = sensors.get('door_open') + if door_ts is not None and (lamp_ts is None or door_ts > lamp_ts): + if door_open is True: + sensors['lamp'] = 'On' + elif door_open is False: + sensors['lamp'] = 'Off' return sensors @@ -749,11 +789,16 @@ def command_handlers(): # Empirical ceiling on this firmware is ~8 req/s (probe_poll_rate_combined.py # 2026-06-03). Hot tier covers what changes mid-cook; doors get the tightest # cadence because door open/close needs sub-second freshness in HA. +# Per-tier timeouts are scaled to cadence: hot tier retries every 1s, so +# a tight 2s ceiling caps the cascade damage from one wedged poll. Warm +# and cold tiers have more headroom; sweep is multi-block Block2 and +# tolerates ~15s. OVEN_POLL_TIERS = [ PollTier( name='hot', interval_s=1.0, active_interval_s=0.5, + timeout_s=2.0, paths=( ('operational', 'state', 'vs', '0'), ('doors', 'vs', '0'), @@ -764,6 +809,7 @@ OVEN_POLL_TIERS = [ PollTier( name='warm', interval_s=30.0, + timeout_s=4.0, paths=( ('power', 'vs', '0'), ('kidslock', 'vs', '0'), @@ -776,6 +822,7 @@ OVEN_POLL_TIERS = [ PollTier( name='cold', interval_s=600.0, + timeout_s=6.0, paths=( ('otninformation', 'vs', '0'), ), @@ -783,6 +830,7 @@ OVEN_POLL_TIERS = [ PollTier( name='sweep', interval_s=300.0, + timeout_s=15.0, paths=(('device', '0'),), is_sweep=True, ), diff --git a/samsung_appliance/bridge.py b/samsung_appliance/bridge.py index 1e5d2e4..d7d66c8 100644 --- a/samsung_appliance/bridge.py +++ b/samsung_appliance/bridge.py @@ -17,6 +17,7 @@ MQTT client; each owns one DTLS session, one cache, one scheduler. """ import json import os +import random import threading import time @@ -27,6 +28,7 @@ from .coap_dtls import DtlsCoapSession, fmt_code from .config import ApplianceConfig, SharedConfig from .keepalive import KeepaliveTask from .logger import bridge_logger +from .observe_refresh import ObserveRefreshTask from .poll_scheduler import PollScheduler from .sensors import index_links from .state_cache import StateCache @@ -42,6 +44,21 @@ def _href_to_segs(href: str) -> list[str]: SERIAL_PATH = '/information/vs/0' SERIAL_FIELD = 'x.com.samsung.da.serialNum' +# If the keepalive watchdog flags the device unreachable for this long, +# force a session reconnect. Catches the half-open case where the DTLS +# socket is still writable but the peer has gone silent — without this, +# the bridge sits in offline state until the reader thread dies on its +# own (which may not happen at all if the OS sees no socket errors). +UNREACHABLE_RECONNECT_S = 120.0 + +# Periodic OBSERVE re-subscribe interval. Safety net for the case where +# the device stays reachable on the DTLS layer but Samsung's RT-OCF +# clears its observer table (e.g. during cloud auth blips). Without +# this, push delivery stays dead even after upstream connectivity +# recovers, since nothing triggers a fresh subscribe on the existing +# session. +OBSERVE_REFRESH_INTERVAL_S = 6 * 3600.0 + class PushBridge: @@ -63,6 +80,7 @@ class PushBridge: self.session: DtlsCoapSession | None = None self.scheduler: PollScheduler | None = None self.keepalive: KeepaliveTask | None = None + self.observe_refresh: ObserveRefreshTask | None = None self.cache = StateCache(descriptor) self.cache.set_on_change(self._on_cache_change) @@ -83,6 +101,11 @@ class PushBridge: self._last_change_source: str | None = None self._last_observe_change_ts: float | None = None self._last_push_active_pub: str | None = None + # Wall-clock timestamp the keepalive watchdog first reported the + # device unreachable on the current session. Cleared on recovery + # or session start. Drives the force-reconnect watchdog below. + self._unreachable_since: float | None = None + self._force_close_in_flight: bool = False # Push is considered "active" if an OBSERVE-sourced change # arrived within this window. Long enough that a quiet but @@ -222,6 +245,8 @@ class PushBridge: self.connect_count += 1 self.cache.descriptor_state.clear() self._publish_gate = False + self._unreachable_since = None + self._force_close_in_flight = False self.log.info("DTLS connected — subscribing %d paths", len(self.descriptor.observe_paths)) @@ -274,6 +299,10 @@ class PushBridge: is_active_fn=self.descriptor.is_active, logger=self.log, ) + # Half-open detection: if no successful poll lands inside this + # window the session is wedged, regardless of whether ping sends + # leave the socket. 60s gives ~60 hot-tier cycles of margin. + liveness_window_s = 60.0 keepalive = KeepaliveTask( sess, interval_s=float(self.shared.PING_INTERVAL_S), @@ -281,9 +310,18 @@ class PushBridge: on_reachable=self._on_reachable, on_unreachable=self._on_unreachable, logger=self.log, + liveness_fn=lambda: (time.monotonic() - scheduler.last_success_ts + ) < liveness_window_s, + ) + observe_refresh = ObserveRefreshTask( + sess, + paths=self.descriptor.observe_paths, + interval_s=OBSERVE_REFRESH_INTERVAL_S, + logger=self.log, ) self.scheduler = scheduler self.keepalive = keepalive + self.observe_refresh = observe_refresh sched_t = threading.Thread( target=scheduler.run_forever, args=(self.stop,), @@ -291,14 +329,19 @@ class PushBridge: ka_t = threading.Thread( target=keepalive.run_forever, args=(self.stop,), daemon=True, name=f'{self.app.klass}-ping') + ref_t = threading.Thread( + target=observe_refresh.run_forever, args=(self.stop,), + daemon=True, name=f'{self.app.klass}-obsref') sched_t.start() ka_t.start() + ref_t.start() try: sess.join() finally: self.scheduler = None self.keepalive = None + self.observe_refresh = None def _seed_from_device0(self, sess): code, pl = sess.get(self.descriptor.seed_path, timeout=15.0) @@ -340,12 +383,37 @@ class PushBridge: self.log.warning("oic/res get: %s", e) def _on_reachable(self) -> None: + self._unreachable_since = None self.set_availability(True) self.reassert_availability() def _on_unreachable(self) -> None: + if self._unreachable_since is None: + self._unreachable_since = time.time() self.set_availability(False) + def _maybe_force_reconnect(self) -> None: + """If the device has been unreachable for UNREACHABLE_RECONNECT_S, + close the DTLS session to break run_forever's session_once() out + of its sess.join() and trigger a fresh connect. Without this, a + half-open session (writable socket, silent peer) holds the bridge + in offline limbo until the OS surfaces a socket error.""" + if self._unreachable_since is None or self._force_close_in_flight: + return + elapsed = time.time() - self._unreachable_since + if elapsed < UNREACHABLE_RECONNECT_S: + return + sess = self.session + if sess is None: + return + self.log.warning( + "unreachable for %.0fs — forcing session reconnect", elapsed) + self._force_close_in_flight = True + try: + sess.close() + except Exception as e: + self.log.warning("force-close: %s", e) + # ---- MQTT publishing -------------------------------------------- def maybe_publish_state(self, force=False): @@ -519,8 +587,8 @@ class PushBridge: self._win_prev_poll = poll self._win_prev_poll_err = poll_err self._win_prev_ping_fail = ping_fail - win_max_rtt, win_slow = (sched.take_window_stats() - if sched else (0.0, 0)) + win_max_rtt, win_slow, win_timeouts = ( + sched.take_window_stats() if sched else (0.0, 0, 0)) window_polls_ok = max(0, d_poll - d_err) h = { @@ -536,6 +604,7 @@ class PushBridge: 'poll_window_errors': d_err, 'poll_window_max_rtt_ms': round(win_max_rtt, 0), 'poll_window_slow_count': win_slow, + 'poll_window_timeout_count': win_timeouts, 'ping_count': ka.ping_count if ka else 0, 'ping_fail_count': ping_fail, 'reachable': ka.reachable if ka else None, @@ -561,13 +630,15 @@ class PushBridge: self.log.warning("health publish: %s", e) self.publish_push_active(push_active) + self._maybe_force_reconnect() if d_poll > 0 or d_err > 0 or d_ping_fail > 0: self.log.info( "poll-window: %d ok, %d err, %d ping-fail, " - "p_max=%.0fms, slow=%d (%ds)", + "p_max=%.0fms, slow=%d, timeouts=%d (%ds)", window_polls_ok, d_err, d_ping_fail, - win_max_rtt, win_slow, self.shared.HEALTH_INTERVAL_S) + win_max_rtt, win_slow, win_timeouts, + self.shared.HEALTH_INTERVAL_S) def publish_push_active(self, active: bool, force: bool = False) -> None: value = 'online' if active else 'offline' @@ -601,8 +672,12 @@ class PushBridge: self.session_started_ts = None if self.stop.is_set(): break - wait = min(backoff, 30.0) - self.log.info("reconnect in %.0fs", wait) + # Jitter the backoff so multiple bridges (dryer + oven) don't + # reconnect in lockstep after a router blip — synchronized + # storms make the broker / DTLS layer flap harder than need + # be. ±30% noise spreads the retry attempts. + wait = min(backoff, 30.0) * random.uniform(0.7, 1.3) + self.log.info("reconnect in %.1fs", wait) if self.stop.wait(wait): break backoff = min(backoff * 2, 30.0) diff --git a/samsung_appliance/coap_dtls.py b/samsung_appliance/coap_dtls.py index 396d333..de64044 100644 --- a/samsung_appliance/coap_dtls.py +++ b/samsung_appliance/coap_dtls.py @@ -2,7 +2,7 @@ Replaces the TLS-over-TCP transport used in the original dryer bridge. Both the oven (UDP/49154) and the dryer (UDP/49155) speak CoAP-over-DTLS -with the ECDHE-ECDSA-AES128-GCM-SHA256 cipher and ab0b0ac4 client cert. +with the ECDHE-ECDSA-AES128-GCM-SHA256 cipher and a client cert. Wire-level details that matter (from local-tools/oven-findings.md §17): * DTLS ciphertext MTU must be 1200; otherwise OpenSSL fragments the @@ -226,10 +226,8 @@ class DtlsCoapSession: ConnectionError / TimeoutError on failure.""" ctx = SSL.Context(SSL.DTLS_METHOD) ctx.set_verify(SSL.VERIFY_NONE, lambda *_: True) - # @SECLEVEL=0 needed because the AC14K_M-rooted ab0b0ac4 chain - # is SHA-1 signed, which OpenSSL 3.x's default security level - # rejects. The chain comes from Samsung's leaked CA so the - # signature algorithm isn't ours to change. + # @SECLEVEL=0 — the AC14K_M-rooted chain is SHA-1 signed, which + # OpenSSL 3.x's default security level rejects. ctx.set_cipher_list(b'ECDHE-ECDSA-AES128-GCM-SHA256:@SECLEVEL=0') ctx.use_certificate_chain_file(self.cert_path) ctx.use_privatekey_file(self.key_path) @@ -592,25 +590,51 @@ class DtlsCoapSession: def ping(self): """RFC 7252 §4.4 CoAP Ping — empty CON, no token, no payload. - Peer MUST respond with an RST sharing the same Message ID. - We don't block waiting for the RST here; the reader thread - sees it, finds nothing in _pending/_observe_tokens for an - empty token, and drops it quietly — which is the documented - behaviour for unsolicited RST. + Fire-and-forget: we do not wait for the matching RST because + Samsung's RT-OCF doesn't reliably emit one (verified + 2026-06-04: every sync ping timed out while polls succeeded + at 200+/window). The send itself is the keepalive — it + tickles Samsung's observer state so OBSERVE subscriptions + aren't aged out. - Sole purpose is keepalive: gives Samsung's RT-OCF stack a - visible "client still here" signal so it doesn't expire our - OBSERVE subscriptions (symptom of expiry: POSTs still get - 2.04 but OBSERVE notifications stop arriving).""" + Real half-open-session detection lives in PollScheduler's + `last_success_ts`, surfaced through KeepaliveTask's + `liveness_fn`.""" if self.conn is None: raise ConnectionError("DTLS session closed") mid = self._next_mid() - # Empty message: ver=01, type=CON, tkl=0, code=0.00, mid, no - # token, no options, no payload. build_coap with empty token - # and no options yields exactly that. self._send_dgram(build_coap(TYPE_CON, 0, mid, b'', [])) return mid + def refresh_observes(self, paths): + """Drop all current OBSERVE registrations and re-subscribe to + the given paths. Used as a periodic safety net — CoAP OBSERVE + has no built-in TTL but Samsung's RT-OCF can age out its + observer table during cloud blips even while the DTLS session + stays healthy. Without this, internet recovery on a still- + reachable device leaves push permanently dead. + + Best-effort: dereg failures are logged and we still bind fresh + tokens via subscribe. Brief race window where a notify on the + old token gets dropped as 'stale' — acceptable for a 6h-scale + safety net.""" + if self.conn is None: + raise ConnectionError("DTLS session closed") + for tok, href in list(self._observe_tokens.items()): + segs = [s for s in href.split('/') if s] + try: + self._send_observe_dereg(tok, segs) + except Exception as e: + logger.warning("refresh dereg %s: %s", href, e) + self._observe_tokens.clear() + time.sleep(0.1) + for path in paths: + try: + self.subscribe(list(path)) + time.sleep(0.05) + except Exception as e: + logger.warning("refresh subscribe %s: %s", path, e) + def subscribe(self, path_segs): """Register an OBSERVE on the given path. The initial 2.05 notification and all subsequent state-change notifications diff --git a/samsung_appliance/keepalive.py b/samsung_appliance/keepalive.py index f704c40..9e47ac3 100644 --- a/samsung_appliance/keepalive.py +++ b/samsung_appliance/keepalive.py @@ -1,8 +1,19 @@ -"""DTLS-layer liveness via CoAP empty-CON ping. +"""DTLS-layer liveness via CoAP empty-CON ping + poll-success watchdog. -Pings the appliance every interval_s. After fail_threshold consecutive -failures, fires on_unreachable. First success after a fail streak fires -on_reachable. Bridge wires these to MQTT availability. +Each interval_s the task: + 1. Sends a CoAP ping. This is fire-and-forget — Samsung's RT-OCF + doesn't reliably reply with an RST, so the send itself is the + keepalive (it tickles Samsung's observer state). The send only + fails if the underlying socket is gone, in which case the failure + counts toward fail_threshold. + 2. Calls liveness_fn() if provided. This is the real half-open + detection: PollScheduler exposes last_success_ts, and the bridge + wraps it as "did we get a 2.05 in the last 60s". If not, count + the tick as a failure even though the ping send succeeded. + +After fail_threshold consecutive failures, fires on_unreachable. First +success after a fail streak fires on_reachable. Bridge wires these to +MQTT availability. """ from __future__ import annotations @@ -20,13 +31,15 @@ class KeepaliveTask: fail_threshold: int = 3, on_reachable: Optional[Callable[[], None]] = None, on_unreachable: Optional[Callable[[], None]] = None, - logger=None): + logger=None, + liveness_fn: Optional[Callable[[], bool]] = None): self.session = session self.interval_s = interval_s self.fail_threshold = fail_threshold self.on_reachable = on_reachable self.on_unreachable = on_unreachable self.log = logger + self.liveness_fn = liveness_fn self._fail_streak = 0 self._reachable = True @@ -56,6 +69,21 @@ class KeepaliveTask: ok = True except Exception as e: if self.log: self.log.warning("ping: %s", e) + # Real half-open detection: ping sends can succeed against a + # silently-wedged peer, but polls won't. If the scheduler + # hasn't recorded a 2.05 inside the liveness window, treat + # this tick as a failure even though the ping itself went out. + if ok and self.liveness_fn is not None: + try: + alive = bool(self.liveness_fn()) + except Exception as e: + if self.log: self.log.warning("liveness_fn: %s", e) + alive = True + if not alive: + if self.log: + self.log.warning("liveness: no successful poll " + "in the liveness window") + ok = False self._ping_count += 1 if ok: if not self._reachable: diff --git a/samsung_appliance/logger.py b/samsung_appliance/logger.py index 7661908..55d6439 100644 --- a/samsung_appliance/logger.py +++ b/samsung_appliance/logger.py @@ -8,17 +8,55 @@ Two logger families share the root handler: via bridge_logger() once the seed reveals the appliance's serial. Both trees propagate to the root logger, which is the one with the -StreamHandler, so the same format applies everywhere.""" +StreamHandler, so the same format applies everywhere. + +WARNING lines render yellow, ERROR red, when colour is enabled. Set +`NO_COLOR=1` in the environment to fall back to plain text (e.g. when +writing logs to a file).""" import logging +import os import sys -logging.basicConfig( - level=logging.INFO, - format='%(asctime)s %(levelname)-5s %(name)-26s %(message)s', +_LEVEL_COLOURS = { + logging.WARNING: '\033[33m', # yellow + logging.ERROR: '\033[31m', # red + logging.CRITICAL: '\033[1;31m', # bold red +} +_RESET = '\033[0m' + + +def _colour_enabled() -> bool: + return 'NO_COLOR' not in os.environ + + +class _LevelColourFormatter(logging.Formatter): + """Wraps the rendered line in an ANSI colour for WARNING+ levels. + INFO/DEBUG pass through unchanged so the bulk of the log stays + readable and a yellow line draws the eye.""" + + def __init__(self, fmt: str, datefmt: str, use_colour: bool): + super().__init__(fmt=fmt, datefmt=datefmt) + self.use_colour = use_colour + + def format(self, record: logging.LogRecord) -> str: + line = super().format(record) + if not self.use_colour: + return line + colour = _LEVEL_COLOURS.get(record.levelno) + if colour is None: + return line + return f"{colour}{line}{_RESET}" + + +_handler = logging.StreamHandler(sys.stdout) +_handler.setFormatter(_LevelColourFormatter( + fmt='%(asctime)s %(levelname)-5s %(name)-26s %(message)s', datefmt='%H:%M:%S', - handlers=[logging.StreamHandler(sys.stdout)], -) + use_colour=_colour_enabled(), +)) + +logging.basicConfig(level=logging.INFO, handlers=[_handler], force=True) logger = logging.getLogger("samsung_appliance") diff --git a/samsung_appliance/observe_refresh.py b/samsung_appliance/observe_refresh.py new file mode 100644 index 0000000..bb8c24d --- /dev/null +++ b/samsung_appliance/observe_refresh.py @@ -0,0 +1,50 @@ +"""Periodic OBSERVE re-subscribe. + +CoAP OBSERVE (RFC 7641) has no built-in TTL, but real-world peers age +out observer state on their own schedule — Samsung's RT-OCF is known +to silently drop notify delivery during cloud auth blips even though +the DTLS session stays healthy. Without a re-subscribe, recovery from +such a blip requires a full session reconnect. + +This task derregisters the current observer tokens and re-subscribes +every `interval_s`. Cheap (one register CON per path), idempotent +(Samsung silently no-ops a register on an already-active token), and +resilient — individual subscribe failures are logged but don't abort +the task. +""" +from __future__ import annotations + +import threading +from typing import Optional + +from .coap_dtls import DtlsCoapSession + + +class ObserveRefreshTask: + + def __init__(self, + session: DtlsCoapSession, + paths, + interval_s: float = 6 * 3600.0, + logger=None): + self.session = session + self.paths = [list(p) for p in paths] + self.interval_s = interval_s + self.log = logger + self._refresh_count = 0 + + @property + def refresh_count(self) -> int: + return self._refresh_count + + def run_forever(self, stop: threading.Event) -> None: + while not stop.wait(self.interval_s): + try: + self.session.refresh_observes(self.paths) + self._refresh_count += 1 + if self.log: + self.log.info("OBSERVE refresh #%d (%d paths)", + self._refresh_count, len(self.paths)) + except Exception as e: + if self.log: + self.log.warning("OBSERVE refresh: %s", e) diff --git a/samsung_appliance/poll_scheduler.py b/samsung_appliance/poll_scheduler.py index 2023817..a18ab15 100644 --- a/samsung_appliance/poll_scheduler.py +++ b/samsung_appliance/poll_scheduler.py @@ -8,6 +8,20 @@ result into many href reps. Adaptive cadence: when descriptor.is_active(cache.links) returns True and tier.active_interval_s is set, that tier uses the tighter cadence. +If the previous health window saw `active_throttle_threshold` timeouts, +the throttle drops back to idle cadence even when active=True — the +RT-OCF stack wedges under load and stacking poll attempts only makes +it worse. + +Per-tier timeouts: tier.timeout_s overrides the scheduler default. Hot +tiers want a tight ceiling (e.g. 2s) so one wedged path can't eat +several poll cycles. + +Cooldown on timeout: when a path times out, it's deferred for ~3 of +its tier's intervals (clamped 5–60s) via the same _defer_until mechanism +used for write_in_progress. Breaks the cluster-cascade where the next +tier tick fires immediately after an 8s wedge and reattempts the same +stalled path. Post-write defer: bridge calls write_in_progress(href) before POSTing a write; the scheduler skips that href for settle_s to avoid Samsung's @@ -35,6 +49,10 @@ class PollTier: paths: tuple[tuple[str, ...], ...] active_interval_s: Optional[float] = None is_sweep: bool = False + # Per-tier CoAP request timeout. Falls back to PollScheduler.timeout_s + # when None. Hot tiers want a tight ceiling (e.g. 2s) so one wedged + # path can't eat several poll cycles; sweep tiers tolerate longer. + timeout_s: Optional[float] = None class PollScheduler: @@ -46,7 +64,8 @@ class PollScheduler: sweep_index_fn: Callable[[object], dict[str, dict]], is_active_fn: Optional[Callable[[dict[str, dict]], bool]] = None, logger=None, - timeout_s: float = 8.0): + timeout_s: float = 8.0, + active_throttle_timeout_threshold: int = 3): self.session = session self.cache = cache self.tiers = tiers @@ -54,6 +73,7 @@ class PollScheduler: self.is_active_fn = is_active_fn self.log = logger self.timeout_s = timeout_s + self.active_throttle_threshold = active_throttle_timeout_threshold now = time.monotonic() self._next_due: dict[str, float] = {t.name: now for t in tiers} @@ -62,12 +82,24 @@ class PollScheduler: self._poll_count = 0 self._poll_error_count = 0 self._last_active: Optional[bool] = None + self._last_throttled: bool = False + # Real-liveness signal for KeepaliveTask: updated on every 2.05 + # we receive from the wire. Initialized to "now" so the first + # keepalive tick after start doesn't fire a false unreachable. + self._last_success_ts: float = now # Per-window tail-latency tracking. Bridge consumes-and-resets # these via take_window_stats() once per HEALTH_INTERVAL_S. + # _window_max_rtt_ms tracks SUCCESSFUL polls only — timeouts go + # into _window_timeout_count so the dashboard sees the real tail + # instead of an 8000ms wall. self._stats_lock = threading.Lock() self._window_max_rtt_ms = 0.0 self._window_slow_count = 0 + self._window_timeout_count = 0 + # Snapshot of the previous window's timeout count, used by the + # active-window throttle to back off when polls are wedging. + self._last_window_timeouts = 0 self.slow_threshold_ms = 1000.0 def write_in_progress(self, href: str, settle_s: float = 4.0) -> None: @@ -89,17 +121,34 @@ class PollScheduler: def poll_error_count(self) -> int: return self._poll_error_count - def take_window_stats(self) -> tuple[float, int]: - """Return (max RTT ms, slow-poll count) seen since the last call, - and reset both. Slow threshold is `self.slow_threshold_ms`.""" + @property + def last_success_ts(self) -> float: + """Monotonic timestamp of the most recent 2.05 response from any + tier. KeepaliveTask uses this as its half-open-detection signal: + if no 2.05 has landed in `liveness_window_s`, the session is + wedged regardless of whether ping sends succeed.""" + return self._last_success_ts + + def take_window_stats(self) -> tuple[float, int, int]: + """Return (max RTT ms over successful polls, slow-poll count, + timeout count) seen since the last call, and reset all three. + Slow threshold is `self.slow_threshold_ms`. The timeout count + is snapshotted into `_last_window_timeouts` for the throttle.""" with self._stats_lock: - out = (self._window_max_rtt_ms, self._window_slow_count) + out = (self._window_max_rtt_ms, + self._window_slow_count, + self._window_timeout_count) + self._last_window_timeouts = self._window_timeout_count self._window_max_rtt_ms = 0.0 self._window_slow_count = 0 + self._window_timeout_count = 0 return out - def _record_rtt(self, rtt_ms: float) -> None: + def _record_rtt(self, rtt_ms: float, *, timed_out: bool = False) -> None: with self._stats_lock: + if timed_out: + self._window_timeout_count += 1 + return if rtt_ms > self._window_max_rtt_ms: self._window_max_rtt_ms = rtt_ms if rtt_ms >= self.slow_threshold_ms: @@ -120,11 +169,29 @@ class PollScheduler: if self.log and self._last_active is not None: self.log.info("active=%s", active) self._last_active = active + # Active-window throttle: if the previous health window saw a + # cluster of timeouts (RT-OCF wedging under load), drop back to + # idle cadence even when is_active=True. Lets the device breathe + # instead of stacking poll attempts on a stalled responder. + with self._stats_lock: + recent_to = self._last_window_timeouts + throttled = (active + and recent_to >= self.active_throttle_threshold) + if throttled != self._last_throttled: + if self.log: + if throttled: + self.log.warning( + "active-throttle ON (%d timeouts last window) — " + "using idle cadence", recent_to) + else: + self.log.info("active-throttle OFF") + self._last_throttled = throttled + effective_active = active and not throttled for tier in self.tiers: if self._next_due[tier.name] > now: continue interval = (tier.active_interval_s - if (active and tier.active_interval_s is not None) + if (effective_active and tier.active_interval_s is not None) else tier.interval_s) self._next_due[tier.name] = now + interval try: @@ -136,7 +203,23 @@ class PollScheduler: self._poll_error_count += 1 if self.log: self.log.warning("tier %s: %s", tier.name, e) + def _tier_timeout(self, tier: PollTier) -> float: + return tier.timeout_s if tier.timeout_s is not None else self.timeout_s + + def _cooldown_for(self, tier: PollTier) -> float: + # On timeout, defer the wedged href for ~3 cycles (clamped to a + # 5–60s band) so repeated tier ticks don't stack attempts on a + # stalled responder. Breaks the cluster-cascade we see in the + # Poll Max RTT chart during heavy device use. + return max(5.0, min(60.0, tier.interval_s * 3.0)) + + def _set_cooldown(self, href: str, cooldown_s: float) -> None: + with self._defer_lock: + self._defer_until[href] = time.monotonic() + cooldown_s + def _do_tier(self, tier: PollTier) -> None: + timeout = self._tier_timeout(tier) + cooldown = self._cooldown_for(tier) for path in tier.paths: href = '/' + '/'.join(path) with self._defer_lock: @@ -145,7 +228,15 @@ class PollScheduler: self._poll_count += 1 t0 = time.monotonic() try: - code, body = self.session.get(list(path), timeout=self.timeout_s) + code, body = self.session.get(list(path), timeout=timeout) + except TimeoutError: + self._poll_error_count += 1 + self._record_rtt(0.0, timed_out=True) + self._set_cooldown(href, cooldown) + if self.log: + self.log.warning("poll %s timeout (cooldown %.0fs)", + href, cooldown) + return except Exception as e: self._poll_error_count += 1 self._record_rtt((time.monotonic() - t0) * 1000.0) @@ -156,6 +247,7 @@ class PollScheduler: self._poll_error_count += 1 if self.log: self.log.warning("poll %s -> %s", href, fmt_code(code)) continue + self._last_success_ts = time.monotonic() try: rep = cbor2.loads(body) except Exception as e: @@ -166,11 +258,22 @@ class PollScheduler: self.cache.apply_rep(href, rep, source='poll') def _do_sweep(self, tier: PollTier) -> None: + timeout = self._tier_timeout(tier) + cooldown = self._cooldown_for(tier) path = list(tier.paths[0]) + href = '/' + '/'.join(path) t0 = time.monotonic() self._poll_count += 1 try: - code, body = self.session.get(path, timeout=self.timeout_s) + code, body = self.session.get(path, timeout=timeout) + except TimeoutError: + self._poll_error_count += 1 + self._record_rtt(0.0, timed_out=True) + self._set_cooldown(href, cooldown) + if self.log: + self.log.warning("sweep %s timeout (cooldown %.0fs)", + path, cooldown) + return except Exception as e: self._poll_error_count += 1 self._record_rtt((time.monotonic() - t0) * 1000.0) @@ -181,6 +284,7 @@ class PollScheduler: self._poll_error_count += 1 if self.log: self.log.warning("sweep -> %s", fmt_code(code)) return + self._last_success_ts = time.monotonic() try: tree = cbor2.loads(body) except Exception as e: