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.
This commit is contained in:
@@ -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'),
|
||||
|
||||
@@ -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,
|
||||
),
|
||||
|
||||
@@ -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,
|
||||
),
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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")
|
||||
|
||||
|
||||
@@ -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)
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user