Refactor to polling-first architecture with OBSERVE as accelerator

State freshness now comes from a tiered PollScheduler over the persistent
DTLS session; OBSERVE registrations are kept as an opportunistic
acceleration layer. Behaviour is identical online vs air-gapped except
for worst-case freshness latency.

Adds three modules:
- StateCache: single source of truth, source-tagged change events
- PollScheduler: hot/warm/cold + sweep tiers, write-defer past the
  fetchback-revert window, per-window RTT/slow-poll tracking
- KeepaliveTask: CoAP empty-CON ping with consecutive-fail detection
  driving MQTT availability

Bridge publishes per-appliance diagnostic entities (Push Active, Last
Update Source, Poll Max RTT, Slow Polls, Poll Errors, Stalest Resource
Age, Last OBSERVE Age) under HA's Diagnostic section. Tier cadences
are descriptor-declared, calibrated against measured per-firmware
ceilings (dryer ~14 req/s, oven ~8 req/s via probe_poll_rate_combined.py).

Drops HEARTBEAT_INTERVAL_S in favour of the descriptor-declared sweep
tier; PING_INTERVAL_S now consumed by KeepaliveTask inside the bridge
rather than driven from main.py.

README explains the push/poll split and what happens when the appliance
is blocked from internet.
This commit is contained in:
Jack Nagy
2026-06-03 18:38:04 +01:00
parent dbc9a57f1b
commit b00c2fd90b
11 changed files with 866 additions and 324 deletions
+8 -7
View File
@@ -43,14 +43,15 @@ MQTT_PASS=
HA_DISCOVERY_PREFIX=homeassistant
# Bridge timers (seconds).
# HEALTH_INTERVAL_S — how often <prefix>/bridge/health republishes.
# HEARTBEAT_INTERVAL_S — periodic full /device/0 re-seed. Refreshes
# every resource (observed too) — useful for
# appliances like the oven that don't reliably
# push OBSERVE on /mode/vs/0 option changes. 0
# disables.
# HEALTH_INTERVAL_S — how often <prefix>/bridge/health republishes.
# PING_INTERVAL_S — CoAP empty-CON ping cadence (DTLS-layer
# liveness). Three consecutive failures publish
# availability=offline.
# State freshness itself comes from the in-bridge PollScheduler whose
# tier cadences are declared in the appliance descriptor — there is
# no top-level heartbeat env var to tune.
HEALTH_INTERVAL_S=60
HEARTBEAT_INTERVAL_S=600
PING_INTERVAL_S=25
# Container TZ.
TZ=Europe/London
+31 -19
View File
@@ -16,9 +16,9 @@
### What you get
- **Multi-appliance, one container.** Single Docker service holds N DTLS sessions in parallel, one per appliance, sharing one MQTT client. Adding an appliance class is ~150 lines and one descriptor file.
- **Sub-second push for state changes.** Cycle starts, pauses, ends, course changes, door opens, lamp toggles — Home Assistant reflects it within ~1 second on appliances that push OBSERVE notifications, or after the 3-second post-POST verify on appliances that don't.
- **Bounded state latency.** Hot-tier resources (job state, door, operational state) refresh on a sub-second cadence regardless of whether the appliance has internet. Worst-case lag is the tier interval (≤1s idle, ≤500ms during an active cycle on the dryer).
- **Writes that work**: dryer Start/Pause/Stop, course selection, wrinkle prevent; oven lamp (light entity), sound, fast preheat, setpoint slider, mode select, stop.
- **Optimistic publish + verify**: HA sees the new value the instant the device 2.04-confirms the write; a Block2 fetch-back 3 seconds later corrects if the device silently coerced or rejected.
- **Optimistic publish + verify**: HA sees the new value the instant the device 2.04-confirms the write; the PollScheduler verifies on its next tier tick (after a 4s defer past Samsung's fetchback-revert window).
- **HA Energy Dashboard ready** (dryer): live watts + cumulative kWh as `total_increasing`.
- **Bridge logs tagged per-appliance** with `<class>.<serial>` once each device's serial is read on connect — `dryer.<serial>` vs `oven.<serial>` interleaved in the same log stream, easy to grep.
- **Zero HA YAML.** Every entity is auto-discovered via MQTT discovery.
@@ -26,7 +26,7 @@
### Under the hood
Each appliance runs an independent push-mode bridge: one sustained DTLS session, CoAP OBSERVE (RFC 7641) on ~11 of the appliance's `/<x>/vs/0` resources, token-stable Block2 (RFC 7959) for the multi-block reads, optimistic state publish + Block2 fetch-back verification after every write. Reconnect with exponential backoff on session errors.
Each appliance runs an independent bridge built around three coordinated pieces over one persistent DTLS session: a `StateCache` (single source of truth for all reps), a `PollScheduler` (tiered adaptive polling — hot/warm/cold + a periodic `/device/0` sweep), and a `KeepaliveTask` (CoAP empty-CON ping for DTLS-layer liveness, with consecutive-failure detection for MQTT availability). Tier cadences are descriptor-declared and were calibrated against the empirically-measured per-firmware ceilings (`local-tools/probe_poll_rate_combined.py`): dryer ~14 req/s, oven ~8 req/s. OBSERVE registrations (RFC 7641) are kept as an opportunistic freshness accelerator — when the appliance has internet and emits notifications, the cache absorbs them and the next-poll timer is reset for that resource; when it's air-gapped, polling alone carries the UX with no other code change. Token-stable Block2 (RFC 7959) handles multi-block reads. Writes are optimistically merged into the cache the moment the device 2.04-confirms, with the scheduler deferring that resource's next poll past the fetchback-revert window. Reconnect with exponential backoff on session errors.
Authentication uses **Samsung's publicly-published cloud-bridge identity** (UUID `ab0b0ac4-…`), present in every Samsung Tizen/RT-OCF appliance's factory ACL with `perm=31` (full CRUDN) on `href=*`. One cert chain works across the whole fleet. Setup is one Python script.
@@ -50,13 +50,26 @@ Read the result:
| Appliance class | Model family | Confirmed |
|---|---|---|
| Dryer | DV5000T (`DA_WM_TP2_20_COMMON`, `mnid=0AJT`) | All entities, sub-second OBSERVE push |
| Oven | NV7000BS-class (`TP1X_DA-KS-OVEN-0107X`, `mnid=0AJT`) | All entities; OBSERVE-push lazy on options-array writes (see "Per-appliance notes" below) |
| Dryer | DV5000T (`DA_WM_TP2_20_COMMON`, `mnid=0AJT`) | All entities, ≤1s hot-tier poll (OBSERVE accelerates when online) |
| Oven | NV7000BS-class (`TP1X_DA-KS-OVEN-0107X`, `mnid=0AJT`) | All entities; hot-tier poll covers door + operational state regardless of cloud reachability |
Other appliances on the same firmware family (washers, dishwashers, AC units) almost certainly speak the same protocol — the auth path and read primitives are common. You'd write one new descriptor in `samsung_appliance/appliances/`.
---
## How the app keeps in sync with the appliance
There are two parallel paths between the appliance and the app over the local CoAP-DTLS socket:
- **Push (OBSERVE).** When the appliance can reach Samsung's cloud, it emits a CoAP OBSERVE notification on the LAN socket within ~100ms of any state change — cycle start, door open, mode flip. The notification travels over the LAN; nothing about the push itself routes via Samsung. **But** the appliance's decision to emit it at all is gated inside its cloud-publish thread. Block the appliance from the internet and the LAN OBSERVE pushes stop, even though the LAN path itself is unaffected and the appliance still answers reads + accepts writes normally.
- **Polling.** The app always polls a small tier of hot resources (operational state, door, etc.) on a sub-second cadence, a warmer tier (mode, kidslock, alarms, …) every 15–30 s, and a full `/device/0` sweep every 5 minutes. This carries the UX regardless of whether OBSERVE is firing.
In normal operation both happen at once: an OBSERVE notification arrives first, the cache absorbs it, and the next-poll timer for that resource is reset. In an air-gapped LAN the app keeps working — only the worst-case freshness changes (from ~100 ms with push to ≤1 s on hot-tier resources via polling). Reads, writes, and HA entities behave identically.
Which path is doing the work is visible in Home Assistant. The bridge publishes per-appliance diagnostic entities including **Push Active** (on while OBSERVE is firing), **Last Update Source** (`observe` / `poll` / `sweep` / `optimistic`), **Last OBSERVE Age**, **Poll Max RTT**, **Slow Polls (window)**, **Poll Errors (window)**, and **Stalest Resource Age** — all under each device's Diagnostic section.
---
## Part 2 — Auth: get the cloud-identity cert
The bridge authenticates with a **client cert** signed by `AC14K_M` (Samsung's leaked diagnostic intermediate CA — used inside Samsung tooling and still trusted by current firmware). The cert's Subject DN contains the cloud-bridge UUID Samsung publishes on its wildcard cloud TLS cert at `*.samsungiotcloud.com`.
@@ -214,7 +227,7 @@ In HA: **Settings → Devices & Services → MQTT** should show both devices pop
| Power on/off | ❌ | Accepted (2.04) but reverts within seconds — hardware-mirrored |
| Child Lock / Remote Control toggle | ❌ | Same — hardware-mirrored physical buttons |
The dryer pushes OBSERVE notifications on every state-changing write within ~100ms. State propagation is sub-second.
The dryer's `/operational/state/vs/0` is on the bridge's hot poll tier (1s idle / 0.5s while a cycle is active) and also accepts OBSERVE registration. When the appliance has internet it pushes notifications within ~100ms of any state change and the cache absorbs them as fast freshness; when air-gapped the hot-tier poll carries the same UX with worst-case lag of one tier interval.
### Oven
@@ -222,17 +235,16 @@ The dryer pushes OBSERVE notifications on every state-changing write within ~100
|---|---|---|
| Read state | ✅ | Cavity state, current/target temp, door, mode, alarms, firmware-update-available |
| Lamp (light entity) | ✅ | Binary On/Off only — High/Low/Dim values are accepted (2.04) but silently coerced back. Works regardless of Remote Control. |
| Sound, Fast preheat | ⚠️ | Wired but untested write-side; RC-gated as a safety. |
| Setpoint slider | ⚠️ | Wired but untested mid-cook behaviour. RC-gated. |
| Mode select | ⚠️ | Wired but untested mid-cook behaviour. RC-gated. |
| Stop button | ⚠️ | Wired but untested. **Not** RC-gated (the SmartThings app stops without Remote Control on, so we don't gate either). |
| Power on/off as a switch | ❌ | Not exposed as a writeable entity — cold-start panel is a physical action. Read-only sensor only. |
| Sound, Fast preheat | ⚠️ | Wired but untested RC-gated. |
| Setpoint slider | ⚠️ | Wired but untested RC-gated. |
| Mode select | ⚠️ | Wired but untested RC-gated. |
| Stop button | ✅ | |
| **Kitchen timer (`⏲` icon)** | ❌ | **The oven's panel kitchen timer is not exposed via CoAP at all.** Confirmed by full `/device/0` dump — `UpperTimer*` fields in `/mode/vs/0` only populate when set via the API, not from the panel. |
**The oven doesn't push OBSERVE on `/mode/vs/0` writes** (the dryer does). The bridge defends with:
1. **Optimistic publish** — the moment a POST returns 2.04, the bridge merges the write body into the local state and publishes to MQTT. HA reflects the new value instantly.
2. **Fetch-back verification** — 3 seconds later, the bridge does a token-stable Block2 GET of the just-written resource. If the device's actual state differs from optimistic (silently coerced), the corrected state is republished and HA reverts.
3. **Periodic heartbeat** — every `HEARTBEAT_INTERVAL_S` (default 600s), the bridge re-fetches `/device/0` and refreshes ALL resources (including observed ones), bounding worst-case drift.
**The oven doesn't push OBSERVE on `/mode/vs/0` writes** (the dryer does). The bridge handles this transparently because state freshness comes from polling rather than from OBSERVE:
1. **Optimistic publish** — the moment a POST returns 2.04, the bridge merges the write body into the cache and publishes to MQTT. HA reflects the new value instantly.
2. **Scheduler reconciliation** — the PollScheduler defers polling the just-written resource for ~4s (past Samsung's fetchback-revert window), then refreshes it on its tier cadence. If the device silently coerced the value, the corrected state is republished and HA reverts.
3. **Periodic `/device/0` sweep** — every 5 minutes the scheduler's sweep tier re-fetches the whole device tree, bounding worst-case drift on any resource the per-tier polls don't cover.
---
@@ -252,7 +264,7 @@ The dryer pushes OBSERVE notifications on every state-changing write within ~100
| `HA_DISCOVERY_PREFIX` | HA discovery topic root (default `homeassistant`) |
| `CERT_PATH` / `KEY_PATH` | Override cert lookup (auto-detects `/config/` then `./certs/`) |
| `HEALTH_INTERVAL_S` | Seconds between `<prefix>/bridge/health` publishes (default 60) |
| `HEARTBEAT_INTERVAL_S` | Seconds between full `/device/0` re-seeds; `0` disables (default 600) |
| `PING_INTERVAL_S` | CoAP empty-CON ping cadence; three consecutive failures publish `availability=offline` (default 25). Tier polling cadences are descriptor-declared, not env-tunable. |
| `SSH_HOST` / `REMOTE_DIR` / `APPDATA_DIR` | Used by `deploy.sh` only |
### MQTT topics — outgoing (bridge → broker)
@@ -264,7 +276,7 @@ Per appliance, where `<prefix>` is its `APPLIANCE_<n>_TOPIC`.
| `<prefix>/availability` | ✓ | `online` after seed; `offline` on disconnect (LWT for appliance #1) |
| `<prefix>/remote_available` | ✓ | `online` iff bridge is up AND Remote Control on the appliance is on. Gates the control entities. |
| `<prefix>/state` | ✓ | JSON sensor dict; published only when sensors actually diff |
| `<prefix>/bridge/health` | ✓ | Every `HEALTH_INTERVAL_S` — connect_count, error_count, notif_count, last_change_age_s, session_age_s, serial |
| `<prefix>/bridge/health` | ✓ | Every `HEALTH_INTERVAL_S` — connect_count, error_count, notif_count, poll_count, poll_error_count, ping_count, ping_fail_count, reachable, last_change_age_s, last_seed_age_s, session_age_s, stalest_href, stalest_age_s, serial |
| `<ha_prefix>/{sensor,binary_sensor,switch,light,number,select,button}/<prefix>/.../config` | ✓ | HA MQTT discovery, republished on every MQTT (re)connect |
### MQTT topics — incoming (bridge subscribes)
@@ -348,7 +360,7 @@ The bridge is appliance-agnostic. Adding e.g. a washer is mechanical:
3. Add `WASHER` to `DESCRIPTORS` in `samsung_appliance/appliances/__init__.py`.
4. Add `APPLIANCE_<n>_CLASS=washer` to `.env`, bump `APPLIANCE_COUNT`, redeploy.
The descriptor pattern handles everything else — DTLS, MQTT, HA discovery, optimistic+verify writes, Block2 reads, OBSERVE notifications, reconnect, periodic heartbeat.
For a new appliance class also declare a `poll_tiers: list[PollTier]` and (optionally) an `is_active(links) -> bool` predicate on the descriptor. Hot-tier resources are whatever needs sub-second freshness for HA UX; warm covers everything else that's not static; the `/device/0` sweep tier catches anything you forgot. The descriptor pattern handles everything else — DTLS, MQTT, HA discovery, optimistic writes, Block2 reads, OBSERVE accelerator, reconnect, liveness pings.
---
@@ -357,7 +369,7 @@ The descriptor pattern handles everything else — DTLS, MQTT, HA discovery, opt
These each looked like obvious improvements at some point. Each one broke something.
- **Don't add OBSERVE subscriptions on OCF-standard `/<x>/0` paths.** They register successfully but never push. Use the Samsung `/<x>/vs/0` siblings (which do).
- **Don't half-block the cloud.** Either let the appliance reach Samsung normally (rock-solid local session, sub-second push) or fully block it (the local session tears down every ~30s; bridge reconnects). Don't sinkhole DNS while letting IPs resolve to unreachable hosts — the appliance holds a stable local session but stops emitting OBSERVE pushes entirely. Worst of both worlds.
- **Don't assume OBSERVE silence means the appliance is broken.** When the appliance can't reach Samsung's cloud, its OBSERVE notify dispatch goes quiet even though the local DTLS session, GETs, POSTs, and the cache continue to work normally (measured at `~14 req/s` dryer / `~8 req/s` oven with 200/200 GETs successful while firewalled — see `local-tools/probe_poll_rate_combined.py`). The polling tiers are the structural answer to this; treat OBSERVE strictly as an optional accelerator.
- **Don't touch `/oic/sec/*` (doxm, pstat, cred, acl).** The bridge doesn't, and you shouldn't from helper scripts either — those resources have wedge/brick risk on Samsung's RT-OCF security stack. The bridge surfaces are strictly `/<x>/vs/0` and `/device/0`.
- **Don't run two clients against the same appliance simultaneously.** Samsung's RT-OCF DTLS allows one active session per peer; a second handshake will get the device to drop the new socket. If HA seems to flap, check whether you've got `main.py` running locally AND the Docker container up.
- **Don't expect parity from every write surface.** Samsung's firmware accepts a lot of writes with `2.04 Changed` but only some of them stick — power, child-lock, and remote-control writes are accepted-then-reverted because they're hardware-mirrored. The bridge's optimistic-publish-then-verify pattern handles this transparently: HA briefly shows the new value, the 3s fetch-back republishes the actual value, HA reverts.
+4 -26
View File
@@ -3,13 +3,15 @@
One process supervises N Samsung appliances over their OCF CoAP-DTLS
local APIs, publishing state + HA discovery to MQTT. Each appliance
runs its own DTLS session in its own thread. MQTT is shared.
runs its own DTLS session in its own thread, with an in-session
PollScheduler driving state freshness and a KeepaliveTask driving
DTLS-layer liveness. MQTT is shared.
Config is env-var driven:
* APPLIANCE_COUNT plus APPLIANCE_<n>_{CLASS,IP,OCF_PORT,TOPIC,NAME}
define the appliances to bridge.
* Shared keys (MQTT_*, HA_DISCOVERY_PREFIX, CERT_PATH, KEY_PATH,
HEALTH_INTERVAL_S, HEARTBEAT_INTERVAL_S) apply to all.
HEALTH_INTERVAL_S, PING_INTERVAL_S) apply to all.
Reconnects on session errors per-appliance; shuts down cleanly on
SIGINT / SIGTERM."""
@@ -146,36 +148,12 @@ def main():
break
return loop
def make_heartbeat(b: PushBridge):
def loop():
while not b.stop.is_set():
if b.stop.wait(shared.HEARTBEAT_INTERVAL_S):
break
b.heartbeat()
return loop
def make_ping(b: PushBridge):
def loop():
while not b.stop.is_set():
if b.stop.wait(shared.PING_INTERVAL_S):
break
b.ping_once()
return loop
for b in bridges:
tag = b.app.klass
threads.append(threading.Thread(
target=b.run_forever, daemon=True, name=f'{tag}-session'))
threads.append(threading.Thread(
target=make_health(b), daemon=True, name=f'{tag}-health'))
if shared.HEARTBEAT_INTERVAL_S > 0:
threads.append(threading.Thread(
target=make_heartbeat(b), daemon=True,
name=f'{tag}-heartbeat'))
if shared.PING_INTERVAL_S > 0:
threads.append(threading.Thread(
target=make_ping(b), daemon=True,
name=f'{tag}-ping'))
stopping = threading.Event()
+88 -1
View File
@@ -24,7 +24,10 @@ from __future__ import annotations
import json
from dataclasses import dataclass, field
from typing import Callable, Optional
from typing import Callable, Optional, TYPE_CHECKING
if TYPE_CHECKING:
from ..poll_scheduler import PollTier
@dataclass
@@ -72,6 +75,18 @@ class ApplianceDescriptor:
# the freshly-projected sensors dict; returns a short string.
log_state_change: Optional[Callable[[dict], str]] = None
# Tiered polling cadence. Hot tier resources are the user-visible
# state that needs sub-second freshness; warm/cold/sweep cover the
# rest. Empirically measured per-device ceilings in
# local-tools/probe_poll_rate_combined.py (dryer ~14 req/s, oven
# ~8 req/s) inform the defaults each descriptor sets.
poll_tiers: list['PollTier'] = field(default_factory=list)
# Predicate the PollScheduler calls each tick to decide whether to
# use a tier's active_interval_s. Gets a shallow snapshot of the
# link dict so it can read whichever resource indicates activity.
is_active: Optional[Callable[[dict], bool]] = None
# --- HA-discovery helpers ----------------------------------------------
# Pure builder fns used by descriptor build_discovery() implementations.
@@ -136,3 +151,75 @@ def avail_with_remote_and_cycle(avail_topic: str,
def encode(cfg: dict) -> bytes:
return json.dumps(cfg).encode()
def bridge_diagnostic_discovery(topic_prefix: str,
ha_discovery_prefix: str,
device_name: str,
model: str = 'Bridge') -> list[tuple[str, bytes]]:
"""HA-discovery payloads for the per-bridge diagnostic entities
(push state, polling RTT, freshness). Returned as a list of
(topic, payload-bytes) tuples to be merged with the descriptor's
own build_discovery output."""
avail_topic = f"{topic_prefix}/availability"
health_topic = f"{topic_prefix}/bridge/health"
push_topic = f"{topic_prefix}/bridge/push_active"
device = device_block(topic_prefix, device_name, model)
avail = avail_base(avail_topic)
def sensor(slug: str, name: str, value_template: str,
unit: str | None = None, device_class: str | None = None,
icon: str | None = None) -> tuple[str, bytes]:
cfg = {
'name': name,
'unique_id': f"{topic_prefix}_bridge_{slug}",
'object_id': f"{topic_prefix}_bridge_{slug}",
'state_topic': health_topic,
'value_template': value_template,
'availability': avail,
'device': device,
'entity_category': 'diagnostic',
}
if unit is not None: cfg['unit_of_measurement'] = unit
if device_class is not None: cfg['device_class'] = device_class
if icon is not None: cfg['icon'] = icon
topic = (f"{ha_discovery_prefix}/sensor/"
f"{topic_prefix}/bridge_{slug}/config")
return topic, encode(cfg)
push_active_cfg = {
'name': 'Push Active',
'unique_id': f"{topic_prefix}_bridge_push_active",
'object_id': f"{topic_prefix}_bridge_push_active",
'state_topic': push_topic,
'payload_on': 'online',
'payload_off': 'offline',
'availability': avail,
'device': device,
'entity_category': 'diagnostic',
'icon': 'mdi:rss',
}
push_active_topic = (f"{ha_discovery_prefix}/binary_sensor/"
f"{topic_prefix}/bridge_push_active/config")
return [
(push_active_topic, encode(push_active_cfg)),
sensor('update_source', 'Last Update Source',
"{{ value_json.last_change_source | default('?') }}",
icon='mdi:transit-connection-variant'),
sensor('stalest_age_s', 'Stalest Resource Age',
"{{ value_json.stalest_age_s | default(0) }}",
unit='s', icon='mdi:clock-alert-outline'),
sensor('poll_max_rtt_ms', 'Poll Max RTT',
"{{ value_json.poll_window_max_rtt_ms | default(0) | int }}",
unit='ms', icon='mdi:speedometer'),
sensor('slow_polls', 'Slow Polls (window)',
"{{ value_json.poll_window_slow_count | default(0) }}",
icon='mdi:timer-sand'),
sensor('poll_errors', 'Poll Errors (window)',
"{{ value_json.poll_window_errors | default(0) }}",
icon='mdi:alert-circle-outline'),
sensor('observe_age_s', 'Last OBSERVE Age',
"{{ value_json.last_observe_age_s if value_json.last_observe_age_s is not none else 'never' }}",
icon='mdi:radar'),
]
+47
View File
@@ -13,6 +13,7 @@ from .base import (
device_block,
encode,
)
from ..poll_scheduler import PollTier
# --- OBSERVE paths -----------------------------------------------------
@@ -434,6 +435,50 @@ def command_handlers():
}
# --- Poll tiers --------------------------------------------------------
# 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.
DRYER_POLL_TIERS = [
PollTier(
name='hot',
interval_s=1.0,
active_interval_s=0.5,
paths=(
('operational', 'state', 'vs', '0'),
),
),
PollTier(
name='warm',
interval_s=15.0,
paths=(
('power', 'vs', '0'),
('kidslock', 'vs', '0'),
('remotectrl', 'vs', '0'),
('alarms', 'vs', '0'),
('course', 'vs', '0'),
('washer', 'vs', '0'),
('st', 'dryercourse', 'vs', '0'),
('wm', 'jobbeginingstatus', 'vs', '0'),
('energy', 'consumption', 'vs', '0'),
('diagnosis', 'vs', '0'),
),
),
PollTier(
name='sweep',
interval_s=300.0,
paths=(('device', '0'),),
is_sweep=True,
),
]
def _is_active(links: dict) -> bool:
rep = links.get('/operational/state/vs/0') or {}
sam_state = rep.get('x.com.samsung.da.state')
return _SAMSUNG_STATE_TO_OCF.get(sam_state) == 'active'
# --- Descriptor --------------------------------------------------------
DRYER = ApplianceDescriptor(
name='dryer',
@@ -447,4 +492,6 @@ DRYER = ApplianceDescriptor(
project=project,
remote_available_field='remote_control_binary',
log_state_change=log_state_change,
poll_tiers=DRYER_POLL_TIERS,
is_active=_is_active,
)
+53
View File
@@ -31,6 +31,7 @@ from .base import (
device_block,
encode,
)
from ..poll_scheduler import PollTier
# ---------------------------------------------------------------------
@@ -744,6 +745,56 @@ def command_handlers():
}
# --- Poll tiers --------------------------------------------------------
# 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.
OVEN_POLL_TIERS = [
PollTier(
name='hot',
interval_s=1.0,
active_interval_s=0.5,
paths=(
('operational', 'state', 'vs', '0'),
('doors', 'vs', '0'),
('oven', 'vs', '0'),
('temperatures', 'vs', '0'),
),
),
PollTier(
name='warm',
interval_s=30.0,
paths=(
('power', 'vs', '0'),
('kidslock', 'vs', '0'),
('remotectrl', 'vs', '0'),
('mode', 'vs', '0'),
('alarms', 'vs', '0'),
('connected', 'vs', '0'),
),
),
PollTier(
name='cold',
interval_s=600.0,
paths=(
('otninformation', 'vs', '0'),
),
),
PollTier(
name='sweep',
interval_s=300.0,
paths=(('device', '0'),),
is_sweep=True,
),
]
def _is_active(links: dict) -> bool:
rep = links.get('/operational/state/vs/0') or {}
sam_state = rep.get('x.com.samsung.da.state')
return _SAMSUNG_STATE_TO_OCF.get(sam_state) == 'active'
# ---------------------------------------------------------------------
OVEN = ApplianceDescriptor(
name='oven',
@@ -758,4 +809,6 @@ OVEN = ApplianceDescriptor(
remote_available_field='remote_control_binary',
cycle_active_field='cycle_active',
log_state_change=log_state_change,
poll_tiers=OVEN_POLL_TIERS,
is_active=_is_active,
)
+276 -268
View File
@@ -1,23 +1,19 @@
"""Push-mode bridge: OCF CoAP-DTLS Observe → MQTT.
"""Bridge: OCF CoAP-DTLS appliance → MQTT, polling-first with opportunistic OBSERVE.
Appliance ──CoAP OBSERVE notifications──► PushBridge
│
▼
MQTT broker
│
▼
Home Assistant
Appliance ──CoAP DTLS─► PushBridge ──MQTT──► Home Assistant
State changes push from the appliance over a sustained DTLS session.
The bridge updates an in-memory link dict, recomputes flat sensors via
the appliance descriptor, and publishes to MQTT ONLY when the flat-
sensor dict actually changes.
State freshness comes from a tiered PollScheduler over a persistent DTLS
session. OBSERVE registrations are kept as an opportunistic freshness
accelerator — when the appliance has internet and pushes notifications,
the cache absorbs them; when it's air-gapped, polling carries the UX
unchanged.
The bridge is appliance-class-agnostic — it delegates every
appliance-specific decision to an ApplianceDescriptor.
DTLS-layer liveness is a separate KeepaliveTask (CoAP empty-CON ping
every PING_INTERVAL_S). Three consecutive failures publish MQTT
availability=offline.
Multiple PushBridges run concurrently in a single process — see
main.py. They share one MQTT client; each owns one DTLS session.
Multiple PushBridges run concurrently — see main.py. They share one
MQTT client; each owns one DTLS session, one cache, one scheduler.
"""
import json
import os
@@ -26,43 +22,28 @@ import time
import cbor2
from .appliances.base import ApplianceDescriptor
from .appliances.base import ApplianceDescriptor, bridge_diagnostic_discovery
from .coap_dtls import DtlsCoapSession, fmt_code
from .config import ApplianceConfig, SharedConfig
from .logger import bridge_logger, logger as module_logger
from .keepalive import KeepaliveTask
from .logger import bridge_logger
from .poll_scheduler import PollScheduler
from .sensors import index_links
from .state_cache import StateCache
# Re-enable with `DEBUG_BRIDGE=1` env var. When on, the bridge dumps:
# - every received CoAP frame (in coap_dtls.py)
# - the full /device/0 link tree at seed time
# - /oic/res directory
# - REP changes for /operational/state, /oven, /power, /mode-options
# Designed for reverse-engineering new oven/dryer behaviour next
# session — leave at 0 in production.
DEBUG_BRIDGE = os.environ.get('DEBUG_BRIDGE') == '1'
def _href_to_segs(href: str) -> list[str]:
"""`/mode/vs/0` → `['mode', 'vs', '0']`. Used to translate an
OBSERVE-notification href back into the path-segs the Block2 GET
needs."""
return [s for s in href.split('/') if s]
# Samsung's `/information/vs/0` resource carries a unique serial number.
# Verified on both dryer (DV5000T) and oven (NV7000BS); we use the
# value to tag per-bridge log lines once the seed completes.
SERIAL_PATH = '/information/vs/0'
SERIAL_FIELD = 'x.com.samsung.da.serialNum'
class PushBridge:
"""Single sustained DTLS-CoAP session to one appliance.
Reconnects with exponential backoff on session errors. Publishes
availability=offline when the appliance is unreachable so HA marks
entities unavailable instead of trusting stale state."""
def __init__(self,
shared: SharedConfig,
@@ -74,20 +55,22 @@ class PushBridge:
self.descriptor = descriptor
self.mqtt = mqtt_client
# Bridge-scoped logger; retagged with serial after first seed.
self.log = bridge_logger(app.klass)
self._serial: str | None = None
# Resolve port (descriptor default if unset in env).
self.port = app.ocf_port or descriptor.default_observe_port
self.session: DtlsCoapSession | None = None
self.links: dict[str, dict] = {} # href → rep
self.descriptor_state: dict = {} # descriptor scratch space
self.scheduler: PollScheduler | None = None
self.keepalive: KeepaliveTask | None = None
self.cache = StateCache(descriptor)
self.cache.set_on_change(self._on_cache_change)
self.last_state_pub = None
self.last_remote_pub = None
self.last_cycle_pub = None
self.last_avail_pub: str | None = None
self.stop = threading.Event()
self.started_ts = time.time()
self.session_started_ts = None
@@ -97,14 +80,25 @@ class PushBridge:
self.connect_count = 0
self.error_count = 0
self._publish_gate = False
self._last_change_source: str | None = None
self._last_observe_change_ts: float | None = None
self._last_push_active_pub: str | None = None
# Per-href fetchback generation counter. Every new schedule
# bumps the gen; a fetchback aborts on wake (and again after
# its GET completes) if its captured gen is no longer the
# latest. This coalesces bursts: rapid lamp toggles or slider
# drags result in many scheduled fetchbacks but only the
# latest one actually publishes. Lock guards the dict mutation
# and the gen comparison.
# Push is considered "active" if an OBSERVE-sourced change
# arrived within this window. Long enough that a quiet but
# working appliance doesn't flap to inactive; short enough that
# a genuinely silent push channel is visible within minutes.
self.push_active_window_s = 600.0
# Snapshot counters for the per-health-window deltas surfaced in
# publish_health's log summary.
self._win_prev_poll = 0
self._win_prev_poll_err = 0
self._win_prev_ping_fail = 0
# Per-href fetchback generation counter coalesces bursts of
# OBSERVE-block2 partial notifications: rapid changes scheduled
# many fetchbacks but only the latest actually publishes.
self._fetch_gen: dict[str, int] = {}
self._fetch_lock = threading.Lock()
@@ -114,27 +108,33 @@ class PushBridge:
self.remote_topic = f"{p}/remote_available"
self.cycle_topic = f"{p}/cycle_active"
self.health_topic = f"{p}/bridge/health"
self.push_active_topic = f"{p}/bridge/push_active"
self.cmd_handlers = descriptor.command_handlers()
self.cmd_topic_prefix = f"{p}/cmd/"
# Pre-built HA discovery payloads. Republished on every MQTT
# (re)connect by main.py.
self.discovery_payloads = descriptor.build_discovery(
app.topic_prefix, shared.HA_DISCOVERY_PREFIX, app.device_name)
self.discovery_payloads = (
descriptor.build_discovery(
app.topic_prefix, shared.HA_DISCOVERY_PREFIX, app.device_name)
+ bridge_diagnostic_discovery(
app.topic_prefix, shared.HA_DISCOVERY_PREFIX, app.device_name,
model=descriptor.name.title()))
# ---- DTLS session helpers ---------------------------------------
# ---- cache plumbing ---------------------------------------------
def _on_cache_change(self, changed: bool, source: str) -> None:
if changed:
self.notif_count += 1
self.last_change_ts = time.time()
self._last_change_source = source
if source == 'observe':
self._last_observe_change_ts = self.last_change_ts
self.maybe_publish_state()
def _on_notification(self, href, payload_bytes):
"""Invoked by the DTLS reader thread for OBSERVE notifications.
Resources larger than one CoAP block (notably the oven's
`/mode/vs/0` at ~9KB) arrive truncated: Samsung sends only
block 0 with Block2.M=1 and expects the client to fetch the
rest via Block2 GET. We use cbor decode failure as the
robust "this notification is partial" signal, then spawn a
worker thread to fetch the full resource."""
"""Reader-thread callback for OBSERVE notifications. Large
resources (oven /mode/vs/0 ~9KB) arrive truncated with Block2.M=1
and we use cbor-decode failure as the partial signal."""
if not payload_bytes:
# Empty payload — almost certainly a Block2 announcement.
self._schedule_fetchback(href)
return
try:
@@ -144,53 +144,19 @@ class PushBridge:
return
if not isinstance(rep, dict):
return
self._apply_rep(href, rep)
def _apply_rep(self, href, rep):
"""Update self.links + fire descriptor hooks + maybe publish.
Shared between the OBSERVE path and the Block2 fetch-back path."""
if DEBUG_BRIDGE:
# /mode/vs/0 carries a huge modeSpec JSON we don't want to
# log; surface just modes + options. Small resources get
# full-rep dumps.
if href == '/mode/vs/0' and isinstance(rep, dict):
self.log.info("mode modes=%r options=%r",
rep.get('x.com.samsung.da.modes'),
rep.get('x.com.samsung.da.options'))
elif href in ('/operational/state/vs/0', '/oven/vs/0',
'/power/vs/0'):
self.log.info("REP %s = %r", href, rep)
self.links[href] = rep
hook = self.descriptor.on_observation
if hook is not None:
try:
hook(self.descriptor_state, href, rep)
except Exception as e:
self.log.warning("on_observation %s: %s", href, e)
self.notif_count += 1
self.last_change_ts = time.time()
self.maybe_publish_state()
self._debug_log_rep(href, rep)
self.cache.apply_rep(href, rep, source='observe')
def _apply_optimistic(self, href, body):
"""Optimistically merge a just-POSTed body into the link dict
and republish state. Samsung accepts (2.04) writes whose
bodies are field-replacements — we mirror that semantics here:
each top-level key in `body` overwrites the corresponding key
in the existing rep. The fetchback that follows republishes
the device's real state, which corrects any field where the
write was silently coerced or rejected."""
if not isinstance(body, dict):
return
rep = dict(self.links.get(href) or {})
rep.update(body)
self._apply_rep(href, rep)
def _debug_log_rep(self, href, rep):
if href == '/mode/vs/0' and isinstance(rep, dict):
self.log.info("mode modes=%r options=%r",
rep.get('x.com.samsung.da.modes'),
rep.get('x.com.samsung.da.options'))
elif href in ('/operational/state/vs/0', '/oven/vs/0', '/power/vs/0'):
self.log.info("REP %s = %r", href, rep)
def _schedule_fetchback(self, href, delay_s: float = 0.0):
"""Spawn a worker thread to fetch the full payload of `href`
via Block2 GET. Each schedule bumps a per-href generation
counter — if a newer fetchback is scheduled before this one
fires, this one aborts (so rapid commands coalesce into a
single verification read of the FINAL state)."""
with self._fetch_lock:
gen = self._fetch_gen.get(href, 0) + 1
self._fetch_gen[href] = gen
@@ -202,15 +168,8 @@ class PushBridge:
).start()
def _fetch_back(self, href, delay_s: float, gen: int):
if delay_s > 0:
# Allow the device's read-side to propagate a recent
# write. The oven needs ~1s after a /mode/vs/0 POST;
# the dryer is faster but the delay is harmless there.
if self.stop.wait(delay_s):
return
# Has a newer fetchback been scheduled during our delay?
# If so, abort — our read would publish stale state relative
# to the user's most recent intent.
if delay_s > 0 and self.stop.wait(delay_s):
return
with self._fetch_lock:
if self._fetch_gen.get(href) != gen:
return
@@ -223,31 +182,24 @@ class PushBridge:
except Exception as e:
self.log.warning("fetchback %s: %s", href, e)
return
# Re-check generation after the GET — a new write may have
# come in during the Block2 round-trip, in which case our
# payload is also superseded.
with self._fetch_lock:
if self._fetch_gen.get(href) != gen:
return
if code != 0x45:
self.log.warning("fetchback %s: %s",
href, fmt_code(code))
self.log.warning("fetchback %s: %s", href, fmt_code(code))
return
try:
rep = cbor2.loads(payload) if payload else {}
except Exception as e:
self.log.warning("fetchback %s cbor: %s", href, e)
return
if not isinstance(rep, dict):
return
self._apply_rep(href, rep)
if isinstance(rep, dict):
self.cache.apply_rep(href, rep, source='observe')
def _retag_logger_with_serial(self):
"""Look up the appliance's serial in the seeded link dict and
retarget self.log to a serial-tagged child. Idempotent."""
if self._serial is not None:
return
info = self.links.get(SERIAL_PATH) or {}
info = self.cache.get(SERIAL_PATH) or {}
serial = info.get(SERIAL_FIELD)
if not serial:
return
@@ -258,8 +210,6 @@ class PushBridge:
# ---- session lifecycle ------------------------------------------
def session_once(self):
"""Run one DTLS session end-to-end. Raises on error; the outer
run_forever wraps this with reconnect/backoff."""
sess = DtlsCoapSession(
self.app.ip, self.port,
cert_path=self.shared.CERT_PATH,
@@ -270,7 +220,7 @@ class PushBridge:
self.session = sess
self.session_started_ts = time.time()
self.connect_count += 1
self.descriptor_state = {}
self.cache.descriptor_state.clear()
self._publish_gate = False
self.log.info("DTLS connected — subscribing %d paths",
@@ -278,10 +228,6 @@ class PushBridge:
sess.start_reader()
# When the bridge is asked to stop, close the session — which
# fires the OBSERVE-dereg sequence and tears DTLS down. Without
# this, session_once would block forever in sess.join() because
# the reader only exits on stop/socket-death.
session_ended = threading.Event()
def _stop_watcher():
@@ -299,18 +245,62 @@ class PushBridge:
try:
self._run_session_inner(sess)
finally:
# Always release the watcher so it doesn't sit pinned on
# self.stop forever after the session ends.
session_ended.set()
def _run_session_inner(self, sess):
"""Body of session_once after start_reader() — split out so
session_once's try/finally cleanly bounds the stop-watcher
thread's lifetime."""
for path in self.descriptor.observe_paths:
sess.subscribe(path)
time.sleep(0.05)
# Inline seed so the publish gate opens before the scheduler's
# first tick. The scheduler's sweep tier will refresh /device/0
# on its own cadence afterwards.
self._seed_from_device0(sess)
self._retag_logger_with_serial()
if DEBUG_BRIDGE:
self._debug_dump_links(sess)
self._publish_gate = True
self.maybe_publish_state(force=True)
self.set_availability(True)
self.log.info("seeded → %d links; sensors live",
len(self.cache.links))
scheduler = PollScheduler(
sess, self.cache,
tiers=self.descriptor.poll_tiers,
sweep_index_fn=index_links,
is_active_fn=self.descriptor.is_active,
logger=self.log,
)
keepalive = KeepaliveTask(
sess,
interval_s=float(self.shared.PING_INTERVAL_S),
fail_threshold=3,
on_reachable=self._on_reachable,
on_unreachable=self._on_unreachable,
logger=self.log,
)
self.scheduler = scheduler
self.keepalive = keepalive
sched_t = threading.Thread(
target=scheduler.run_forever, args=(self.stop,),
daemon=True, name=f'{self.app.klass}-poll')
ka_t = threading.Thread(
target=keepalive.run_forever, args=(self.stop,),
daemon=True, name=f'{self.app.klass}-ping')
sched_t.start()
ka_t.start()
try:
sess.join()
finally:
self.scheduler = None
self.keepalive = None
def _seed_from_device0(self, sess):
code, pl = sess.get(self.descriptor.seed_path, timeout=15.0)
if code != 0x45:
raise RuntimeError(
@@ -321,100 +311,63 @@ class PushBridge:
raise RuntimeError(
f"/{'/'.join(self.descriptor.seed_path)} cbor decode: {e}"
) from e
# During the seed we want the cache populated without triggering
# a publish per resource — gate the on_change callback off until
# the publish gate opens just below.
for href, rep in index_links(body).items():
self.links.setdefault(href, rep)
# Once the seed is in, we know the appliance's serial — retag
# the logger so the remaining log lines this session emits are
# serial-tagged.
self._retag_logger_with_serial()
if DEBUG_BRIDGE:
# Dump every link's rep (skipping the huge modeSpec on
# /mode/vs/0) + the /oic/res directory. Useful for finding
# new writable resources next session.
for href, rep in sorted(self.links.items()):
if href == '/mode/vs/0':
short = {k: v for k, v in rep.items()
if k not in (
'x.com.samsung.da.modeSpec',
'x.com.samsung.da.supportedModes',
)}
self.log.info("LINK %s = %r", href, short)
else:
self.log.info("LINK %s = %r", href, rep)
try:
code, pl = sess.get(['oic', 'res'], timeout=10.0)
if code == 0x45 and pl:
self.log.info("OIC_RES = %r", cbor2.loads(pl))
else:
self.log.info("oic/res → %s", fmt_code(code))
except Exception as e:
self.log.warning("oic/res get: %s", e)
hook = self.descriptor.on_observation
if hook is not None:
for href, rep in self.links.items():
try:
hook(self.descriptor_state, href, rep)
except Exception as e:
self.log.warning("seed on_observation %s: %s", href, e)
if href not in self.cache.links:
self.cache.apply_rep(href, rep, source='seed')
self.last_seed_ts = time.time()
self._publish_gate = True
self.maybe_publish_state(force=True)
def _debug_dump_links(self, sess):
for href, rep in sorted(self.cache.links.items()):
if href == '/mode/vs/0':
short = {k: v for k, v in rep.items()
if k not in (
'x.com.samsung.da.modeSpec',
'x.com.samsung.da.supportedModes',
)}
self.log.info("LINK %s = %r", href, short)
else:
self.log.info("LINK %s = %r", href, rep)
try:
code, pl = sess.get(['oic', 'res'], timeout=10.0)
if code == 0x45 and pl:
self.log.info("OIC_RES = %r", cbor2.loads(pl))
else:
self.log.info("oic/res → %s", fmt_code(code))
except Exception as e:
self.log.warning("oic/res get: %s", e)
def _on_reachable(self) -> None:
self.set_availability(True)
self.log.info("seeded → %d links; sensors live", len(self.links))
self.reassert_availability()
sess.join()
def heartbeat(self):
sess = self.session
if sess is None:
return
try:
code, pl = sess.get(self.descriptor.seed_path, timeout=15.0)
except Exception as e:
self.log.warning("heartbeat seed: %s", e)
return
if code != 0x45:
self.log.warning("heartbeat seed: %s", fmt_code(code))
return
try:
body = cbor2.loads(pl)
except Exception as e:
self.log.warning("heartbeat seed cbor: %s", e)
return
# Refresh ALL resources, not just non-observed ones. We used
# to skip observed resources on the assumption OBSERVE kept
# them fresh, but the oven doesn't reliably push OBSERVE on
# /mode/vs/0 option changes (timer / lamp / sound), so we'd
# be stuck with stale values until the next user-driven POST
# triggered a fetchback. Refreshing everything bounds HA's
# divergence to HEARTBEAT_INTERVAL_S in the worst case.
hook = self.descriptor.on_observation
for href, rep in index_links(body).items():
self.links[href] = rep
if hook is not None:
try:
hook(self.descriptor_state, href, rep)
except Exception as e:
self.log.warning("heartbeat hook %s: %s", href, e)
self.last_seed_ts = time.time()
self.maybe_publish_state()
def _on_unreachable(self) -> None:
self.set_availability(False)
# ---- MQTT publishing --------------------------------------------
def maybe_publish_state(self, force=False):
if not force and not self._publish_gate:
return
sensors = self.descriptor.flatten(self.links)
snap = self.cache.snapshot()
sensors = self.descriptor.flatten(snap)
project = self.descriptor.project
if project is not None:
sensors = project(self.descriptor_state, sensors)
sensors = project(self.cache.descriptor_state, sensors)
if not force and sensors == self.last_state_pub:
return
if DEBUG_BRIDGE and self.last_state_pub is not None:
diffs = {k: (self.last_state_pub.get(k), v)
for k, v in sensors.items()
if self.last_state_pub.get(k) != v}
diffs.update({k: (self.last_state_pub.get(k), None)
for k in self.last_state_pub
if k not in sensors})
if diffs:
self.log.info("sensor diff: %s",
{k: f"{a!r} → {b!r}" for k, (a, b) in diffs.items()})
self.last_state_pub = sensors
self.mqtt.publish(self.state_topic,
json.dumps(sensors).encode(),
@@ -428,7 +381,8 @@ class PushBridge:
if not force:
log_fn = self.descriptor.log_state_change
extra = log_fn(sensors) if log_fn is not None else ''
self.log.info("state changed (%s notif#%d)",
self.log.info("state changed [%s] (%s notif#%d)",
self._last_change_source or '?',
extra or 'descriptor-no-log', self.notif_count)
def publish_remote_available(self, remote_on, force=False):
@@ -465,12 +419,18 @@ class PushBridge:
if cycle_field is not None:
self.publish_cycle_active(
self.last_state_pub.get(cycle_field), force=True)
now = time.time()
active = (self._last_observe_change_ts is not None
and (now - self._last_observe_change_ts) <= self.push_active_window_s)
self.publish_push_active(active, force=True)
def set_availability(self, online):
value = 'online' if online else 'offline'
if value == self.last_avail_pub:
return
self.last_avail_pub = value
try:
self.mqtt.publish(self.avail_topic,
'online' if online else 'offline',
qos=1, retain=True)
self.mqtt.publish(self.avail_topic, value, qos=1, retain=True)
except Exception as e:
self.log.warning("avail publish: %s", e)
if not online and self.descriptor.remote_available_field is not None:
@@ -487,6 +447,13 @@ class PushBridge:
qos=1, retain=True)
except Exception:
pass
if not online:
self._last_push_active_pub = None
try:
self.mqtt.publish(self.push_active_topic, 'offline',
qos=1, retain=True)
except Exception:
pass
# ---- MQTT command handling --------------------------------------
@@ -498,12 +465,9 @@ class PushBridge:
if handler is None:
self.log.warning("unknown command topic: %s", topic)
return
# Shallow-snapshot self.links so the handler sees a consistent
# view across the read-modify-write it may need to perform
# (e.g. oven lamp / sound / fastpreheat all RMW /mode/vs/0
# options). Inner reps are mutated only by the OBSERVE reader
# thread; handlers that mutate items must deep-copy themselves.
result = handler(payload, dict(self.links))
# Handler gets a links snapshot so its read-modify-write sees a
# consistent view across the multi-field operation.
result = handler(payload, self.cache.snapshot())
if result is None:
self.log.warning("rejected command %s payload=%r",
topic, payload)
@@ -513,69 +477,113 @@ class PushBridge:
if sess is None:
self.log.warning("command %s: no DTLS session", topic)
return
href = '/' + '/'.join(path_segs)
sched = self.scheduler
defer_s = 4.0
if sched is not None:
sched.write_in_progress(href, settle_s=defer_s)
try:
code, _ = sess.post(path_segs, cbor2.dumps(body), timeout=8.0)
except Exception as e:
self.log.warning("command %s POST failed: %s", topic, e)
return
self.log.info("command %s payload=%r → %s",
suffix, payload, fmt_code(code))
defer_note = f" (poll-defer {href} {defer_s:.0f}s)" if sched is not None else ''
self.log.info("command %s payload=%r → %s%s",
suffix, payload, fmt_code(code), defer_note)
if code >> 5 == 2:
href = '/' + '/'.join(path_segs)
# Optimistic publish — apply the write to our local state.
# No fetchback: empirically (2026-05-31) the post-write
# Block2 GET was causing the appliance to roll our values
# back ~3s later, on every writable resource. OBSERVE
# pushes keep HA in sync without polling. (Was the root
# cause of the "mid-cycle setpoint reverts" symptom.)
self._apply_optimistic(href, body)
# Optimistic local merge so HA sees the write reflected
# immediately. No Block2 fetchback — that triggers Samsung's
# 3-second revert (project_fetchback_revert_root_cause.md).
# The PollScheduler will reconcile on its next tier tick
# after the write_in_progress settle window expires.
self.cache.apply_optimistic(href, body)
def publish_health(self):
now = time.time()
sched = self.scheduler
ka = self.keepalive
last_obs_age = (round(now - self._last_observe_change_ts, 1)
if self._last_observe_change_ts else None)
push_active = (last_obs_age is not None
and last_obs_age <= self.push_active_window_s)
# Per-window deltas for the log summary AND the health topic
# (HA gets the same numbers without doing template arithmetic).
poll = sched.poll_count if sched else 0
poll_err = sched.poll_error_count if sched else 0
ping_fail = ka.ping_fail_count if ka else 0
d_poll = poll - self._win_prev_poll
d_err = poll_err - self._win_prev_poll_err
d_ping_fail = ping_fail - self._win_prev_ping_fail
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))
window_polls_ok = max(0, d_poll - d_err)
h = {
'mode': 'push',
'device_class': self.descriptor.name,
'serial': self._serial,
'connect_count': self.connect_count,
'error_count': self.error_count,
'notif_count': self.notif_count,
'last_change_age_s': (round(now - self.last_change_ts, 1)
if self.last_change_ts else None),
'last_seed_age_s': (round(now - self.last_seed_ts, 1)
if self.last_seed_ts else None),
'session_age_s': (round(now - self.session_started_ts, 1)
if self.session_started_ts else None),
'uptime_seconds': round(now - self.started_ts, 0),
'mode': 'poll+observe',
'device_class': self.descriptor.name,
'serial': self._serial,
'connect_count': self.connect_count,
'error_count': self.error_count,
'notif_count': self.notif_count,
'poll_count': poll,
'poll_error_count': poll_err,
'poll_window_ok': window_polls_ok,
'poll_window_errors': d_err,
'poll_window_max_rtt_ms': round(win_max_rtt, 0),
'poll_window_slow_count': win_slow,
'ping_count': ka.ping_count if ka else 0,
'ping_fail_count': ping_fail,
'reachable': ka.reachable if ka else None,
'last_change_source': self._last_change_source,
'last_observe_age_s': last_obs_age,
'push_active': push_active,
'last_change_age_s': (round(now - self.last_change_ts, 1)
if self.last_change_ts else None),
'last_seed_age_s': (round(now - self.last_seed_ts, 1)
if self.last_seed_ts else None),
'session_age_s': (round(now - self.session_started_ts, 1)
if self.session_started_ts else None),
'uptime_seconds': round(now - self.started_ts, 0),
}
stalest = self.cache.stalest()
if stalest is not None:
h['stalest_href'] = stalest[0]
h['stalest_age_s'] = round(stalest[1], 1)
try:
self.mqtt.publish(self.health_topic, json.dumps(h).encode(),
qos=0, retain=True)
except Exception as e:
self.log.warning("health publish: %s", e)
self.publish_push_active(push_active)
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)",
window_polls_ok, d_err, d_ping_fail,
win_max_rtt, win_slow, self.shared.HEALTH_INTERVAL_S)
def publish_push_active(self, active: bool, force: bool = False) -> None:
value = 'online' if active else 'offline'
if not force and value == self._last_push_active_pub:
return
self._last_push_active_pub = value
try:
self.mqtt.publish(self.push_active_topic, value,
qos=1, retain=True)
self.log.info("push_active → %s", value)
except Exception as e:
self.log.warning("push_active publish: %s", e)
# ---- top-level loop ---------------------------------------------
def _publish_tick_loop(self):
while not self.stop.wait(15.0):
try:
self.maybe_publish_state()
except Exception as e:
self.log.warning("publish tick: %s", e)
def ping_once(self):
"""Send one CoAP Ping. Caller (main.py's per-bridge ping
thread) handles the cadence. No-op if no live session."""
sess = self.session
if sess is None:
return
try:
sess.ping()
except Exception as e:
self.log.warning("ping: %s", e)
def run_forever(self):
threading.Thread(target=self._publish_tick_loop, daemon=True,
name=f'{self.app.klass}-tick').start()
backoff = 1.0
while not self.stop.is_set():
try:
-3
View File
@@ -57,7 +57,6 @@ class SharedConfig:
MQTT_PASS: Optional[str]
HA_DISCOVERY_PREFIX: str
HEALTH_INTERVAL_S: int
HEARTBEAT_INTERVAL_S: int
PING_INTERVAL_S: int
@classmethod
@@ -72,8 +71,6 @@ class SharedConfig:
HA_DISCOVERY_PREFIX=os.getenv('HA_DISCOVERY_PREFIX',
'homeassistant'),
HEALTH_INTERVAL_S=int(os.getenv('HEALTH_INTERVAL_S', '60')),
HEARTBEAT_INTERVAL_S=int(os.getenv('HEARTBEAT_INTERVAL_S',
'600')),
PING_INTERVAL_S=int(os.getenv('PING_INTERVAL_S', '25')),
)
+82
View File
@@ -0,0 +1,82 @@
"""DTLS-layer liveness via CoAP empty-CON ping.
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.
"""
from __future__ import annotations
import threading
from typing import Callable, Optional
from .coap_dtls import DtlsCoapSession
class KeepaliveTask:
def __init__(self,
session: DtlsCoapSession,
interval_s: float = 25.0,
fail_threshold: int = 3,
on_reachable: Optional[Callable[[], None]] = None,
on_unreachable: Optional[Callable[[], None]] = None,
logger=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._fail_streak = 0
self._reachable = True
self._ping_count = 0
self._ping_fail_count = 0
@property
def ping_count(self) -> int:
return self._ping_count
@property
def ping_fail_count(self) -> int:
return self._ping_fail_count
@property
def reachable(self) -> bool:
return self._reachable
def run_forever(self, stop: threading.Event) -> None:
while not stop.wait(self.interval_s):
self._tick()
def _tick(self) -> None:
ok = False
try:
self.session.ping()
ok = True
except Exception as e:
if self.log: self.log.warning("ping: %s", e)
self._ping_count += 1
if ok:
if not self._reachable:
if self.log:
self.log.info("ping recovered after %d fails",
self._fail_streak)
self._reachable = True
if self.on_reachable is not None:
try: self.on_reachable()
except Exception as e:
if self.log: self.log.warning("on_reachable: %s", e)
self._fail_streak = 0
return
self._ping_fail_count += 1
self._fail_streak += 1
if self._reachable and self._fail_streak >= self.fail_threshold:
if self.log:
self.log.warning("device unreachable after %d ping failures",
self._fail_streak)
self._reachable = False
if self.on_unreachable is not None:
try: self.on_unreachable()
except Exception as e:
if self.log: self.log.warning("on_unreachable: %s", e)
+199
View File
@@ -0,0 +1,199 @@
"""Tiered adaptive polling against a DtlsCoapSession.
Tiers are descriptor-declared (hot/warm/cold + sweep). Per tick:
each tier whose deadline has passed polls all its paths sequentially
on the shared session, writing into the StateCache. The sweep tier
issues one Block2 GET of /device/0 and uses index_links to fan its
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.
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
fetchback-revert bug.
"""
from __future__ import annotations
import threading
import time
from dataclasses import dataclass
from typing import Callable, Optional, TYPE_CHECKING
import cbor2
from .coap_dtls import DtlsCoapSession, fmt_code
if TYPE_CHECKING:
from .state_cache import StateCache
@dataclass(frozen=True)
class PollTier:
name: str
interval_s: float
paths: tuple[tuple[str, ...], ...]
active_interval_s: Optional[float] = None
is_sweep: bool = False
class PollScheduler:
def __init__(self,
session: DtlsCoapSession,
cache: 'StateCache',
tiers: list[PollTier],
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):
self.session = session
self.cache = cache
self.tiers = tiers
self.sweep_index = sweep_index_fn
self.is_active_fn = is_active_fn
self.log = logger
self.timeout_s = timeout_s
now = time.monotonic()
self._next_due: dict[str, float] = {t.name: now for t in tiers}
self._defer_until: dict[str, float] = {}
self._defer_lock = threading.Lock()
self._poll_count = 0
self._poll_error_count = 0
self._last_active: Optional[bool] = None
# Per-window tail-latency tracking. Bridge consumes-and-resets
# these via take_window_stats() once per HEALTH_INTERVAL_S.
self._stats_lock = threading.Lock()
self._window_max_rtt_ms = 0.0
self._window_slow_count = 0
self.slow_threshold_ms = 1000.0
def write_in_progress(self, href: str, settle_s: float = 4.0) -> None:
with self._defer_lock:
self._defer_until[href] = time.monotonic() + settle_s
def run_forever(self, stop: threading.Event) -> None:
while not stop.is_set():
self._run_due_tiers()
sleep_for = max(0.05, min(1.0, self._earliest_deadline() - time.monotonic()))
if stop.wait(sleep_for):
return
@property
def poll_count(self) -> int:
return self._poll_count
@property
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`."""
with self._stats_lock:
out = (self._window_max_rtt_ms, self._window_slow_count)
self._window_max_rtt_ms = 0.0
self._window_slow_count = 0
return out
def _record_rtt(self, rtt_ms: float) -> None:
with self._stats_lock:
if rtt_ms > self._window_max_rtt_ms:
self._window_max_rtt_ms = rtt_ms
if rtt_ms >= self.slow_threshold_ms:
self._window_slow_count += 1
def _earliest_deadline(self) -> float:
return min(self._next_due.values())
def _run_due_tiers(self) -> None:
now = time.monotonic()
active = False
if self.is_active_fn is not None:
try:
active = bool(self.is_active_fn(self.cache.snapshot()))
except Exception as e:
if self.log: self.log.warning("is_active: %s", e)
if active != self._last_active:
if self.log and self._last_active is not None:
self.log.info("active=%s", active)
self._last_active = active
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)
else tier.interval_s)
self._next_due[tier.name] = now + interval
try:
if tier.is_sweep:
self._do_sweep(tier)
else:
self._do_tier(tier)
except Exception as e:
self._poll_error_count += 1
if self.log: self.log.warning("tier %s: %s", tier.name, e)
def _do_tier(self, tier: PollTier) -> None:
for path in tier.paths:
href = '/' + '/'.join(path)
with self._defer_lock:
if self._defer_until.get(href, 0) > time.monotonic():
continue
self._poll_count += 1
t0 = time.monotonic()
try:
code, body = self.session.get(list(path), timeout=self.timeout_s)
except Exception as e:
self._poll_error_count += 1
self._record_rtt((time.monotonic() - t0) * 1000.0)
if self.log: self.log.warning("poll %s: %s", href, e)
return
self._record_rtt((time.monotonic() - t0) * 1000.0)
if code != 0x45 or not body:
self._poll_error_count += 1
if self.log: self.log.warning("poll %s -> %s", href, fmt_code(code))
continue
try:
rep = cbor2.loads(body)
except Exception as e:
self._poll_error_count += 1
if self.log: self.log.warning("poll %s cbor: %s", href, e)
continue
if isinstance(rep, dict):
self.cache.apply_rep(href, rep, source='poll')
def _do_sweep(self, tier: PollTier) -> None:
path = list(tier.paths[0])
t0 = time.monotonic()
self._poll_count += 1
try:
code, body = self.session.get(path, timeout=self.timeout_s)
except Exception as e:
self._poll_error_count += 1
self._record_rtt((time.monotonic() - t0) * 1000.0)
if self.log: self.log.warning("sweep %s: %s", path, e)
return
self._record_rtt((time.monotonic() - t0) * 1000.0)
if code != 0x45 or not body:
self._poll_error_count += 1
if self.log: self.log.warning("sweep -> %s", fmt_code(code))
return
try:
tree = cbor2.loads(body)
except Exception as e:
self._poll_error_count += 1
if self.log: self.log.warning("sweep cbor: %s", e)
return
indexed = self.sweep_index(tree)
for href, rep in indexed.items():
with self._defer_lock:
if self._defer_until.get(href, 0) > time.monotonic():
continue
self.cache.apply_rep(href, rep, source='sweep')
if self.log:
elapsed_ms = (time.monotonic() - t0) * 1000.0
self.log.info("sweep complete (%d links, %.0fms)",
len(indexed), elapsed_ms)
+78
View File
@@ -0,0 +1,78 @@
"""Single source of truth for one appliance's state.
All writers (OBSERVE notify, poll, seed, optimistic) call apply_rep().
A registered on_change callback fires after any apply that mutated the
cache, which the bridge wires to its MQTT publish gate.
"""
from __future__ import annotations
import threading
import time
from typing import Callable, Optional, TYPE_CHECKING
if TYPE_CHECKING:
from .appliances.base import ApplianceDescriptor
class StateCache:
def __init__(self, descriptor: 'ApplianceDescriptor'):
self.descriptor = descriptor
self.links: dict[str, dict] = {}
self.last_updated: dict[str, float] = {}
self.source: dict[str, str] = {}
self.descriptor_state: dict = {}
self._on_change: Optional[Callable[[bool, str], None]] = None
self._lock = threading.RLock()
def set_on_change(self, cb: Callable[[bool, str], None]) -> None:
self._on_change = cb
def apply_rep(self, href: str, rep: dict, source: str) -> bool:
if not isinstance(rep, dict):
return False
with self._lock:
prior = self.links.get(href)
changed = prior != rep
self.links[href] = rep
self.last_updated[href] = time.time()
self.source[href] = source
hook = self.descriptor.on_observation
if hook is not None:
try:
hook(self.descriptor_state, href, rep)
except Exception:
pass
if self._on_change is not None:
try:
self._on_change(changed, source)
except Exception:
pass
return changed
def apply_optimistic(self, href: str, body: dict) -> bool:
if not isinstance(body, dict):
return False
with self._lock:
merged = dict(self.links.get(href) or {})
merged.update(body)
return self.apply_rep(href, merged, source='optimistic')
def get(self, href: str) -> Optional[dict]:
with self._lock:
return self.links.get(href)
def snapshot(self) -> dict[str, dict]:
with self._lock:
return dict(self.links)
def freshness_s(self, href: str) -> Optional[float]:
ts = self.last_updated.get(href)
return None if ts is None else (time.time() - ts)
def stalest(self) -> Optional[tuple[str, float]]:
with self._lock:
if not self.last_updated:
return None
href = min(self.last_updated, key=self.last_updated.get)
return href, time.time() - self.last_updated[href]