commit c99ef324fc1ab7cea777d1198e5d7970bafe1219 Author: Jack Nagy Date: Sun May 31 15:50:15 2026 +0100 Initial commit diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..5bd04b1 --- /dev/null +++ b/.env.example @@ -0,0 +1,62 @@ +# SmartThings-Local Bridge config. +# Copy to `.env` and fill in. Never commit `.env`. + +# ============================================================= +# Appliances — one process supervises N appliances over DTLS. +# ============================================================= +# APPLIANCE_COUNT defines how many entries to read. Per-appliance +# keys are 1-indexed (APPLIANCE_1_*, APPLIANCE_2_*, …). +APPLIANCE_COUNT=1 + +# Appliance 1 — Samsung dryer +APPLIANCE_1_CLASS=dryer +APPLIANCE_1_IP=192.168.1.100 +# Leave OCF_PORT blank to inherit the descriptor's default +# (dryer=49155, oven=49154). +APPLIANCE_1_OCF_PORT= +APPLIANCE_1_TOPIC=samsung_dryer +APPLIANCE_1_NAME=Samsung Dryer + +# Future: +# APPLIANCE_2_CLASS=oven +# APPLIANCE_2_IP=192.168.1.101 +# APPLIANCE_2_OCF_PORT= +# APPLIANCE_2_TOPIC=samsung_oven +# APPLIANCE_2_NAME=Samsung Oven +# (Don't forget to bump APPLIANCE_COUNT=2.) + +# --- Cert paths --- +# Defaults work for Docker (mount as /config) and bare-metal (drop +# into ./certs). The ab0b0ac4 admin-override cert + key are built by +# local-tools/setup_samsung_cloud_cert.py. +# CERT_PATH=./certs/ab0b0ac4_fullchain.pem +# KEY_PATH=./certs/ab0b0ac4.key + +# --- MQTT broker (HA Mosquitto add-on or any broker) --- +MQTT_BROKER=192.168.1.5 +MQTT_PORT=1883 +MQTT_USER=samsung_bridge +MQTT_PASS= + +# HA discovery prefix — must match the MQTT integration's setting in HA +# (default `homeassistant`). +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=60 +HEARTBEAT_INTERVAL_S=600 + +# Container TZ. +TZ=Europe/London + + +# --- Deploy (deploy.sh — tar + ssh docker compose) --- +SSH_HOST=user@your-server +REMOTE_DIR=/mnt/user/compose/smartthings-local +APPDATA_DIR=/mnt/user/appdata/smartthings-local diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..e43687b --- /dev/null +++ b/.gitignore @@ -0,0 +1,35 @@ +# Secrets and per-instance config — never commit +.env +.DS_Store + +# Certificates and keys — per-instance, never push +certs/ +*.pem +*.crt +*.p12 +*.pfx +*.key +token.txt + +# Python artifacts +__pycache__/ +*.pyc +.venv/ +venv/ + +# Bridge runtime +bridge.log +course_mapper.log + +# Capture / dump files +*.pcap +*.har +*.mitm +capture/ +flows/ + +# Local-only research tools (not part of the bridge runtime) +local-tools/ + +# Claude Code per-project state (permissions allowlist, etc.) +.claude/ diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..ddad334 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,31 @@ +FROM python:3.11-slim + +WORKDIR /app + +# Python deps first so layer cache survives code changes +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +# Application code +COPY main.py . +COPY samsung_appliance/ ./samsung_appliance/ + +# /config holds the ab0b0ac4 client cert + key. Mount from the host so +# secrets aren't baked into the image. +RUN mkdir -p /config + +# Unbuffered stdout so docker logs is live +ENV PYTHONUNBUFFERED=1 + +# Defaults — override in .env or `docker run -e …`. Per-appliance +# keys (APPLIANCE_COUNT, APPLIANCE__*) have no universal default +# and must be set in .env. +ENV CERT_PATH=/config/ab0b0ac4_fullchain.pem \ + KEY_PATH=/config/ab0b0ac4.key \ + HA_DISCOVERY_PREFIX=homeassistant \ + HEALTH_INTERVAL_S=60 \ + HEARTBEAT_INTERVAL_S=600 + +# No port — bridge is outbound-only (DTLS UDP to appliance, MQTT to broker). + +CMD ["python", "main.py"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..152c60d --- /dev/null +++ b/README.md @@ -0,0 +1,387 @@ +# SmartThings-Local + +**Local-first Home Assistant integration for newer-generation Samsung connected appliances.** One process supervises multiple appliances (dryer + oven currently), each over its own CoAP-DTLS session, publishing state + writes through MQTT with HA auto-discovery — no SmartThings cloud round-trip for any of it. + +> ### Proof of concept — collaborators wanted +> +> This is working code running in my home and I rely on it daily, but it's a **proof of concept**, not a polished product. No unit tests; one person's hardware as the validation set (one dryer model, one oven model); hand-rolled MQTT-based integration instead of a proper HA custom component; "wired-but-untested" comments scattered through the oven descriptor; brittle to per-firmware quirks (the "oven doesn't push OBSERVE on options writes" finding is the kind of thing that needs ongoing care). +> +> **I would love for someone to take this further and build a proper HA integration out of it.** All the protocol research is done — DTLS auth via Samsung's published cloud identity, token-stable Block2 reads, OBSERVE-then-fetchback notifications, write semantics, the optimistic-publish-then-verify pattern, brick-avoiding resource boundaries — and the descriptor pattern is the seed of a clean per-appliance abstraction. The HA-side polish that's missing is custom-component shape: config flow, native entity classes, async-Python DTLS instead of MQTT round-trips, error surfacing into HA's notification system, support across more firmware versions, and someone who actually lives in the HA codebase. +> +> If you're that person, get in touch — happy to co-author, hand off, or hand over entirely. + +### 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. +- **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. +- **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. +- **Your state stays on your LAN.** Bridge → broker → HA. Samsung's cloud sees nothing from HA. *(The appliance still maintains its own TLS session to Samsung — appliance design, not ours.)* + +### 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. + +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. + +--- + +## Part 1 — Is your appliance compatible? + +Check before anything else; if it's older firmware, this project doesn't target it. + +```sh +# UDP scan for DTLS-CoAP ports +nmap -Pn -sU -p 49152-49160 "$APPLIANCE_IP" +``` + +Read the result: + +- **`49154/udp` (or similar 4915x) open|filtered with a DTLS handshake responding** → newer firmware (Tizen RT 3.x with DAWIT 3.0). This is what the bridge talks to. +- **Only `8888/tcp` open (token-based HTTPS)** → older firmware (~2018–2022). **Not supported here.** + +### Tested combinations + +| 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) | + +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/`. + +--- + +## 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`. + +You can verify the UUID yourself with one OpenSSL command: + +```sh +openssl s_client -connect connect-v2.samsungiotcloud.com:443 \ + -servername connect-v2.samsungiotcloud.com \ + -showcerts < /dev/null 2>/dev/null \ + | openssl x509 -noout -subject +# subject=C=KR, O=Samsung Electronics, OU=uuid:, CN=*.samsungiotcloud.com +``` + +The UUID lives in `OU=uuid:`. Samsung's cert is valid through **2035-04-09**. + +This README deliberately doesn't pin the literal UUID — the setup script extracts it live each run, so it self-updates if Samsung ever rotates. + +### Why this works + +- Every Samsung Tizen/RT-OCF appliance has a **factory-baked ACE** in `/oic/sec/acl` granting this UUID `perm=31` on `href=*`. It's the identity Samsung's own cloud-bridge daemon uses when forwarding cloud-issued commands to the on-device OCF stack. +- TizenRT iotivity derives peerId from `memmem(subject_dn, "uuid:")` — RDN-agnostic. A cert with the UUID in CN authenticates the same as one with it in OU. +- We don't have Samsung's matching private key (HSM-bound on their cloud) but we don't need it — we mint our own key and have `AC14K_M` sign our leaf. Different key, same identity, same access. + +### One-command setup + +You need `AC14K_M.pem`, its key, and the three upstream chain certs (`cert_1.pem`…`cert_4.pem`). These are published in [cicciovo/homebridge-samsung-airconditioner](https://github.com/cicciovo/homebridge-samsung-airconditioner). Drop them into `./certs/`. + +```sh +AC14K_M_CERT=./certs/ac14k_m.pem \ +AC14K_M_KEY=./certs/ac14k_m.key \ +CHAIN_DIR=./certs/ \ +OUT_DIR=./certs/ \ +TARGET_IP=$APPLIANCE_IP TARGET_PORT=49154 \ +python local-tools/setup_samsung_cloud_cert.py --test +``` + +What it does: + +1. **Live-fetches** Samsung's wildcard cloud cert and extracts the current cloud-bridge UUID. +2. Generates a fresh RSA-2048 key pair you own. +3. Builds a CSR with the UUID in OU + CN + SAN, signs it with `AC14K_M` (SHA-1). +4. Concatenates `leaf + AC14K_M + 3 upstream CAs` into `fullchain.pem`. +5. With `--test`: opens a DTLS handshake against `$TARGET_IP:$TARGET_PORT` and GETs `/oic/sec/acl` — a `2.05` reply proves the cert authenticated as the cloud-identity peer (anonymous peers get `4.01` on that resource). + +Output: `ab0b0ac4_fullchain.pem` + `ab0b0ac4.key` (filename matches the UUID prefix as a convention; the actual UUID is whatever was published live). Drop them in `./certs/`. + +The UUID is **not hardcoded** anywhere in the script or this README. If the live fetch fails (restricted network), `UUID= python setup_samsung_cloud_cert.py …` lets you supply it manually; the docstring documents the openssl-extract one-liner. + +### How durable is this? + +Rotating the cloud-bridge UUID is roughly equivalent to Samsung re-issuing TLS certs across their entire IoT cloud AND pushing new ACLs to every device in the field AND updating the on-device cloud-bridge daemon's identity — a multi-quarter project with a months-long backwards-compat window. The `AC14K_M` signing CA has been publicly leaked for years and still appears in 2026 firmware trust stores. Our access is roughly as durable as SmartThings cloud control of these appliances. + +> **Legacy path:** earlier versions of this project used a per-hub-UUID cert via an anonymous `/oic/sec/doxm` read escalation. That still works on the dryer-family firmware but isn't necessary — the ab0b0ac4 cert is one identity that authenticates against every appliance, factory ACL, and survives device resets. `bootstrap.py` in the repo automates the legacy path if you'd rather; otherwise ignore it. + +--- + +## Part 3 — Configure your appliances + +Copy `.env.example` to `.env` and fill in. + +### Layered envs + +The bridge config splits into: + +- **Shared keys** (one per process): MQTT broker + creds, HA discovery prefix, cert paths, timer intervals. +- **Per-appliance keys** (one block per appliance) under `APPLIANCE__*` (1-indexed). + +`APPLIANCE_COUNT` tells the bridge how many indexed blocks to read. Bump it as you add appliances. + +```bash +APPLIANCE_COUNT=2 + +# Appliance 1 — dryer +APPLIANCE_1_CLASS=dryer +APPLIANCE_1_IP=192.168.1.100 +APPLIANCE_1_OCF_PORT= # blank → descriptor default (49155 for dryer) +APPLIANCE_1_TOPIC=samsung_dryer +APPLIANCE_1_NAME=Samsung Dryer + +# Appliance 2 — oven +APPLIANCE_2_CLASS=oven +APPLIANCE_2_IP=192.168.1.101 +APPLIANCE_2_OCF_PORT= # blank → descriptor default (49154 for oven) +APPLIANCE_2_TOPIC=samsung_oven +APPLIANCE_2_NAME=Samsung Oven +``` + +Each `APPLIANCE__CLASS` must match a descriptor key in `samsung_appliance/appliances/__init__.py::DESCRIPTORS` — currently `dryer` and `oven`. + +--- + +## Part 4 — Run it + +### Docker (the real deployment) + +```sh +docker compose up -d --build +docker compose logs -f +``` + +Container name `smartthings-local`. Outbound-only — no ports exposed. Needs egress to each appliance's IP/port (UDP) and to your MQTT broker. The certs in `./certs/` (or whatever `APPDATA_DIR` points to via the volume mount) are read-only mounted at `/config`. + +### Deploying to a remote Linux host (Unraid, etc.) + +```sh +# Once: upload the cert + key onto the remote. +ssh "$SSH_HOST" mkdir -p "$APPDATA_DIR" +scp certs/ab0b0ac4_fullchain.pem certs/ab0b0ac4.key "$SSH_HOST:$APPDATA_DIR/" + +# Each deploy: ship source + .env, rebuild container on the host. +./deploy.sh +``` + +Set `SSH_HOST`, `REMOTE_DIR`, `APPDATA_DIR` in `.env`. `deploy.sh` extracts those three keys via `grep` rather than `source .env`, so values containing spaces (like `APPLIANCE_1_NAME=Samsung Dryer`) don't break it. + +### Bare metal (first test / debugging) + +```sh +python3 -m venv .venv +.venv/bin/pip install -r requirements.txt +.venv/bin/python main.py +``` + +### Expected first-run logs + +``` +14:08:42 INFO samsung_appliance SmartThings-Local Bridge starting (2 appliances) +14:08:42 INFO samsung_appliance broker = :1883 (user=) +14:08:42 INFO samsung_appliance [1] dryer @ :49155 (DTLS) → topic samsung_dryer/* +14:08:42 INFO samsung_appliance [2] oven @ :49154 (DTLS) → topic samsung_oven/* +14:08:42 INFO samsung_appliance MQTT connected → :1883 +14:08:43 INFO dryer DTLS connected — subscribing 11 paths +14:08:44 INFO dryer. identified — serial=… +14:08:44 INFO dryer. seeded → 25 links; sensors live +14:08:44 INFO oven DTLS connected — subscribing 11 paths +14:08:46 INFO oven. identified — serial=… +14:08:46 INFO oven. seeded → 16 links; sensors live +``` + +In HA: **Settings → Devices & Services → MQTT** should show both devices populated. + +--- + +## Per-appliance notes + +### Dryer + +| Capability | Works? | Notes | +|---|---|---| +| Read all state | ✅ | Machine state, job state, energy (W + kWh), course, dry level, completion time, remote control, child lock, alarms | +| Wrinkle Prevent toggle | ✅ | Persists | +| Start / Pause / Stop | ✅ | Via `/operational/state/vs/0`; needs Remote Control on | +| Change course | ✅ | Via `/st/dryercourse/vs/0`; needs Remote Control on. **Not exposed by the SmartThings cloud HA integration.** | +| 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. + +### Oven + +| Capability | Works? | Notes | +|---|---|---| +| 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. | +| **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. + +--- + +## Reference + +### Config keys + +| Key | Meaning | +|---|---| +| `APPLIANCE_COUNT` | Number of `APPLIANCE__*` blocks to read (1-indexed) | +| `APPLIANCE__CLASS` | Descriptor name: `dryer`, `oven` | +| `APPLIANCE__IP` | LAN IP of the appliance | +| `APPLIANCE__OCF_PORT` | Optional override (blank → descriptor default: dryer=49155, oven=49154) | +| `APPLIANCE__TOPIC` | MQTT topic prefix (also the HA device identifier — changing it re-keys the device) | +| `APPLIANCE__NAME` | Friendly name on the HA device card | +| `MQTT_BROKER` / `MQTT_PORT` / `MQTT_USER` / `MQTT_PASS` | Broker config | +| `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) | +| `SSH_HOST` / `REMOTE_DIR` / `APPDATA_DIR` | Used by `deploy.sh` only | + +### MQTT topics — outgoing (bridge → broker) + +Per appliance, where `` is its `APPLIANCE__TOPIC`. + +| Topic | Retain | When | +|---|---|---| +| `/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 | +| `/{sensor,binary_sensor,switch,light,number,select,button}//.../config` | ✓ | HA MQTT discovery, republished on every MQTT (re)connect | + +### MQTT topics — incoming (bridge subscribes) + +`/cmd/#`. **The MQTT user must have READ permission on this subtree** — without it the broker silently drops the TCP connection shortly after SUBSCRIBE. Check broker logs if writes never land. + +Dryer: + +| Suffix | Payloads | Effect | +|---|---|---| +| `cmd/wrinkle_prevent` | `On`, `Off` | POST `/washer/vs/0` | +| `cmd/operational_state` | `Run`, `Pause`, `Ready` | POST `/operational/state/vs/0` — requires RC | +| `cmd/dryer_mode` | Course name (e.g. `Cotton`) | Translated to `Course_HH` then POST `/st/dryercourse/vs/0` — requires RC | + +Oven: + +| Suffix | Payloads | Effect | +|---|---|---| +| `cmd/lamp` | `On`, `Off` | RMW of `/mode/vs/0 .options[UpperLamp_*]` | +| `cmd/sound` | `On`, `Off` | RMW of `/mode/vs/0 .options[Sound_*]` | +| `cmd/fastpreheat` | `On`, `Off` | RMW of `/mode/vs/0 .options[fastpreheat_*]` | +| `cmd/setpoint` | Integer °C (30–270, step 5) | RMW of `/temperatures/vs/0 .items[0].desired` — requires RC | +| `cmd/mode` | Mode name (e.g. `Convection`, `LargeGrill`) | POST `/mode/vs/0 {modes: []}` — requires RC | +| `cmd/stop` | (button press) | POST `/operational/state/vs/0 {state: Ready}` | + +### Entity counts (approximate, per appliance) + +| Type | Dryer | Oven | +|---|---|---| +| `sensor` | 17 | 17 | +| `binary_sensor` | 4 | 7 | +| `switch` | 1 (wrinkle) | 2 (sound, fastpreheat) | +| `light` | — | 1 (lamp) | +| `number` | — | 1 (setpoint slider) | +| `select` | 1 (course) | 1 (mode) | +| `button` | 3 (start/pause/stop) | 1 (stop) | + +Gated control entities use HA's `availability_mode: all` against `/availability` AND `/remote_available`. Flip Remote Control on the appliance's front panel and those entities un-grey in HA. + +### Repo layout + +``` +main.py Entry point — loads config, spawns one PushBridge per appliance +samsung_appliance/ The bridge package + __init__.py + config.py SharedConfig + ApplianceConfig dataclasses + logger.py Tagged logger helpers + bridge.py PushBridge — one DTLS session per appliance, descriptor-driven + coap_dtls.py DTLS-CoAP session: handshake, token-stable Block2 GET, POST, OBSERVE + sensors.py /device/0 link-dict indexer (shared util) + appliances/ + __init__.py DESCRIPTORS registry + get_descriptor() + base.py ApplianceDescriptor dataclass + HA discovery helpers + dryer.py Dryer descriptor (paths, flatten, discovery, commands) + oven.py Oven descriptor +Dockerfile Container build (python:3.11-slim + 3 deps) +docker-compose.yml One service: smartthings-local +deploy.sh tar + ssh + docker compose up --build +.env.example Template — copy to .env, fill in +local-tools/ Research/probes — gitignored + setup_samsung_cloud_cert.py One-shot cert minting script + probe_oven_*.py DTLS probes for the oven (lamp, OBSERVE, full /device/0 fetch) + comparisons/ Per-appliance /device/0 dumps + diff +``` + +`certs/` is gitignored. Drop the privileged client cert + key there; the container mounts that directory read-only at `/config`. + +--- + +## Adding a new appliance class + +The bridge is appliance-agnostic. Adding e.g. a washer is mechanical: + +1. Capture the appliance's `/device/0` to see what resources/fields it exposes. The setup script's `--test` mode is a good start; for the full dump use `local-tools/probe_oven_full_fetch.py` as a template. +2. Create `samsung_appliance/appliances/washer.py` with: + - `OBSERVE_PATHS` — list of `[seg, …]` paths to subscribe to (only `//vs/0` resources push; the OCF-standard `//0` siblings register but never fire) + - `flatten(links) -> dict` — map link dict to the flat sensor dict that goes on MQTT + - `build_discovery(prefix, ha_prefix, name) -> [(topic, payload), …]` — HA discovery configs + - `command_handlers() -> {suffix: fn(payload, links)}` — MQTT commands → `(path_segs, body_dict)` + - A module-level `WASHER = ApplianceDescriptor(name='washer', default_observe_port=…, …)` +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. + +--- + +## Traps to avoid + +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 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. + +--- + +## Known DTLS flakiness + +Samsung's RT-OCF DTLS stack occasionally closes sessions actively — usually right after a Block2 GET or in the seconds after a POST. The bridge handles this with exponential reconnect (1s → 30s) and a re-seed on each new session. From HA's perspective the entity briefly goes offline then comes back; from the bridge's perspective you'll see lines like: + +``` +oven.… DTLS recv: Unexpected EOF +oven.… reconnect in 1s +oven.… DTLS connected — subscribing 11 paths +oven.… seeded → 16 links; sensors live +``` + +If reconnects become persistent (e.g. >10 in a minute) something's actually wrong — check the appliance's Wi-Fi link first, then look for a competing DTLS client on the LAN. + +--- + +## Contributing + +Patches welcome — especially: + +- New appliance descriptors (washer, dishwasher, AC, fridge, etc.) on the same Tizen RT 3.x firmware family. +- Confirmation/refutation on additional dryer or oven models. `nmap` + `/device/0` dump + `/oic/d` GET is enough to know if you're on the same firmware family. +- A proper HA custom component wrapping the bridge so there's a config flow instead of YAML/env editing. + +If you submit a PR, please don't include real device UUIDs, MACs, serials, IPs, or bearer tokens — use the placeholders from `.env.example`. diff --git a/bootstrap.py b/bootstrap.py new file mode 100644 index 0000000..f225ea8 --- /dev/null +++ b/bootstrap.py @@ -0,0 +1,480 @@ +#!/usr/bin/env python3 +"""Interactive setup for samsung-appliance-local. + +Run this once before `main.py`. It will: + + 1. Ask for your dryer's IP and OCF port; verify the port is reachable. + 2. Locate Samsung's AC14K_M intermediate CA cert + key on disk + (you have to fetch these yourself — see the README link). + 3. Try to discover your SmartThings hub UUID anonymously from the + dryer's /oic/sec/acl. If that fails, ask you for it. + 4. Generate a leaf cert (SHA-1 RSA, Samsung iot-Identity + role OIDs, + Subject CN=urn:uuid:) signed by AC14K_M, and write + certs/mega.key + certs/mega_chain.pem. + 5. Offer to populate .env from .env.example with the IP/port. + +This script is setup-only — `cryptography` is not a runtime dep. Install +into a venv: + + python -m venv .venv + .venv/bin/pip install -r requirements-bootstrap.txt + .venv/bin/python bootstrap.py +""" +import os +import shutil +import socket +import ssl +import subprocess +import sys +import tempfile +from pathlib import Path + +try: + import cbor2 +except ImportError: + sys.exit("cbor2 not installed — pip install -r requirements-bootstrap.txt") + +from samsung_dryer.coap import ( + URI_PATH, CSM, enc_opts, enc_tcp, read_tcp, fmt_code, +) + + +REPO_ROOT = Path(__file__).resolve().parent +CERTS_DIR = REPO_ROOT / 'certs' + +# Samsung-specific OIDs the dryer firmware looks for in the leaf. +SAMSUNG_IOT_IDENTITY_OID = '1.3.6.1.4.1.51414.0.1.2' +SAMSUNG_ROLE_OID = '1.3.6.1.4.1.51414.1.3' + +# AC14K_M cert link — used in user-facing error messages so the recipe +# is self-contained. +AC14K_M_SOURCE = ( + 'https://github.com/cicciovo/homebridge-samsung-airconditioner ' + '(see ac14k_m.pem and the matching key)' +) + + +# ---------- tiny UX helpers ------------------------------------------------ + +BOLD = '\033[1m' +DIM = '\033[2m' +GREEN = '\033[32m' +RED = '\033[31m' +YEL = '\033[33m' +END = '\033[0m' + +def _tty(): + return sys.stdout.isatty() + +def info(msg): print(f"{BOLD}»{END} {msg}" if _tty() else f"» {msg}") +def ok(msg): print(f"{GREEN}✓{END} {msg}" if _tty() else f"OK {msg}") +def warn(msg): print(f"{YEL}!{END} {msg}" if _tty() else f"! {msg}") +def fail(msg): print(f"{RED}✗{END} {msg}" if _tty() else f"FAIL {msg}") +def dim(msg): print(f"{DIM}{msg}{END}" if _tty() else msg) + +def prompt(question, default=None): + suffix = f" [{default}]" if default is not None else "" + while True: + try: + ans = input(f" {question}{suffix}: ").strip() + except EOFError: + print(); sys.exit(130) + if ans: + return ans + if default is not None: + return default + +def confirm(question, default=True): + suffix = ' [Y/n]' if default else ' [y/N]' + while True: + try: + ans = input(f" {question}{suffix}: ").strip().lower() + except EOFError: + print(); sys.exit(130) + if not ans: + return default + if ans in ('y', 'yes'): return True + if ans in ('n', 'no'): return False + + +# ---------- step 1: AC14K_M discovery ------------------------------------- +# Note: we deliberately do NOT do a bare TCP reachability probe before +# the real TLS handshake. The dryer's OCF stack treats a plain +# TCP-open-then-close (no TLS) as anomalous and enters a defensive state +# that closes subsequent handshakes' sockets immediately after CSM. +# Empirically observed; see commit history. Reachability is checked +# implicitly when we open TLS in step 3. + +def find_ac14km(): + """Look in ./certs/ for the AC14K_M cert + key under any of the + common filenames. Returns (cert_path, key_path) or (None, None).""" + cert_candidates = ['ac14k_m.pem', 'AC14K_M.pem', 'cert_1.pem'] + key_candidates = ['ac14k_m.key', 'AC14K_M.key', 'key.pem', 'ac14k_m_key.pem'] + cert = next((CERTS_DIR / n for n in cert_candidates if (CERTS_DIR / n).exists()), None) + key = next((CERTS_DIR / n for n in key_candidates if (CERTS_DIR / n).exists()), None) + return cert, key + + +def check_openssl(): + """Bootstrap shells out to openssl for cert generation — SHA-1 signing + was removed from python-cryptography in v43, and openssl is ubiquitous + enough that requiring it is reasonable.""" + if shutil.which('openssl') is None: + fail("openssl not found in PATH — required for cert generation") + return False + return True + + +def _run(cmd, **kw): + """Wrapper that surfaces stderr on failure.""" + res = subprocess.run(cmd, capture_output=True, text=True, **kw) + if res.returncode != 0: + raise RuntimeError( + f"`{' '.join(cmd)}` failed:\n{res.stderr.strip() or res.stdout.strip()}" + ) + return res + + +def _openssl_config(common_name, hub_uuid=None, include_samsung_role=True): + """Return an OpenSSL config snippet matching the proven canonical recipe + used to generate the original working `mega_chain.pem` for this project + (see spoof/mega_ext.cnf). All four SAN entries and the `clientAuth, + serverAuth` EKU values are defensive — the dryer's `memmem` scan only + cares about the Subject DN, but adjacent tooling reads the rest.""" + v3_lines = [ + "basicConstraints = CA:FALSE", + "keyUsage = digitalSignature, keyEncipherment", + f"extendedKeyUsage = clientAuth, serverAuth, {SAMSUNG_IOT_IDENTITY_OID}", + ] + if hub_uuid: + v3_lines.append("subjectAltName = @alt_names") + if include_samsung_role: + v3_lines.append( + f"{SAMSUNG_ROLE_OID} = ASN1:UTF8String:samsung.role.hub") + sections = [ + "[ req ]", + "distinguished_name = dn", + "prompt = no", + "req_extensions = v3", + "", + "[ dn ]", + f"CN = {common_name}", + "O = Samsung Electronics", + "C = KR", + "", + "[ v3 ]", + *v3_lines, + ] + if hub_uuid: + # Belt-and-braces SAN entries — three URI forms and a DNS name. + # Matches the canonical mega_ext.cnf exactly so the leaf is + # bit-for-bit equivalent to the cert known to authenticate. + sections += [ + "", + "[ alt_names ]", + f"URI.1 = urn:uuid:{hub_uuid}", + f"URI.2 = uri:uuid:{hub_uuid}", + f"URI.3 = uuid:{hub_uuid}", + f"DNS.1 = {hub_uuid}", + ] + return "\n".join(sections) + "\n" + + +def _generate_signed_cert(*, common_name, hub_uuid, include_samsung_role, + ca_cert, ca_key, out_key, out_cert, days): + """Generate an RSA-2048 key + SHA-1 signed cert via openssl.""" + with tempfile.TemporaryDirectory() as td: + tdp = Path(td) + conf = tdp / 'leaf.cnf' + csr = tdp / 'leaf.csr' + conf.write_text(_openssl_config(common_name, hub_uuid, + include_samsung_role)) + # 1) key + CSR with extensions baked into req_extensions + _run(['openssl', 'req', '-new', '-newkey', 'rsa:2048', '-nodes', + '-keyout', str(out_key), '-out', str(csr), '-config', str(conf)]) + # 2) sign with AC14K_M, SHA-1, copy the v3 extensions through + _run(['openssl', 'x509', '-req', '-in', str(csr), + '-CA', str(ca_cert), '-CAkey', str(ca_key), + '-CAcreateserial', '-out', str(out_cert), + '-days', str(days), '-sha1', + '-extfile', str(conf), '-extensions', 'v3']) + os.chmod(out_key, 0o600) + + +def generate_leaf(hub_uuid, ca_cert, ca_key, out_dir): + """The real leaf — Subject CN contains `urn:uuid:` so the + dryer's `memmem` scan recognises us as the SmartThings hub. Writes + mega.key and mega_chain.pem (leaf || AC14K_M).""" + subject_uri = f"urn:uuid:{hub_uuid}" + out_key = out_dir / 'mega.key' + out_leaf = out_dir / 'mega_leaf.pem' + out_chain = out_dir / 'mega_chain.pem' + _generate_signed_cert( + common_name=subject_uri, + hub_uuid=hub_uuid, + include_samsung_role=True, + ca_cert=ca_cert, ca_key=ca_key, + out_key=out_key, out_cert=out_leaf, + days=365 * 5, + ) + # Concatenate leaf || AC14K_M for the bridge's load_cert_chain. + out_chain.write_bytes(out_leaf.read_bytes() + Path(ca_cert).read_bytes()) + out_leaf.unlink() + return out_key, out_chain + + +def generate_probe(ca_cert, ca_key, tmp_dir): + """Throwaway leaf with NO `uuid:` in the Subject DN — the dryer treats + us as an anonymous-but-CA-trusted peer. Used once to attempt the + anonymous ACL read; never written to disk outside tmp_dir.""" + out_key = tmp_dir / 'probe.key' + out_leaf = tmp_dir / 'probe.pem' + out_chain = tmp_dir / 'probe_chain.pem' + _generate_signed_cert( + common_name='samsung-local-bootstrap-probe', + hub_uuid=None, + include_samsung_role=False, + ca_cert=ca_cert, ca_key=ca_key, + out_key=out_key, out_cert=out_leaf, + days=30, + ) + out_chain.write_bytes(out_leaf.read_bytes() + Path(ca_cert).read_bytes()) + out_leaf.unlink() + return out_key, out_chain + + +# ---------- step 4: anonymous ACL read ------------------------------------ + +def open_tls(host, port, cert_path, key_path, timeout=8): + """Same pattern as samsung_dryer.bridge._open_tls — drop OpenSSL 3.x + security level so SHA-1 leaves are accepted.""" + ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT) + ctx.check_hostname = False + ctx.verify_mode = ssl.CERT_NONE + try: + ctx.set_ciphers('DEFAULT:@SECLEVEL=0') + except ssl.SSLError: + pass + ctx.load_cert_chain(certfile=str(cert_path), keyfile=str(key_path)) + raw = socket.create_connection((host, port), timeout=timeout) + sock = ctx.wrap_socket(raw) + sock.send(CSM) + sock.settimeout(2) + try: read_tcp(sock) + except (socket.timeout, ConnectionError): pass + sock.settimeout(timeout) + return sock + + +def coap_get(sock, path_segs, token=b'\x01\x02\x03\x04'): + opts = [(URI_PATH, s.encode()) for s in path_segs] + sock.send(enc_tcp(0x01, token=token, opts_b=enc_opts(opts))) + code, _tok, _opts, pl = read_tcp(sock) + return code, pl + + +def extract_hub_uuid_from_doxm(doxm_payload): + """Parse the CBOR-encoded /oic/sec/doxm response and return the hub + UUID. On this firmware, `devowneruuid` and `rowneruuid` both carry + the SmartThings hub's UUID — they're the same value in practice and + we prefer devowneruuid (the OCF spec field for the device's owner).""" + try: + doc = cbor2.loads(doxm_payload) + except Exception as e: + warn(f"doxm CBOR decode failed: {e}") + return None + if not isinstance(doc, dict): + warn(f"doxm decoded to {type(doc).__name__}, expected dict") + return None + for key in ('devowneruuid', 'rowneruuid'): + val = doc.get(key) + if isinstance(val, str) and looks_like_uuid(val): + return val + warn(f"doxm payload had no devowneruuid/rowneruuid (keys: " + f"{list(doc.keys())})") + return None + + +def try_anonymous_doxm_read(host, port, ca_cert, ca_key): + """Discover the hub UUID by reading /oic/sec/doxm anonymously. + + Mechanism: the dryer's baseline ACL contains a wildcard ACE + (`subjectuuid=*` perm=2) granting any authenticated peer read access + to /oic/sec/doxm. We don't need to be the hub — we just need to + complete a chain-valid TLS handshake. doxm.devowneruuid is the + SmartThings hub's UUID.""" + with tempfile.TemporaryDirectory() as td: + tdp = Path(td) + try: + key_path, chain_path = generate_probe(ca_cert, ca_key, tdp) + except RuntimeError as e: + warn(f"probe cert generation failed: {e}") + return None + try: + sock = open_tls(host, port, chain_path, key_path) + except ConnectionRefusedError: + fail(f"connection refused at {host}:{port} — wrong port, or " + f"the dryer isn't on the LAN.") + return None + except (ssl.SSLError, OSError) as e: + warn(f"anonymous TLS handshake failed: {e}") + return None + try: + code, pl = coap_get(sock, ['oic', 'sec', 'doxm']) + except ConnectionError as e: + warn(f"dryer closed the CoAP session immediately after CSM: {e}") + dim(" This usually means the dryer's OCF stack is in a " + "defensive cooldown — typically caused by a concurrent " + "TLS session (the bridge running) or rapid recent probes. " + "Stop main.py / the bridge container, wait ~60s, then re-run.") + return None + finally: + try: sock.close() + except Exception: pass + if code != 0x45: + warn(f"GET /oic/sec/doxm → {fmt_code(code)} (expected 2.05) " + f"— switching to manual entry") + return None + return extract_hub_uuid_from_doxm(pl) + + +# ---------- step 5: .env --------------------------------------------------- + +def maybe_write_env(appliance_ip, appliance_port): + env_path = REPO_ROOT / '.env' + example = REPO_ROOT / '.env.example' + if not example.exists(): + warn(".env.example missing — skipping .env generation") + return + if env_path.exists(): + if not confirm("Overwrite existing .env with new IP/port? (other " + "values preserved)", default=False): + dim(" leaving .env untouched") + return + text = example.read_text() + text = _replace_kv(text, 'APPLIANCE_IP', appliance_ip) + text = _replace_kv(text, 'APPLIANCE_OCF_PORT', str(appliance_port)) + env_path.write_text(text) + ok(f"wrote {env_path} — fill in MQTT_BROKER / MQTT_USER / MQTT_PASS before running main.py") + + +def _replace_kv(text, key, value): + out = [] + for line in text.splitlines(): + if line.startswith(f"{key}="): + out.append(f"{key}={value}") + else: + out.append(line) + return '\n'.join(out) + ('\n' if text.endswith('\n') else '') + + +# ---------- step 6: hub UUID validation ----------------------------------- + +def looks_like_uuid(s): + import re + return bool(re.fullmatch( + r'[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-' + r'[0-9a-fA-F]{4}-[0-9a-fA-F]{12}', s.strip())) + + +# ---------- main ---------------------------------------------------------- + +def main(): + print() + print(f"{BOLD}samsung-appliance-local — bootstrap{END}" if _tty() + else "samsung-appliance-local — bootstrap") + print(f"{DIM}This will discover your dryer, locate your CA cert, and " + f"generate the leaf used to authenticate as the SmartThings hub.{END}" + if _tty() else + "This will discover your dryer, locate your CA cert, and generate " + "the leaf used to authenticate as the SmartThings hub.") + print() + + # --- 1. dryer location --- + # Reachability is verified implicitly by the TLS handshake in step 3. + # We can't do a bare TCP probe here — that knocks the dryer's OCF + # session into a defensive state and breaks the subsequent TLS attempt. + info("Step 1 — dryer location") + appliance_ip = prompt("Dryer IP on your LAN", default=None) + appliance_port = int(prompt("OCF port (newer firmware uses 49154)", + default='49154')) + dim(f" Will connect to {appliance_ip}:{appliance_port} once we have " + f"a probe cert.") + print() + + # --- 2. AC14K_M --- + info("Step 2 — locate Samsung's AC14K_M intermediate CA") + CERTS_DIR.mkdir(parents=True, exist_ok=True) + cert_path, key_path = find_ac14km() + if cert_path is None or key_path is None: + fail(f"AC14K_M cert + key not found in {CERTS_DIR}/") + dim(f" Fetch them from: {AC14K_M_SOURCE}") + dim(f" Place as: {CERTS_DIR}/ac14k_m.pem and " + f"{CERTS_DIR}/ac14k_m.key (other common names accepted)") + return 2 + ok(f"found CA cert: {cert_path.name}") + ok(f"found CA key: {key_path.name}") + if not check_openssl(): + return 2 + print() + + # --- 3. hub UUID --- + info("Step 3 — discover your SmartThings hub UUID") + dim(" Reading /oic/sec/doxm anonymously — the dryer's baseline ACL") + dim(" allows any authenticated peer to read it (wildcard ACE).") + hub_uuid = try_anonymous_doxm_read(appliance_ip, appliance_port, + cert_path, key_path) + if hub_uuid: + ok(f"discovered hub UUID from /oic/sec/doxm: {hub_uuid}") + if not confirm("Use this UUID?", default=True): + hub_uuid = None + if not hub_uuid: + warn("Falling back to manual entry. Options B/C in the README " + "describe how to obtain it.") + while True: + hub_uuid = prompt("Hub UUID (8-4-4-4-12 hex)", default=None) + if looks_like_uuid(hub_uuid): + hub_uuid = hub_uuid.strip().lower() + break + warn("That doesn't look like a UUID. Format: " + "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx") + print() + + # --- 4. leaf --- + info("Step 4 — generate the leaf cert (mega.key + mega_chain.pem)") + mega_key = CERTS_DIR / 'mega.key' + mega_chain = CERTS_DIR / 'mega_chain.pem' + if mega_key.exists() or mega_chain.exists(): + warn(f"existing leaf cert detected in {CERTS_DIR}/") + if not confirm("Overwrite?", default=False): + dim(" leaving existing leaf in place — skipping generation") + print() + maybe_write_env(appliance_ip, appliance_port) + print() + ok("Done.") + return 0 + try: + key_out, chain_out = generate_leaf(hub_uuid, cert_path, key_path, + CERTS_DIR) + except RuntimeError as e: + fail(f"leaf cert generation failed: {e}") + return 2 + ok(f"wrote {key_out}") + ok(f"wrote {chain_out}") + print() + + # --- 5. .env --- + info("Step 5 — populate .env") + maybe_write_env(appliance_ip, appliance_port) + print() + + ok("Done. Next: edit .env to fill in MQTT_BROKER / MQTT_USER / " + "MQTT_PASS, then run main.py.") + return 0 + + +if __name__ == '__main__': + try: + sys.exit(main()) + except KeyboardInterrupt: + print(); sys.exit(130) diff --git a/deploy.sh b/deploy.sh new file mode 100755 index 0000000..de36c29 --- /dev/null +++ b/deploy.sh @@ -0,0 +1,80 @@ +#!/bin/bash +# Sync source + .env to the remote and rebuild the container. +# +# Two host paths are used: +# REMOTE_DIR — compose project (source code, .env, docker-compose.yml) +# Convention: /mnt/user/compose/samsung-bridge/ +# APPDATA_DIR — bind-mount source for /config inside the container +# (ab0b0ac4 client cert + key live here). +# Convention: /mnt/user/appdata/samsung-bridge/ +# +# The remote must already have the certs in $APPDATA_DIR. Run once +# before the first deploy: +# +# source .env +# ssh "$SSH_HOST" mkdir -p "$APPDATA_DIR" +# scp certs/ab0b0ac4_fullchain.pem certs/ab0b0ac4.key \ +# "$SSH_HOST:$APPDATA_DIR/" +# +# Subsequent deploys (this script) ship source code + .env only; the +# certs in $APPDATA_DIR are preserved. +set -e + +if [ ! -f .env ]; then + echo "Error: .env file not found. Copy .env.example to .env and configure it." + exit 1 +fi + +# Pull only the keys deploy.sh actually needs, without sourcing .env. +# Sourcing would tokenize unquoted spaces in values (e.g. +# `APPLIANCE_1_NAME=Samsung Dryer`) as shell commands. +get_env() { + grep -E "^${1}=" .env | head -1 | cut -d= -f2- +} +SSH_HOST=$(get_env SSH_HOST) +REMOTE_DIR=$(get_env REMOTE_DIR) +APPDATA_DIR=$(get_env APPDATA_DIR) + +: "${SSH_HOST:?SSH_HOST not set in .env}" +: "${REMOTE_DIR:?REMOTE_DIR not set in .env}" +: "${APPDATA_DIR:?APPDATA_DIR not set in .env}" + +echo "Deploying to ${SSH_HOST}:${REMOTE_DIR}…" +ssh "${SSH_HOST}" mkdir -p "${REMOTE_DIR}" "${APPDATA_DIR}" + +# Source code — explicit allowlist instead of an excludelist. Anything +# else in the repo (research files, certs, logs, the .git dir) stays +# local. +COPYFILE_DISABLE=1 tar cz \ + main.py \ + samsung_appliance/ \ + Dockerfile \ + docker-compose.yml \ + requirements.txt \ + deploy.sh \ + README.md \ + .env.example \ + .gitignore \ + | ssh "${SSH_HOST}" "cd ${REMOTE_DIR} && tar xz && find . -name '._*' -delete" + +# Ship .env separately and lock it down on the remote. +scp .env "${SSH_HOST}:${REMOTE_DIR}/.env" +ssh "${SSH_HOST}" "chmod 600 ${REMOTE_DIR}/.env" + +# Verify certs are present on the remote — they have to be uploaded +# once before the first build. +if ! ssh "${SSH_HOST}" "test -s ${APPDATA_DIR}/ab0b0ac4_fullchain.pem && test -s ${APPDATA_DIR}/ab0b0ac4.key"; then + echo + echo "WARNING: ${APPDATA_DIR}/ab0b0ac4_fullchain.pem and ab0b0ac4.key not" + echo "found on the remote. The container will start but fail to" + echo "connect to the appliance until you upload them, e.g.:" + echo " ssh ${SSH_HOST} mkdir -p ${APPDATA_DIR}" + echo " scp certs/ab0b0ac4_fullchain.pem certs/ab0b0ac4.key ${SSH_HOST}:${APPDATA_DIR}/" + echo +fi + +echo "Rebuilding container…" +ssh "${SSH_HOST}" "cd ${REMOTE_DIR} && docker compose up -d --build" + +echo "Done." +echo "Logs: ssh ${SSH_HOST} 'cd ${REMOTE_DIR} && docker compose logs -f'" diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..6a0b127 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,21 @@ +services: + smartthings-local: + build: . + container_name: smartthings-local + restart: unless-stopped + + # Bridge is outbound-only (DTLS UDP to each appliance, MQTT to the + # broker on 1883). No ports to expose. + + volumes: + # Holds the ab0b0ac4 client cert + key. APPDATA_DIR comes from + # .env; on Unraid this is typically + # /mnt/user/appdata/smartthings-local/. Bare-metal dev falls + # back to ./certs alongside this compose file. + - ${APPDATA_DIR:-./certs}:/config:ro + + # All runtime config is in .env. env_file passes every variable + # straight into the container, so adding a new appliance is a + # .env edit only — no compose change. + env_file: + - .env diff --git a/main.py b/main.py new file mode 100644 index 0000000..8e0700e --- /dev/null +++ b/main.py @@ -0,0 +1,203 @@ +#!/usr/bin/env python3 +"""SmartThings-Local Bridge — entry point. + +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. + +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. + +Reconnects on session errors per-appliance; shuts down cleanly on +SIGINT / SIGTERM.""" +import logging +import os +import signal +import sys +import threading + +import paho.mqtt.client as mqtt + +from samsung_appliance.appliances import get_descriptor +from samsung_appliance.bridge import PushBridge +from samsung_appliance.config import SharedConfig, load_appliances +from samsung_appliance.logger import logger + + +def main(): + shared = SharedConfig.from_env() + + try: + appliances = load_appliances() + except ValueError as e: + logger.error("config: %s", e) + return 2 + + if not shared.MQTT_BROKER: + logger.error("config: MQTT_BROKER not set") + return 2 + for path in (shared.CERT_PATH, shared.KEY_PATH): + if not path.exists(): + logger.error("client cert/key not found: %s", path) + return 2 + + # Resolve descriptors up front — bad APPLIANCE__CLASS should fail + # at startup, not 10s into the first DTLS attempt. + pairs = [] + for app in appliances: + try: + desc = get_descriptor(app.klass) + except ValueError as e: + logger.error("APPLIANCE_%d_CLASS: %s", app.index, e) + return 2 + pairs.append((app, desc)) + + logger.info("SmartThings-Local Bridge starting (%d appliance%s)", + len(appliances), '' if len(appliances) == 1 else 's') + logger.info(" broker = %s:%d (user=%s)", + shared.MQTT_BROKER, shared.MQTT_PORT, + shared.MQTT_USER or '') + for app, desc in pairs: + port = app.ocf_port or desc.default_observe_port + logger.info(" [%d] %s @ %s:%d (DTLS) → topic %s/*", + app.index, app.klass, app.ip, port, app.topic_prefix) + + # --- MQTT client (shared) --- + cli = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, + client_id='smartthings_local_bridge') + if shared.MQTT_USER: + cli.username_pw_set(shared.MQTT_USER, shared.MQTT_PASS) + # Headroom for the on_connect burst across all appliances (each + # publishes ~30 retained QoS-1 messages: discovery + state + avail). + # 100 + 50*N keeps us comfortably ahead of paho's default 20. + cli.max_inflight_messages_set(100 + 50 * len(appliances)) + + # Use the FIRST appliance's availability topic for LWT — paho only + # supports one will message. If a second appliance is added later, + # its availability is managed via explicit publishes on disconnect + # rather than LWT. For Phase 1 (dryer only) this is exact. + first_app = appliances[0] + cli.will_set(f"{first_app.topic_prefix}/availability", + payload='offline', qos=1, retain=True) + + # Build bridges. Each builds its own discovery payloads in __init__. + bridges: list[PushBridge] = [PushBridge(shared, app, desc, cli) + for app, desc in pairs] + by_prefix = {b.cmd_topic_prefix.rstrip('/'): b for b in bridges} + + def on_connect(client, userdata, flags, rc, props=None): + if rc != 0: + logger.warning("MQTT connect rc=%s", rc) + return + logger.info("MQTT connected → %s:%d", + shared.MQTT_BROKER, shared.MQTT_PORT) + for b in bridges: + for topic, payload in b.discovery_payloads: + client.publish(topic, payload, qos=1, retain=True) + cmd_wildcard = f"{b.app.topic_prefix}/cmd/#" + client.subscribe(cmd_wildcard, qos=1) + b.reassert_availability() + logger.info("subscribed to %d cmd wildcards", len(bridges)) + + def on_disconnect(client, userdata, flags, rc, props=None): + logger.warning("MQTT disconnected rc=%s", rc) + + def on_message(client, userdata, msg): + # Route by topic prefix. Each appliance owns a distinct + # `/cmd/*` namespace, so the prefix-match is unambiguous. + for prefix, bridge in by_prefix.items(): + if msg.topic.startswith(prefix + '/'): + try: + payload = msg.payload.decode('utf-8', + errors='replace').strip() + except Exception: + return + bridge.handle_command(msg.topic, payload) + return + + cli.on_connect = on_connect + cli.on_disconnect = on_disconnect + cli.on_message = on_message + + if os.getenv('PAHO_DEBUG'): + paho_logger = logging.getLogger('paho.mqtt.client') + paho_logger.setLevel(logging.DEBUG) + paho_handler = logging.StreamHandler(sys.stdout) + paho_handler.setFormatter(logging.Formatter( + '%(asctime)s PAHO %(message)s', datefmt='%H:%M:%S')) + paho_logger.addHandler(paho_handler) + paho_logger.propagate = False + cli.enable_logger(paho_logger) + + cli.connect_async(shared.MQTT_BROKER, shared.MQTT_PORT, keepalive=60) + cli.loop_start() + + # Per-bridge runner / health / heartbeat threads. + threads: list[threading.Thread] = [] + + def make_health(b: PushBridge): + def loop(): + while not b.stop.is_set(): + b.publish_health() + if b.stop.wait(shared.HEALTH_INTERVAL_S): + 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 + + 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')) + + stopping = threading.Event() + + def shutdown(*_): + if stopping.is_set(): + return + stopping.set() + logger.info("shutting down…") + for b in bridges: + b.stop.set() + try: b.set_availability(False) + except Exception: pass + + signal.signal(signal.SIGINT, shutdown) + signal.signal(signal.SIGTERM, shutdown) + + for t in threads: + t.start() + + # Wait for the session threads (the only non-daemon-equivalent + # loops). They exit when their bridge's stop event is set. + try: + for t in threads: + if t.name.endswith('-session'): + t.join() + finally: + try: + cli.loop_stop() + cli.disconnect() + except Exception: + pass + logger.info("stopped") + return 0 + + +if __name__ == '__main__': + sys.exit(main()) diff --git a/requirements-bootstrap.txt b/requirements-bootstrap.txt new file mode 100644 index 0000000..952be50 --- /dev/null +++ b/requirements-bootstrap.txt @@ -0,0 +1,4 @@ +# Setup-only deps. bootstrap.py reuses cbor2 to parse the dryer's ACL +# response and shells out to `openssl` for cert generation (so SHA-1 +# signing keeps working independent of python-cryptography's policy). +-r requirements.txt diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..0ee25c0 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,9 @@ +cbor2>=5.6 +paho-mqtt>=2.0 +# pyOpenSSL 23.x added the set_ciphertext_mtu API the DTLS handshake +# needs; older versions fragment the client cert across two UDP +# datagrams and TizenRT drops the second one. +pyOpenSSL>=23.0 +# scapy is only needed for the offline analyze_pcap.py helper; keep +# optional to avoid pulling it into the bridge container. +# scapy>=2.5 diff --git a/samsung_appliance/__init__.py b/samsung_appliance/__init__.py new file mode 100644 index 0000000..e84bf95 --- /dev/null +++ b/samsung_appliance/__init__.py @@ -0,0 +1,5 @@ +"""Samsung appliance local-API → MQTT bridge with HA discovery. + +Multi-device: dryer, oven, etc. The device class is selected at startup +via the DEVICE_CLASS env var (see samsung_appliance.appliances).""" +__version__ = "2.0.0" diff --git a/samsung_appliance/appliances/__init__.py b/samsung_appliance/appliances/__init__.py new file mode 100644 index 0000000..75a4046 --- /dev/null +++ b/samsung_appliance/appliances/__init__.py @@ -0,0 +1,32 @@ +"""Appliance descriptors registry. + +Adding a new appliance class: + 1. Write `appliances/.py` with an ApplianceDescriptor named + after the class (uppercase, e.g. OVEN). + 2. Add it to DESCRIPTORS below. + 3. Set DEVICE_CLASS= in the per-device .env. + +main.py imports get_descriptor(name) to look up the descriptor at +startup; the bridge itself stays class-agnostic. +""" +from .base import ApplianceDescriptor +from .dryer import DRYER +from .oven import OVEN + + +DESCRIPTORS: dict[str, ApplianceDescriptor] = { + DRYER.name: DRYER, + OVEN.name: OVEN, +} + + +def get_descriptor(name: str) -> ApplianceDescriptor: + try: + return DESCRIPTORS[name] + except KeyError: + raise ValueError( + f"unknown DEVICE_CLASS={name!r}; " + f"available: {sorted(DESCRIPTORS)}") from None + + +__all__ = ['ApplianceDescriptor', 'DESCRIPTORS', 'get_descriptor'] diff --git a/samsung_appliance/appliances/base.py b/samsung_appliance/appliances/base.py new file mode 100644 index 0000000..96dd667 --- /dev/null +++ b/samsung_appliance/appliances/base.py @@ -0,0 +1,103 @@ +"""ApplianceDescriptor — the per-device-class abstraction. + +The bridge is appliance-class-agnostic: it owns the DTLS session, the +OBSERVE-token bookkeeping, and MQTT publish gating. Each appliance +class (dryer, oven, …) provides a descriptor that supplies: + + * observe_paths — which CoAP resources to OBSERVE + * seed_path — the resource to fetch on connect (usually + /device/0) to populate the link dict + * flatten(links) — links → flat-dict that lands on MQTT + * build_discovery(…) — list of (HA-discovery topic, payload) + * command_handlers() — MQTT command-suffix → (path_segs, body_dict) + +Optional hooks let an appliance hold transient state across pushes: + + * on_observation(state, href, rep) — capture anchors, e.g. for + time extrapolation + * project(state, sensors) — fill in extrapolated fields + on publish + +state is a free-form dict the bridge owns and threads into both hooks. +""" +from __future__ import annotations + +import json +from dataclasses import dataclass, field +from typing import Callable, Optional + + +@dataclass +class ApplianceDescriptor: + """Static description of a Samsung appliance class. + + `name` is the value DEVICE_CLASS resolves to (e.g. 'dryer'). + `default_observe_port` documents the UDP port this firmware + exposes its DTLS-CoAP endpoint on, for the .env templates and + docs — config.APPLIANCE_OCF_PORT still wins at runtime.""" + + name: str + default_observe_port: int + + observe_paths: list[list[str]] + seed_path: list[str] + + flatten: Callable[[dict], dict] + build_discovery: Callable[[str, str, str], list[tuple[str, bytes]]] + # command_handlers() returns {topic_suffix: fn(payload_str, links_snapshot)} + # where the fn returns (path_segs, body_dict) | None. Handlers receive + # a snapshot of the bridge's link dict so they can do read-modify-write + # on resources like `/mode/vs/0` options or `/temperatures/vs/0` items. + command_handlers: Callable[[], dict[str, Callable[[str, dict], Optional[tuple]]]] + + # Optional behavioural hooks. state is a mutable dict the bridge + # threads in; the descriptor decides what keys to put in it. + on_observation: Optional[Callable[[dict, str, dict], None]] = None + project: Optional[Callable[[dict, dict], dict]] = None + + # If set, the bridge maintains a second availability topic + # `/remote_available` derived from sensors[remote_field]. + # HA entities that the appliance only honours with Remote-Control + # enabled gate themselves on it. + remote_available_field: Optional[str] = None + + # Optional log-line callback for state-change notifications. Gets + # the freshly-projected sensors dict; returns a short string. + log_state_change: Optional[Callable[[dict], str]] = None + + +# --- HA-discovery helpers ---------------------------------------------- +# Pure builder fns used by descriptor build_discovery() implementations. +# Kept here so the per-appliance modules stay focused on their entity +# inventory. + +def device_block(topic_prefix: str, device_name: str, + model: str) -> dict: + return { + 'identifiers': [topic_prefix], + 'name': device_name, + 'manufacturer': 'Samsung', + 'model': model, + } + + +def avail_base(avail_topic: str) -> list[dict]: + return [{'topic': avail_topic, + 'payload_available': 'online', + 'payload_not_available': 'offline'}] + + +def avail_with_remote(avail_topic: str, + remote_topic: str) -> list[dict]: + return [ + {'topic': avail_topic, + 'payload_available': 'online', + 'payload_not_available': 'offline'}, + {'topic': remote_topic, + 'payload_available': 'online', + 'payload_not_available': 'offline'}, + ] + + +def encode(cfg: dict) -> bytes: + return json.dumps(cfg).encode() diff --git a/samsung_appliance/appliances/dryer.py b/samsung_appliance/appliances/dryer.py new file mode 100644 index 0000000..d9d965c --- /dev/null +++ b/samsung_appliance/appliances/dryer.py @@ -0,0 +1,450 @@ +"""Dryer descriptor. + +Lifts the dryer-specific OBSERVE paths, sensor flattening, HA discovery +inventory, and MQTT command handlers out of the original +samsung_dryer/{bridge,sensors,discovery}.py modules into one place. +""" +import time + +from .base import ( + ApplianceDescriptor, + avail_base, + avail_with_remote, + device_block, + encode, +) + + +# --- OBSERVE paths ----------------------------------------------------- +# Only Samsung's `//vs/0` siblings actually push notifications; the +# OCF-standard `//0` paths accept registration silently but never +# fire. flatten() derives the OCF-shaped values from the live /vs/0 +# strings. +OBSERVE_PATHS = [ + ['operational', 'state', 'vs', '0'], # state, remainingTime, progress + ['power', 'vs', '0'], # power on/off + ['kidslock', 'vs', '0'], # child lock + ['remotectrl', 'vs', '0'], # remote control enabled + ['energy', 'consumption', 'vs', '0'], + ['course', 'vs', '0'], + ['washer', 'vs', '0'], # dryLevel, dryTime, type + ['diagnosis', 'vs', '0'], + ['alarms', 'vs', '0'], + ['st', 'dryercourse', 'vs', '0'], + ['wm', 'jobbeginingstatus', 'vs', '0'], +] + + +# --- Course table ------------------------------------------------------ +# Captured 2026-05-29 by dialing every course on a +# DA_WM_TP2_20_COMMON_DV5000T dryer. Other Samsung dryers may report a +# different Table_NN; capture a fresh table for them with +# local-tools/course_mapper.py. +COURSE_NAMES = { + 'Table_03': { + 0x16: 'Cotton', + 0x18: 'Synthetics', + 0x19: 'Delicates', + 0x1A: 'Wool', + 0x1B: 'Bedding', + 0x1C: 'Shirts', + 0x1D: 'Towels', + 0x1E: 'Outdoor', + 0x1F: 'Mixed Load', + 0x20: 'Iron Dry', + 0x23: 'Quick Dry 35', + 0x24: 'Cool Air', + 0x25: 'Warm Air', + 0x27: 'Time Dry', + }, +} + +_COURSE_CODE_BY_NAME = { + name: code + for table_codes in COURSE_NAMES.values() + for code, name in table_codes.items() +} + + +def _decode_course(s): + """`Table_03_Course_16` → `Cotton`. Pass through verbatim if the + table or code isn't in our lookup.""" + if not isinstance(s, str) or '_Course_' not in s: + return s + table_part, _, code_str = s.partition('_Course_') + table = COURSE_NAMES.get(table_part) + if not table: + return s + try: + code = int(code_str, 16) + except ValueError: + return s + return table.get(code, s) + + +def _encode_course(name): + """`Cotton` → `Course_16`. Returns None for unknown names so the + caller refuses rather than POST garbage.""" + code = _COURSE_CODE_BY_NAME.get(name) + if code is None: + return None + return f"Course_{code:02X}" + + +def _course_options(): + """Stable-sorted human course names for the HA select dropdown.""" + return sorted(_COURSE_CODE_BY_NAME.keys()) + + +# --- Samsung-state → OCF currentMachineState --------------------------- +_SAMSUNG_STATE_TO_OCF = { + 'Ready': 'idle', + 'Run': 'active', + 'Running': 'active', + 'Pause': 'pause', + 'Paused': 'pause', + 'End': 'idle', +} + + +def _num(v): + try: + return float(v) + except (TypeError, ValueError): + return None + + +def _int(v): + try: + return int(v) + except (TypeError, ValueError): + return None + + +# --- flatten ----------------------------------------------------------- +def flatten(links): + """Map a /device/0 link dict to the flat sensor dict that's + published to MQTT. Every field reads from `//vs/0` paths so push + updates immediately drive every entity.""" + g = lambda href, k, default=None: (links.get(href) or {}).get(k, default) + + inst_w = _num(g('/energy/consumption/vs/0', + 'x.com.samsung.da.instantaneousPower')) + cum_wh = _num(g('/energy/consumption/vs/0', + 'x.com.samsung.da.cumulativePower')) + if inst_w is not None and inst_w < 0: + # The dryer reports a phantom -500W when idle; HA energy + # dashboard hates negatives. + inst_w = 0.0 + + sam_state = g('/operational/state/vs/0', 'x.com.samsung.da.state') + machine_state = (_SAMSUNG_STATE_TO_OCF.get(sam_state, sam_state) + if sam_state is not None + else g('/operational/state/0', 'currentMachineState')) + + progress = g('/operational/state/vs/0', 'x.com.samsung.da.progress') + job_state = progress or g('/operational/state/0', 'currentJobState') + # HA's value_template treats the literal "None" as null (renders as + # "Unknown"). Substitute something we can render verbatim. + if job_state in (None, 'None'): + job_state = 'Idle' + if progress in (None, 'None'): + progress = 'Idle' + + remaining = (g('/operational/state/vs/0', + 'x.com.samsung.da.remainingTime') + or g('/operational/state/0', 'remainingTime')) + rem_min = None + if remaining: + try: + h, m, s = remaining.split(':') + rem_min = int(h) * 60 + int(m) + (1 if int(s) > 0 else 0) + except Exception: + pass + + sam_power = g('/power/vs/0', 'x.com.samsung.da.power') + sam_kids = g('/kidslock/vs/0', 'x.com.samsung.da.kidsLock') + sam_rc = g('/remotectrl/vs/0', + 'x.com.samsung.da.remoteControlEnabled') + power_bin = (sam_power == 'On') if sam_power is not None else None + kids_bin = (sam_kids != 'Ready') if sam_kids is not None else None + rc_bin = (str(sam_rc).lower() == 'true') if sam_rc is not None else None + + return { + 'machine_state': machine_state, + 'job_state': job_state, + 'progress': progress, + 'progress_percentage': _int(g('/operational/state/vs/0', + 'x.com.samsung.da.progressPercentage') + or g('/operational/state/0', + 'progressPercentage')), + 'completion_time': remaining, + 'completion_minutes': rem_min, + 'delay_end_time': g('/operational/state/vs/0', + 'x.com.samsung.da.delayEndTime'), + 'power_state': sam_power, + 'power_state_binary': power_bin, + 'child_lock': sam_kids, + 'child_lock_binary': kids_bin, + 'remote_control': sam_rc, + 'remote_control_binary': rc_bin, + 'power_watts': inst_w, + 'energy_kwh': round(cum_wh / 1000.0, 2) + if cum_wh is not None else None, + 'energy_wh_cumulative': int(cum_wh) if cum_wh is not None else None, + 'dryer_mode': _decode_course( + g('/st/dryercourse/vs/0', + 'x.com.samsung.da.st.dryerMode')), + 'dry_level': _int(g('/washer/vs/0', + 'x.com.samsung.da.dryLevel')), + 'dry_time': g('/washer/vs/0', + 'x.com.samsung.da.dryTime'), + 'dryer_type': g('/washer/vs/0', + 'x.com.samsung.da.dryerType'), + 'wrinkle_prevent': g('/washer/vs/0', + 'x.com.samsung.da.wrinklePrevent'), + 'diagnosis': g('/diagnosis/vs/0', + 'x.com.samsung.da.diagnosisStart'), + 'country_code': g('/configuration/vs/0', + 'x.com.samsung.da.countryCode'), + } + + +# --- Remaining-time anchor + extrapolation ---------------------------- +# The dryer pushes /operational/state/vs/0 on state transitions but not +# on remainingTime ticks. Anchor = (timestamp, total_seconds) at last +# push; project() extrapolates downward while machine_state == 'active'. + +def on_observation(state, href, rep): + if href != '/operational/state/vs/0': + return + rem = rep.get('x.com.samsung.da.remainingTime') + if not isinstance(rem, str): + return + try: + h, m, s = rem.split(':') + state['remaining_anchor'] = (time.time(), + int(h) * 3600 + int(m) * 60 + int(s)) + except (ValueError, AttributeError): + pass + + +def project(state, sensors): + anchor = state.get('remaining_anchor') + if sensors.get('machine_state') != 'active' or anchor is None: + return sensors + ts, total = anchor + remaining = max(0, int(total - (time.time() - ts))) + h, rest = divmod(remaining, 3600) + m, s = divmod(rest, 60) + sensors = dict(sensors) + sensors['completion_time'] = f"{h}:{m:02d}:{s:02d}" + sensors['completion_minutes'] = h * 60 + m + (1 if s > 0 else 0) + return sensors + + +# --- Log-line ---------------------------------------------------------- +def log_state_change(sensors): + return (f"machine={sensors.get('machine_state')} " + f"power={sensors.get('power_watts')}W " + f"energy={sensors.get('energy_kwh')}kWh") + + +# --- HA discovery ------------------------------------------------------ +MODEL = 'OCF dryer (TizenRT-iotivity)' + +# (key, friendly name, extra-config-dict) +_SENSORS = [ + ('machine_state', 'Machine state', {'icon': 'mdi:tumble-dryer'}), + ('job_state', 'Job state', {}), + ('progress', 'Progress', {}), + ('progress_percentage', 'Progress percent', + {'unit_of_measurement': '%', 'state_class': 'measurement'}), + ('completion_time', 'Completion time', {'icon': 'mdi:timer-sand'}), + ('completion_minutes', 'Remaining minutes', + {'unit_of_measurement': 'min', 'device_class': 'duration', + 'state_class': 'measurement'}), + ('delay_end_time', 'Delay end time', {'icon': 'mdi:timer'}), + ('power_state', 'Power state', {}), + ('power_watts', 'Power', + {'unit_of_measurement': 'W', 'device_class': 'power', + 'state_class': 'measurement'}), + ('energy_kwh', 'Energy', + {'unit_of_measurement': 'kWh', 'device_class': 'energy', + 'state_class': 'total_increasing'}), + ('dryer_mode', 'Dryer mode', {}), + ('dry_level', 'Dry level', {}), + ('dry_time', 'Dry time', {}), + ('dryer_type', 'Dryer type', {}), + ('wrinkle_prevent', 'Wrinkle prevent', {}), + ('diagnosis', 'Diagnosis', {}), + ('country_code', 'Country code', {}), +] + +# (key, friendly name, value_template, device_class) +_BINARY_SENSORS = [ + ('running', 'Running', + "{{ 'ON' if value_json.machine_state == 'active' else 'OFF' }}", + 'running'), + ('power_switch', 'Power switch', + "{{ 'ON' if value_json.power_state_binary else 'OFF' }}", + 'power'), + ('child_lock_active', 'Child lock', + "{{ 'ON' if value_json.child_lock_binary else 'OFF' }}", + 'lock'), + ('remote_control_enabled', 'Remote control', + "{{ 'ON' if value_json.remote_control_binary else 'OFF' }}", + 'connectivity'), +] + +# MQTT command-topic suffixes. The bridge subscribes to /cmd/# +# and dispatches by suffix. +CMD_WRINKLE_PREVENT = 'cmd/wrinkle_prevent' +CMD_OPERATIONAL = 'cmd/operational_state' +CMD_DRYER_MODE = 'cmd/dryer_mode' + + +def build_discovery(topic_prefix, ha_prefix, device_name): + """Return list of (discovery_topic, payload_bytes) tuples ready to + publish (retained) on MQTT connect.""" + state_topic = f"{topic_prefix}/state" + avail_topic = f"{topic_prefix}/availability" + remote_topic = f"{topic_prefix}/remote_available" + dev = device_block(topic_prefix, device_name, MODEL) + out = [] + + # read-only sensors + for key, name, extra in _SENSORS: + cfg = { + 'name': name, + 'unique_id': f"{topic_prefix}_{key}", + 'object_id': f"{topic_prefix}_{key}", + 'state_topic': state_topic, + 'value_template': f"{{{{ value_json.{key} }}}}", + 'availability': avail_base(avail_topic), + 'device': dev, + } + cfg.update(extra) + out.append((f"{ha_prefix}/sensor/{topic_prefix}/{key}/config", + encode(cfg))) + + for key, name, template, dclass in _BINARY_SENSORS: + cfg = { + 'name': name, + 'unique_id': f"{topic_prefix}_{key}", + 'object_id': f"{topic_prefix}_{key}", + 'state_topic': state_topic, + 'value_template': template, + 'payload_on': 'ON', + 'payload_off': 'OFF', + 'device_class': dclass, + 'availability': avail_base(avail_topic), + 'device': dev, + } + out.append((f"{ha_prefix}/binary_sensor/{topic_prefix}/{key}/config", + encode(cfg))) + + # switch: wrinkle prevent (always available) + cfg = { + 'name': 'Wrinkle prevent', + 'unique_id': f"{topic_prefix}_wrinkle_prevent_switch", + 'object_id': f"{topic_prefix}_wrinkle_prevent_switch", + 'state_topic': state_topic, + 'value_template': '{{ value_json.wrinkle_prevent }}', + 'state_on': 'On', + 'state_off': 'Off', + 'command_topic': f"{topic_prefix}/{CMD_WRINKLE_PREVENT}", + 'payload_on': 'On', + 'payload_off': 'Off', + 'icon': 'mdi:iron', + 'availability': avail_base(avail_topic), + 'device': dev, + } + out.append((f"{ha_prefix}/switch/{topic_prefix}/wrinkle_prevent/config", + encode(cfg))) + + # buttons: Start / Pause / Stop (gated on remote control) + buttons = [ + ('start', 'Start cycle', 'Run', 'mdi:play'), + ('pause', 'Pause cycle', 'Pause', 'mdi:pause'), + ('stop', 'Stop cycle', 'Ready', 'mdi:stop'), + ] + for key, name, payload_press, icon in buttons: + cfg = { + 'name': name, + 'unique_id': f"{topic_prefix}_{key}", + 'object_id': f"{topic_prefix}_{key}", + 'command_topic': f"{topic_prefix}/{CMD_OPERATIONAL}", + 'payload_press': payload_press, + 'icon': icon, + 'availability': avail_with_remote(avail_topic, remote_topic), + 'availability_mode': 'all', + 'device': dev, + } + out.append((f"{ha_prefix}/button/{topic_prefix}/{key}/config", + encode(cfg))) + + # select: course (gated on remote control) + cfg = { + 'name': 'Course', + 'unique_id': f"{topic_prefix}_course_select", + 'object_id': f"{topic_prefix}_course_select", + 'state_topic': state_topic, + 'value_template': '{{ value_json.dryer_mode }}', + 'command_topic': f"{topic_prefix}/{CMD_DRYER_MODE}", + 'options': _course_options(), + 'icon': 'mdi:tumble-dryer', + 'availability': avail_with_remote(avail_topic, remote_topic), + 'availability_mode': 'all', + 'device': dev, + } + out.append((f"{ha_prefix}/select/{topic_prefix}/course/config", + encode(cfg))) + + return out + + +# --- MQTT command handlers -------------------------------------------- +def command_handlers(): + """topic_suffix → fn(payload, links) → (path_segs, body_dict) | None. + + `None` means refuse the command (caller logs & drops). Dryer + handlers don't need the links snapshot — they're all single-field + writes.""" + def _wrinkle(p, _links): + if p not in ('On', 'Off'): + return None + return ['washer', 'vs', '0'], {'x.com.samsung.da.wrinklePrevent': p} + + def _operational(p, _links): + if p not in ('Run', 'Pause', 'Ready'): + return None + return ['operational', 'state', 'vs', '0'], {'x.com.samsung.da.state': p} + + def _course(p, _links): + code = _encode_course(p) + if code is None: + return None + return ['st', 'dryercourse', 'vs', '0'], {'x.com.samsung.da.st.dryerMode': code} + + return { + CMD_WRINKLE_PREVENT: _wrinkle, + CMD_OPERATIONAL: _operational, + CMD_DRYER_MODE: _course, + } + + +# --- Descriptor -------------------------------------------------------- +DRYER = ApplianceDescriptor( + name='dryer', + default_observe_port=49155, + observe_paths=OBSERVE_PATHS, + seed_path=['device', '0'], + flatten=flatten, + build_discovery=build_discovery, + command_handlers=command_handlers, + on_observation=on_observation, + project=project, + remote_available_field='remote_control_binary', + log_state_change=log_state_change, +) diff --git a/samsung_appliance/appliances/oven.py b/samsung_appliance/appliances/oven.py new file mode 100644 index 0000000..5de6ccd --- /dev/null +++ b/samsung_appliance/appliances/oven.py @@ -0,0 +1,665 @@ +"""Oven descriptor (Samsung NV7000BS-class). + +Resource map captured 2026-05-31 via DTLS-CoAP with the ab0b0ac4 cert. +See `local-tools/comparisons/oven-tree.md` for the full field reference. + +Write surfaces this descriptor exposes: + + proven: + * UpperLamp via /mode/vs/0 options RMW (probe_oven_lamp_toggle.py) + — works even with Remote Control off. + + unproven (first HA use is also the test): + * Sound, FastPreheat — same RMW pattern as lamp. + * Setpoint via /temperatures/vs/0 items RMW. Mid-cook write may + or may not retune the element (plan §K-U #2). + * Mode select via /mode/vs/0 .modes — mid-cook acceptance unknown + (plan §K-U #3). + * Power on/off via /power/vs/0. + * Stop via /operational/state/vs/0 (dryer convention; oven may + use a different state value). + +Untested writes are gated behind /remote_available so HA +disables them in the UI when the oven's Remote Control switch is off.""" +import time + +from .base import ( + ApplianceDescriptor, + avail_base, + avail_with_remote, + device_block, + encode, +) + + +# --------------------------------------------------------------------- +# OBSERVE paths — every push-eligible //vs/0 resource on the oven. +# Same wedge-safety story as the dryer: only `//vs/0` siblings push; +# OCF-standard `//0` paths accept registration but never fire. +# Security paths (/oic/sec/{doxm,pstat,acl,cred}) are deliberately +# EXCLUDED — those are the surfaces that nearly bricked the oven in +# prior sessions. The bridge has no reason to touch them. +# --------------------------------------------------------------------- +OBSERVE_PATHS = [ + ['operational', 'state', 'vs', '0'], # state, time, progress + ['power', 'vs', '0'], # power On/Off + ['oven', 'vs', '0'], # cavity state (Cooking, Idle, …) + ['temperatures','vs', '0'], # current + desired temp + ['doors', 'vs', '0'], # openState + ['kidslock', 'vs', '0'], # child lock + ['remotectrl', 'vs', '0'], # remote control enabled + ['mode', 'vs', '0'], # cooking mode + options array + ['alarms', 'vs', '0'], # alarm code (OV_E_OFF etc.) + ['connected', 'vs', '0'], # cloud connectivity status + ['otninformation', 'vs', '0'], # firmware-update flags +] + + +# --------------------------------------------------------------------- +# Mode dropdown — supportedModes minus the explicitly NotSupported one. +# Order matches the oven's own supportedModes list so the dropdown +# matches the device UI's order. +# --------------------------------------------------------------------- +SUPPORTED_MODES = [ + 'Autocook', + 'Convection', + 'TopHeatPluseConvection', + 'Conventional', + 'LargeGrill', + 'SmallGrill', + 'BottomHeatPluseConvection', + 'PlateWarm', + 'KeepWarm', + 'Bottom', + 'EcoConvection', + 'FanGrill', + 'Defrost', + # 'SteamClean', # control=NotSupported per modeSpec — exclude. +] + + +# Setpoint bounds — union across modeSpec entries on this oven. Per-mode +# bounds (e.g. PlateWarm 30–80) tighten this; the firmware will refuse +# out-of-range writes for the active mode and the HA UI will surface +# the resulting 4.xx in the bridge log. +SETPOINT_MIN_C = 30 +SETPOINT_MAX_C = 270 +SETPOINT_STEP_C = 5 + + +# Samsung's operational state strings → OCF currentMachineState shape. +_SAMSUNG_STATE_TO_OCF = { + 'Ready': 'idle', + 'Run': 'active', + 'Running': 'active', + 'Pause': 'pause', + 'Paused': 'pause', + 'End': 'idle', + 'Stop': 'idle', +} + + +def _num(v): + try: return float(v) + except (TypeError, ValueError): return None + + +def _int(v): + try: return int(v) + except (TypeError, ValueError): return None + + +def _option_value(options, prefix, default=None): + """Find `_` in an options array and return .""" + for o in options: + if o.startswith(prefix + '_'): + return o.split('_', 1)[1] + return default + + +def _replace_in_options(options, prefix, new_value): + """Return a new options array with any `_*` entry replaced + by `_`. Caller must verify `options` is the live + options array first (Samsung uses replace-not-merge on this field).""" + return [f"{prefix}_{new_value}" if o.startswith(prefix + '_') else o + for o in options] + + +def _fmt_hms(seconds): + """Format an integer second count as `H:MM:SS`. Returns None on + bad input so callers can leave the field null rather than emitting + a misleading `0:00:00`.""" + try: + s = int(seconds) + except (TypeError, ValueError): + return None + if s < 0: + s = 0 + h, rest = divmod(s, 3600) + m, sec = divmod(rest, 60) + return f"{h}:{m:02d}:{sec:02d}" + + +# --------------------------------------------------------------------- +# flatten — Samsung /device/0 links → HA-flavoured sensor dict. +# Every field reads from `//vs/0` paths so push updates immediately +# drive every entity. Where a field is settable (lamp, mode, setpoint), +# we publish it as a read-side sensor here AND as a writeable entity +# in build_discovery; the read side closes the HA UI feedback loop. +# --------------------------------------------------------------------- +def flatten(links): + g = lambda href, k, default=None: (links.get(href) or {}).get(k, default) + + # Operational + sam_state = g('/operational/state/vs/0', 'x.com.samsung.da.state') + machine_state = (_SAMSUNG_STATE_TO_OCF.get(sam_state, sam_state) + if sam_state is not None else None) + + operation_time = g('/operational/state/vs/0', + 'x.com.samsung.da.operationTime') + remaining = g('/operational/state/vs/0', + 'x.com.samsung.da.remainingTime') + rem_min = None + if remaining: + try: + h, m, s = remaining.split(':') + rem_min = int(h) * 60 + int(m) + (1 if int(s) > 0 else 0) + except Exception: + pass + + # Cavity state — Cooking, Idle, Preheating, … + oven_state = g('/oven/vs/0', 'x.com.samsung.da.state') + + # Temperatures + temps_items = (g('/temperatures/vs/0', + 'x.com.samsung.da.items') or []) + cur_c = des_c = None + if temps_items: + cur_c = _int(temps_items[0].get('x.com.samsung.da.current')) + des_c = _int(temps_items[0].get('x.com.samsung.da.desired')) + + # Door + doors_items = g('/doors/vs/0', 'x.com.samsung.da.items') or [] + door = doors_items[0].get('x.com.samsung.da.openState') if doors_items else None + door_open = (door == 'Open') if door is not None else None + + # Power + sam_power = g('/power/vs/0', 'x.com.samsung.da.power') + power_bin = (sam_power == 'On') if sam_power is not None else None + + # Kidslock + Remote + sam_kids = g('/kidslock/vs/0', 'x.com.samsung.da.kidsLock') + kids_bin = (sam_kids != 'Ready') if sam_kids is not None else None + sam_rc = g('/remotectrl/vs/0', + 'x.com.samsung.da.remoteControlEnabled') + rc_bin = (str(sam_rc).lower() == 'true') if sam_rc is not None else None + + # Mode + options + modes = g('/mode/vs/0', 'x.com.samsung.da.modes') or [] + current_mode = modes[0] if modes else None + options = g('/mode/vs/0', 'x.com.samsung.da.options') or [] + lamp = _option_value(options, 'UpperLamp') # 'On' / 'Off' + sound = _option_value(options, 'Sound') # 'On' / 'Off' + fastpreheat = _option_value(options, 'fastpreheat') # 'On' / 'Off' + timer_state = _option_value(options, 'UpperTimerState') # 'Ready' / 'Running' + # UpperTimerCurrent/UpperTimerSet are integer seconds. Format as + # H:MM:SS for HA display so users see "1:10:00", not "4200". + timer_current_raw = _option_value(options, 'UpperTimerCurrent') + timer_set_raw = _option_value(options, 'UpperTimerSet') + timer_current = _fmt_hms(timer_current_raw) + timer_set = _fmt_hms(timer_set_raw) + timer_current_seconds = _int(timer_current_raw) + timer_set_seconds = _int(timer_set_raw) + + # Alarms + alarm_items = g('/alarms/vs/0', 'x.com.samsung.da.items') or [] + alarm_code = (alarm_items[0].get('x.com.samsung.da.code') + if alarm_items else None) + alarm_time = (alarm_items[0].get('x.com.samsung.da.triggeredTime') + if alarm_items else None) + # OV_E_OFF appears when the oven is off / no alarm; treat as inactive. + alarm_active = bool(alarm_code) and alarm_code != 'OV_E_OFF' + + # Connectivity / firmware + sam_connected = g('/connected/vs/0', 'x.com.samsung.da.connected') + connected_bin = (sam_connected == 'On') if sam_connected is not None else None + fw_update_available = g('/otninformation/vs/0', + 'x.com.samsung.da.newVersionAvailable') + fw_update_bin = (str(fw_update_available).lower() == 'true' + if fw_update_available is not None else None) + + return { + 'machine_state': machine_state, + 'oven_state': oven_state, + 'progress_percentage': _int(g('/operational/state/vs/0', + 'x.com.samsung.da.progressPercentage')), + 'operation_time': operation_time, + 'completion_time': remaining, + 'completion_minutes': rem_min, + 'current_temp_c': cur_c, + 'target_temp_c': des_c, + 'door': door, + 'door_open': door_open, + 'power_state': sam_power, + 'power_state_binary': power_bin, + 'child_lock': sam_kids, + 'child_lock_binary': kids_bin, + 'remote_control': sam_rc, + 'remote_control_binary': rc_bin, + 'mode': current_mode, + 'lamp': lamp, + 'sound': sound, + 'fastpreheat': fastpreheat, + 'timer_state': timer_state, + 'timer_current': timer_current, + 'timer_set': timer_set, + 'timer_current_seconds': timer_current_seconds, + 'timer_set_seconds': timer_set_seconds, + 'alarm_code': alarm_code, + 'alarm_time': alarm_time, + 'alarm_active': alarm_active, + 'connected': sam_connected, + 'connected_binary': connected_bin, + 'firmware_update_available': fw_update_bin, + } + + +# --------------------------------------------------------------------- +# Remaining-time anchor + projection. The oven pushes /operational/state +# on state transitions but probably not on remainingTime ticks (matches +# dryer behaviour). Capture (ts, total_seconds) at each push and +# extrapolate downward while machine_state == active. +# --------------------------------------------------------------------- +def on_observation(state, href, rep): + if href != '/operational/state/vs/0': + return + rem = rep.get('x.com.samsung.da.remainingTime') + if not isinstance(rem, str): + return + try: + h, m, s = rem.split(':') + state['remaining_anchor'] = (time.time(), + int(h) * 3600 + int(m) * 60 + int(s)) + except (ValueError, AttributeError): + pass + + +def project(state, sensors): + anchor = state.get('remaining_anchor') + if sensors.get('machine_state') != 'active' or anchor is None: + return sensors + ts, total = anchor + remaining = max(0, int(total - (time.time() - ts))) + h, rest = divmod(remaining, 3600) + m, s = divmod(rest, 60) + sensors = dict(sensors) + sensors['completion_time'] = f"{h}:{m:02d}:{s:02d}" + sensors['completion_minutes'] = h * 60 + m + (1 if s > 0 else 0) + return sensors + + +def log_state_change(sensors): + return (f"machine={sensors.get('machine_state')} " + f"oven={sensors.get('oven_state')} " + f"temp={sensors.get('current_temp_c')}/" + f"{sensors.get('target_temp_c')}°C " + f"mode={sensors.get('mode')}") + + +# --------------------------------------------------------------------- +# HA discovery inventory +# --------------------------------------------------------------------- +MODEL = 'OCF oven (TizenRT-iotivity, NV7000BS-class)' + +# (key, friendly name, extra config) +# +# Only read-only sensors live here. Fields that ALSO have an +# interactive entity (light, switch, number, select) are removed — +# the interactive entity already surfaces the live state, so a +# duplicate read-only "Lamp state" / "Fast preheat state" / etc. +# sensor would just clutter the device card with the same value +# twice. +_SENSORS = [ + ('machine_state', 'Machine state', {'icon': 'mdi:stove'}), + ('oven_state', 'Cavity state', {}), + ('progress_percentage', 'Progress percent', + {'unit_of_measurement': '%', 'state_class': 'measurement'}), + ('operation_time', 'Elapsed time', {'icon': 'mdi:timer'}), + ('completion_time', 'Completion time', {'icon': 'mdi:timer-sand'}), + ('completion_minutes', 'Remaining minutes', + {'unit_of_measurement': 'min', 'device_class': 'duration', + 'state_class': 'measurement'}), + ('current_temp_c', 'Temperature', + {'unit_of_measurement': '°C', 'device_class': 'temperature', + 'state_class': 'measurement'}), + # target_temp_c is also exposed as a Number entity for editing, + # but the Number is RC-gated. The sensor stays always-visible so + # the user can see the current setpoint even with Remote Control + # off at the oven. + ('target_temp_c', 'Setpoint', + {'unit_of_measurement': '°C', 'device_class': 'temperature', + 'state_class': 'measurement', 'icon': 'mdi:thermometer-chevron-up'}), + # power_state: read-only. The oven doesn't expose a meaningful + # POST /power/vs/0 from cold — turning the unit on at the panel + # is a physical action — so we don't ship a Power switch entity. + ('power_state', 'Power state', {'icon': 'mdi:power'}), + ('door', 'Door state', {}), + ('child_lock', 'Child lock state', {}), + ('remote_control', 'Remote control state', {}), + ('timer_state', 'Timer state', {}), + ('timer_current', 'Timer remaining', {'icon': 'mdi:timer-sand'}), + ('timer_set', 'Timer set', {'icon': 'mdi:timer'}), + ('alarm_code', 'Alarm code', + {'icon': 'mdi:alert', 'entity_category': 'diagnostic'}), + ('alarm_time', 'Alarm time', + {'icon': 'mdi:clock-alert', 'entity_category': 'diagnostic'}), + ('connected', 'Cloud connectivity', + {'entity_category': 'diagnostic'}), +] + +# (key, friendly, value_template, device_class, extras) +_BINARY_SENSORS = [ + ('running', 'Running', + "{{ 'ON' if value_json.machine_state == 'active' else 'OFF' }}", + 'running', {}), + ('door_open', 'Door', + "{{ 'ON' if value_json.door_open else 'OFF' }}", + 'door', {}), + # `power_switch` binary_sensor would duplicate the Power switch + # entity below; the switch already shows on/off state. + ('child_lock_active', 'Child lock', + "{{ 'ON' if value_json.child_lock_binary else 'OFF' }}", + 'lock', {}), + ('remote_control_enabled', 'Remote control', + "{{ 'ON' if value_json.remote_control_binary else 'OFF' }}", + 'connectivity', {}), + ('alarm_active', 'Alarm active', + "{{ 'ON' if value_json.alarm_active else 'OFF' }}", + 'problem', {}), + ('connected_bin', 'Connected', + "{{ 'ON' if value_json.connected_binary else 'OFF' }}", + 'connectivity', {'entity_category': 'diagnostic'}), + ('firmware_update_available', 'Firmware update available', + "{{ 'ON' if value_json.firmware_update_available else 'OFF' }}", + 'update', {'entity_category': 'diagnostic'}), +] + + +# MQTT command-topic suffixes (under /cmd/…) +CMD_LAMP = 'cmd/lamp' +CMD_SOUND = 'cmd/sound' +CMD_FASTPREHEAT = 'cmd/fastpreheat' +CMD_POWER = 'cmd/power' +CMD_STOP = 'cmd/stop' +CMD_MODE = 'cmd/mode' +CMD_SETPOINT = 'cmd/setpoint' + + +def build_discovery(topic_prefix, ha_prefix, device_name): + state_topic = f"{topic_prefix}/state" + avail_topic = f"{topic_prefix}/availability" + remote_topic = f"{topic_prefix}/remote_available" + dev = device_block(topic_prefix, device_name, MODEL) + out = [] + + # --- read-only sensors ------------------------------------------- + for key, name, extra in _SENSORS: + cfg = { + 'name': name, + 'unique_id': f"{topic_prefix}_{key}", + 'object_id': f"{topic_prefix}_{key}", + 'state_topic': state_topic, + 'value_template': f"{{{{ value_json.{key} }}}}", + 'availability': avail_base(avail_topic), + 'device': dev, + } + cfg.update(extra) + out.append((f"{ha_prefix}/sensor/{topic_prefix}/{key}/config", + encode(cfg))) + + for key, name, template, dclass, extra in _BINARY_SENSORS: + cfg = { + 'name': name, + 'unique_id': f"{topic_prefix}_{key}", + 'object_id': f"{topic_prefix}_{key}", + 'state_topic': state_topic, + 'value_template': template, + 'payload_on': 'ON', + 'payload_off': 'OFF', + 'device_class': dclass, + 'availability': avail_base(avail_topic), + 'device': dev, + } + cfg.update(extra) + out.append((f"{ha_prefix}/binary_sensor/{topic_prefix}/{key}/config", + encode(cfg))) + + # --- light: oven lamp (proven via probe_oven_lamp_states.py; + # binary On/Off only — High/Low/Dim coerce back to previous + # state. Works regardless of Remote Control switch, so we only + # gate on base availability). For the MQTT light default schema, + # state_value_template's output must match payload_on/payload_off + # exactly (case-sensitive) for HA to recognise the state. + cfg = { + 'name': 'Lamp', + 'unique_id': f"{topic_prefix}_lamp_light", + 'object_id': f"{topic_prefix}_lamp_light", + 'state_topic': state_topic, + 'state_value_template': "{{ value_json.lamp }}", + 'command_topic': f"{topic_prefix}/{CMD_LAMP}", + 'payload_on': 'On', + 'payload_off': 'Off', + 'icon': 'mdi:track-light', + 'availability': avail_base(avail_topic), + 'device': dev, + } + out.append((f"{ha_prefix}/light/{topic_prefix}/lamp/config", + encode(cfg))) + + # --- switches (RC-gated; untested mid-cook). Sound + fastpreheat + # are options-array writes (same RMW path as lamp); Power is a + # /power/vs/0 single-field write. -------------------------------- + untested_switches = [ + ('sound', 'Sound', '{{ value_json.sound }}', CMD_SOUND, 'mdi:volume-high'), + ('fastpreheat', 'Fast preheat', '{{ value_json.fastpreheat }}', CMD_FASTPREHEAT, 'mdi:fire'), + # Power deliberately omitted: turning the oven on at the + # cold-start panel is a physical action; the read-only + # power_state sensor (above) reflects its state. + ] + for key, name, tpl, cmd, icon in untested_switches: + cfg = { + 'name': name, + 'unique_id': f"{topic_prefix}_{key}_switch", + 'object_id': f"{topic_prefix}_{key}_switch", + 'state_topic': state_topic, + 'value_template': tpl, + 'state_on': 'On', + 'state_off': 'Off', + 'command_topic': f"{topic_prefix}/{cmd}", + 'payload_on': 'On', + 'payload_off': 'Off', + 'icon': icon, + 'availability': avail_with_remote(avail_topic, remote_topic), + 'availability_mode': 'all', + 'device': dev, + } + out.append((f"{ha_prefix}/switch/{topic_prefix}/{key}/config", + encode(cfg))) + + # --- number: setpoint (RC-gated, slider input) ------------------ + cfg = { + 'name': 'Setpoint', + 'unique_id': f"{topic_prefix}_setpoint", + 'object_id': f"{topic_prefix}_setpoint", + 'state_topic': state_topic, + 'value_template': '{{ value_json.target_temp_c }}', + 'command_topic': f"{topic_prefix}/{CMD_SETPOINT}", + 'min': SETPOINT_MIN_C, + 'max': SETPOINT_MAX_C, + 'step': SETPOINT_STEP_C, + 'unit_of_measurement': '°C', + 'device_class': 'temperature', + 'mode': 'slider', + 'icon': 'mdi:thermometer-chevron-up', + 'availability': avail_with_remote(avail_topic, remote_topic), + 'availability_mode': 'all', + 'device': dev, + } + out.append((f"{ha_prefix}/number/{topic_prefix}/setpoint/config", + encode(cfg))) + + # --- select: mode (RC-gated) ------------------------------------ + cfg = { + 'name': 'Cooking mode', + 'unique_id': f"{topic_prefix}_mode_select", + 'object_id': f"{topic_prefix}_mode_select", + 'state_topic': state_topic, + 'value_template': '{{ value_json.mode }}', + 'command_topic': f"{topic_prefix}/{CMD_MODE}", + 'options': SUPPORTED_MODES, + 'icon': 'mdi:tune', + 'availability': avail_with_remote(avail_topic, remote_topic), + 'availability_mode': 'all', + 'device': dev, + } + out.append((f"{ha_prefix}/select/{topic_prefix}/mode/config", + encode(cfg))) + + # --- button: Stop cycle ---------------------------------------- + # NOT gated on remote_available: the SmartThings app stops the + # oven regardless of the Remote Control switch state, so the + # device clearly honours Stop without that gate. Only requires + # the bridge itself to be online. + cfg = { + 'name': 'Stop cycle', + 'unique_id': f"{topic_prefix}_stop", + 'object_id': f"{topic_prefix}_stop", + 'command_topic': f"{topic_prefix}/{CMD_STOP}", + 'payload_press': 'Stop', + 'icon': 'mdi:stop', + 'availability': avail_base(avail_topic), + 'device': dev, + } + out.append((f"{ha_prefix}/button/{topic_prefix}/stop/config", + encode(cfg))) + + return out + + +# --------------------------------------------------------------------- +# Command handlers — fn(payload, links) → (path_segs, body_dict) | None. +# Read-modify-write handlers (lamp/sound/fastpreheat) snapshot the +# `/mode/vs/0` options array and replace just their slot. /temperatures +# is also RMW because Samsung's write semantics on the items array are +# replace-not-merge. +# --------------------------------------------------------------------- +def _mode_options(links): + """Return the live `/mode/vs/0` options array (a copy), or None + if /mode/vs/0 isn't seeded yet.""" + rep = links.get('/mode/vs/0') or {} + opts = rep.get('x.com.samsung.da.options') + if not opts: + return None + return list(opts) + + +def _temps_items(links): + """Return a deep-ish copy of the /temperatures/vs/0 items array.""" + rep = links.get('/temperatures/vs/0') or {} + items = rep.get('x.com.samsung.da.items') or [] + return [dict(it) for it in items] if items else None + + +def command_handlers(): + def _lamp(p, links): + if p not in ('On', 'Off'): + return None + opts = _mode_options(links) + if opts is None: + return None + return ['mode', 'vs', '0'], { + 'x.com.samsung.da.options': _replace_in_options(opts, 'UpperLamp', p), + } + + def _sound(p, links): + if p not in ('On', 'Off'): + return None + opts = _mode_options(links) + if opts is None: + return None + return ['mode', 'vs', '0'], { + 'x.com.samsung.da.options': _replace_in_options(opts, 'Sound', p), + } + + def _fastpreheat(p, links): + if p not in ('On', 'Off'): + return None + opts = _mode_options(links) + if opts is None: + return None + return ['mode', 'vs', '0'], { + 'x.com.samsung.da.options': _replace_in_options( + opts, 'fastpreheat', p), + } + + def _power(p, _links): + if p not in ('On', 'Off'): + return None + return ['power', 'vs', '0'], {'x.com.samsung.da.power': p} + + def _stop(_p, _links): + # Untested for the oven. Dryer convention is state='Ready' to + # leave the cycle in idle. If this turns out to be wrong, the + # bridge will log the 4.xx but the oven won't be harmed — + # /operational/state/vs/0 isn't a wedge-trigger surface. + return ['operational', 'state', 'vs', '0'], { + 'x.com.samsung.da.state': 'Ready', + } + + def _mode(p, _links): + if p not in SUPPORTED_MODES: + return None + return ['mode', 'vs', '0'], {'x.com.samsung.da.modes': [p]} + + def _setpoint(p, links): + try: + temp = float(p) + except (TypeError, ValueError): + return None + # Snap to step and bounds. + temp_i = int(round(temp / SETPOINT_STEP_C) * SETPOINT_STEP_C) + if not (SETPOINT_MIN_C <= temp_i <= SETPOINT_MAX_C): + return None + items = _temps_items(links) + if items is None: + return None + items[0]['x.com.samsung.da.desired'] = str(temp_i) + return ['temperatures', 'vs', '0'], { + 'x.com.samsung.da.items': items, + } + + return { + CMD_LAMP: _lamp, + CMD_SOUND: _sound, + CMD_FASTPREHEAT: _fastpreheat, + CMD_POWER: _power, + CMD_STOP: _stop, + CMD_MODE: _mode, + CMD_SETPOINT: _setpoint, + } + + +# --------------------------------------------------------------------- +OVEN = ApplianceDescriptor( + name='oven', + default_observe_port=49154, + observe_paths=OBSERVE_PATHS, + seed_path=['device', '0'], + flatten=flatten, + build_discovery=build_discovery, + command_handlers=command_handlers, + on_observation=on_observation, + project=project, + remote_available_field='remote_control_binary', + log_state_change=log_state_change, +) diff --git a/samsung_appliance/bridge.py b/samsung_appliance/bridge.py new file mode 100644 index 0000000..5b54ccc --- /dev/null +++ b/samsung_appliance/bridge.py @@ -0,0 +1,502 @@ +"""Push-mode bridge: OCF CoAP-DTLS Observe → MQTT. + + Appliance ──CoAP OBSERVE notifications──► PushBridge + │ + ▼ + MQTT broker + │ + ▼ + 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. + +The bridge is appliance-class-agnostic — it delegates every +appliance-specific decision to an ApplianceDescriptor. + +Multiple PushBridges run concurrently in a single process — see +main.py. They share one MQTT client; each owns one DTLS session. +""" +import json +import threading +import time + +import cbor2 + +from .appliances.base import ApplianceDescriptor +from .coap_dtls import DtlsCoapSession, fmt_code +from .config import ApplianceConfig, SharedConfig +from .logger import bridge_logger, logger as module_logger +from .sensors import index_links + + +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, + app: ApplianceConfig, + descriptor: ApplianceDescriptor, + mqtt_client): + self.shared = shared + self.app = app + 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.last_state_pub = None + self.last_remote_pub = None + self.stop = threading.Event() + self.started_ts = time.time() + self.session_started_ts = None + self.last_change_ts = None + self.last_seed_ts = None + self.notif_count = 0 + self.connect_count = 0 + self.error_count = 0 + self._publish_gate = False + + # 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. + self._fetch_gen: dict[str, int] = {} + self._fetch_lock = threading.Lock() + + p = app.topic_prefix + self.state_topic = f"{p}/state" + self.avail_topic = f"{p}/availability" + self.remote_topic = f"{p}/remote_available" + self.health_topic = f"{p}/bridge/health" + 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) + + # ---- DTLS session helpers --------------------------------------- + + 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.""" + if not payload_bytes: + # Empty payload — almost certainly a Block2 announcement. + self._schedule_fetchback(href) + return + try: + rep = cbor2.loads(payload_bytes) + except Exception: + self._schedule_fetchback(href) + 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.""" + self.links[href] = rep + hook = self.descriptor.on_observation + if hook is not None: + try: + hook(self.descriptor_state, href, rep) + except Exception as e: + self.log.warning("on_observation %s: %s", href, e) + self.notif_count += 1 + self.last_change_ts = time.time() + self.maybe_publish_state() + + def _apply_optimistic(self, href, body): + """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 _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 + threading.Thread( + target=self._fetch_back, + args=(href, delay_s, gen), + daemon=True, + name=f'fetch{href}', + ).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. + with self._fetch_lock: + if self._fetch_gen.get(href) != gen: + return + sess = self.session + if sess is None: + return + segs = _href_to_segs(href) + try: + code, payload = sess.get(segs, timeout=15.0) + 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)) + 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) + + 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 {} + serial = info.get(SERIAL_FIELD) + if not serial: + return + self._serial = serial + self.log = bridge_logger(self.app.klass, serial) + self.log.info("identified — serial=%s", serial) + + # ---- 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, + key_path=self.shared.KEY_PATH, + on_notification=self._on_notification, + ) + sess.connect() + self.session = sess + self.session_started_ts = time.time() + self.connect_count += 1 + self.descriptor_state = {} + self._publish_gate = False + + self.log.info("DTLS connected — subscribing %d paths", + len(self.descriptor.observe_paths)) + + sess.start_reader() + + for path in self.descriptor.observe_paths: + sess.subscribe(path) + time.sleep(0.05) + + code, pl = sess.get(self.descriptor.seed_path, timeout=15.0) + if code != 0x45: + raise RuntimeError( + f"/{'/'.join(self.descriptor.seed_path)} -> {fmt_code(code)}") + try: + body = cbor2.loads(pl) + except Exception as e: + raise RuntimeError( + f"/{'/'.join(self.descriptor.seed_path)} cbor decode: {e}" + ) from e + 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() + + hook = self.descriptor.on_observation + if hook is not None: + for href, rep in self.links.items(): + try: + hook(self.descriptor_state, href, rep) + except Exception as e: + self.log.warning("seed on_observation %s: %s", href, e) + self.last_seed_ts = time.time() + + self._publish_gate = True + self.maybe_publish_state(force=True) + self.set_availability(True) + self.log.info("seeded → %d links; sensors live", len(self.links)) + + 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() + + # ---- MQTT publishing -------------------------------------------- + + def maybe_publish_state(self, force=False): + if not force and not self._publish_gate: + return + sensors = self.descriptor.flatten(self.links) + project = self.descriptor.project + if project is not None: + sensors = project(self.descriptor_state, sensors) + if not force and sensors == self.last_state_pub: + return + self.last_state_pub = sensors + self.mqtt.publish(self.state_topic, + json.dumps(sensors).encode(), + qos=1, retain=True) + field = self.descriptor.remote_available_field + if field is not None: + self.publish_remote_available(sensors.get(field)) + 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)", + extra or 'descriptor-no-log', self.notif_count) + + def publish_remote_available(self, remote_on, force=False): + value = 'online' if remote_on else 'offline' + if not force and value == self.last_remote_pub: + return + self.last_remote_pub = value + try: + self.mqtt.publish(self.remote_topic, value, qos=1, retain=True) + self.log.info("remote_available → %s", value) + except Exception as e: + self.log.warning("remote_available publish: %s", e) + + def reassert_availability(self): + if self.session is None or self.last_state_pub is None: + return + self.set_availability(True) + field = self.descriptor.remote_available_field + if field is not None: + self.publish_remote_available( + self.last_state_pub.get(field), force=True) + + def set_availability(self, online): + try: + self.mqtt.publish(self.avail_topic, + 'online' if online else 'offline', + 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: + self.last_remote_pub = None + try: + self.mqtt.publish(self.remote_topic, 'offline', + qos=1, retain=True) + except Exception: + pass + + # ---- MQTT command handling -------------------------------------- + + def handle_command(self, topic, payload): + if not topic.startswith(self.cmd_topic_prefix): + return + suffix = topic[len(self.cmd_topic_prefix) - len('cmd/'):] + handler = self.cmd_handlers.get(suffix) + 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)) + if result is None: + self.log.warning("rejected command %s payload=%r", + topic, payload) + return + path_segs, body = result + sess = self.session + if sess is None: + self.log.warning("command %s: no DTLS session", topic) + return + 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)) + # Defensive re-read of the just-written resource. The dryer + # pushes an OBSERVE notification within ~100ms of a 2.xx write + # and we'd see the new state anyway, but the oven doesn't + # push on options-only writes (lamp/sound/fastpreheat). Without + # this fetchback, HA would only see the new state on the next + # heartbeat (10 min by default). + if code >> 5 == 2: + href = '/' + '/'.join(path_segs) + # Optimistic publish: apply the write to our local state + # and republish the MQTT state immediately. HA sees the + # new value with no UI flash. The fetchback that follows + # acts as verification — if the appliance didn't actually + # honour the write (silent coerce), the fetchback's + # publish will revert HA to the device's true state. + self._apply_optimistic(href, body) + # 3s settling window before the verification read. + # 1.5s was sometimes too short — the oven's read-side + # propagation lags more than that, and an early Block2 + # GET on /mode/vs/0 right after a POST appears to be one + # of the triggers for the oven actively closing DTLS. + # The fetchback runs in its own worker thread, so this + # delay is non-blocking; the optimistic publish has + # already given HA the new state. + self._schedule_fetchback(href, delay_s=3.0) + + def publish_health(self): + now = time.time() + 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), + } + 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) + + # ---- 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 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: + self.session_once() + backoff = 1.0 + except Exception as e: + self.error_count += 1 + self.log.warning("session error: %s", e) + sess = self.session + self.session = None + if sess is not None: + try: sess.close() + except Exception: pass + self.set_availability(False) + self.session_started_ts = None + if self.stop.is_set(): + break + wait = min(backoff, 30.0) + self.log.info("reconnect in %.0fs", wait) + if self.stop.wait(wait): + break + backoff = min(backoff * 2, 30.0) diff --git a/samsung_appliance/coap_dtls.py b/samsung_appliance/coap_dtls.py new file mode 100644 index 0000000..9f84056 --- /dev/null +++ b/samsung_appliance/coap_dtls.py @@ -0,0 +1,536 @@ +"""CoAP-over-DTLS client for Samsung RT-OCF appliances (RFC 7252 + 6347). + +Replaces the TLS-over-TCP transport used in the original dryer bridge. +Both the oven (UDP/49154) and the dryer (UDP/49155) speak CoAP-over-DTLS +with the ECDHE-ECDSA-AES128-GCM-SHA256 cipher and ab0b0ac4 client cert. + +Wire-level details that matter (from local-tools/oven-findings.md §17): + * DTLS ciphertext MTU must be 1200; otherwise OpenSSL fragments the + client cert across two datagrams and TizenRT drops the second. + * Samsung's RT-OCF uses ACK+separate-CON for the larger responses. + The reader MUST correlate by (token, mid) — not arrival order — + or interleaved one-shot / OBSERVE traffic mis-attributes. + * Multi-block GET requires the SAME CoAP token across every block + of the response ("token-stable Block2"). Fresh-token-per-block + is silently dropped by the server. + +Reader thread owns the UDP socket. Callers issue get()/post() and block +on a per-token Event the reader signals. OBSERVE notifications are +delivered via the on_notification callback. +""" +import socket +import struct +import threading +import time + +from OpenSSL import SSL + +from .logger import logger + + +# CoAP option numbers (RFC 7252 + 7641 + 7959) +URI_PATH = 11 +URI_QUERY = 15 +OBSERVE = 6 +CONTENT_FORMAT = 12 +ACCEPT = 17 +BLOCK2 = 23 +SIZE2 = 28 + +# CoAP message types +TYPE_CON = 0 +TYPE_NON = 1 +TYPE_ACK = 2 +TYPE_RST = 3 + +# CoAP method codes +METHOD_GET = 0x01 +METHOD_POST = 0x02 + +# CoAP content-format value for application/cbor +CF_CBOR = b'\x3c' + +# OBSERVE option values (RFC 7641 §2) +OBSERVE_REGISTER = b'' # register / refresh +OBSERVE_DEREGISTER = bytes([1]) # deregister + +# Block2 SZX=6 → 1024-byte blocks. The largest size Samsung's RT-OCF +# will honour and the only one the probes have validated end-to-end. +BLOCK_SZX = 6 + + +def _vlen(v): + """Variable-length integer encoder used in option deltas + lengths.""" + if v < 13: return v, b'' + if v < 269: return 13, bytes([v - 13]) + return 14, struct.pack('>H', v - 269) + + +def encode_options(opts): + """Encode a list of (option_number, value_bytes) tuples.""" + out = b'' + prev = 0 + for n, val in sorted(opts, key=lambda x: x[0]): + d, dx = _vlen(n - prev) + l, lx = _vlen(len(val)) + out += bytes([(d << 4) | l]) + dx + lx + val + prev = n + return out + + +def parse_coap(data): + """Decode a CoAP datagram. Returns (mtype, code, mid, token, + options, payload). options is a list of (num, value_bytes).""" + mt = (data[0] >> 4) & 0x03 + tkl = data[0] & 0x0F + code = data[1] + mid = int.from_bytes(data[2:4], 'big') + tok = data[4:4 + tkl] + i = 4 + tkl + opts = [] + prev = 0 + payload = b'' + while i < len(data): + b = data[i] + if b == 0xFF: + payload = data[i + 1:] + break + d_nib, l_nib = b >> 4, b & 0x0F + i += 1 + if d_nib == 13: + delta = 13 + data[i]; i += 1 + elif d_nib == 14: + delta = 269 + int.from_bytes(data[i:i + 2], 'big'); i += 2 + elif d_nib == 15: + raise ValueError("reserved option delta nibble 15") + else: + delta = d_nib + if l_nib == 13: + length = 13 + data[i]; i += 1 + elif l_nib == 14: + length = 269 + int.from_bytes(data[i:i + 2], 'big'); i += 2 + elif l_nib == 15: + raise ValueError("reserved option length nibble 15") + else: + length = l_nib + num = prev + delta + opts.append((num, data[i:i + length])) + i += length + prev = num + return mt, code, mid, tok, opts, payload + + +def build_coap(mtype, code, mid, token, options, payload=b''): + """Build a CoAP datagram. mtype: CON/NON/ACK/RST. token: bytes (may + be empty for ACK). options: list of (num, value_bytes).""" + tkl = len(token) + hdr = bytes([(1 << 6) | (mtype << 4) | tkl, code, + (mid >> 8) & 0xFF, mid & 0xFF]) + body = hdr + token + encode_options(options) + if payload: + body += b'\xFF' + payload + return body + + +def block_value(num, more, szx): + """Encode a CoAP Block-N option value.""" + v = (num << 4) | ((more & 1) << 3) | (szx & 7) + if v <= 0xFF: return bytes([v]) + if v <= 0xFFFF: return struct.pack('>H', v) + return struct.pack('>I', v)[1:] + + +def fmt_code(c): + """0x45 → '2.05', 0x84 → '4.04'. Used in log lines.""" + return f"{c >> 5}.{c & 0x1F:02d}" + + +def _split_dtls(buf): + """Split a UDP datagram that contains one-or-more DTLS records. + OpenSSL sometimes hands the BIO multiple records back-to-back; we + must send each as its own UDP datagram or TizenRT drops them.""" + o, out = 0, [] + while o + 13 <= len(buf): + L = int.from_bytes(buf[o + 11:o + 13], 'big') + end = o + 13 + L + if end > len(buf): + break + out.append(buf[o:end]) + o = end + return out + + +class DtlsCoapSession: + """Single sustained DTLS-CoAP session. + + Caller drives lifecycle: + sess = DtlsCoapSession(host, port, cert, key) + sess.connect() + sess.start_reader() + sess.subscribe([...], on_notification=cb) # OBSERVE + code, body = sess.get(['device', '0']) # Block2 fetch + code, _ = sess.post(['mode','vs','0'], cbor) + sess.close() + """ + + HANDSHAKE_TIMEOUT_S = 12.0 + READER_RECV_TIMEOUT_S = 1.0 # short so stop_event propagates quickly + MAX_BLOCKS = 32 # safety bound for Block2 fetches + + def __init__(self, host, port, cert_path, key_path, + on_notification=None, mtu=1200): + self.host = host + self.port = port + self.cert_path = str(cert_path) + self.key_path = str(key_path) + self.on_notification = on_notification # fn(href, payload_bytes) + self.mtu = mtu + + self.sock = None + self.conn = None + self.dest = None + + self._send_lock = threading.Lock() + self._mid = 0x5000 + self._tok_counter = 0 + # token (bytes) → (Event, container_dict) + self._pending = {} + # token (bytes) → href (str) + self._observe_tokens = {} + + self._stop = threading.Event() + self._reader_thread = None + + # ---- lifecycle --------------------------------------------------- + + def connect(self): + """DTLS handshake. Blocks up to HANDSHAKE_TIMEOUT_S. Raises + ConnectionError / TimeoutError on failure.""" + ctx = SSL.Context(SSL.DTLS_METHOD) + ctx.set_verify(SSL.VERIFY_NONE, lambda *_: True) + # @SECLEVEL=0 needed because the AC14K_M-rooted ab0b0ac4 chain + # is SHA-1 signed, which OpenSSL 3.x's default security level + # rejects. The chain comes from Samsung's leaked CA so the + # signature algorithm isn't ours to change. + ctx.set_cipher_list(b'ECDHE-ECDSA-AES128-GCM-SHA256:@SECLEVEL=0') + ctx.use_certificate_chain_file(self.cert_path) + ctx.use_privatekey_file(self.key_path) + ctx.check_privatekey() + + conn = SSL.Connection(ctx, None) + conn.set_connect_state() + conn.set_ciphertext_mtu(self.mtu) + + sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + sock.settimeout(2.0) + dest = (self.host, self.port) + + t0 = time.time() + while time.time() - t0 < self.HANDSHAKE_TIMEOUT_S: + try: + conn.do_handshake() + break + except SSL.WantReadError: + pass + except SSL.Error as e: + sock.close() + raise ConnectionError(f"DTLS handshake error: {e}") from e + try: + o = conn.bio_read(65535) + if o: + for r in _split_dtls(o): + sock.sendto(r, dest) + except SSL.WantReadError: + pass + try: + d, _ = sock.recvfrom(65535) + if d: + conn.bio_write(d) + except socket.timeout: + pass + time.sleep(0.05) + else: + sock.close() + raise TimeoutError( + f"DTLS handshake timeout to {self.host}:{self.port}") + + self.sock = sock + self.conn = conn + self.dest = dest + self._stop.clear() + + def start_reader(self): + """Spawn the reader thread. Must be called after connect().""" + if self.sock is None: + raise RuntimeError("connect() before start_reader()") + t = threading.Thread(target=self._reader_loop, + daemon=True, name='dtls-reader') + t.start() + self._reader_thread = t + + def join(self): + """Block until the reader thread exits (i.e. socket dies).""" + if self._reader_thread is not None: + self._reader_thread.join() + + def close(self): + """Tear down session. Signals reader_loop and wakes pending + waiters with an error.""" + self._stop.set() + if self.conn is not None: + try: + self.conn.shutdown() + except Exception: + pass + if self.sock is not None: + try: + self.sock.close() + except Exception: + pass + for tok, (ev, container) in list(self._pending.items()): + container.setdefault('err', 'socket closed') + ev.set() + self._pending.clear() + self._observe_tokens.clear() + self.sock = None + self.conn = None + + # ---- send / receive plumbing ------------------------------------- + + def _next_mid(self): + self._mid = (self._mid + 1) & 0xFFFF + return self._mid + + def _next_tok(self): + self._tok_counter = (self._tok_counter + 1) & 0xFFFFFFFF + # 4-byte tokens — fits within tkl=8 cap with headroom and + # avoids collisions across long-running OBSERVE subscriptions. + return self._tok_counter.to_bytes(4, 'big') + + def _send_dgram(self, datagram): + """Send a CoAP datagram. Holds the send lock for the + BIO-drain so two writers can't interleave records.""" + with self._send_lock: + if self.conn is None: + raise ConnectionError("DTLS session closed") + self.conn.send(datagram) + try: + while True: + o = self.conn.bio_read(65535) + if not o: + break + for r in _split_dtls(o): + self.sock.sendto(r, self.dest) + except SSL.WantReadError: + pass + + def _reader_loop(self): + """Pump UDP socket → DTLS BIO → CoAP parser. Demuxes to pending + / observe handlers. Exits on socket error or stop event.""" + sock = self.sock + conn = self.conn + sock.settimeout(self.READER_RECV_TIMEOUT_S) + try: + while not self._stop.is_set(): + try: + d, _ = sock.recvfrom(65535) + except socket.timeout: + continue + except (OSError, ValueError): + return + if not d: + continue + try: + conn.bio_write(d) + except SSL.Error as e: + logger.warning("DTLS bio_write: %s", e) + return + # Drain all app data the DTLS conn has buffered. One + # UDP datagram can yield zero, one, or several CoAP + # records depending on how mbedtls packed them. + while True: + try: + pl = conn.recv(65535) + except SSL.WantReadError: + break + except SSL.ZeroReturnError: + logger.info("DTLS peer closed connection") + return + except SSL.Error as e: + logger.warning("DTLS recv: %s", e) + return + if not pl: + break + try: + self._dispatch_coap(pl) + except Exception as e: + logger.warning("dispatch: %s", e) + finally: + # Make sure pending waiters don't hang if the reader dies. + for tok, (ev, container) in list(self._pending.items()): + container.setdefault('err', 'reader exited') + ev.set() + + def _dispatch_coap(self, datagram): + try: + mt, code, mid, tok, ropts, payload = parse_coap(datagram) + except Exception as e: + logger.debug("malformed CoAP: %s", e) + return + + # ACK back any CON from the device to suppress retransmits. + # RFC 7252 §4.2 — ACK is a bare frame (token len 0, code 0). + if mt == TYPE_CON: + try: + self._send_dgram(build_coap(TYPE_ACK, 0, mid, b'', [])) + except Exception as e: + logger.warning("ACK send: %s", e) + + # Empty ACK with no options & no payload = "separate response + # coming" — used by Samsung's RT-OCF for the larger reads. Stop + # the retransmit timer on the client side and wait for the CON. + if mt == TYPE_ACK and code == 0 and not payload and not ropts: + return + + # Pending one-shot? Resolve and return. + rec = self._pending.get(tok) + if rec is not None: + ev, container = rec + container['code'] = code + container['mtype'] = mt + container['mid'] = mid + container['options'] = ropts + container['payload'] = payload + ev.set() + return + + # OBSERVE notification? + href = self._observe_tokens.get(tok) + if href is not None: + if code != 0x45: + logger.warning("observe %s: non-2.05 %s", + href, fmt_code(code)) + return + cb = self.on_notification + if cb is not None: + try: + cb(href, payload) + except Exception as e: + logger.warning("notification callback %s: %s", + href, e) + return + + # Stale token (post-reconnect or unknown) — drop quietly. + + # ---- request primitives ------------------------------------------ + + def get(self, path_segs, query=(), timeout=10.0): + """Token-stable Block2 GET. Returns (code, payload_bytes). + + Reuses one CoAP token across every block of a multi-block + response — Samsung's server keys per-transfer state on the + token, and dropping a fresh token on block 1+ silently drops + the request.""" + if self.conn is None: + raise ConnectionError("DTLS session closed") + tok = self._next_tok() + blob = b'' + num = 0 + last_code = None + last_opts = [] + deadline = time.time() + timeout + while True: + ev = threading.Event() + container = {} + self._pending[tok] = (ev, container) + try: + mid = self._next_mid() + opts = [(URI_PATH, s.encode()) for s in path_segs] + for q in query: + opts.append((URI_QUERY, q.encode())) + opts.append((ACCEPT, CF_CBOR)) + if num > 0: + opts.append((BLOCK2, block_value(num, 0, BLOCK_SZX))) + self._send_dgram( + build_coap(TYPE_CON, METHOD_GET, mid, tok, opts)) + wait = max(0.1, deadline - time.time()) + if not ev.wait(wait): + raise TimeoutError( + f"GET /{'/'.join(path_segs)} block {num} timeout") + if 'err' in container: + raise ConnectionError(container['err']) + finally: + self._pending.pop(tok, None) + + code = container['code'] + payload = container['payload'] + ropts = container['options'] + last_code = code + last_opts = ropts + blob += payload + # 4.xx / 5.xx responses don't carry Block2 continuation — + # bail with whatever we got. Caller decides if 4.xx is fatal. + if code >> 5 != 2: + return code, blob + b2 = [v for n, v in ropts if n == BLOCK2] + more = 0 + if b2: + bv = int.from_bytes(b2[0], 'big') + more = (bv >> 3) & 1 + if not more: + break + num += 1 + if num > self.MAX_BLOCKS: + raise ConnectionError( + f"GET /{'/'.join(path_segs)}: >{self.MAX_BLOCKS} " + f"blocks, aborting") + return last_code, blob + + def post(self, path_segs, body_cbor, timeout=8.0): + """Single-frame POST with a CBOR-encoded body. Returns + (code, payload_bytes). body_cbor must already be encoded.""" + if self.conn is None: + raise ConnectionError("DTLS session closed") + tok = self._next_tok() + mid = self._next_mid() + opts = [(URI_PATH, s.encode()) for s in path_segs] + opts.append((CONTENT_FORMAT, CF_CBOR)) + opts.append((ACCEPT, CF_CBOR)) + datagram = build_coap(TYPE_CON, METHOD_POST, mid, tok, opts, + body_cbor) + ev = threading.Event() + container = {} + self._pending[tok] = (ev, container) + try: + self._send_dgram(datagram) + if not ev.wait(timeout): + raise TimeoutError( + f"POST /{'/'.join(path_segs)} timeout") + if 'err' in container: + raise ConnectionError(container['err']) + return container['code'], container['payload'] + finally: + self._pending.pop(tok, None) + + def subscribe(self, path_segs): + """Register an OBSERVE on the given path. The initial 2.05 + notification and all subsequent state-change notifications + will fire on_notification(href, payload_bytes). + + Returns the token used (in case the caller wants to deregister + later).""" + if self.conn is None: + raise ConnectionError("DTLS session closed") + tok = self._next_tok() + href = '/' + '/'.join(path_segs) + # Register the token BEFORE sending — otherwise the device + # could respond between send() and the dict insert, and the + # reader thread would drop the initial 2.05 as "stale". + self._observe_tokens[tok] = href + mid = self._next_mid() + opts = [(URI_PATH, s.encode()) for s in path_segs] + opts.append((OBSERVE, OBSERVE_REGISTER)) + opts.append((ACCEPT, CF_CBOR)) + self._send_dgram( + build_coap(TYPE_CON, METHOD_GET, mid, tok, opts)) + return tok diff --git a/samsung_appliance/config.py b/samsung_appliance/config.py new file mode 100644 index 0000000..7a15468 --- /dev/null +++ b/samsung_appliance/config.py @@ -0,0 +1,118 @@ +"""Configuration — env-var driven. + +Two flavours: + * SharedConfig — MQTT broker, cert paths, HA prefix, timers. One per + process. + * ApplianceConfig — one per appliance the bridge is supervising. Keys + come from `APPLIANCE__*` env vars (1-indexed). The list of + appliances is `APPLIANCE_COUNT` entries long. + +.env in cwd hydrates os.environ at import time; docker-compose env wins +over file contents.""" +import os +from dataclasses import dataclass +from pathlib import Path +from typing import Optional + + +def _load_env_file(filename='.env'): + """If a .env-style file is present in cwd, hydrate os.environ from + it. Existing env wins so docker-compose `environment:` overrides + file contents.""" + env_path = Path.cwd() / filename + if not env_path.exists(): + return + for line in env_path.read_text().splitlines(): + line = line.strip() + if not line or line.startswith('#') or '=' not in line: + continue + k, v = line.split('=', 1) + k = k.strip(); v = v.strip() + if k and k not in os.environ: + os.environ[k] = v + + +_load_env_file('.env') + + +def _resolve_cert(env_key, basename): + """Cert lookup: explicit env > /config/ (Docker mount) > + ./certs/ (bare-metal dev).""" + if os.getenv(env_key): + return Path(os.environ[env_key]) + docker_path = Path('/config') / basename + if docker_path.exists(): + return docker_path + return Path.cwd() / 'certs' / basename + + +@dataclass(frozen=True) +class SharedConfig: + """Process-wide config (MQTT, cert paths, intervals).""" + CERT_PATH: Path + KEY_PATH: Path + MQTT_BROKER: Optional[str] + MQTT_PORT: int + MQTT_USER: Optional[str] + MQTT_PASS: Optional[str] + HA_DISCOVERY_PREFIX: str + HEALTH_INTERVAL_S: int + HEARTBEAT_INTERVAL_S: int + + @classmethod + def from_env(cls) -> 'SharedConfig': + return cls( + CERT_PATH=_resolve_cert('CERT_PATH', 'ab0b0ac4_fullchain.pem'), + KEY_PATH=_resolve_cert('KEY_PATH', 'ab0b0ac4.key'), + MQTT_BROKER=os.getenv('MQTT_BROKER'), + MQTT_PORT=int(os.getenv('MQTT_PORT', '1883')), + MQTT_USER=os.getenv('MQTT_USER') or None, + MQTT_PASS=os.getenv('MQTT_PASS') or None, + 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')), + ) + + +@dataclass(frozen=True) +class ApplianceConfig: + """Per-appliance runtime config. + + Sourced from `APPLIANCE__*` env vars (1-indexed). `index` is + just a stable identifier for logs; it does not appear in MQTT + topics or HA discovery (those are keyed off `topic_prefix`).""" + index: int + klass: str # 'dryer', 'oven', … + ip: str + ocf_port: Optional[int] # None → descriptor.default_observe_port + topic_prefix: str + device_name: str + + @classmethod + def from_env(cls, index: int) -> 'ApplianceConfig': + prefix = f'APPLIANCE_{index}_' + klass = os.getenv(prefix + 'CLASS') + if not klass: + raise ValueError(f"{prefix}CLASS not set") + ip = os.getenv(prefix + 'IP') + if not ip: + raise ValueError(f"{prefix}IP not set") + port_env = os.getenv(prefix + 'OCF_PORT') + port = int(port_env) if port_env else None + topic = os.getenv(prefix + 'TOPIC') or f'samsung_{klass}' + name = os.getenv(prefix + 'NAME') or f'Samsung {klass.title()}' + return cls( + index=index, klass=klass, ip=ip, ocf_port=port, + topic_prefix=topic, device_name=name, + ) + + +def load_appliances() -> list[ApplianceConfig]: + """Read APPLIANCE_COUNT and build the appliance list. At least one + is required.""" + count = int(os.getenv('APPLIANCE_COUNT', '1')) + if count < 1: + raise ValueError("APPLIANCE_COUNT must be >= 1") + return [ApplianceConfig.from_env(i + 1) for i in range(count)] diff --git a/samsung_appliance/logger.py b/samsung_appliance/logger.py new file mode 100644 index 0000000..7661908 --- /dev/null +++ b/samsung_appliance/logger.py @@ -0,0 +1,32 @@ +"""Stdout logging — works correctly under Docker's PYTHONUNBUFFERED=1. + +Two logger families share the root handler: + * `samsung_appliance` (and its children) for module-level lines — + DTLS warnings, generic startup chatter, MQTT plumbing. + * `.` for per-appliance bridge lines — e.g. + `dryer.`. Each PushBridge gets its own such logger + via bridge_logger() once the seed reveals the appliance's serial. + +Both trees propagate to the root logger, which is the one with the +StreamHandler, so the same format applies everywhere.""" +import logging +import sys + + +logging.basicConfig( + level=logging.INFO, + format='%(asctime)s %(levelname)-5s %(name)-26s %(message)s', + datefmt='%H:%M:%S', + handlers=[logging.StreamHandler(sys.stdout)], +) + +logger = logging.getLogger("samsung_appliance") + + +def bridge_logger(klass: str, serial: str | None = None) -> logging.Logger: + """Return a top-level logger tagged with the appliance class and, + once known, the serial. Pre-seed callers pass serial=None and get + e.g. `dryer`; post-seed callers pass the serial and get e.g. + `dryer.`.""" + name = f"{klass}.{serial}" if serial else klass + return logging.getLogger(name) diff --git a/samsung_appliance/sensors.py b/samsung_appliance/sensors.py new file mode 100644 index 0000000..5344492 --- /dev/null +++ b/samsung_appliance/sensors.py @@ -0,0 +1,19 @@ +"""Shared sensor helpers. + +Appliance-specific flattening lives in samsung_appliance/appliances/*.py; +this module only carries utilities that every descriptor uses (currently +the /device/0 link-dict indexer). +""" + + +def index_links(device0_body): + """Turn the /device/0 CBOR list-of-{href, rep} into a dict keyed + by href. The first list entry is the device-level rep itself and + isn't useful here, so skip it.""" + out = {} + if not isinstance(device0_body, list): + return out + for entry in device0_body[1:]: + if isinstance(entry, dict) and 'href' in entry: + out[entry['href']] = entry.get('rep') or {} + return out