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