feat(protocol): surface OBSERVE refetch outcomes under DEBUG_BRIDGE

The refetch path logged only at debug, and the bridge configures logging
at INFO, so a successful re-read and a total failure produced identical
output: nothing. That makes the hardware validation for #39 impossible
to read.

One line per refetch, promoted to INFO when DEBUG_BRIDGE=1 and left at
debug otherwise, naming the href, the one-shot token, the block count,
and the reassembled size. The token is the part that matters: it is what
shows the re-read used a fresh 4-byte token rather than the observation's
1-byte one, which is the assumption the whole design rests on.

Gating on DEBUG_BRIDGE rather than raising the logger keeps the per-block
retransmit lines out of the way, and matches how the module already gates
its frame dump.
This commit is contained in:
Jack Nagy
2026-08-15 13:28:00 +01:00
parent 231e88a8c6
commit e63acb759f
2 changed files with 69 additions and 8 deletions
+30 -8
View File
@@ -717,6 +717,15 @@ class DtlsCoapSession:
# ---- OBSERVE refetch ---------------------------------------------
@staticmethod
def _log_refetch(msg, *args):
"""Refetch outcomes are debug-level in normal operation, which is
below the bridge's INFO default, so a healthy session stays quiet.
DEBUG_BRIDGE=1 promotes them to INFO for hardware validation:
that shows which token the re-read used and whether it completed,
without also turning on every per-block retransmit line."""
(logger.info if DEBUG_BRIDGE else logger.debug)(msg, *args)
def _queue_refetch(self, href):
"""Queue a blockwise notification for re-reading.
@@ -728,7 +737,9 @@ class DtlsCoapSession:
with self._refetch_cond:
if (href not in self._refetch_pending
and len(self._refetch_pending) >= _MAX_PENDING_REFETCH):
logger.debug("observe %s: refetch queue full, dropping", href)
self._log_refetch(
"refetch %s dropped: queue full (%d pending)",
href, len(self._refetch_pending))
return
self._refetch_seq += 1
self._refetch_pending[href] = self._refetch_seq
@@ -779,24 +790,28 @@ class DtlsCoapSession:
self.pace()
segs = [s for s in href.split('/') if s]
try:
code, payload = self._blockwise_get(
code, payload, blocks, tok = self._blockwise_get(
segs, (), _REFETCH_TIMEOUT_S)
except Exception as e:
# Device silent, session gone, ETag never settled, block cap
# hit. Whatever the reason, dropping the notification is the
# contract: the poll tiers still carry freshness, and handing
# over the first block is the bug this replaced.
logger.debug("observe %s: refetch failed (%s)", href, e)
self._log_refetch("refetch %s failed: %s", href, e)
return
if code != 0x45:
logger.debug("observe %s: refetch returned %s",
href, fmt_code(code))
self._log_refetch("refetch %s returned %s", href, fmt_code(code))
return
with self._refetch_cond:
# A newer notification landed while we were reading. That one
# has its own refetch queued, so this result is already stale.
if self._refetch_pending.get(href, 0) > seq:
self._log_refetch(
"refetch %s tok=%s blocks=%d bytes=%d superseded",
href, tok.hex(), blocks, len(payload))
return
self._log_refetch("refetch %s tok=%s blocks=%d bytes=%d ok",
href, tok.hex(), blocks, len(payload))
cb = self.on_notification
if cb is not None:
try:
@@ -814,11 +829,16 @@ class DtlsCoapSession:
token, and dropping a fresh token on block 1+ silently drops
the request."""
self._check_live()
return self._blockwise_get(path_segs, query, timeout)
code, blob, _blocks, _tok = self._blockwise_get(
path_segs, query, timeout)
return code, blob
def _blockwise_get(self, path_segs, query=(), timeout=10.0):
"""Shared token-stable Block2 reassembly (RFC 7959 §2.4).
Returns (code, payload, block_count, token). The last two are
diagnostics for the refetch log; get() drops them.
Mints one fresh 4-byte token and holds it across every block of
the transfer. Also the notification-refetch primitive: RFC 7959
§3.4 forbids continuing a blockwise notification on the
@@ -850,6 +870,7 @@ class DtlsCoapSession:
tok = self._next_tok()
blob = b''
num = 0
blocks = 0
last_code = None
etag = None
deadline = time.time() + timeout
@@ -861,6 +882,7 @@ class DtlsCoapSession:
tok, path_segs, query, num, szx, deadline)
if 'err' in container:
raise container['err']
blocks += 1
code = container['code']
payload = container['payload']
@@ -869,7 +891,7 @@ class DtlsCoapSession:
# 4.xx / 5.xx responses don't carry Block2 continuation —
# bail with whatever we got. Caller decides if 4.xx is fatal.
if code >> 5 != 2:
return code, blob
return code, blob, blocks, tok
# RFC 7959 §2.4: compare ETags across blocks, or we splice
# two versions of the resource into one buffer.
@@ -897,7 +919,7 @@ class DtlsCoapSession:
num += 1
if num > self.MAX_BLOCKS:
raise BlockwiseError()
return last_code, blob
return last_code, blob, blocks, tok
def _exchange_block(self, tok, path_segs, query, num, szx, deadline):
"""Send one block request under `tok` and return its response
+39
View File
@@ -12,6 +12,7 @@ reasons, both recorded on #39: RFC 7959 §3.4 forbids continuing on the
observation's token, and Samsung's RT-OCF drops a transfer that opens
at NUM>0 under a token it has not seen.
"""
import logging
import socket
import threading
import time
@@ -20,6 +21,7 @@ import pytest
from OpenSSL import SSL
from smartthings_local.errors import BlockwiseError
from smartthings_local.protocol import dtls_session
from smartthings_local.protocol.coap import (
BLOCK2, ETAG, METHOD_GET, OBSERVE, TYPE_ACK, TYPE_NON,
block_fields, block_value, build_coap, parse_coap,
@@ -27,6 +29,7 @@ from smartthings_local.protocol.coap import (
from smartthings_local.protocol.dtls_session import DtlsCoapSession
SZX = 6 # 1024-byte blocks, the only size these appliances honour
_LOGGER_NAME = "smartthings_local.protocol.dtls_session"
class _NullAuth:
@@ -371,6 +374,42 @@ def test_refetch_worker_exits_when_the_reader_dies():
sess.close()
def test_debug_bridge_promotes_the_refetch_outcome_to_info(monkeypatch, caplog):
"""The hardware validation for #39 reads this line to confirm which
token the re-read used, so it has to survive the bridge's INFO
default. Without DEBUG_BRIDGE it stays at debug."""
monkeypatch.setattr(dtls_session, "DEBUG_BRIDGE", True)
blocks = [b"A" * 1024, b"B" * 40]
def responder(request):
_mtype, _code, mid, tok, opts, _ = request
if any(n == OBSERVE for n, _ in opts):
return []
num, _ = _requested_block(request)
more = 1 if num + 1 < len(blocks) else 0
return [_content(tok, mid, blocks[num],
block2=block_value(num, more, SZX))]
sess, calls = _make_session(responder)
try:
with caplog.at_level(logging.INFO, logger=_LOGGER_NAME):
tok = sess.subscribe(["mode", "vs", "0"])
sess.conn.inject(
_notification(tok, blocks[0], block2=block_value(0, 1, SZX)))
assert _wait_for(lambda: calls)
line = next((r.getMessage() for r in caplog.records
if r.getMessage().startswith("refetch /mode/vs/0")), None)
assert line is not None, "no refetch line at INFO"
assert "blocks=2" in line
assert f"bytes={sum(len(b) for b in blocks)}" in line
assert line.endswith("ok")
# The token in the line is the one-shot token, not the observe one.
assert f"tok={tok.hex()} " not in line
finally:
_close(sess)
# --------------------------------------------------------------------
# Shared Block2 loop hardening