diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..fffbd75 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,15 @@ +# Build context is the repo root (see mqtt_demo/docker-compose.yml's +# `context: ..`). The Dockerfile COPYs whole directories (protocol/, +# ocf/, mqtt_demo/), so anything under them that shouldn't land in the +# image has to be excluded here explicitly. + +# Secrets — never bake these into the image. Mounted at runtime via +# the /config volume instead. +**/.env +certs/ + +# Local dev cruft. +**/__pycache__/ +**/*.pyc +.git/ +.venv/ diff --git a/README.md b/README.md index 0a64f18..07428c5 100644 --- a/README.md +++ b/README.md @@ -4,6 +4,12 @@ image +> **Looking to control your Samsung appliance from Home Assistant?** +> Use [localthings](https://github.com/mbillow/localthings) — a Home +> Assistant custom component built on this repo's `protocol/` + `ocf/` +> layers. This repo is the protocol research project and a +> self-contained MQTT bridge demo; new appliance support (capability +> mappings, HA entities) should go to localthings, not here. > ### Proof of concept — collaborators wanted > @@ -54,11 +60,11 @@ Read the result: | Oven | NV7000BS-class (`TP1X_DA-KS-OVEN-0107X`, `mnid=0AJT`) | All entities; hot-tier poll covers door + operational state regardless of cloud reachability | | Fridge | ARTIK051_REF_17K (`DA-REF-ART-COMMON-1_20201124`) | Contributed by [@aminorjourney](https://github.com/aminorjourney) (PR #1). Older firmware family; port 49155, minimal `/oic/res` with full tree under `/device/0` | -Other appliances on the same firmware family (washers, dishwashers, AC units) almost certainly speak the same protocol — the auth path and read primitives are common. You'd write one new descriptor in `samsung_appliance/appliances/`. +Other appliances on the same firmware family (washers, dishwashers, AC units) almost certainly speak the same protocol — the auth path and read primitives are common. You'd write one new descriptor for the `localthings` registry. ### Firmware families — a limitation -Descriptors are firmware-family-specific. Each file in `samsung_appliance/appliances/` hardcodes the resource layout of one firmware family: which hrefs it polls, which fields it reads, which write surfaces it exposes. There's no runtime feature detection. +Descriptors are firmware-family-specific. Each descriptor hardcodes the resource layout of one firmware family: which hrefs it polls, which fields it reads, which write surfaces it exposes. There's no runtime feature detection. The three sample descriptors here (`mqtt_demo/samples/`) are frozen references. **What this means in practice:** if you set `APPLIANCE__CLASS=fridge` on a fridge that speaks a different firmware family than the one this descriptor was built for, the bridge will start and connect fine, but many sensors will publish as unknown and some controls won't work. Nothing catastrophic — you just get a half-broken HA device card. @@ -162,7 +168,7 @@ 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`. +Each `APPLIANCE__CLASS` must match a descriptor key in `mqtt_demo/samples/__init__.py::DESCRIPTORS` — currently `dryer`, `oven`, and `fridge`. --- @@ -194,18 +200,18 @@ Set `SSH_HOST`, `REMOTE_DIR`, `APPDATA_DIR` in `.env`. `deploy.sh` extracts thos ```sh python3 -m venv .venv -.venv/bin/pip install -r requirements.txt -.venv/bin/python main.py +.venv/bin/pip install -r mqtt_demo/requirements.txt +.venv/bin/python -m mqtt_demo ``` ### 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:42 INFO mqtt_demo SmartThings-Local Bridge starting (2 appliances) +14:08:42 INFO mqtt_demo broker = :1883 (user=) +14:08:42 INFO mqtt_demo [1] dryer @ :49155 (DTLS) → topic samsung_dryer/* +14:08:42 INFO mqtt_demo [2] oven @ :49154 (DTLS) → topic samsung_oven/* +14:08:42 INFO mqtt_demo 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 @@ -344,45 +350,52 @@ Gated control entities use HA's `availability_mode: all` against `/avail ### Repo layout ``` -main.py Entry point — loads config, spawns one PushBridge per appliance setup_cert.py One-shot cert minting script (live-fetches AC14K_M + UUID) -samsung_appliance/ The bridge package +protocol/ DTLS-CoAP protocol layer (reusable for non-MQTT bridges) __init__.py + auth.py DTLS client cert setup + authentication + dtls_session.py DTLS session management, handshake, liveness + coap.py CoAP wire protocol: message encode/decode, token handling +ocf/ OCF resource + state management (reusable layer) + __init__.py + state_cache.py StateCache — single source of truth for appliance state + poll_scheduler.py Tiered adaptive polling (hot/warm/cold + sweep) + keepalive.py CoAP liveness checks (empty-CON pings) + observe_refresh.py OBSERVE registration management +mqtt_demo/ MQTT bridge demo (uses protocol/ + ocf/) + __init__.py + __main__.py Entry point — loads config, spawns one bridge per appliance 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 + bridge.py Bridge — one DTLS session per appliance, descriptor-driven + descriptor.py ApplianceDescriptor dataclass + HA discovery helpers + samples/ + __init__.py Sample DESCRIPTORS registry (frozen reference implementations) 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 + fridge.py Fridge descriptor (ARTIK051 firmware family) + Dockerfile Container build (python:3.11-slim + 3 deps) + docker-compose.yml One service: smartthings-local + deploy.sh tar + ssh + docker compose up --build + requirements.txt Python dependencies for the bridge + .env.example Template — copy to .env, fill in ``` -`certs/` is gitignored. Drop the privileged client cert + key there; the container mounts that directory read-only at `/config`. +`certs/` is gitignored. Drop the privileged client cert + key there; the container mounts that directory read-only at `/config`. See [`localthings`](https://github.com/mbillow/localthings) for production HA integration. --- -## Adding a new appliance class +## Adding appliance support -The bridge is appliance-agnostic. Adding e.g. a washer is mechanical: +The three descriptors in `mqtt_demo/samples/` (dryer, oven, fridge) are +frozen reference implementations — enough to exercise both the newer +Tizen RT 3.x family and the older ARTIK051 family, proving the +`protocol/` + `ocf/` layers generalize across firmware generations. +They are not updated for new appliance models. -1. Capture the appliance's `/device/0` to see what resources/fields it exposes. Use the GET helpers in `samsung_appliance/coap_dtls.py` against your authenticated session. -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. - -For a new appliance class also declare a `poll_tiers: list[PollTier]` and (optionally) an `is_active(links) -> bool` predicate on the descriptor. Hot-tier resources are whatever needs sub-second freshness for HA UX; warm covers everything else that's not static; the `/device/0` sweep tier catches anything you forgot. The descriptor pattern handles everything else — DTLS, MQTT, HA discovery, optimistic writes, Block2 reads, OBSERVE accelerator, reconnect, liveness pings. +**To add support for a new appliance, submit it to +[localthings](https://github.com/mbillow/localthings)**, which owns +the capability registry and Home Assistant integration. --- @@ -393,7 +406,7 @@ These each looked like obvious improvements at some point. Each one broke someth - **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 assume OBSERVE silence means the appliance is broken.** When the appliance can't reach Samsung's cloud, its OBSERVE notify dispatch goes quiet even though the local DTLS session, GETs, POSTs, and the cache continue to work normally (measured at `~14 req/s` dryer / `~8 req/s` oven with 200/200 GETs successful while firewalled). The polling tiers are the structural answer to this; treat OBSERVE strictly as an optional accelerator. - **Don't touch `/oic/sec/*` (doxm, pstat, cred, acl).** The bridge doesn't, and you shouldn't from helper scripts either — those resources have wedge/brick risk on Samsung's RT-OCF security stack. The bridge surfaces are strictly `//vs/0` and `/device/0`. -- **Don't run two clients against the same appliance simultaneously.** Samsung's RT-OCF DTLS allows one active session per peer; a second handshake will get the device to drop the new socket. If HA seems to flap, check whether you've got `main.py` running locally AND the Docker container up. +- **Don't 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 `python -m mqtt_demo` running locally AND the Docker container up. - **Don't expect parity from every write surface.** Samsung's firmware accepts a lot of writes with `2.04 Changed` but only some of them stick — power, child-lock, and remote-control writes are accepted-then-reverted because they're hardware-mirrored. The bridge's optimistic-publish-then-verify pattern handles this transparently: HA briefly shows the new value, the 3s fetch-back republishes the actual value, HA reverts. --- diff --git a/.env.example b/mqtt_demo/.env.example similarity index 100% rename from .env.example rename to mqtt_demo/.env.example diff --git a/Dockerfile b/mqtt_demo/Dockerfile similarity index 69% rename from Dockerfile rename to mqtt_demo/Dockerfile index ddad334..8510c75 100644 --- a/Dockerfile +++ b/mqtt_demo/Dockerfile @@ -3,12 +3,15 @@ FROM python:3.11-slim WORKDIR /app # Python deps first so layer cache survives code changes -COPY requirements.txt . +COPY mqtt_demo/requirements.txt . RUN pip install --no-cache-dir -r requirements.txt -# Application code -COPY main.py . -COPY samsung_appliance/ ./samsung_appliance/ +# Application code — build context is the repo root (see +# docker-compose.yml's `context: ..`) since mqtt_demo/ depends on the +# protocol/ and ocf/ library packages that live alongside it. +COPY protocol/ ./protocol/ +COPY ocf/ ./ocf/ +COPY mqtt_demo/ ./mqtt_demo/ # /config holds the ab0b0ac4 client cert + key. Mount from the host so # secrets aren't baked into the image. @@ -28,4 +31,4 @@ ENV CERT_PATH=/config/ab0b0ac4_fullchain.pem \ # No port — bridge is outbound-only (DTLS UDP to appliance, MQTT to broker). -CMD ["python", "main.py"] +CMD ["python", "-m", "mqtt_demo"] diff --git a/mqtt_demo/__init__.py b/mqtt_demo/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/main.py b/mqtt_demo/__main__.py similarity index 96% rename from main.py rename to mqtt_demo/__main__.py index 6826cd7..570a10f 100644 --- a/main.py +++ b/mqtt_demo/__main__.py @@ -23,10 +23,10 @@ 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 +from mqtt_demo.samples import get_descriptor +from mqtt_demo.bridge import PushBridge +from mqtt_demo.config import SharedConfig, load_appliances +from mqtt_demo.logger import logger def main(): diff --git a/samsung_appliance/bridge.py b/mqtt_demo/bridge.py similarity index 98% rename from samsung_appliance/bridge.py rename to mqtt_demo/bridge.py index 43dec1a..36b8747 100644 --- a/samsung_appliance/bridge.py +++ b/mqtt_demo/bridge.py @@ -23,15 +23,16 @@ import time import cbor2 -from .appliances.base import ApplianceDescriptor, bridge_diagnostic_discovery -from .coap_dtls import DtlsCoapSession, fmt_code +from protocol.dtls_session import DtlsCoapSession, fmt_code + +from ocf.keepalive import KeepaliveTask +from ocf.observe_refresh import ObserveRefreshTask +from ocf.poll_scheduler import PollScheduler +from ocf.state_cache import StateCache + +from .descriptor import ApplianceDescriptor, bridge_diagnostic_discovery from .config import ApplianceConfig, SharedConfig -from .keepalive import KeepaliveTask from .logger import bridge_logger -from .observe_refresh import ObserveRefreshTask -from .poll_scheduler import PollScheduler -from .sensors import index_links -from .state_cache import StateCache DEBUG_BRIDGE = os.environ.get('DEBUG_BRIDGE') == '1' @@ -295,7 +296,6 @@ class PushBridge: scheduler = PollScheduler( sess, self.cache, tiers=self.descriptor.poll_tiers, - sweep_index_fn=index_links, is_active_fn=self.descriptor.is_active, logger=self.log, ) @@ -357,7 +357,7 @@ class PushBridge: # During the seed we want the cache populated without triggering # a publish per resource — gate the on_change callback off until # the publish gate opens just below. - for href, rep in index_links(body).items(): + for href, rep in StateCache.index_device_tree(body).items(): if href not in self.cache.links: self.cache.apply_rep(href, rep, source='seed') self.last_seed_ts = time.time() diff --git a/samsung_appliance/config.py b/mqtt_demo/config.py similarity index 100% rename from samsung_appliance/config.py rename to mqtt_demo/config.py diff --git a/deploy.sh b/mqtt_demo/deploy.sh similarity index 67% rename from deploy.sh rename to mqtt_demo/deploy.sh index 1ae387c..c160a49 100755 --- a/deploy.sh +++ b/mqtt_demo/deploy.sh @@ -20,8 +20,15 @@ # 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." +# This script lives in mqtt_demo/ but the build context is the repo +# root (mqtt_demo/docker-compose.yml uses `context: ..`, since the +# image needs the protocol/ and ocf/ library packages alongside +# mqtt_demo/). Run everything from the repo root so the tar allowlist +# and remote layout line up with that context. +cd "$(dirname "$0")/.." + +if [ ! -f mqtt_demo/.env ]; then + echo "Error: mqtt_demo/.env file not found. Copy mqtt_demo/.env.example to mqtt_demo/.env and configure it." exit 1 fi @@ -29,7 +36,7 @@ fi # 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- + grep -E "^${1}=" mqtt_demo/.env | head -1 | cut -d= -f2- } SSH_HOST=$(get_env SSH_HOST) REMOTE_DIR=$(get_env REMOTE_DIR) @@ -44,22 +51,21 @@ 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. +# local. protocol/ and ocf/ are the library packages mqtt_demo/ imports +# from; they need to land as REMOTE_DIR's siblings of mqtt_demo/ so the +# compose file's `context: ..` resolves the same way it does locally. COPYFILE_DISABLE=1 tar cz \ - main.py \ - samsung_appliance/ \ - Dockerfile \ - docker-compose.yml \ - requirements.txt \ - deploy.sh \ + protocol/ \ + ocf/ \ + mqtt_demo/ \ README.md \ - .env.example \ .gitignore \ + .dockerignore \ | 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" +scp mqtt_demo/.env "${SSH_HOST}:${REMOTE_DIR}/mqtt_demo/.env" +ssh "${SSH_HOST}" "chmod 600 ${REMOTE_DIR}/mqtt_demo/.env" # Verify certs are present on the remote — they have to be uploaded # once before the first build. @@ -74,7 +80,7 @@ if ! ssh "${SSH_HOST}" "test -s ${APPDATA_DIR}/client_fullchain.pem && test -s $ fi echo "Rebuilding container…" -ssh "${SSH_HOST}" "cd ${REMOTE_DIR} && docker compose up -d --build" +ssh "${SSH_HOST}" "cd ${REMOTE_DIR}/mqtt_demo && docker compose up -d --build" echo "Done." -echo "Logs: ssh ${SSH_HOST} 'cd ${REMOTE_DIR} && docker compose logs -f'" +echo "Logs: ssh ${SSH_HOST} 'cd ${REMOTE_DIR}/mqtt_demo && docker compose logs -f'" diff --git a/samsung_appliance/appliances/base.py b/mqtt_demo/descriptor.py similarity index 99% rename from samsung_appliance/appliances/base.py rename to mqtt_demo/descriptor.py index fe3977b..ed335d2 100644 --- a/samsung_appliance/appliances/base.py +++ b/mqtt_demo/descriptor.py @@ -27,7 +27,7 @@ from dataclasses import dataclass, field from typing import Callable, Optional, TYPE_CHECKING if TYPE_CHECKING: - from ..poll_scheduler import PollTier + from ocf.poll_scheduler import PollTier @dataclass diff --git a/docker-compose.yml b/mqtt_demo/docker-compose.yml similarity index 91% rename from docker-compose.yml rename to mqtt_demo/docker-compose.yml index 769bb08..aa1a471 100644 --- a/docker-compose.yml +++ b/mqtt_demo/docker-compose.yml @@ -1,6 +1,8 @@ services: smartthings-local: - build: . + build: + context: .. + dockerfile: mqtt_demo/Dockerfile container_name: smartthings-local restart: unless-stopped diff --git a/samsung_appliance/logger.py b/mqtt_demo/logger.py similarity index 95% rename from samsung_appliance/logger.py rename to mqtt_demo/logger.py index 55d6439..f287c9c 100644 --- a/samsung_appliance/logger.py +++ b/mqtt_demo/logger.py @@ -1,7 +1,7 @@ """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 — + * `mqtt_demo` (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 @@ -58,7 +58,7 @@ _handler.setFormatter(_LevelColourFormatter( logging.basicConfig(level=logging.INFO, handlers=[_handler], force=True) -logger = logging.getLogger("samsung_appliance") +logger = logging.getLogger("mqtt_demo") def bridge_logger(klass: str, serial: str | None = None) -> logging.Logger: diff --git a/requirements.txt b/mqtt_demo/requirements.txt similarity index 100% rename from requirements.txt rename to mqtt_demo/requirements.txt diff --git a/samsung_appliance/appliances/__init__.py b/mqtt_demo/samples/__init__.py similarity index 69% rename from samsung_appliance/appliances/__init__.py rename to mqtt_demo/samples/__init__.py index 20657a9..ec8c8ee 100644 --- a/samsung_appliance/appliances/__init__.py +++ b/mqtt_demo/samples/__init__.py @@ -1,15 +1,15 @@ """Appliance descriptors registry. Adding a new appliance class: - 1. Write `appliances/.py` with an ApplianceDescriptor named - after the class (uppercase, e.g. OVEN). + 1. Write `mqtt_demo/samples/.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. +mqtt_demo/__main__.py imports get_descriptor(name) to look up the +descriptor at startup; the bridge itself stays class-agnostic. """ -from .base import ApplianceDescriptor +from ..descriptor import ApplianceDescriptor from .dryer import DRYER from .oven import OVEN from .fridge import FRIDGE diff --git a/samsung_appliance/appliances/dryer.py b/mqtt_demo/samples/dryer.py similarity index 99% rename from samsung_appliance/appliances/dryer.py rename to mqtt_demo/samples/dryer.py index 0186680..533b97d 100644 --- a/samsung_appliance/appliances/dryer.py +++ b/mqtt_demo/samples/dryer.py @@ -6,14 +6,14 @@ samsung_dryer/{bridge,sensors,discovery}.py modules into one place. """ import time -from .base import ( +from ..descriptor import ( ApplianceDescriptor, avail_base, avail_with_remote, device_block, encode, ) -from ..poll_scheduler import PollTier +from ocf.poll_scheduler import PollTier # --- OBSERVE paths ----------------------------------------------------- diff --git a/samsung_appliance/appliances/fridge.py b/mqtt_demo/samples/fridge.py similarity index 99% rename from samsung_appliance/appliances/fridge.py rename to mqtt_demo/samples/fridge.py index 4687eec..49b48c3 100644 --- a/samsung_appliance/appliances/fridge.py +++ b/mqtt_demo/samples/fridge.py @@ -24,13 +24,13 @@ Key differences from newer Tizen RT firmware: """ import time -from .base import ( +from ..descriptor import ( ApplianceDescriptor, avail_base, device_block, encode, ) -from ..poll_scheduler import PollTier +from ocf.poll_scheduler import PollTier MODEL = 'ARTIK051_REF_17K' diff --git a/samsung_appliance/appliances/oven.py b/mqtt_demo/samples/oven.py similarity index 99% rename from samsung_appliance/appliances/oven.py rename to mqtt_demo/samples/oven.py index 9da542c..48c6ce8 100644 --- a/samsung_appliance/appliances/oven.py +++ b/mqtt_demo/samples/oven.py @@ -23,7 +23,7 @@ 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 ( +from ..descriptor import ( ApplianceDescriptor, avail_base, avail_with_cycle, @@ -31,7 +31,7 @@ from .base import ( device_block, encode, ) -from ..poll_scheduler import PollTier +from ocf.poll_scheduler import PollTier # --------------------------------------------------------------------- diff --git a/ocf/__init__.py b/ocf/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/samsung_appliance/keepalive.py b/ocf/keepalive.py similarity index 98% rename from samsung_appliance/keepalive.py rename to ocf/keepalive.py index 11fe464..9543770 100644 --- a/samsung_appliance/keepalive.py +++ b/ocf/keepalive.py @@ -20,7 +20,7 @@ from __future__ import annotations import threading from typing import Callable, Optional -from .coap_dtls import DtlsCoapSession +from protocol.dtls_session import DtlsCoapSession class KeepaliveTask: diff --git a/samsung_appliance/observe_refresh.py b/ocf/observe_refresh.py similarity index 97% rename from samsung_appliance/observe_refresh.py rename to ocf/observe_refresh.py index bb8c24d..15a9ca0 100644 --- a/samsung_appliance/observe_refresh.py +++ b/ocf/observe_refresh.py @@ -17,7 +17,7 @@ from __future__ import annotations import threading from typing import Optional -from .coap_dtls import DtlsCoapSession +from protocol.dtls_session import DtlsCoapSession class ObserveRefreshTask: diff --git a/samsung_appliance/poll_scheduler.py b/ocf/poll_scheduler.py similarity index 97% rename from samsung_appliance/poll_scheduler.py rename to ocf/poll_scheduler.py index ff8e6ba..6634957 100644 --- a/samsung_appliance/poll_scheduler.py +++ b/ocf/poll_scheduler.py @@ -3,8 +3,8 @@ Tiers are descriptor-declared (hot/warm/cold + sweep). Per tick: each tier whose deadline has passed polls all its paths sequentially on the shared session, writing into the StateCache. The sweep tier -issues one Block2 GET of /device/0 and uses index_links to fan its -result into many href reps. +issues one Block2 GET of /device/0 and uses +StateCache.index_device_tree to fan its result into many href reps. Adaptive cadence: when descriptor.is_active(cache.links) returns True and tier.active_interval_s is set, that tier uses the tighter cadence. @@ -36,7 +36,7 @@ from typing import Callable, Optional, TYPE_CHECKING import cbor2 -from .coap_dtls import DtlsCoapSession, fmt_code +from protocol.dtls_session import DtlsCoapSession, fmt_code if TYPE_CHECKING: from .state_cache import StateCache @@ -61,7 +61,6 @@ class PollScheduler: session: DtlsCoapSession, cache: 'StateCache', tiers: list[PollTier], - sweep_index_fn: Callable[[object], dict[str, dict]], is_active_fn: Optional[Callable[[dict[str, dict]], bool]] = None, logger=None, timeout_s: float = 8.0, @@ -69,7 +68,6 @@ class PollScheduler: self.session = session self.cache = cache self.tiers = tiers - self.sweep_index = sweep_index_fn self.is_active_fn = is_active_fn self.log = logger self.timeout_s = timeout_s @@ -220,11 +218,13 @@ class PollScheduler: def _do_tier(self, tier: PollTier) -> None: timeout = self._tier_timeout(tier) cooldown = self._cooldown_for(tier) - for path in tier.paths: + for i, path in enumerate(tier.paths): href = '/' + '/'.join(path) with self._defer_lock: if self._defer_until.get(href, 0) > time.monotonic(): continue + if i > 0: + self.session.pace() self._poll_count += 1 t0 = time.monotonic() try: @@ -301,7 +301,7 @@ class PollScheduler: self._poll_error_count += 1 if self.log: self.log.warning("sweep cbor: %s", e) return - indexed = self.sweep_index(tree) + indexed = self.cache.index_device_tree(tree) for href, rep in indexed.items(): with self._defer_lock: if self._defer_until.get(href, 0) > time.monotonic(): diff --git a/samsung_appliance/state_cache.py b/ocf/state_cache.py similarity index 71% rename from samsung_appliance/state_cache.py rename to ocf/state_cache.py index 852c34f..7aa19fa 100644 --- a/samsung_appliance/state_cache.py +++ b/ocf/state_cache.py @@ -8,15 +8,16 @@ from __future__ import annotations import threading import time -from typing import Callable, Optional, TYPE_CHECKING +from typing import Callable, Optional, Protocol -if TYPE_CHECKING: - from .appliances.base import ApplianceDescriptor + +class _ObservationHook(Protocol): + def on_observation(self, state: dict, href: str, rep: dict) -> None: ... class StateCache: - def __init__(self, descriptor: 'ApplianceDescriptor'): + def __init__(self, descriptor: '_ObservationHook'): self.descriptor = descriptor self.links: dict[str, dict] = {} self.last_updated: dict[str, float] = {} @@ -58,6 +59,23 @@ class StateCache: merged.update(body) return self.apply_rep(href, merged, source='optimistic') + @staticmethod + def index_device_tree(device0_body) -> dict[str, dict]: + """Turn a /device/0 CBOR list-of-{href, rep} sweep response into + a dict keyed by href. Entry [0] is the device-level rep itself + and isn't useful here, so it's skipped. + + Replaces the old standalone sensors.index_links — folded in + here because every current and future caller immediately feeds + the result into apply_rep on this same cache.""" + out: dict[str, dict] = {} + 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 + def get(self, href: str) -> Optional[dict]: with self._lock: return self.links.get(href) diff --git a/protocol/__init__.py b/protocol/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/protocol/coap.py b/protocol/coap.py new file mode 100644 index 0000000..f20287b --- /dev/null +++ b/protocol/coap.py @@ -0,0 +1,139 @@ +"""CoAP wire encoding/decoding (RFC 7252 + 7641 + 7959). + +Pure functions — no sockets, no DTLS. Split out of the original +coap_dtls.py so protocol/dtls_session.py (the stateful session) and +this module (stateless wire format) can be reasoned about and tested +independently. +""" +import struct + +# 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 diff --git a/samsung_appliance/coap_dtls.py b/protocol/dtls_session.py similarity index 80% rename from samsung_appliance/coap_dtls.py rename to protocol/dtls_session.py index de64044..014f242 100644 --- a/samsung_appliance/coap_dtls.py +++ b/protocol/dtls_session.py @@ -20,13 +20,25 @@ delivered via the on_notification callback. """ import os import socket -import struct import threading import time +from pathlib import Path from OpenSSL import SSL -from .logger import logger +from .coap import ( + URI_PATH, URI_QUERY, OBSERVE, CONTENT_FORMAT, ACCEPT, BLOCK2, SIZE2, + TYPE_CON, TYPE_NON, TYPE_ACK, TYPE_RST, + METHOD_GET, METHOD_POST, CF_CBOR, + OBSERVE_REGISTER, OBSERVE_DEREGISTER, BLOCK_SZX, + encode_options, parse_coap, build_coap, block_value, fmt_code, + split_dtls as _split_dtls, +) +import logging + +logger = logging.getLogger(__name__) + +_OCF_ROOT_CA = str(Path(__file__).parent / 'ocf_root_ca.pem') # Diagnostic logging — when DEBUG_BRIDGE=1 in env, the bridge dumps @@ -36,137 +48,18 @@ from .logger import logger # new resources and field semantics; otherwise quiet. DEBUG_BRIDGE = os.environ.get('DEBUG_BRIDGE') == '1' +# Per-block retransmission: send up to this many times before giving up. +# Each attempt waits at most _BLOCK_ACK_TIMEOUT seconds (capped by the +# overall deadline). Matches RFC 7252 CON retransmit behaviour. +_BLOCK_MAX_ATTEMPTS = 3 +_BLOCK_ACK_TIMEOUT = 4.0 -# 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 +# Inter-request pacing: minimum seconds between CoAP CON sends on one session. +# Samsung's RT-OCF stacks drop requests when hit faster than their firmware +# ceiling (dryer ~14 req/s, oven ~8 req/s, dishwasher unknown). 5 req/s +# (200 ms) is conservative enough for all tested devices; tune per device +# once the ceiling is measured empirically. +_DEFAULT_RATE_LIMIT_RPS = 5.0 class DtlsCoapSession: @@ -187,13 +80,15 @@ class DtlsCoapSession: MAX_BLOCKS = 32 # safety bound for Block2 fetches def __init__(self, host, port, cert_path, key_path, - on_notification=None, mtu=1200): + on_notification=None, mtu=1200, + rate_limit_rps: float = _DEFAULT_RATE_LIMIT_RPS): 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._min_req_interval = 1.0 / rate_limit_rps self.sock = None self.conn = None @@ -218,6 +113,14 @@ class DtlsCoapSession: self._stop = threading.Event() self._reader_thread = None + self._last_send_ts = 0.0 + + def pace(self) -> None: + """Sleep only the part of the rate-limit interval not already consumed + since the last real send. Uses _stop so session teardown wakes it.""" + remaining = self._min_req_interval - (time.monotonic() - self._last_send_ts) + if remaining > 0: + self._stop.wait(remaining) # ---- lifecycle --------------------------------------------------- @@ -225,9 +128,13 @@ class DtlsCoapSession: """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 — the AC14K_M-rooted chain is SHA-1 signed, which - # OpenSSL 3.x's default security level rejects. + + ctx.load_verify_locations(_OCF_ROOT_CA) + ctx.set_verify(SSL.VERIFY_PEER, lambda conn, cert, err, depth, ok: ok) + # @SECLEVEL=0 permits SHA-1 in Samsung's server cert chain (AC14K_M + # intermediate is SHA-1 signed). This is the only channel that reaches + # the OpenSSL instance cryptography bundles — ctypes and cffi bindings + # do not expose SSL_CTX_set_security_level on this build. 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) @@ -368,8 +275,9 @@ class DtlsCoapSession: with self._send_lock: if self.conn is None: raise ConnectionError("DTLS session closed") - self.conn.send(datagram) try: + self.conn.send(datagram) + self._last_send_ts = time.monotonic() while True: o = self.conn.bio_read(65535) if not o: @@ -515,28 +423,46 @@ class DtlsCoapSession: last_code = None last_opts = [] deadline = time.time() + timeout + szx = BLOCK_SZX # server may negotiate down; track per-transfer while True: - ev = threading.Event() + if num > 0: + self.pace() 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) + for attempt in range(_BLOCK_MAX_ATTEMPTS): + 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, szx))) + self._send_dgram( + build_coap(TYPE_CON, METHOD_GET, mid, tok, opts)) + per_wait = min(_BLOCK_ACK_TIMEOUT, + max(0.1, deadline - time.time())) + if ev.wait(per_wait): + break # got a response + remaining = deadline - time.time() + if remaining <= 0 or attempt == _BLOCK_MAX_ATTEMPTS - 1: + logger.debug( + "GET %s /%s block %d: timed out after %d attempt(s)", + self.host, '/'.join(path_segs), num, attempt + 1, + ) + raise TimeoutError( + f"GET /{'/'.join(path_segs)} block {num} timeout") + logger.debug( + "GET %s /%s block %d: attempt %d/%d timeout, retrying", + self.host, '/'.join(path_segs), num, + attempt + 1, _BLOCK_MAX_ATTEMPTS, + ) + finally: + self._pending.pop(tok, None) + if 'err' in container: + raise ConnectionError(container['err']) code = container['code'] payload = container['payload'] @@ -553,6 +479,9 @@ class DtlsCoapSession: if b2: bv = int.from_bytes(b2[0], 'big') more = (bv >> 3) & 1 + server_szx = bv & 0x07 + if server_szx != szx: + szx = server_szx if not more: break num += 1 diff --git a/protocol/ocf_root_ca.pem b/protocol/ocf_root_ca.pem new file mode 100644 index 0000000..43f34a2 --- /dev/null +++ b/protocol/ocf_root_ca.pem @@ -0,0 +1,15 @@ +-----BEGIN CERTIFICATE----- +MIICXTCCAgGgAwIBAgIBATAMBggqhkjOPQQDAgUAMGsxKDAmBgNVBAMTH1NhbXN1 +bmcgRWxlY3Ryb25pY3MgT0NGIFJvb3QgQ0ExFDASBgNVBAsTC09DRiBSb290IENB +MRwwGgYDVQQKExNTYW1zdW5nIEVsZWN0cm9uaWNzMQswCQYDVQQGEwJLUjAgFw0x +NjExMjQwMjU1MTFaGA8yMDY5MTIzMTE0NTk1OVowazEoMCYGA1UEAxMfU2Ftc3Vu +ZyBFbGVjdHJvbmljcyBPQ0YgUm9vdCBDQTEUMBIGA1UECxMLT0NGIFJvb3QgQ0Ex +HDAaBgNVBAoTE1NhbXN1bmcgRWxlY3Ryb25pY3MxCzAJBgNVBAYTAktSMFkwEwYH +KoZIzj0CAQYIKoZIzj0DAQcDQgAEYsf/Kx+sUFBQESbuytTDPwLPIe0X/8/B1L7b +3abxE9w0gQZAfI8WYUkKfNfP7HXh1M5SCnOkfwWraltGOKTeX6OBkTCBjjAOBgNV +HQ8BAf8EBAMCAcYwMgYDVR0fBCswKTAnoCWgI4YhaHR0cDovL3Byb2RjYS5zYW1z +dW5naW90cy5jb20vY3JsMA8GA1UdEwEB/wQFMAMBAf8wNwYIKwYBBQUHAQEEKzAp +MCcGCCsGAQUFBzABhhtodHRwOi8vb2NzcC5zYW1zdW5naW90cy5jb20wDAYIKoZI +zj0EAwIFAANIADBFAiARY9aSE30q31q5v8B4sJczBqOp7AsD9o8ZIuNmH7IwSwIh +ALfR56jcXoFis/nDx0tQ2xTI/f0b7F6tWPqj2vyKQeNR +-----END CERTIFICATE----- diff --git a/requirements-bootstrap.txt b/requirements-bootstrap.txt index 2360d27..a749cc8 100644 --- a/requirements-bootstrap.txt +++ b/requirements-bootstrap.txt @@ -1,4 +1,4 @@ # Setup-only deps for setup_cert.py: shells out to `openssl` for SHA-1 # signing (independent of python-cryptography's policy) and uses # pyOpenSSL for the optional --test DTLS handshake. --r requirements.txt +-r mqtt_demo/requirements.txt diff --git a/requirements-dev.txt b/requirements-dev.txt new file mode 100644 index 0000000..039d26e --- /dev/null +++ b/requirements-dev.txt @@ -0,0 +1 @@ +pytest>=8.0 diff --git a/samsung_appliance/__init__.py b/samsung_appliance/__init__.py deleted file mode 100644 index e84bf95..0000000 --- a/samsung_appliance/__init__.py +++ /dev/null @@ -1,5 +0,0 @@ -"""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/sensors.py b/samsung_appliance/sensors.py deleted file mode 100644 index 5344492..0000000 --- a/samsung_appliance/sensors.py +++ /dev/null @@ -1,19 +0,0 @@ -"""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 diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..f9953a2 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,7 @@ +import sys +from pathlib import Path + +# Ensure repo root is on sys.path so that protocol/ and ocf/ can be imported +repo_root = Path(__file__).resolve().parent.parent +if str(repo_root) not in sys.path: + sys.path.insert(0, str(repo_root)) diff --git a/tests/test_coap_wire.py b/tests/test_coap_wire.py new file mode 100644 index 0000000..e76fa98 --- /dev/null +++ b/tests/test_coap_wire.py @@ -0,0 +1,45 @@ +from protocol.coap import ( + build_coap, parse_coap, encode_options, block_value, fmt_code, + TYPE_CON, METHOD_GET, URI_PATH, ACCEPT, CF_CBOR, BLOCK2, +) + + +def test_build_then_parse_roundtrip_no_payload(): + opts = [(URI_PATH, b'device'), (URI_PATH, b'0'), (ACCEPT, CF_CBOR)] + datagram = build_coap(TYPE_CON, METHOD_GET, 0xABCD, b'\x01\x02', opts) + mtype, code, mid, tok, parsed_opts, payload = parse_coap(datagram) + assert mtype == TYPE_CON + assert code == METHOD_GET + assert mid == 0xABCD + assert tok == b'\x01\x02' + assert payload == b'' + assert sorted(parsed_opts) == sorted(opts) + + +def test_build_then_parse_roundtrip_with_payload(): + datagram = build_coap(TYPE_CON, 0x45, 1, b'\xff', [], payload=b'\xa1\x01\x02') + _, code, _, _, _, payload = parse_coap(datagram) + assert code == 0x45 + assert payload == b'\xa1\x01\x02' + + +def test_encode_options_orders_by_option_number(): + # ACCEPT (17) must be encoded after URI_PATH (11) regardless of input order + encoded_in_order = encode_options([(URI_PATH, b'x'), (ACCEPT, CF_CBOR)]) + encoded_reversed = encode_options([(ACCEPT, CF_CBOR), (URI_PATH, b'x')]) + assert encoded_in_order == encoded_reversed + + +def test_block_value_encodes_num_more_szx(): + # num=2, more=1, szx=6 -> (2<<4)|(1<<3)|6 = 0x2E + assert block_value(2, 1, 6) == bytes([0x2E]) + + +def test_block_value_promotes_to_two_bytes_when_num_is_large(): + v = block_value(num=0xFFF, more=0, szx=0) + assert len(v) == 2 + + +def test_fmt_code_formats_class_dot_detail(): + assert fmt_code(0x45) == '2.05' + assert fmt_code(0x84) == '4.04' diff --git a/tests/test_import_isolation.py b/tests/test_import_isolation.py new file mode 100644 index 0000000..d2d6119 --- /dev/null +++ b/tests/test_import_isolation.py @@ -0,0 +1,42 @@ +"""protocol/ and ocf/ must be vendorable on their own — no dependency +on mqtt_demo/. This copies just those two directories into an empty +temp dir and imports every module in them there, so a stray +`from mqtt_demo... import ...` fails loudly instead of silently +passing because mqtt_demo/ happens to also be on sys.path in-repo.""" +import os +import shutil +import subprocess +import sys +from pathlib import Path + +REPO_ROOT = Path(__file__).resolve().parent.parent + + +def test_protocol_and_ocf_import_without_mqtt_demo_present(tmp_path): + for pkg in ('protocol', 'ocf'): + shutil.copytree(REPO_ROOT / pkg, tmp_path / pkg) + + import_lines = [ + "import protocol.coap", + "import protocol.dtls_session", + "import ocf.state_cache", + "import ocf.poll_scheduler", + "import ocf.keepalive", + "import ocf.observe_refresh", + ] + script = "\n".join(import_lines) + "\nprint('OK')\n" + + env = dict(os.environ) + env["PYTHONPATH"] = str(tmp_path) + + result = subprocess.run( + [sys.executable, "-c", script], + cwd=str(tmp_path), + env=env, + capture_output=True, text=True, + ) + assert result.returncode == 0, ( + f"protocol/ocf failed to import without mqtt_demo/ present:\n" + f"stdout: {result.stdout}\nstderr: {result.stderr}" + ) + assert "OK" in result.stdout diff --git a/tests/test_state_cache.py b/tests/test_state_cache.py new file mode 100644 index 0000000..a505176 --- /dev/null +++ b/tests/test_state_cache.py @@ -0,0 +1,40 @@ +# tests/test_state_cache.py +from ocf.state_cache import StateCache + + +class _FakeDescriptor: + on_observation = None + + +def test_index_device_tree_skips_device_level_entry_at_index_zero(): + tree = [ + {'href': '/device/0', 'rep': {'n': 'Dryer'}}, # device-level, skipped + {'href': '/mode/vs/0', 'rep': {'x': 1}}, + {'href': '/power/vs/0', 'rep': {'y': 2}}, + ] + indexed = StateCache.index_device_tree(tree) + assert indexed == {'/mode/vs/0': {'x': 1}, '/power/vs/0': {'y': 2}} + assert '/device/0' not in indexed + + +def test_index_device_tree_stub_entry_becomes_empty_dict(): + tree = [ + {'href': '/device/0', 'rep': {}}, + {'href': '/oven/vs/0'}, # no 'rep' key at all — a stub resource + ] + indexed = StateCache.index_device_tree(tree) + assert indexed == {'/oven/vs/0': {}} + + +def test_index_device_tree_non_list_input_returns_empty(): + assert StateCache.index_device_tree({'not': 'a list'}) == {} + assert StateCache.index_device_tree(None) == {} + + +def test_apply_rep_reports_change_and_updates_cache(): + cache = StateCache(_FakeDescriptor()) + changed = cache.apply_rep('/mode/vs/0', {'x': 1}, source='poll') + assert changed is True + assert cache.get('/mode/vs/0') == {'x': 1} + unchanged = cache.apply_rep('/mode/vs/0', {'x': 1}, source='poll') + assert unchanged is False