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:
+8
-7
@@ -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
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|
||||||
|
|||||||
@@ -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'),
|
||||||
|
]
|
||||||
|
|||||||
@@ -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,
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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
@@ -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:
|
||||||
|
|||||||
@@ -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')),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
@@ -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)
|
||||||
@@ -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]
|
||||||
Reference in New Issue
Block a user