Add oven controls and fix Samsung-OCF write semantics

Major session of local-OCF reverse engineering against the NV7000BS
oven and DV5000T dryer. Surfaces a working set of HA entities for the
oven and resolves several Samsung-quirk regressions in the bridge's
write path.

Key behavioural fixes:
- OBSERVE registrations now use single-byte tokens. Samsung RT-OCF
  silently drops registrations with TKL>1; same 4-byte tokens work
  fine for GET/POST. Symptom was that writes returned 2.04 but the
  appliance never pushed state changes.
- Per-session random starting tokens + MID. Samsung retains observer
  state across DTLS reconnects from the same cert; reusing tokens on
  reconnect silently no-ops.
- OBSERVE deregister sent on DtlsCoapSession.close(), with a stop-
  watcher thread in PushBridge.session_once() so SIGINT/SIGTERM
  actually reaches close() instead of hanging in sess.join().
- pyOpenSSL is not thread-safe — reader-loop conn.* calls now hold
  the same _send_lock the sender uses, dispatching decrypted packets
  outside the lock so the auto-ACK send doesn't deadlock.
- Periodic CoAP Ping (RFC 7252 §4.4) keepalive to keep DTLS warm.
- Post-write Block2 fetchback REMOVED. It was the root cause of
  every "setpoint/operationTime/modes revert ~3s after write"
  symptom — Samsung's stack treats a read on a freshly-written
  resource as a signal to invalidate that write. OBSERVE pushes
  keep HA in sync without the verification GET.

HA-facing changes (oven):
- New entities: Lamp (light), Sound (switch), Fast preheat, Natural
  steam, Setpoint (number), Cook time (number), Stop cycle (button).
- Cooking mode surfaced as a read-only sensor — the oven owns the
  modes field once a cycle is active and rolls local writes back.
- Cook time writes operationTime + remainingTime on
  /operational/state/vs/0 (discovered via OBSERVE capture of
  SmartThings mid-cycle changes — UpperTimerSet on /mode/vs/0
  options is vestigial and doesn't drive the running cycle).
- New cycle_active MQTT availability topic. Writes the oven only
  honours mid-cycle (setpoint, cook time, fast preheat, natural
  steam, stop) gate on it via avail_with_cycle / avail_with_remote_
  and_cycle. Sound + Lamp remain always-available.
- Cycle Start deliberately NOT exposed. Every byte-level approxi-
  mation of SmartThings's working start sequence is rejected at
  the firmware level. Empty discovery payloads remove the previous
  Start button and Cooking-mode select cleanly from HA.

Diagnostics:
- DEBUG_BRIDGE=1 env var enables verbose tracing (rx CON/NON/ACK/
  RST per frame, full link-tree dump at seed, /oic/res directory,
  REP changes on /operational/state, /oven, /power, mode options).
  Quiet in production.
This commit is contained in:
Jack Nagy
2026-05-31 20:43:47 +01:00
parent 85580d81b7
commit 8775f7a930
6 changed files with 475 additions and 133 deletions
+12
View File
@@ -154,6 +154,14 @@ def main():
b.heartbeat()
return loop
def make_ping(b: PushBridge):
def loop():
while not b.stop.is_set():
if b.stop.wait(shared.PING_INTERVAL_S):
break
b.ping_once()
return loop
for b in bridges:
tag = b.app.klass
threads.append(threading.Thread(
@@ -164,6 +172,10 @@ def main():
threads.append(threading.Thread(
target=make_heartbeat(b), daemon=True,
name=f'{tag}-heartbeat'))
if shared.PING_INTERVAL_S > 0:
threads.append(threading.Thread(
target=make_ping(b), daemon=True,
name=f'{tag}-ping'))
stopping = threading.Event()
+35
View File
@@ -61,6 +61,13 @@ class ApplianceDescriptor:
# enabled gate themselves on it.
remote_available_field: Optional[str] = None
# If set, the bridge maintains a third availability topic
# `<prefix>/cycle_active` derived from sensors[cycle_active_field]
# (boolean). HA entities that the appliance only accepts writes for
# while a cycle is running (e.g. oven setpoint, cook time, options
# toggles) gate themselves on it.
cycle_active_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
@@ -99,5 +106,33 @@ def avail_with_remote(avail_topic: str,
]
def avail_with_cycle(avail_topic: str,
cycle_topic: str) -> list[dict]:
return [
{'topic': avail_topic,
'payload_available': 'online',
'payload_not_available': 'offline'},
{'topic': cycle_topic,
'payload_available': 'online',
'payload_not_available': 'offline'},
]
def avail_with_remote_and_cycle(avail_topic: str,
remote_topic: str,
cycle_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'},
{'topic': cycle_topic,
'payload_available': 'online',
'payload_not_available': 'offline'},
]
def encode(cfg: dict) -> bytes:
return json.dumps(cfg).encode()
+185 -89
View File
@@ -26,7 +26,8 @@ import time
from .base import (
ApplianceDescriptor,
avail_base,
avail_with_remote,
avail_with_cycle,
avail_with_remote_and_cycle,
device_block,
encode,
)
@@ -55,29 +56,6 @@ OBSERVE_PATHS = [
]
# ---------------------------------------------------------------------
# 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
@@ -166,6 +144,16 @@ def flatten(links):
rem_min = int(h) * 60 + int(m) + (1 if int(s) > 0 else 0)
except Exception:
pass
# operationTime parsed as minutes — the source of truth for "Cook
# time" in HA (mid-cycle SmartThings updates land here, not in
# /mode/vs/0 UpperTimerSet which is vestigial).
op_min = None
if operation_time:
try:
h, m, s = operation_time.split(':')
op_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')
@@ -201,6 +189,12 @@ def flatten(links):
lamp = _option_value(options, 'UpperLamp') # 'On' / 'Off'
sound = _option_value(options, 'Sound') # 'On' / 'Off'
fastpreheat = _option_value(options, 'fastpreheat') # 'On' / 'Off'
# NaturalSteam only appears in the options array after it's been
# touched in the SmartThings app at least once. Until then it's
# absent, so `_option_value` returns None — surface that as None
# (HA renders as "Unknown") rather than "Off", which would falsely
# imply we know it's disabled.
natural_steam = _option_value(options, 'NaturalSteam') # 'On' / 'Off' / None
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".
@@ -230,10 +224,16 @@ def flatten(links):
return {
'machine_state': machine_state,
# `cycle_active` gates the writable controls in HA. The oven
# only honours setpoint / cook-time / option writes (and Stop)
# while a cycle is active — outside an active cycle, writes
# return 2.04 but get rolled back within ~3s.
'cycle_active': machine_state == 'active',
'oven_state': oven_state,
'progress_percentage': _int(g('/operational/state/vs/0',
'x.com.samsung.da.progressPercentage')),
'operation_time': operation_time,
'operation_time_minutes': op_min,
'completion_time': remaining,
'completion_minutes': rem_min,
'current_temp_c': cur_c,
@@ -250,6 +250,7 @@ def flatten(links):
'lamp': lamp,
'sound': sound,
'fastpreheat': fastpreheat,
'natural_steam': natural_steam,
'timer_state': timer_state,
'timer_current': timer_current,
'timer_set': timer_set,
@@ -303,7 +304,9 @@ def log_state_change(sensors):
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')}")
f"mode={sensors.get('mode')} "
f"timer_set={sensors.get('timer_set_seconds')} "
f"timer_cur={sensors.get('timer_current_seconds')}")
# ---------------------------------------------------------------------
@@ -322,6 +325,9 @@ MODEL = 'OCF oven (TizenRT-iotivity, NV7000BS-class)'
_SENSORS = [
('machine_state', 'Machine state', {'icon': 'mdi:stove'}),
('oven_state', 'Cavity state', {}),
# Cooking mode is read-only via local OCF — the oven owns the
# `modes` field once a cycle is active and rolls back any writes.
('mode', 'Cooking mode', {'icon': 'mdi:tune'}),
('progress_percentage', 'Progress percent',
{'unit_of_measurement': '%', 'state_class': 'measurement'}),
('operation_time', 'Elapsed time', {'icon': 'mdi:timer'}),
@@ -386,19 +392,39 @@ _BINARY_SENSORS = [
# MQTT command-topic suffixes (under <prefix>/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'
CMD_LAMP = 'cmd/lamp'
CMD_SOUND = 'cmd/sound'
CMD_FASTPREHEAT = 'cmd/fastpreheat'
CMD_NATURALSTEAM = 'cmd/naturalsteam'
CMD_POWER = 'cmd/power'
CMD_STOP = 'cmd/stop'
CMD_SETPOINT = 'cmd/setpoint'
CMD_COOK_TIME = 'cmd/cook_time'
# NOTE — no CMD_START or CMD_MODE. Reverse-engineered 2026-05-31:
# * `state='Run'` writes to /operational/state/vs/0 are accepted
# (2.04) and machine briefly goes active, but the oven cavity
# stays Ready (no Preheat) and the cycle self-cancels within
# ~3s. Tried every byte-level approximation of SmartThings's
# working start (matching all four fields on /operational/state,
# +operationTime, +remainingTime, +progressPercentage='1', plus
# /temperatures/vs/0 desired, with and without /mode/vs/0 modes,
# with PUT vs POST, paced 1s apart, with OCF-version-options
# 2049/2053, with Samsung vendor-option 65524=0xc0) — none of
# these engage the cavity. The differentiator must be something
# invisible at the OBSERVE-push level (likely a cloud-mediated
# auth path the SmartThings app uses). See project_oven_remote
# _start_open.md for full notes.
# * `modes=['Convection']` writes to /mode/vs/0 succeed (2.04)
# but the oven owns the field once a cycle is active and rolls
# local writes back to ['NoOperation']. mode is surfaced as a
# read-only sensor instead.
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"
cycle_topic = f"{topic_prefix}/cycle_active"
dev = device_block(topic_prefix, device_name, MODEL)
out = []
@@ -456,17 +482,35 @@ def build_discovery(topic_prefix, ha_prefix, device_name):
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.
# --- switches. Sound is always-available — independent of cycle
# state, no RC required. Fast preheat + Natural steam are
# options-array writes the oven only honours mid-cycle, so they
# gate on RC + cycle_active. Power deliberately omitted: turning
# the oven on is a physical-panel action; read-only power_state
# sensor reflects its state.
cfg = {
'name': 'Sound',
'unique_id': f"{topic_prefix}_sound_switch",
'object_id': f"{topic_prefix}_sound_switch",
'state_topic': state_topic,
'value_template': '{{ value_json.sound }}',
'state_on': 'On',
'state_off': 'Off',
'command_topic': f"{topic_prefix}/{CMD_SOUND}",
'payload_on': 'On',
'payload_off': 'Off',
'icon': 'mdi:volume-high',
'availability': avail_base(avail_topic),
'device': dev,
}
out.append((f"{ha_prefix}/switch/{topic_prefix}/sound/config",
encode(cfg)))
cycle_switches = [
('fastpreheat', 'Fast preheat', '{{ value_json.fastpreheat }}', CMD_FASTPREHEAT, 'mdi:fire'),
('natural_steam', 'Natural steam', '{{ value_json.natural_steam }}', CMD_NATURALSTEAM, 'mdi:kettle-steam'),
]
for key, name, tpl, cmd, icon in untested_switches:
for key, name, tpl, cmd, icon in cycle_switches:
cfg = {
'name': name,
'unique_id': f"{topic_prefix}_{key}_switch",
@@ -479,7 +523,8 @@ def build_discovery(topic_prefix, ha_prefix, device_name):
'payload_on': 'On',
'payload_off': 'Off',
'icon': icon,
'availability': avail_with_remote(avail_topic, remote_topic),
'availability': avail_with_remote_and_cycle(
avail_topic, remote_topic, cycle_topic),
'availability_mode': 'all',
'device': dev,
}
@@ -501,48 +546,70 @@ def build_discovery(topic_prefix, ha_prefix, device_name):
'device_class': 'temperature',
'mode': 'slider',
'icon': 'mdi:thermometer-chevron-up',
'availability': avail_with_remote(avail_topic, remote_topic),
# RC + cycle_active gated — Samsung's local-OCF surface only
# honours setpoint changes while a cycle is actually running
# (idle writes get rolled back within ~3s).
'availability': avail_with_remote_and_cycle(
avail_topic, remote_topic, cycle_topic),
'availability_mode': 'all',
'device': dev,
}
out.append((f"{ha_prefix}/number/{topic_prefix}/setpoint/config",
encode(cfg)))
# --- select: mode (RC-gated) ------------------------------------
# --- button: Stop cycle ----------------------------------------
# Gated on cycle_active — there's nothing to stop when idle.
# There is no Start button: local-OCF cycle start is not
# reproducible on this firmware (see project_oven_remote_start
# _open.md memory note for the full investigation). Cooking mode
# is similarly omitted — read-only via local OCF, surfaced as a
# sensor.
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),
'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_with_cycle(avail_topic, cycle_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)))
# --- number: Cook time in minutes (RC + cycle gated; the oven
# only honours operationTime writes while running). Source of
# truth is `operationTime` on /operational/state/vs/0;
# SmartThings's mid-cycle time changes land in that same field.
cfg = {
'name': 'Cook time',
'unique_id': f"{topic_prefix}_cook_time",
'object_id': f"{topic_prefix}_cook_time",
'state_topic': state_topic,
'value_template': '{{ value_json.operation_time_minutes | int(0) }}',
'command_topic': f"{topic_prefix}/{CMD_COOK_TIME}",
'min': 0,
'max': 1439, # 23:59 — matches modeSpec timeMax
'step': 1,
'unit_of_measurement': 'min',
'mode': 'box',
'icon': 'mdi:timer',
'availability': avail_with_remote_and_cycle(
avail_topic, remote_topic, cycle_topic),
'availability_mode': 'all',
'device': dev,
}
out.append((f"{ha_prefix}/number/{topic_prefix}/cook_time/config",
encode(cfg)))
# --- removal: publish empty payload to the discovery topics of
# entities we used to expose. HA treats an empty retained payload
# on a discovery topic as "delete this entity", so previously-set
# up Start buttons and Cooking-mode selects disappear cleanly.
out.append((f"{ha_prefix}/button/{topic_prefix}/start/config", b''))
out.append((f"{ha_prefix}/select/{topic_prefix}/mode/config", b''))
return out
@@ -602,31 +669,39 @@ def command_handlers():
opts, 'fastpreheat', p),
}
def _naturalsteam(p, links):
# NaturalSteam_* only appears in the options array after the
# SmartThings app has touched it once. If absent, append the
# slot — the oven creates it on first write, so the bridge
# doesn't need a "prime via app" dance.
if p not in ('On', 'Off'):
return None
opts = _mode_options(links)
if opts is None:
return None
if not any(o.startswith('NaturalSteam_') for o in opts):
opts = opts + [f'NaturalSteam_{p}']
else:
opts = _replace_in_options(opts, 'NaturalSteam', p)
return ['mode', 'vs', '0'], {
'x.com.samsung.da.options': opts,
}
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
@@ -638,14 +713,34 @@ def command_handlers():
'x.com.samsung.da.items': items,
}
def _cook_time(p, links):
# HA sends minutes; oven cycle duration lives in
# /operational/state/vs/0 as `operationTime` / `remainingTime`
# (H:MM:SS strings). Writing both — mirroring SmartThings's
# observed behaviour, which resets the live countdown to the
# new duration. Clamp to modeSpec's 0..23:59.
try:
minutes = int(round(float(p)))
except (TypeError, ValueError):
return None
if not (0 <= minutes <= 1439):
return None
h, m = divmod(minutes, 60)
hms = f"{h:02d}:{m:02d}:00"
return ['operational', 'state', 'vs', '0'], {
'x.com.samsung.da.operationTime': hms,
'x.com.samsung.da.remainingTime': hms,
}
return {
CMD_LAMP: _lamp,
CMD_SOUND: _sound,
CMD_FASTPREHEAT: _fastpreheat,
CMD_POWER: _power,
CMD_STOP: _stop,
CMD_MODE: _mode,
CMD_SETPOINT: _setpoint,
CMD_LAMP: _lamp,
CMD_SOUND: _sound,
CMD_FASTPREHEAT: _fastpreheat,
CMD_NATURALSTEAM: _naturalsteam,
CMD_POWER: _power,
CMD_STOP: _stop,
CMD_SETPOINT: _setpoint,
CMD_COOK_TIME: _cook_time,
}
@@ -661,5 +756,6 @@ OVEN = ApplianceDescriptor(
on_observation=on_observation,
project=project,
remote_available_field='remote_control_binary',
cycle_active_field='cycle_active',
log_state_change=log_state_change,
)
+119 -21
View File
@@ -20,6 +20,7 @@ Multiple PushBridges run concurrently in a single process — see
main.py. They share one MQTT client; each owns one DTLS session.
"""
import json
import os
import threading
import time
@@ -32,6 +33,16 @@ from .logger import bridge_logger, logger as module_logger
from .sensors import index_links
# Re-enable with `DEBUG_BRIDGE=1` env var. When on, the bridge dumps:
# - every received CoAP frame (in coap_dtls.py)
# - the full /device/0 link tree at seed time
# - /oic/res directory
# - REP changes for /operational/state, /oven, /power, /mode-options
# Designed for reverse-engineering new oven/dryer behaviour next
# session — leave at 0 in production.
DEBUG_BRIDGE = os.environ.get('DEBUG_BRIDGE') == '1'
def _href_to_segs(href: str) -> list[str]:
"""`/mode/vs/0` → `['mode', 'vs', '0']`. Used to translate an
OBSERVE-notification href back into the path-segs the Block2 GET
@@ -76,6 +87,7 @@ class PushBridge:
self.last_state_pub = None
self.last_remote_pub = None
self.last_cycle_pub = None
self.stop = threading.Event()
self.started_ts = time.time()
self.session_started_ts = None
@@ -100,6 +112,7 @@ class PushBridge:
self.state_topic = f"{p}/state"
self.avail_topic = f"{p}/availability"
self.remote_topic = f"{p}/remote_available"
self.cycle_topic = f"{p}/cycle_active"
self.health_topic = f"{p}/bridge/health"
self.cmd_handlers = descriptor.command_handlers()
self.cmd_topic_prefix = f"{p}/cmd/"
@@ -136,6 +149,17 @@ class PushBridge:
def _apply_rep(self, href, rep):
"""Update self.links + fire descriptor hooks + maybe publish.
Shared between the OBSERVE path and the Block2 fetch-back path."""
if DEBUG_BRIDGE:
# /mode/vs/0 carries a huge modeSpec JSON we don't want to
# log; surface just modes + options. Small resources get
# full-rep dumps.
if href == '/mode/vs/0' and isinstance(rep, dict):
self.log.info("mode modes=%r options=%r",
rep.get('x.com.samsung.da.modes'),
rep.get('x.com.samsung.da.options'))
elif href in ('/operational/state/vs/0', '/oven/vs/0',
'/power/vs/0'):
self.log.info("REP %s = %r", href, rep)
self.links[href] = rep
hook = self.descriptor.on_observation
if hook is not None:
@@ -254,6 +278,35 @@ class PushBridge:
sess.start_reader()
# When the bridge is asked to stop, close the session — which
# fires the OBSERVE-dereg sequence and tears DTLS down. Without
# this, session_once would block forever in sess.join() because
# the reader only exits on stop/socket-death.
session_ended = threading.Event()
def _stop_watcher():
while not session_ended.is_set():
if self.stop.wait(1.0):
try:
sess.close()
except Exception as e:
self.log.warning("stop close: %s", e)
return
threading.Thread(target=_stop_watcher, daemon=True,
name=f'{self.app.klass}-stopw').start()
try:
self._run_session_inner(sess)
finally:
# Always release the watcher so it doesn't sit pinned on
# self.stop forever after the session ends.
session_ended.set()
def _run_session_inner(self, sess):
"""Body of session_once after start_reader() — split out so
session_once's try/finally cleanly bounds the stop-watcher
thread's lifetime."""
for path in self.descriptor.observe_paths:
sess.subscribe(path)
time.sleep(0.05)
@@ -276,6 +329,30 @@ class PushBridge:
# serial-tagged.
self._retag_logger_with_serial()
if DEBUG_BRIDGE:
# Dump every link's rep (skipping the huge modeSpec on
# /mode/vs/0) + the /oic/res directory. Useful for finding
# new writable resources next session.
for href, rep in sorted(self.links.items()):
if href == '/mode/vs/0':
short = {k: v for k, v in rep.items()
if k not in (
'x.com.samsung.da.modeSpec',
'x.com.samsung.da.supportedModes',
)}
self.log.info("LINK %s = %r", href, short)
else:
self.log.info("LINK %s = %r", href, rep)
try:
code, pl = sess.get(['oic', 'res'], timeout=10.0)
if code == 0x45 and pl:
self.log.info("OIC_RES = %r", cbor2.loads(pl))
else:
self.log.info("oic/res → %s", fmt_code(code))
except Exception as e:
self.log.warning("oic/res get: %s", e)
hook = self.descriptor.on_observation
if hook is not None:
for href, rep in self.links.items():
@@ -345,6 +422,9 @@ class PushBridge:
field = self.descriptor.remote_available_field
if field is not None:
self.publish_remote_available(sensors.get(field))
cycle_field = self.descriptor.cycle_active_field
if cycle_field is not None:
self.publish_cycle_active(sensors.get(cycle_field))
if not force:
log_fn = self.descriptor.log_state_change
extra = log_fn(sensors) if log_fn is not None else ''
@@ -362,6 +442,17 @@ class PushBridge:
except Exception as e:
self.log.warning("remote_available publish: %s", e)
def publish_cycle_active(self, cycle_on, force=False):
value = 'online' if cycle_on else 'offline'
if not force and value == self.last_cycle_pub:
return
self.last_cycle_pub = value
try:
self.mqtt.publish(self.cycle_topic, value, qos=1, retain=True)
self.log.info("cycle_active → %s", value)
except Exception as e:
self.log.warning("cycle_active publish: %s", e)
def reassert_availability(self):
if self.session is None or self.last_state_pub is None:
return
@@ -370,6 +461,10 @@ class PushBridge:
if field is not None:
self.publish_remote_available(
self.last_state_pub.get(field), force=True)
cycle_field = self.descriptor.cycle_active_field
if cycle_field is not None:
self.publish_cycle_active(
self.last_state_pub.get(cycle_field), force=True)
def set_availability(self, online):
try:
@@ -385,6 +480,13 @@ class PushBridge:
qos=1, retain=True)
except Exception:
pass
if not online and self.descriptor.cycle_active_field is not None:
self.last_cycle_pub = None
try:
self.mqtt.publish(self.cycle_topic, 'offline',
qos=1, retain=True)
except Exception:
pass
# ---- MQTT command handling --------------------------------------
@@ -418,30 +520,15 @@ class PushBridge:
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.
# Optimistic publish — apply the write to our local state.
# No fetchback: empirically (2026-05-31) the post-write
# Block2 GET was causing the appliance to roll our values
# back ~3s later, on every writable resource. OBSERVE
# pushes keep HA in sync without polling. (Was the root
# cause of the "mid-cycle setpoint reverts" symptom.)
self._apply_optimistic(href, body)
# 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()
@@ -475,6 +562,17 @@ class PushBridge:
except Exception as e:
self.log.warning("publish tick: %s", e)
def ping_once(self):
"""Send one CoAP Ping. Caller (main.py's per-bridge ping
thread) handles the cadence. No-op if no live session."""
sess = self.session
if sess is None:
return
try:
sess.ping()
except Exception as e:
self.log.warning("ping: %s", e)
def run_forever(self):
threading.Thread(target=self._publish_tick_loop, daemon=True,
name=f'{self.app.klass}-tick').start()
+122 -23
View File
@@ -18,6 +18,7 @@ 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 os
import socket
import struct
import threading
@@ -28,6 +29,14 @@ from OpenSSL import SSL
from .logger import logger
# Diagnostic logging — when DEBUG_BRIDGE=1 in env, the bridge dumps
# every received CoAP frame, every /operational/state/vs/0 + /oven/vs/0
# + /power/vs/0 + /mode/vs/0-options rep change, the full link tree at
# seed time, and the /oic/res directory. Useful for reverse-engineering
# new resources and field semantics; otherwise quiet.
DEBUG_BRIDGE = os.environ.get('DEBUG_BRIDGE') == '1'
# CoAP option numbers (RFC 7252 + 7641 + 7959)
URI_PATH = 11
URI_QUERY = 15
@@ -191,8 +200,17 @@ class DtlsCoapSession:
self.dest = None
self._send_lock = threading.Lock()
self._mid = 0x5000
self._tok_counter = 0
# Randomize MID and token counter starting points so reconnects
# don't reuse identifiers from previous sessions — Samsung's
# RT-OCF appears to remember observer state across DTLS
# sessions, and re-registering with a token it still thinks is
# active is silently no-ops.
self._mid = int.from_bytes(os.urandom(2), 'big')
self._tok_counter = int.from_bytes(os.urandom(4), 'big')
# OBSERVE tokens are 1-byte (Samsung silently drops TKL>1
# OBSERVE registrations). Pick a random starting byte in the
# 0x40..0xff range so each session uses fresh values.
self._observe_tok_counter = 0x40 + (os.urandom(1)[0] & 0xBF)
# token (bytes) → (Event, container_dict)
self._pending = {}
# token (bytes) → href (str)
@@ -273,9 +291,35 @@ class DtlsCoapSession:
if self._reader_thread is not None:
self._reader_thread.join()
def _send_observe_dereg(self, tok, path_segs):
"""Send a single OBSERVE deregister GET (Observe option = 1)
on the existing token. Best-effort — caller swallows errors."""
if self.conn is None:
return
mid = self._next_mid()
opts = [(URI_PATH, s.encode()) for s in path_segs]
opts.append((OBSERVE, OBSERVE_DEREGISTER))
opts.append((ACCEPT, CF_CBOR))
self._send_dgram(
build_coap(TYPE_CON, METHOD_GET, mid, tok, opts))
def close(self):
"""Tear down session. Signals reader_loop and wakes pending
waiters with an error."""
"""Tear down session. Sends best-effort OBSERVE deregisters
first so Samsung's RT-OCF cleans up its observer table —
without this, the per-cert observer state survives DTLS close
and a quick reconnect with the same tokens silently no-ops."""
# Send dereg for every active observation while the conn is
# still healthy. Tiny sleep lets the records reach the wire
# before we shut DTLS down.
if self.conn is not None and self._observe_tokens:
for tok, href in list(self._observe_tokens.items()):
segs = [s for s in href.split('/') if s]
try:
self._send_observe_dereg(tok, segs)
except Exception as e:
logger.warning("dereg %s: %s", href, e)
time.sleep(0.1)
self._stop.set()
if self.conn is not None:
try:
@@ -307,6 +351,19 @@ class DtlsCoapSession:
# avoids collisions across long-running OBSERVE subscriptions.
return self._tok_counter.to_bytes(4, 'big')
def _next_observe_tok(self):
# Single-byte tokens for OBSERVE registrations. Samsung
# RT-OCF accepts these but silently drops TKL=4 OBSERVE
# registrations. Counter is randomly seeded per session so
# reconnects don't collide with stale observer state Samsung
# may still be holding from the previous run.
self._observe_tok_counter = (self._observe_tok_counter + 1) & 0xFF
# Avoid 0x00 — some CoAP stacks treat an all-zero token as
# equivalent to "no token" / empty (TKL=0).
if self._observe_tok_counter == 0:
self._observe_tok_counter = 1
return bytes([self._observe_tok_counter])
def _send_dgram(self, datagram):
"""Send a CoAP datagram. Holds the send lock for the
BIO-drain so two writers can't interleave records."""
@@ -340,31 +397,46 @@ class DtlsCoapSession:
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:
# pyOpenSSL's SSL.Connection is not thread-safe — the
# same SSL object must not be touched by multiple
# threads concurrently. Drain decrypted records into a
# local list under _send_lock so the reader never races
# a sender's conn.send()/bio_read(). Dispatch happens
# AFTER releasing the lock because _dispatch_coap may
# call _send_dgram (auto-ACK for CON frames), which
# re-acquires the lock — holding it across dispatch
# would deadlock.
packets = []
exit_reader = False
with self._send_lock:
try:
pl = conn.recv(65535)
except SSL.WantReadError:
break
except SSL.ZeroReturnError:
logger.info("DTLS peer closed connection")
return
conn.bio_write(d)
except SSL.Error as e:
logger.warning("DTLS recv: %s", e)
logger.warning("DTLS bio_write: %s", e)
return
if not pl:
break
while True:
try:
pl = conn.recv(65535)
except SSL.WantReadError:
break
except SSL.ZeroReturnError:
logger.info("DTLS peer closed connection")
exit_reader = True
break
except SSL.Error as e:
logger.warning("DTLS recv: %s", e)
exit_reader = True
break
if not pl:
break
packets.append(pl)
for pl in packets:
try:
self._dispatch_coap(pl)
except Exception as e:
logger.warning("dispatch: %s", e)
if exit_reader:
return
finally:
# Make sure pending waiters don't hang if the reader dies.
for tok, (ev, container) in list(self._pending.items()):
@@ -378,6 +450,12 @@ class DtlsCoapSession:
logger.debug("malformed CoAP: %s", e)
return
if DEBUG_BRIDGE:
kind = ['CON', 'NON', 'ACK', 'RST'][mt]
logger.info("rx %s code=%s mid=%04x tok=%s opts=%d pl=%d",
kind, fmt_code(code), mid, tok.hex() or '-',
len(ropts), len(payload))
# 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:
@@ -512,6 +590,27 @@ class DtlsCoapSession:
finally:
self._pending.pop(tok, None)
def ping(self):
"""RFC 7252 §4.4 CoAP Ping — empty CON, no token, no payload.
Peer MUST respond with an RST sharing the same Message ID.
We don't block waiting for the RST here; the reader thread
sees it, finds nothing in _pending/_observe_tokens for an
empty token, and drops it quietly — which is the documented
behaviour for unsolicited RST.
Sole purpose is keepalive: gives Samsung's RT-OCF stack a
visible "client still here" signal so it doesn't expire our
OBSERVE subscriptions (symptom of expiry: POSTs still get
2.04 but OBSERVE notifications stop arriving)."""
if self.conn is None:
raise ConnectionError("DTLS session closed")
mid = self._next_mid()
# Empty message: ver=01, type=CON, tkl=0, code=0.00, mid, no
# token, no options, no payload. build_coap with empty token
# and no options yields exactly that.
self._send_dgram(build_coap(TYPE_CON, 0, mid, b'', []))
return mid
def subscribe(self, path_segs):
"""Register an OBSERVE on the given path. The initial 2.05
notification and all subsequent state-change notifications
@@ -521,7 +620,7 @@ class DtlsCoapSession:
later)."""
if self.conn is None:
raise ConnectionError("DTLS session closed")
tok = self._next_tok()
tok = self._next_observe_tok()
href = '/' + '/'.join(path_segs)
# Register the token BEFORE sending — otherwise the device
# could respond between send() and the dict insert, and the
+2
View File
@@ -58,6 +58,7 @@ class SharedConfig:
HA_DISCOVERY_PREFIX: str
HEALTH_INTERVAL_S: int
HEARTBEAT_INTERVAL_S: int
PING_INTERVAL_S: int
@classmethod
def from_env(cls) -> 'SharedConfig':
@@ -73,6 +74,7 @@ class SharedConfig:
HEALTH_INTERVAL_S=int(os.getenv('HEALTH_INTERVAL_S', '60')),
HEARTBEAT_INTERVAL_S=int(os.getenv('HEARTBEAT_INTERVAL_S',
'600')),
PING_INTERVAL_S=int(os.getenv('PING_INTERVAL_S', '25')),
)