Merge pull request #8 from mbillow/protocol-library-reorg

Reorg into protocol/ + ocf/ + mqtt_demo/, port 3 DTLS reliability fixes
This commit is contained in:
Quite Yellow
2026-07-06 17:33:00 +01:00
committed by GitHub
34 changed files with 531 additions and 280 deletions
+15
View File
@@ -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/
+50 -37
View File
@@ -4,6 +4,12 @@
<img width="778" height="367" alt="image" src="https://github.com/user-attachments/assets/cc1dca15-f272-4625-a13c-2dc82283ff95" />
> **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_<n>_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_<n>_CLASS` must match a descriptor key in `samsung_appliance/appliances/__init__.py::DESCRIPTORS` — currently `dryer` and `oven`.
Each `APPLIANCE_<n>_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 = <broker-ip>:1883 (user=<mqtt-user>)
14:08:42 INFO samsung_appliance [1] dryer @ <dryer-ip>:49155 (DTLS) → topic samsung_dryer/*
14:08:42 INFO samsung_appliance [2] oven @ <oven-ip>:49154 (DTLS) → topic samsung_oven/*
14:08:42 INFO samsung_appliance MQTT connected → <broker-ip>:1883
14:08:42 INFO mqtt_demo SmartThings-Local Bridge starting (2 appliances)
14:08:42 INFO mqtt_demo broker = <broker-ip>:1883 (user=<mqtt-user>)
14:08:42 INFO mqtt_demo [1] dryer @ <dryer-ip>:49155 (DTLS) → topic samsung_dryer/*
14:08:42 INFO mqtt_demo [2] oven @ <oven-ip>:49154 (DTLS) → topic samsung_oven/*
14:08:42 INFO mqtt_demo MQTT connected → <broker-ip>:1883
14:08:43 INFO dryer DTLS connected — subscribing 11 paths
14:08:44 INFO dryer.<dryer-serial> identified — serial=…
14:08:44 INFO dryer.<dryer-serial> seeded → 25 links; sensors live
@@ -344,45 +350,52 @@ Gated control entities use HA's `availability_mode: all` against `<prefix>/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 `/<x>/vs/0` resources push; the OCF-standard `/<x>/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_<n>_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 `/<x>/0` paths.** They register successfully but never push. Use the Samsung `/<x>/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 `/<x>/vs/0` and `/device/0`.
- **Don't run two clients against the same appliance simultaneously.** Samsung's RT-OCF DTLS allows one active session per peer; a second handshake will get the device to drop the new socket. If HA seems to flap, check whether you've got `main.py` running locally AND the Docker container up.
- **Don't run two clients against the same appliance simultaneously.** Samsung's RT-OCF DTLS allows one active session per peer; a second handshake will get the device to drop the new socket. If HA seems to flap, check whether you've got `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.
---
+8 -5
View File
@@ -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"]
View File
+4 -4
View File
@@ -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():
@@ -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()
+21 -15
View File
@@ -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'"
@@ -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
@@ -1,6 +1,8 @@
services:
smartthings-local:
build: .
build:
context: ..
dockerfile: mqtt_demo/Dockerfile
container_name: smartthings-local
restart: unless-stopped
@@ -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.
* `<class>.<serial>` for per-appliance bridge lines — e.g.
`dryer.<serial>`. 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:
@@ -1,15 +1,15 @@
"""Appliance descriptors registry.
Adding a new appliance class:
1. Write `appliances/<class>.py` with an ApplianceDescriptor named
after the class (uppercase, e.g. OVEN).
1. Write `mqtt_demo/samples/<class>.py` with an ApplianceDescriptor
named after the class (uppercase, e.g. OVEN).
2. Add it to DESCRIPTORS below.
3. Set DEVICE_CLASS=<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
@@ -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 -----------------------------------------------------
@@ -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'
@@ -23,7 +23,7 @@ Untested writes are gated behind <prefix>/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
# ---------------------------------------------------------------------
View File
@@ -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:
@@ -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:
@@ -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():
@@ -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)
View File
+139
View File
@@ -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
@@ -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
+15
View File
@@ -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-----
+1 -1
View File
@@ -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
+1
View File
@@ -0,0 +1 @@
pytest>=8.0
-5
View File
@@ -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"
-19
View File
@@ -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
+7
View File
@@ -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))
+45
View File
@@ -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'
+42
View File
@@ -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
+40
View File
@@ -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