diff --git a/main.py b/main.py index 8e0700e..830fa38 100644 --- a/main.py +++ b/main.py @@ -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() diff --git a/samsung_appliance/appliances/base.py b/samsung_appliance/appliances/base.py index 96dd667..d363949 100644 --- a/samsung_appliance/appliances/base.py +++ b/samsung_appliance/appliances/base.py @@ -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 + # `/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() diff --git a/samsung_appliance/appliances/oven.py b/samsung_appliance/appliances/oven.py index 5de6ccd..91ac812 100644 --- a/samsung_appliance/appliances/oven.py +++ b/samsung_appliance/appliances/oven.py @@ -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 /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, ) diff --git a/samsung_appliance/bridge.py b/samsung_appliance/bridge.py index 5b54ccc..7dae272 100644 --- a/samsung_appliance/bridge.py +++ b/samsung_appliance/bridge.py @@ -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() diff --git a/samsung_appliance/coap_dtls.py b/samsung_appliance/coap_dtls.py index 9f84056..396d333 100644 --- a/samsung_appliance/coap_dtls.py +++ b/samsung_appliance/coap_dtls.py @@ -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 diff --git a/samsung_appliance/config.py b/samsung_appliance/config.py index 7a15468..6c46188 100644 --- a/samsung_appliance/config.py +++ b/samsung_appliance/config.py @@ -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')), )