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
151 lines
4.3 KiB
Python
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
|