Files
Jack Nagy 231e88a8c6 fix(protocol): reassemble blockwise OBSERVE notifications
A notification carries only the first block of a large representation
(RFC 7959 §2.6). _dispatch_coap handed that block straight to
on_notification, so consumers decoded a truncated CBOR buffer. Reported
twice on /mode/vs/0: #37 and mbillow/localthings#361.

The recovery is a re-read from block 0 on a fresh 4-byte one-shot token,
not a §2.6 continuation. §3.4 rules out reusing the observation's token,
and this server drops a transfer that opens at NUM>0 under a token it
has not seen, so a continuation is the one shape that cannot work here.
A truncated notification is now withheld and queued to a worker thread
that re-reads the resource and delivers the reassembled representation.
When the re-read fails the notification is dropped at debug level and
the poll tiers carry freshness, which is what they already did.

The re-read has to run off the reader thread: _dispatch_coap runs there
and the transfer waits on an event only that same thread can set. The
worker is serialized and paces between transfers, so a notification
storm stays under the firmware request ceiling.

Also in the Block2 loop, now extracted and shared by both paths:

- compare the response's Block2 NUM against the one requested, so a
  retransmitted block is no longer concatenated as if it were the next
- compare ETags across blocks (§2.4) and restart once when the
  representation changes mid-transfer
- recompute the next block number from the accumulated byte offset when
  the server negotiates the block size down
- re-check reader liveness while waiting on a block, so a mid-transfer
  reader death fails fast instead of burning the whole timeout
- guard the token counters and _pending with a lock, now that the
  session issues concurrent reads of its own

Closes #39
2026-08-15 12:39:48 +01:00

151 lines
4.3 KiB
Python

"""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
from ..errors import MalformedMessageError
# CoAP option numbers (RFC 7252 + 7641 + 7959)
URI_PATH = 11
URI_QUERY = 15
OBSERVE = 6
ETAG = 4
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 MalformedMessageError()
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 MalformedMessageError()
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 block_fields(value):
"""Decode a CoAP Block-N option value. Inverse of block_value().
Returns (num, more, szx). An empty value means block 0, no more,
SZX=0 — RFC 7959 §2.2 allows a zero-length option to elide it."""
v = int.from_bytes(value, 'big')
return v >> 4, (v >> 3) & 1, v & 0x07
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