Skip to content

feat(mqtt): encode embedded packets into one bbqueue write ring (design 054 §4.7) - #283

Merged
lxsaah merged 2 commits into
feat/054-connector-boundaryfrom
feat/054-s08-write-ring
Oct 4, 2026
Merged

lxsaah merged 2 commits into
feat/054-connector-boundaryfrom
feat/054-s08-write-ring

Conversation

@lxsaah

@lxsaah lxsaah commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

Stage 8 of the 054 implementation plan, the first connector stage. The embedded MQTT session no longer allocates a Vec per outbound packet: every packet is encoded straight into one bbqueue ring, allocated once per connector. The action and event channels stay for now (stages 9 and 10).

Change

aimdb-mqtt-connector/src/embedded/write_ring.rs (new): WriteRing

  • A BBQueue<BoxedSlice, AtomicCoord, Polling>, plus two embassy-sync AtomicWakers:
    • data: a commit wakes write_out.
    • room: a release wakes whatever waits for space.
      bbqueue's polling notifier wakes nothing, so these carry the signals.
  • Default 4,096 bytes and CONTROL_RESERVE 64, which gives max_publish() 1,984 bytes. That is half the ring minus the reserve: the most a bipbuffer can always grant contiguously, whatever its offset.
  • Writing:
    • put(packet, reserve): sizes the packet with MqttLenWriter, waits for grant_exact(len + reserve), encodes with MqttBufWriter over the grant, and commits only len.
    • try_put: the lossy variant, for pings.
  • Room checks: has_room(n) is a probe grant dropped uncommitted. poll_room and wait_room are the waker-registering forms.
  • drain() discards leftover bytes. write_out(tx) writes read grants, flushes, releases them and wakes room.

session_loop.rs

  • The Channel<Vec<u8>, 4>, encode, queue, queue_lossy and the old write_out are gone. CONNECT and SUBSCRIBE use put(…, 0), PUBLISH uses put(…, CONTROL_RESERVE), and PINGREQ uses try_put (still lossy).
  • run_session takes &WriteRing and drains it first, so no bytes from an old session reach a new socket ahead of its CONNECT.
  • The action arm's gate: room for max_publish() + CONTROL_RESERVE, checked on every poll through the room waker instead of !outbound.is_full() computed once.
  • Oversized PUBLISH: a frame over max_publish() is skipped with a defmt warning before the client state commits to it.
  • PUBACK: drain_packets waits for PUBACK_ROOM (16 bytes) before parsing each packet, then encodes the PUBACK with try_put before state.receive takes &mut.
    • This deviates from awaiting put inside the parse scope. That version kept the parsed PacketGeneric (32 property slots) inside the session future and grew it to 10,488 bytes, which failed the existing ≤ 8,192 size test.
    • A new test pins the PUBACK mountain-mqtt produces (no properties) at ≤ PUBACK_ROOM.

session.rs, tls.rs: each creates the ring once, before its reconnect loop, and passes it to every session.

Cargo.toml: bbqueue = { version = "0.7", default-features = false, features = ["alloc"] } on the embedded feature (MIT OR Apache-2.0; cargo deny check licenses ok).

Tests

New unit tests in write_ring.rs:

  • A probe that wraps changes no data.
  • An empty ring grants half its capacity at every offset (0–63 on a 64-byte ring), and one byte more fails when drained at the middle (the bipbuffer half-capacity case).
  • max_publish plus the reserve is exactly half the default ring.
  • drain discards an old session's bytes.
  • A PUBACK waits for room during a burst, is woken by write_out's release, and goes out in order.
  • try_put drops instead of waiting.

Existing suites, unchanged and passing:

Suite Passed
session_loop (idle ping cadence, outbound under inbound flood) 6
tokio_broker (reconnect and resubscribe) 3
embassy_broker 1
tls_session 3
tls_broker 3
backend_parity 8
--features std 40
embedded-tls lib tests 36

Review follow-up (dc9a3a4)

  • Fixed here: a hang. put now fails with Overflow when a frame plus its reserve exceeds half the ring. Before, a CONNECT or SUBSCRIBE that large (for example a JWT used as the password) parked forever outside the session's select, with no deadline and no error. Such a packet now ends the session. Refusing it at build is planned for stage 10.
  • fits(len, reserve) and put_sized let perform reuse its encoded length instead of computing it twice.
  • The review's eight proofs are added as tests:
    • The two CONNECT cases are flipped to the new behaviour: a 2,100-byte CONNECT fails at offsets 0 and 2,048, and one that fits goes out from offset 2,048.
    • The other six still describe today's behaviour, and each will flip in the stage that changes it:
      • an oversized publish skipped without a count (stage 10)
      • the 1,984-byte publish cap
      • the 3,584-byte receive cap
      • a 4,000-byte retained message reconnecting forever (fixed by a separate PR to main that advertises Maximum Packet Size)
    • tests/write_ring_proofs.rs (a fake broker) runs in make test and make clippy.
  • Two of this PR's earlier tests filled rings with put beyond half their size; they now fill with try_put.
  • Mutation check: removing the new size guard fails a_connect_over_half_the_ring_fails_at_any_offset.

Verification

  • All of the above.
  • cargo clippy -p aimdb-mqtt-connector --target thumbv7em-none-eabihf --no-default-features --features "embassy-runtime,defmt" -- -D warnings: clean.
  • make clippy (whole workspace) and cargo fmt --all --check: clean.
  • CI does not run on PRs into feat/054-connector-boundary.

🤖 Generated with Claude Code

lxsaah and others added 2 commits October 4, 2026 17:40
…gn 054 §4.7)

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…grant

put() now returns Overflow when a frame plus its reserve exceeds half the
ring, the most a bipbuffer can always grant contiguously. A CONNECT or
SUBSCRIBE that large used to park forever outside the session's select.
perform() reuses its encoded length through put_sized.

Adds the review's size-limit proofs as regression tests, the two CONNECT
cases flipped to the new behaviour, and runs write_ring_proofs in make
test and make clippy.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant