Skip to content

Commit b4cec57

Browse files
committed
blockwise: stream single-record SenML-CBOR vd blobs over Block1 (Step L15)
ThingsBoard wraps every RPC write in SenML-CBOR for a ct="11542 112" client — including multi-kilobyte blobs, which then ride Block1. The Block1 lane was opaque-only (structured formats drew 4.13 by design), so every TB blob write failed on the first block. This teaches the lane the one SenML shape that can stream: a single-record CBOR pack whose only value is a definite-length vd byte string, on a resource-level PUT. parse_cbor_blob_prefix reads the record header off block 0 (all of it must fit there, or 4.13); the blob bytes then flow into write_chunk blob-relative, and max_block1_size bounds the blob, not the envelope. The wire truth, pinned from a live capture against leshan-demo-server: Leshan (and therefore ThingsBoard) emits the value FIRST and the name AFTER the blob — [{vd, n}] — so the record's name is generally unknown until the final block. Fields trailing the blob buffer (<= 512 B, more is not a name) and parse_cbor_blob_trailing resolves the name at transfer end; the finale — write_chunk(last) inside the transaction bracket — runs only after the resolved name matches the request path. A lying name aborts with 4.00 and nothing commits. Name-leading packs (this library's own encoder order) validate on block 0 instead. SenML-JSON, multi-record packs and non-vd values stay 4.13; envelopes that lie about the blob length draw 4.00 mid-transfer. Tests: 14 blockwise scenarios (both field orders, end-validation abort with no transaction opened, trailing spanning blocks with the empty finale, SZX downshift blob-offset math, cap-excludes-header, trailing cap, unstreamable/malformed matrices), 7 parser vectors including the live-captured Leshan prefix bytes, and a Leshan e2e writing a 4 KiB opaque as SENML_CBOR through real Californium Block1 plus a small-write dispatch-path leg (verified against a native Leshan before push). The blockwise fuzz target gained a streaming write_chunk sink and a SenML Block1 seed (2.3M-run smoke clean). Version 0.1.6.
1 parent d745f36 commit b4cec57

11 files changed

Lines changed: 1679 additions & 25 deletions

File tree

‎Cargo.lock‎

Lines changed: 4 additions & 4 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎Cargo.toml‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,18 +8,18 @@ members = [
88
]
99

1010
[workspace.package]
11-
version = "0.1.5"
11+
version = "0.1.6"
1212
edition = "2024"
1313
rust-version = "1.85"
1414
license = "Apache-2.0"
1515
authors = ["Wojtek Siudzinski"]
1616

1717
[workspace.dependencies]
1818
# Internal
19-
lekki = { path = "crates/lekki", version = "0.1.5" }
20-
lekki-mbedtls = { path = "crates/lekki-mbedtls", version = "0.1.5" }
21-
lekki-testkit = { path = "crates/lekki-testkit", version = "0.1.5" }
22-
mbedtls-host-sys = { path = "crates/mbedtls-host-sys", version = "0.1.5" }
19+
lekki = { path = "crates/lekki", version = "0.1.6" }
20+
lekki-mbedtls = { path = "crates/lekki-mbedtls", version = "0.1.6" }
21+
lekki-testkit = { path = "crates/lekki-testkit", version = "0.1.6" }
22+
mbedtls-host-sys = { path = "crates/mbedtls-host-sys", version = "0.1.6" }
2323

2424
coap-lite = "=0.13.3"
2525
minicbor = "=2.2.2"

‎README.md‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,10 @@ Linux hosts and ESP-IDF (ESP32-S3) devices. *Lekki* is Polish for
1616
## Scope (what this deliberately is not)
1717

1818
Single server, binding U only. No bootstrap, no Queue Mode, no SMS/TCP, no
19-
ACLs. Streaming server→client Block1 uploads, client Block2 pull
20-
downloads (firmware), and client Block1 Register/Update for oversized
19+
ACLs. Streaming server→client Block1 uploads (raw opaque, and — under
20+
`lwm2m11` — single-record SenML-CBOR `vd` blobs, the shape ThingsBoard
21+
wraps every blob write in), client Block2 pull downloads (firmware), and
22+
client Block1 Register/Update for oversized
2123
object lists are in scope. Phase A is LwM2M 1.0 (TLV); Phase B adds
2224
1.1 + SenML-CBOR behind the `lwm2m11` feature; Observe/Notify with
2325
pmin/pmax (poll-and-diff, all-CON) lives behind the default-off `observe`

‎crates/lekki-testkit/tests/leshan_e2e.rs‎

Lines changed: 162 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,11 @@ struct DeviceState {
6666
timezone: String,
6767
reboots: u32,
6868
manifest_len: usize,
69+
/// Streaming-write staging for /33099/0/1 (the blob-upload e2e).
70+
blob_staged: Vec<u8>,
71+
blob_committed: Option<Vec<u8>>,
72+
blob_chunks: usize,
73+
blob_finishing: bool,
6974
}
7075

7176
/// Standard-modeled /3 Device facade (Leshan ships the object model, so
@@ -182,6 +187,94 @@ impl Lwm2mObject for ManifestObject {
182187
}
183188
}
184189

190+
/// Deliberately **unmodeled** object (no XML anywhere): Leshan encodes
191+
/// writes to it purely from the REST request's explicit `type`, which is
192+
/// also how it behaves toward any vendor object it has no model for.
193+
/// /33099/0/1 is the streaming blob sink (the dongle's profile-upload
194+
/// shape, miniaturized): offset 0 resets the staging, `last` awaits its
195+
/// transaction Commit.
196+
struct BlobSinkObject {
197+
state: Arc<Mutex<DeviceState>>,
198+
}
199+
200+
impl Lwm2mObject for BlobSinkObject {
201+
fn oid(&self) -> u16 {
202+
33099
203+
}
204+
205+
fn instances(&mut self, out: &mut InstanceLister) {
206+
out.instance(0);
207+
}
208+
209+
fn resources(&mut self, _iid: u16, out: &mut ResourceLister) {
210+
out.resource(1, ResourceKind::W, Presence::Absent);
211+
}
212+
213+
fn read(&mut self, _: u16, _: u16, _: Option<u16>, _: &mut ResourceWriter<'_>) -> DmResult {
214+
Err(CoapStatus::NotFound)
215+
}
216+
217+
fn write(
218+
&mut self,
219+
iid: u16,
220+
rid: u16,
221+
_riid: Option<u16>,
222+
value: &ResourceReader<'_>,
223+
) -> DmResult {
224+
// A single-datagram blob write is the degenerate one-chunk case —
225+
// the dongle's /33001/0/1 handler does exactly this.
226+
if rid != 1 {
227+
return Err(CoapStatus::MethodNotAllowed);
228+
}
229+
let data = value.opaque()?.to_vec();
230+
self.write_chunk(iid, rid, 0, &data, true)
231+
}
232+
233+
fn write_chunk(
234+
&mut self,
235+
_iid: u16,
236+
rid: u16,
237+
offset: usize,
238+
data: &[u8],
239+
last: bool,
240+
) -> DmResult {
241+
if rid != 1 {
242+
return Err(CoapStatus::MethodNotAllowed);
243+
}
244+
let mut state = self.state.lock().unwrap();
245+
if offset == 0 {
246+
state.blob_staged.clear();
247+
}
248+
if offset != state.blob_staged.len() {
249+
return Err(CoapStatus::RequestEntityIncomplete);
250+
}
251+
state.blob_staged.extend_from_slice(data);
252+
state.blob_chunks += 1;
253+
if last {
254+
state.blob_finishing = true;
255+
}
256+
Ok(())
257+
}
258+
259+
fn transaction(&mut self, t: lekki::dm::Transaction) -> DmResult {
260+
let mut state = self.state.lock().unwrap();
261+
match t {
262+
lekki::dm::Transaction::Begin => {}
263+
lekki::dm::Transaction::Commit if state.blob_finishing => {
264+
let staged = core::mem::take(&mut state.blob_staged);
265+
state.blob_committed = Some(staged);
266+
state.blob_finishing = false;
267+
}
268+
lekki::dm::Transaction::Commit => {}
269+
lekki::dm::Transaction::Rollback => {
270+
state.blob_staged.clear();
271+
state.blob_finishing = false;
272+
}
273+
}
274+
Ok(())
275+
}
276+
}
277+
185278
// ---- client thread ----------------------------------------------------------
186279

187280
enum Cmd {
@@ -214,6 +307,10 @@ impl Dut {
214307
timezone: "UTC".into(),
215308
reboots: 0,
216309
manifest_len: 4096,
310+
blob_staged: Vec::new(),
311+
blob_committed: None,
312+
blob_chunks: 0,
313+
blob_finishing: false,
217314
}));
218315
let (cmd_tx, cmd_rx) = mpsc::channel::<Cmd>();
219316
let (event_tx, event_rx) = mpsc::channel::<ClientEvent>();
@@ -235,6 +332,11 @@ impl Dut {
235332
.unwrap();
236333
client
237334
.register_object(Box::new(ManifestObject {
335+
state: Arc::clone(&device_state),
336+
}))
337+
.unwrap();
338+
client
339+
.register_object(Box::new(BlobSinkObject {
238340
state: device_state,
239341
}))
240342
.unwrap();
@@ -783,6 +885,66 @@ fn composite_read_against_leshan() {
783885
dut.shutdown();
784886
}
785887

888+
/// The ThingsBoard blob-write shape against real Leshan/Californium
889+
/// (Step L15): a single-resource **opaque write big enough to ride
890+
/// Block1**, encoded as SENML_CBOR. TB wraps every RPC write in
891+
/// SenML-CBOR for a `ct="11542 112"` client, so a profile upload arrives
892+
/// exactly like this — one `[{n, vd}]` record chunked by the server.
893+
/// The client must parse the record header off block 0 and stream the
894+
/// blob bytes into `write_chunk` (this is the wire-format pin the unit
895+
/// tests cannot give: the pack here is built by Leshan's own serializer,
896+
/// not ours).
897+
#[cfg(feature = "lwm2m11")]
898+
#[test]
899+
#[ignore = "requires the Leshan docker harness (see module docs)"]
900+
fn blob_write_over_block1_senml_cbor_against_leshan() {
901+
let endpoint = "lekki-e2e-blob";
902+
provision_psk(endpoint);
903+
let dut = Dut::spawn(endpoint);
904+
dut.wait_event(|e| *e == ClientEvent::Registered, Duration::from_secs(15));
905+
906+
// 4 KiB — well past Californium's MAX_MESSAGE_SIZE (1024), so the
907+
// request arrives block-wise.
908+
let blob: Vec<u8> = (0..4096u32).map(|i| (i % 253) as u8).collect();
909+
let body = rest_put(
910+
&format!("/api/clients/{endpoint}/33099/0/1?format=SENML_CBOR&timeout=30"),
911+
&format!(
912+
"{{\"id\":1,\"kind\":\"singleResource\",\"value\":\"{}\",\"type\":\"opaque\"}}",
913+
hex(&blob)
914+
),
915+
);
916+
assert!(body.contains("CHANGED"), "blob write: {body}");
917+
918+
{
919+
let state = dut.state.lock().unwrap();
920+
assert_eq!(
921+
state.blob_committed.as_deref(),
922+
Some(&blob[..]),
923+
"committed blob differs from what Leshan sent"
924+
);
925+
assert!(
926+
state.blob_chunks > 1,
927+
"expected a chunked (Block1) delivery, got {} chunk(s)",
928+
state.blob_chunks
929+
);
930+
}
931+
932+
// A small opaque write in the same format lands through the plain
933+
// dispatch path (no Block1) — both lanes end in the same staging.
934+
let body = rest_put(
935+
&format!("/api/clients/{endpoint}/33099/0/1?format=SENML_CBOR&timeout=10"),
936+
"{\"id\":1,\"kind\":\"singleResource\",\"value\":\"c0ffee\",\"type\":\"opaque\"}",
937+
);
938+
assert!(body.contains("CHANGED"), "small blob write: {body}");
939+
assert_eq!(
940+
dut.state.lock().unwrap().blob_committed.as_deref(),
941+
Some(&[0xC0u8, 0xFF, 0xEE][..]),
942+
"small SenML opaque write did not land"
943+
);
944+
945+
dut.shutdown();
946+
}
947+
786948
// ---- Observe e2e (Step L13, feature `observe`) ------------------------------
787949

788950
/// The plan's four-beat observe proof against a real Californium:

0 commit comments

Comments
 (0)