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 8775f7a930
commit 709fdf444d
11 changed files with 866 additions and 324 deletions
+8 -7
View File
@@ -43,14 +43,15 @@ MQTT_PASS=
HA_DISCOVERY_PREFIX=homeassistant HA_DISCOVERY_PREFIX=homeassistant
# Bridge timers (seconds). # Bridge timers (seconds).
# HEALTH_INTERVAL_S — how often <prefix>/bridge/health republishes. # HEALTH_INTERVAL_S — how often <prefix>/bridge/health republishes.
# HEARTBEAT_INTERVAL_S — periodic full /device/0 re-seed. Refreshes # PING_INTERVAL_S — CoAP empty-CON ping cadence (DTLS-layer
# every resource (observed too) — useful for # liveness). Three consecutive failures publish
# appliances like the oven that don't reliably # availability=offline.
# push OBSERVE on /mode/vs/0 option changes. 0 # State freshness itself comes from the in-bridge PollScheduler whose
# disables. # tier cadences are declared in the appliance descriptor — there is
# no top-level heartbeat env var to tune.
HEALTH_INTERVAL_S=60 HEALTH_INTERVAL_S=60
HEARTBEAT_INTERVAL_S=600 PING_INTERVAL_S=25
# Container TZ. # Container TZ.
TZ=Europe/London TZ=Europe/London
+31 -19
View File
@@ -16,9 +16,9 @@
### What you get ### 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. - **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. - **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`. - **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. - **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. - **Zero HA YAML.** Every entity is auto-discovered via MQTT discovery.
@@ -26,7 +26,7 @@
### Under the hood ### 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. 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 | | Appliance class | Model family | Confirmed |
|---|---|---| |---|---|---|
| Dryer | DV5000T (`DA_WM_TP2_20_COMMON`, `mnid=0AJT`) | All entities, sub-second OBSERVE push | | 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; OBSERVE-push lazy on options-array writes (see "Per-appliance notes" below) | | 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/`. 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 ## 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`. 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 | | Power on/off | ❌ | Accepted (2.04) but reverts within seconds — hardware-mirrored |
| Child Lock / Remote Control toggle | ❌ | Same — hardware-mirrored physical buttons | | 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 ### 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 | | 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. | | 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. | | Sound, Fast preheat | ⚠️ | Wired but untested RC-gated. |
| Setpoint slider | ⚠️ | Wired but untested mid-cook behaviour. RC-gated. | | Setpoint slider | ⚠️ | Wired but untested RC-gated. |
| Mode select | ⚠️ | Wired but untested mid-cook behaviour. RC-gated. | | Mode select | ⚠️ | Wired but untested RC-gated. |
| Stop button | ⚠️ | Wired but untested. **Not** RC-gated (the SmartThings app stops without Remote Control on, so we don't gate either). | | Stop button | ✅ | |
| Power on/off as a switch | ❌ | Not exposed as a writeable entity — cold-start panel is a physical action. Read-only sensor only. |
| **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. | | **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: **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 local state and publishes to MQTT. HA reflects the new value instantly. 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. **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. 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 heartbeat** — every `HEARTBEAT_INTERVAL_S` (default 600s), the bridge re-fetches `/device/0` and refreshes ALL resources (including observed ones), bounding worst-case drift. 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`) | | `HA_DISCOVERY_PREFIX` | HA discovery topic root (default `homeassistant`) |
| `CERT_PATH` / `KEY_PATH` | Override cert lookup (auto-detects `/config/` then `./certs/`) | | `CERT_PATH` / `KEY_PATH` | Override cert lookup (auto-detects `/config/` then `./certs/`) |
| `HEALTH_INTERVAL_S` | Seconds between `<prefix>/bridge/health` publishes (default 60) | | `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 | | `SSH_HOST` / `REMOTE_DIR` / `APPDATA_DIR` | Used by `deploy.sh` only |
### MQTT topics — outgoing (bridge → broker) ### 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>/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>/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>/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 | | `<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) ### 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`. 3. Add `WASHER` to `DESCRIPTORS` in `samsung_appliance/appliances/__init__.py`.
4. Add `APPLIANCE_<n>_CLASS=washer` to `.env`, bump `APPLIANCE_COUNT`, redeploy. 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. 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 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 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 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. - **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 One process supervises N Samsung appliances over their OCF CoAP-DTLS
local APIs, publishing state + HA discovery to MQTT. Each appliance 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: Config is env-var driven:
* APPLIANCE_COUNT plus APPLIANCE_<n>_{CLASS,IP,OCF_PORT,TOPIC,NAME} * APPLIANCE_COUNT plus APPLIANCE_<n>_{CLASS,IP,OCF_PORT,TOPIC,NAME}
define the appliances to bridge. define the appliances to bridge.
* Shared keys (MQTT_*, HA_DISCOVERY_PREFIX, CERT_PATH, KEY_PATH, * 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 Reconnects on session errors per-appliance; shuts down cleanly on
SIGINT / SIGTERM.""" SIGINT / SIGTERM."""
@@ -146,36 +148,12 @@ def main():
break break
return loop 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: for b in bridges:
tag = b.app.klass tag = b.app.klass
threads.append(threading.Thread( threads.append(threading.Thread(
target=b.run_forever, daemon=True, name=f'{tag}-session')) target=b.run_forever, daemon=True, name=f'{tag}-session'))
threads.append(threading.Thread( threads.append(threading.Thread(
target=make_health(b), daemon=True, name=f'{tag}-health')) 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() stopping = threading.Event()
+88 -1
View File
@@ -24,7 +24,10 @@ from __future__ import annotations
import json import json
from dataclasses import dataclass, field 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 @dataclass
@@ -72,6 +75,18 @@ class ApplianceDescriptor:
# the freshly-projected sensors dict; returns a short string. # the freshly-projected sensors dict; returns a short string.
log_state_change: Optional[Callable[[dict], str]] = None 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 ---------------------------------------------- # --- HA-discovery helpers ----------------------------------------------
# Pure builder fns used by descriptor build_discovery() implementations. # 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: def encode(cfg: dict) -> bytes:
return json.dumps(cfg).encode() 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, device_block,
encode, encode,
) )
from ..poll_scheduler import PollTier
# --- OBSERVE paths ----------------------------------------------------- # --- 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 -------------------------------------------------------- # --- Descriptor --------------------------------------------------------
DRYER = ApplianceDescriptor( DRYER = ApplianceDescriptor(
name='dryer', name='dryer',
@@ -447,4 +492,6 @@ DRYER = ApplianceDescriptor(
project=project, project=project,
remote_available_field='remote_control_binary', remote_available_field='remote_control_binary',
log_state_change=log_state_change, 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, device_block,
encode, 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( OVEN = ApplianceDescriptor(
name='oven', name='oven',
@@ -758,4 +809,6 @@ OVEN = ApplianceDescriptor(
remote_available_field='remote_control_binary', remote_available_field='remote_control_binary',
cycle_active_field='cycle_active', cycle_active_field='cycle_active',
log_state_change=log_state_change, 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 Appliance ──CoAP DTLS─► PushBridge ──MQTT──► Home Assistant
│
▼
MQTT broker
│
▼
Home Assistant
State changes push from the appliance over a sustained DTLS session. State freshness comes from a tiered PollScheduler over a persistent DTLS
The bridge updates an in-memory link dict, recomputes flat sensors via session. OBSERVE registrations are kept as an opportunistic freshness
the appliance descriptor, and publishes to MQTT ONLY when the flat- accelerator — when the appliance has internet and pushes notifications,
sensor dict actually changes. the cache absorbs them; when it's air-gapped, polling carries the UX
unchanged.
The bridge is appliance-class-agnostic — it delegates every DTLS-layer liveness is a separate KeepaliveTask (CoAP empty-CON ping
appliance-specific decision to an ApplianceDescriptor. every PING_INTERVAL_S). Three consecutive failures publish MQTT
availability=offline.
Multiple PushBridges run concurrently in a single process — see Multiple PushBridges run concurrently — see main.py. They share one
main.py. They share one MQTT client; each owns one DTLS session. MQTT client; each owns one DTLS session, one cache, one scheduler.
""" """
import json import json
import os import os
@@ -26,43 +22,28 @@ import time
import cbor2 import cbor2
from .appliances.base import ApplianceDescriptor from .appliances.base import ApplianceDescriptor, bridge_diagnostic_discovery
from .coap_dtls import DtlsCoapSession, fmt_code from .coap_dtls import DtlsCoapSession, fmt_code
from .config import ApplianceConfig, SharedConfig 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 .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' DEBUG_BRIDGE = os.environ.get('DEBUG_BRIDGE') == '1'
def _href_to_segs(href: str) -> list[str]: 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] 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_PATH = '/information/vs/0'
SERIAL_FIELD = 'x.com.samsung.da.serialNum' SERIAL_FIELD = 'x.com.samsung.da.serialNum'
class PushBridge: 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, def __init__(self,
shared: SharedConfig, shared: SharedConfig,
@@ -74,20 +55,22 @@ class PushBridge:
self.descriptor = descriptor self.descriptor = descriptor
self.mqtt = mqtt_client self.mqtt = mqtt_client
# Bridge-scoped logger; retagged with serial after first seed.
self.log = bridge_logger(app.klass) self.log = bridge_logger(app.klass)
self._serial: str | None = None self._serial: str | None = None
# Resolve port (descriptor default if unset in env).
self.port = app.ocf_port or descriptor.default_observe_port self.port = app.ocf_port or descriptor.default_observe_port
self.session: DtlsCoapSession | None = None self.session: DtlsCoapSession | None = None
self.links: dict[str, dict] = {} # href → rep self.scheduler: PollScheduler | None = None
self.descriptor_state: dict = {} # descriptor scratch space 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_state_pub = None
self.last_remote_pub = None self.last_remote_pub = None
self.last_cycle_pub = None self.last_cycle_pub = None
self.last_avail_pub: str | None = None
self.stop = threading.Event() self.stop = threading.Event()
self.started_ts = time.time() self.started_ts = time.time()
self.session_started_ts = None self.session_started_ts = None
@@ -97,14 +80,25 @@ class PushBridge:
self.connect_count = 0 self.connect_count = 0
self.error_count = 0 self.error_count = 0
self._publish_gate = False 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 # Push is considered "active" if an OBSERVE-sourced change
# bumps the gen; a fetchback aborts on wake (and again after # arrived within this window. Long enough that a quiet but
# its GET completes) if its captured gen is no longer the # working appliance doesn't flap to inactive; short enough that
# latest. This coalesces bursts: rapid lamp toggles or slider # a genuinely silent push channel is visible within minutes.
# drags result in many scheduled fetchbacks but only the self.push_active_window_s = 600.0
# latest one actually publishes. Lock guards the dict mutation
# and the gen comparison. # 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_gen: dict[str, int] = {}
self._fetch_lock = threading.Lock() self._fetch_lock = threading.Lock()
@@ -114,27 +108,33 @@ class PushBridge:
self.remote_topic = f"{p}/remote_available" self.remote_topic = f"{p}/remote_available"
self.cycle_topic = f"{p}/cycle_active" self.cycle_topic = f"{p}/cycle_active"
self.health_topic = f"{p}/bridge/health" self.health_topic = f"{p}/bridge/health"
self.push_active_topic = f"{p}/bridge/push_active"
self.cmd_handlers = descriptor.command_handlers() self.cmd_handlers = descriptor.command_handlers()
self.cmd_topic_prefix = f"{p}/cmd/" self.cmd_topic_prefix = f"{p}/cmd/"
# Pre-built HA discovery payloads. Republished on every MQTT self.discovery_payloads = (
# (re)connect by main.py. descriptor.build_discovery(
self.discovery_payloads = descriptor.build_discovery( app.topic_prefix, shared.HA_DISCOVERY_PREFIX, app.device_name)
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): def _on_notification(self, href, payload_bytes):
"""Invoked by the DTLS reader thread for OBSERVE notifications. """Reader-thread callback for OBSERVE notifications. Large
resources (oven /mode/vs/0 ~9KB) arrive truncated with Block2.M=1
Resources larger than one CoAP block (notably the oven's and we use cbor-decode failure as the partial signal."""
`/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."""
if not payload_bytes: if not payload_bytes:
# Empty payload — almost certainly a Block2 announcement.
self._schedule_fetchback(href) self._schedule_fetchback(href)
return return
try: try:
@@ -144,53 +144,19 @@ class PushBridge:
return return
if not isinstance(rep, dict): if not isinstance(rep, dict):
return 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: if DEBUG_BRIDGE:
# /mode/vs/0 carries a huge modeSpec JSON we don't want to self._debug_log_rep(href, rep)
# log; surface just modes + options. Small resources get self.cache.apply_rep(href, rep, source='observe')
# 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()
def _apply_optimistic(self, href, body): def _debug_log_rep(self, href, rep):
"""Optimistically merge a just-POSTed body into the link dict if href == '/mode/vs/0' and isinstance(rep, dict):
and republish state. Samsung accepts (2.04) writes whose self.log.info("mode modes=%r options=%r",
bodies are field-replacements — we mirror that semantics here: rep.get('x.com.samsung.da.modes'),
each top-level key in `body` overwrites the corresponding key rep.get('x.com.samsung.da.options'))
in the existing rep. The fetchback that follows republishes elif href in ('/operational/state/vs/0', '/oven/vs/0', '/power/vs/0'):
the device's real state, which corrects any field where the self.log.info("REP %s = %r", href, rep)
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 _schedule_fetchback(self, href, delay_s: float = 0.0): 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: with self._fetch_lock:
gen = self._fetch_gen.get(href, 0) + 1 gen = self._fetch_gen.get(href, 0) + 1
self._fetch_gen[href] = gen self._fetch_gen[href] = gen
@@ -202,15 +168,8 @@ class PushBridge:
).start() ).start()
def _fetch_back(self, href, delay_s: float, gen: int): def _fetch_back(self, href, delay_s: float, gen: int):
if delay_s > 0: if delay_s > 0 and self.stop.wait(delay_s):
# Allow the device's read-side to propagate a recent return
# 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.
with self._fetch_lock: with self._fetch_lock:
if self._fetch_gen.get(href) != gen: if self._fetch_gen.get(href) != gen:
return return
@@ -223,31 +182,24 @@ class PushBridge:
except Exception as e: except Exception as e:
self.log.warning("fetchback %s: %s", href, e) self.log.warning("fetchback %s: %s", href, e)
return 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: with self._fetch_lock:
if self._fetch_gen.get(href) != gen: if self._fetch_gen.get(href) != gen:
return return
if code != 0x45: if code != 0x45:
self.log.warning("fetchback %s: %s", self.log.warning("fetchback %s: %s", href, fmt_code(code))
href, fmt_code(code))
return return
try: try:
rep = cbor2.loads(payload) if payload else {} rep = cbor2.loads(payload) if payload else {}
except Exception as e: except Exception as e:
self.log.warning("fetchback %s cbor: %s", href, e) self.log.warning("fetchback %s cbor: %s", href, e)
return return
if not isinstance(rep, dict): if isinstance(rep, dict):
return self.cache.apply_rep(href, rep, source='observe')
self._apply_rep(href, rep)
def _retag_logger_with_serial(self): 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: if self._serial is not None:
return return
info = self.links.get(SERIAL_PATH) or {} info = self.cache.get(SERIAL_PATH) or {}
serial = info.get(SERIAL_FIELD) serial = info.get(SERIAL_FIELD)
if not serial: if not serial:
return return
@@ -258,8 +210,6 @@ class PushBridge:
# ---- session lifecycle ------------------------------------------ # ---- session lifecycle ------------------------------------------
def session_once(self): 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( sess = DtlsCoapSession(
self.app.ip, self.port, self.app.ip, self.port,
cert_path=self.shared.CERT_PATH, cert_path=self.shared.CERT_PATH,
@@ -270,7 +220,7 @@ class PushBridge:
self.session = sess self.session = sess
self.session_started_ts = time.time() self.session_started_ts = time.time()
self.connect_count += 1 self.connect_count += 1
self.descriptor_state = {} self.cache.descriptor_state.clear()
self._publish_gate = False self._publish_gate = False
self.log.info("DTLS connected — subscribing %d paths", self.log.info("DTLS connected — subscribing %d paths",
@@ -278,10 +228,6 @@ class PushBridge:
sess.start_reader() 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() session_ended = threading.Event()
def _stop_watcher(): def _stop_watcher():
@@ -299,18 +245,62 @@ class PushBridge:
try: try:
self._run_session_inner(sess) self._run_session_inner(sess)
finally: finally:
# Always release the watcher so it doesn't sit pinned on
# self.stop forever after the session ends.
session_ended.set() session_ended.set()
def _run_session_inner(self, sess): 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: for path in self.descriptor.observe_paths:
sess.subscribe(path) sess.subscribe(path)
time.sleep(0.05) 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) code, pl = sess.get(self.descriptor.seed_path, timeout=15.0)
if code != 0x45: if code != 0x45:
raise RuntimeError( raise RuntimeError(
@@ -321,100 +311,63 @@ class PushBridge:
raise RuntimeError( raise RuntimeError(
f"/{'/'.join(self.descriptor.seed_path)} cbor decode: {e}" f"/{'/'.join(self.descriptor.seed_path)} cbor decode: {e}"
) from 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(): for href, rep in index_links(body).items():
self.links.setdefault(href, rep) if href not in self.cache.links:
self.cache.apply_rep(href, rep, source='seed')
# 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)
self.last_seed_ts = time.time() self.last_seed_ts = time.time()
self._publish_gate = True def _debug_dump_links(self, sess):
self.maybe_publish_state(force=True) 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.set_availability(True)
self.log.info("seeded → %d links; sensors live", len(self.links)) self.reassert_availability()
sess.join() def _on_unreachable(self) -> None:
self.set_availability(False)
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()
# ---- MQTT publishing -------------------------------------------- # ---- MQTT publishing --------------------------------------------
def maybe_publish_state(self, force=False): def maybe_publish_state(self, force=False):
if not force and not self._publish_gate: if not force and not self._publish_gate:
return return
sensors = self.descriptor.flatten(self.links) snap = self.cache.snapshot()
sensors = self.descriptor.flatten(snap)
project = self.descriptor.project project = self.descriptor.project
if project is not None: 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: if not force and sensors == self.last_state_pub:
return 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.last_state_pub = sensors
self.mqtt.publish(self.state_topic, self.mqtt.publish(self.state_topic,
json.dumps(sensors).encode(), json.dumps(sensors).encode(),
@@ -428,7 +381,8 @@ class PushBridge:
if not force: if not force:
log_fn = self.descriptor.log_state_change log_fn = self.descriptor.log_state_change
extra = log_fn(sensors) if log_fn is not None else '' 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) extra or 'descriptor-no-log', self.notif_count)
def publish_remote_available(self, remote_on, force=False): def publish_remote_available(self, remote_on, force=False):
@@ -465,12 +419,18 @@ class PushBridge:
if cycle_field is not None: if cycle_field is not None:
self.publish_cycle_active( self.publish_cycle_active(
self.last_state_pub.get(cycle_field), force=True) 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): def set_availability(self, online):
value = 'online' if online else 'offline'
if value == self.last_avail_pub:
return
self.last_avail_pub = value
try: try:
self.mqtt.publish(self.avail_topic, self.mqtt.publish(self.avail_topic, value, qos=1, retain=True)
'online' if online else 'offline',
qos=1, retain=True)
except Exception as e: except Exception as e:
self.log.warning("avail publish: %s", e) self.log.warning("avail publish: %s", e)
if not online and self.descriptor.remote_available_field is not None: if not online and self.descriptor.remote_available_field is not None:
@@ -487,6 +447,13 @@ class PushBridge:
qos=1, retain=True) qos=1, retain=True)
except Exception: except Exception:
pass 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 -------------------------------------- # ---- MQTT command handling --------------------------------------
@@ -498,12 +465,9 @@ class PushBridge:
if handler is None: if handler is None:
self.log.warning("unknown command topic: %s", topic) self.log.warning("unknown command topic: %s", topic)
return return
# Shallow-snapshot self.links so the handler sees a consistent # Handler gets a links snapshot so its read-modify-write sees a
# view across the read-modify-write it may need to perform # consistent view across the multi-field operation.
# (e.g. oven lamp / sound / fastpreheat all RMW /mode/vs/0 result = handler(payload, self.cache.snapshot())
# 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))
if result is None: if result is None:
self.log.warning("rejected command %s payload=%r", self.log.warning("rejected command %s payload=%r",
topic, payload) topic, payload)
@@ -513,69 +477,113 @@ class PushBridge:
if sess is None: if sess is None:
self.log.warning("command %s: no DTLS session", topic) self.log.warning("command %s: no DTLS session", topic)
return 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: try:
code, _ = sess.post(path_segs, cbor2.dumps(body), timeout=8.0) code, _ = sess.post(path_segs, cbor2.dumps(body), timeout=8.0)
except Exception as e: except Exception as e:
self.log.warning("command %s POST failed: %s", topic, e) self.log.warning("command %s POST failed: %s", topic, e)
return return
self.log.info("command %s payload=%r → %s", defer_note = f" (poll-defer {href} {defer_s:.0f}s)" if sched is not None else ''
suffix, payload, fmt_code(code)) self.log.info("command %s payload=%r → %s%s",
suffix, payload, fmt_code(code), defer_note)
if code >> 5 == 2: if code >> 5 == 2:
href = '/' + '/'.join(path_segs) # Optimistic local merge so HA sees the write reflected
# Optimistic publish — apply the write to our local state. # immediately. No Block2 fetchback — that triggers Samsung's
# No fetchback: empirically (2026-05-31) the post-write # 3-second revert (project_fetchback_revert_root_cause.md).
# Block2 GET was causing the appliance to roll our values # The PollScheduler will reconcile on its next tier tick
# back ~3s later, on every writable resource. OBSERVE # after the write_in_progress settle window expires.
# pushes keep HA in sync without polling. (Was the root self.cache.apply_optimistic(href, body)
# cause of the "mid-cycle setpoint reverts" symptom.)
self._apply_optimistic(href, body)
def publish_health(self): def publish_health(self):
now = time.time() 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 = { h = {
'mode': 'push', 'mode': 'poll+observe',
'device_class': self.descriptor.name, 'device_class': self.descriptor.name,
'serial': self._serial, 'serial': self._serial,
'connect_count': self.connect_count, 'connect_count': self.connect_count,
'error_count': self.error_count, 'error_count': self.error_count,
'notif_count': self.notif_count, 'notif_count': self.notif_count,
'last_change_age_s': (round(now - self.last_change_ts, 1) 'poll_count': poll,
if self.last_change_ts else None), 'poll_error_count': poll_err,
'last_seed_age_s': (round(now - self.last_seed_ts, 1) 'poll_window_ok': window_polls_ok,
if self.last_seed_ts else None), 'poll_window_errors': d_err,
'session_age_s': (round(now - self.session_started_ts, 1) 'poll_window_max_rtt_ms': round(win_max_rtt, 0),
if self.session_started_ts else None), 'poll_window_slow_count': win_slow,
'uptime_seconds': round(now - self.started_ts, 0), '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: try:
self.mqtt.publish(self.health_topic, json.dumps(h).encode(), self.mqtt.publish(self.health_topic, json.dumps(h).encode(),
qos=0, retain=True) qos=0, retain=True)
except Exception as e: except Exception as e:
self.log.warning("health publish: %s", 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 --------------------------------------------- # ---- 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): def run_forever(self):
threading.Thread(target=self._publish_tick_loop, daemon=True,
name=f'{self.app.klass}-tick').start()
backoff = 1.0 backoff = 1.0
while not self.stop.is_set(): while not self.stop.is_set():
try: try:
-3
View File
@@ -57,7 +57,6 @@ class SharedConfig:
MQTT_PASS: Optional[str] MQTT_PASS: Optional[str]
HA_DISCOVERY_PREFIX: str HA_DISCOVERY_PREFIX: str
HEALTH_INTERVAL_S: int HEALTH_INTERVAL_S: int
HEARTBEAT_INTERVAL_S: int
PING_INTERVAL_S: int PING_INTERVAL_S: int
@classmethod @classmethod
@@ -72,8 +71,6 @@ class SharedConfig:
HA_DISCOVERY_PREFIX=os.getenv('HA_DISCOVERY_PREFIX', HA_DISCOVERY_PREFIX=os.getenv('HA_DISCOVERY_PREFIX',
'homeassistant'), 'homeassistant'),
HEALTH_INTERVAL_S=int(os.getenv('HEALTH_INTERVAL_S', '60')), 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')), 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]