diff --git a/.env.example b/.env.example index 5bd04b1..82162b2 100644 --- a/.env.example +++ b/.env.example @@ -43,14 +43,15 @@ MQTT_PASS= HA_DISCOVERY_PREFIX=homeassistant # Bridge timers (seconds). -# HEALTH_INTERVAL_S — how often /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 /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 diff --git a/README.md b/README.md index 6685c0f..c0328fb 100644 --- a/README.md +++ b/README.md @@ -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 `.` once each device's serial is read on connect — `dryer.` vs `oven.` 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 `//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 `/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 `` is its `APPLIANCE__TOPIC`. | `/availability` | ✓ | `online` after seed; `offline` on disconnect (LWT for appliance #1) | | `/remote_available` | ✓ | `online` iff bridge is up AND Remote Control on the appliance is on. Gates the control entities. | | `/state` | ✓ | JSON sensor dict; published only when sensors actually diff | -| `/bridge/health` | ✓ | Every `HEALTH_INTERVAL_S` — connect_count, error_count, notif_count, last_change_age_s, session_age_s, serial | +| `/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 | | `/{sensor,binary_sensor,switch,light,number,select,button}//.../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__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 `//0` paths.** They register successfully but never push. Use the Samsung `//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 `//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. diff --git a/main.py b/main.py index 830fa38..6826cd7 100644 --- a/main.py +++ b/main.py @@ -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__{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() diff --git a/samsung_appliance/appliances/base.py b/samsung_appliance/appliances/base.py index d363949..6f891cb 100644 --- a/samsung_appliance/appliances/base.py +++ b/samsung_appliance/appliances/base.py @@ -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'), + ] diff --git a/samsung_appliance/appliances/dryer.py b/samsung_appliance/appliances/dryer.py index d9d965c..c182289 100644 --- a/samsung_appliance/appliances/dryer.py +++ b/samsung_appliance/appliances/dryer.py @@ -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, ) diff --git a/samsung_appliance/appliances/oven.py b/samsung_appliance/appliances/oven.py index 91ac812..e5eed9e 100644 --- a/samsung_appliance/appliances/oven.py +++ b/samsung_appliance/appliances/oven.py @@ -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, ) diff --git a/samsung_appliance/bridge.py b/samsung_appliance/bridge.py index 7dae272..1e5d2e4 100644 --- a/samsung_appliance/bridge.py +++ b/samsung_appliance/bridge.py @@ -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: diff --git a/samsung_appliance/config.py b/samsung_appliance/config.py index 6c46188..0edb89a 100644 --- a/samsung_appliance/config.py +++ b/samsung_appliance/config.py @@ -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')), ) diff --git a/samsung_appliance/keepalive.py b/samsung_appliance/keepalive.py new file mode 100644 index 0000000..f704c40 --- /dev/null +++ b/samsung_appliance/keepalive.py @@ -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) diff --git a/samsung_appliance/poll_scheduler.py b/samsung_appliance/poll_scheduler.py new file mode 100644 index 0000000..2023817 --- /dev/null +++ b/samsung_appliance/poll_scheduler.py @@ -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) diff --git a/samsung_appliance/state_cache.py b/samsung_appliance/state_cache.py new file mode 100644 index 0000000..852c34f --- /dev/null +++ b/samsung_appliance/state_cache.py @@ -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]