From 51ca2d1ea2584cf70043f27ce0d8c17cb1e84e72 Mon Sep 17 00:00:00 2001 From: "sounds.like.lx" <147444674+lxsaah@users.noreply.github.com> Date: Sun, 4 Oct 2026 21:08:42 +0200 Subject: [PATCH] fix(mqtt): advertise Maximum Packet Size in the embedded CONNECT (#284) --- CHANGELOG.md | 7 ++ aimdb-mqtt-connector/CHANGELOG.md | 10 +++ .../src/embedded/session_loop.rs | 69 +++++++++++++++++- aimdb-mqtt-connector/tests/common/mod.rs | 48 +++++++++++-- aimdb-mqtt-connector/tests/tokio_broker.rs | 72 +++++++++++++++++++ 5 files changed, 200 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 620b1899..cd7465e3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -53,6 +53,13 @@ and `pump_client` take the router. KNX, WebSocket, TCP, UDS and serial use `ExactGrammar`: a `{…}` link on them fails the build. The user-facing link API is unchanged. ([aimdb-core](aimdb-core/CHANGELOG.md)) +### Fixed + +- **The embedded MQTT backend advertises the largest packet it receives** + (`Maximum Packet Size` 3,328 in CONNECT), so an oversized retained message is + withheld by the broker instead of reconnecting the client forever. + ([aimdb-mqtt-connector](aimdb-mqtt-connector/CHANGELOG.md)) + ## [2.0.0] - 2026-09-18 ### Added diff --git a/aimdb-mqtt-connector/CHANGELOG.md b/aimdb-mqtt-connector/CHANGELOG.md index 4f0f62e0..5d499e8d 100644 --- a/aimdb-mqtt-connector/CHANGELOG.md +++ b/aimdb-mqtt-connector/CHANGELOG.md @@ -23,6 +23,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **The `with_qos` doc no longer claims an inbound subscribe QoS.** Inbound subscriptions stay at QoS 1, as before. +### Fixed + +- **An oversized retained message no longer reconnects the `Embedded` backend + forever.** The session receives packets of up to 3,328 bytes (its 3,584-byte + buffer minus one read), but its MQTT 5 CONNECT did not say so. A broker + could send a larger packet, which ends the session, and a retained one is + replayed after every SUBSCRIBE, so the client reconnected once per + reconnection delay. CONNECT now advertises `Maximum Packet Size` 3,328, and + the broker withholds anything larger instead of sending it. + ## [0.7.0] - 2026-09-18 ### Changed (breaking) diff --git a/aimdb-mqtt-connector/src/embedded/session_loop.rs b/aimdb-mqtt-connector/src/embedded/session_loop.rs index ec8975aa..2533670d 100644 --- a/aimdb-mqtt-connector/src/embedded/session_loop.rs +++ b/aimdb-mqtt-connector/src/embedded/session_loop.rs @@ -46,6 +46,13 @@ const RX_CHUNK: usize = 256; /// rather than a fixed buffer. const PACKET_BUFFER_SIZE: usize = BUFFER_SIZE - 2 * RX_CHUNK; +/// The largest packet the reader takes whatever arrives before it: the +/// buffer minus one feed chunk (see [`PacketReader`]). Advertised as the +/// CONNECT's Maximum Packet Size, so a broker never sends a larger packet — +/// one that would end the session, and a retained one would end every +/// session it is replayed into. +const MAX_INBOUND_PACKET: usize = PACKET_BUFFER_SIZE - RX_CHUNK; + /// One chunk of freshly read bytes, in flight from the read half to the loop. type Chunk = heapless::Vec; @@ -176,10 +183,13 @@ async fn client_loop( // Topic aliases are declined: honouring them would mean storing the // server's topic names for the life of the connection. let _ = properties.push(ConnectProperty::TopicAliasMaximum(0.into())); + let _ = properties.push(ConnectProperty::MaximumPacketSize( + (MAX_INBOUND_PACKET as u32).into(), + )); // Ours, not `connection_settings.keep_alive()`: that field has no // setter, so it is always mountain-mqtt's own 60 s constant. The // cadence below is derived from the value we actually send. - let connect: Connect<'_, 1, 0> = Connect::new( + let connect: Connect<'_, 2, 0> = Connect::new( settings.keep_alive_secs, *connection_settings.username(), *connection_settings.password(), @@ -664,6 +674,63 @@ mod tests { ); } + /// A QoS 0 PUBLISH to `t` that is exactly `total` bytes on the wire. + fn publish_of(total: usize) -> Vec { + let varint_len = if total - 2 < 128 { 1 } else { 2 }; + let remaining = total - 1 - varint_len; + let mut bytes = alloc::vec![0x30u8]; + if varint_len == 1 { + bytes.push(remaining as u8); + } else { + bytes.push((remaining % 128) as u8 | 0x80); + bytes.push((remaining / 128) as u8); + } + bytes.extend_from_slice(&[0x00, 0x01, b't', 0x00]); + bytes.resize(total, b'x'); + bytes + } + + /// Feeds a `first`-byte packet, a `second`-byte one and a trailing one in + /// `RX_CHUNK` reads, consuming packets as they complete, as the session + /// does. The trailing packet makes the read that completes `second` a full + /// one that also carries the head of the next packet: the worst case. + fn receive_after(first: usize, second: usize) -> Result<(), PacketReadError> { + let mut stream = publish_of(first); + stream.extend_from_slice(&publish_of(second)); + stream.extend_from_slice(&publish_of(RX_CHUNK)); + let mut reader = PacketReader::::new(); + let mut received = 0; + for chunk in stream.chunks(RX_CHUNK) { + reader.feed(chunk)?; + while let Some(total) = reader.framed_len()? { + reader.consume(total); + received += 1; + } + } + assert!(received >= 2); + Ok(()) + } + + /// The advertised Maximum Packet Size is one the reader takes wherever the + /// packet starts inside a read. The reader's stated limit (buffer minus + /// one read) is conservative by one byte; two bytes more fail at some + /// offset. + #[test] + fn the_advertised_maximum_packet_size_is_always_received() { + assert_eq!(MAX_INBOUND_PACKET, 3328); + let offsets = 8..8 + RX_CHUNK; + for first in offsets.clone() { + assert_eq!( + receive_after(first, MAX_INBOUND_PACKET), + Ok(()), + "after {first} bytes" + ); + } + assert!(offsets + .map(|first| receive_after(first, MAX_INBOUND_PACKET + 2)) + .any(|r| r == Err(PacketReadError::PacketTooLargeForBuffer))); + } + #[test] fn a_deadline_in_the_past_still_sleeps_a_tick() { // Never zero: a zero-length sleep would spin the loop. diff --git a/aimdb-mqtt-connector/tests/common/mod.rs b/aimdb-mqtt-connector/tests/common/mod.rs index e2b35ce3..b6a28d5f 100644 --- a/aimdb-mqtt-connector/tests/common/mod.rs +++ b/aimdb-mqtt-connector/tests/common/mod.rs @@ -26,6 +26,10 @@ pub struct Seen { pub keep_alives: Vec, pub subscribes: Vec>, pub published: Vec<(String, Vec)>, + /// The Maximum Packet Size each MQTT 5 CONNECT advertised, when it did. + pub max_packet_sizes: Vec, + /// Pushes withheld because they exceeded the client's Maximum Packet Size. + pub withheld: usize, } impl Seen { @@ -127,6 +131,30 @@ fn connect_keep_alive(body: &[u8]) -> Option { Some(u16::from_be_bytes([*body.get(8)?, *body.get(9)?])) } +/// The Maximum Packet Size property (0x27) of an MQTT 5 CONNECT, if present. +/// Knows the fixed-size properties a client sends; anything else ends the +/// scan with `None`. +fn connect_max_packet_size(body: &[u8]) -> Option { + let mut i = 10; + let len = take_varint(body, &mut i)?; + let end = i + len; + while i < end { + let id = *body.get(i)?; + i += 1; + match id { + 0x27 => { + let b = body.get(i..i + 4)?; + return Some(u32::from_be_bytes([b[0], b[1], b[2], b[3]])); + } + 0x11 => i += 4, // session expiry interval + 0x21 | 0x22 => i += 2, // receive maximum, topic alias maximum + 0x17 | 0x19 => i += 1, // request problem / response information + _ => return None, + } + } + None +} + /// The identity a CONNECT carries: client id, then the credentials its flags /// advertise. Nothing here sets a will, so the payload fields are contiguous. fn connect_identity(body: &[u8], v5: bool) -> Option<(String, Option<(String, String)>)> { @@ -261,6 +289,9 @@ where { let mut buf = Vec::new(); let mut v5 = true; + // The client's Maximum Packet Size: a broker must not send it anything + // larger, so a push over it is withheld. + let mut client_max: Option = None; loop { let Some((first, body)) = read_packet(socket, &mut buf).await else { @@ -280,6 +311,14 @@ where if let Some(keep_alive) = connect_keep_alive(&body) { seen.keep_alives.push(keep_alive); } + client_max = if v5 { + connect_max_packet_size(&body) + } else { + None + }; + if let Some(max) = client_max { + seen.max_packet_sizes.push(max); + } } let ack: &[u8] = if v5 { &[0x20, 0x03, 0x00, 0x00, 0x00] @@ -299,11 +338,10 @@ where return; } if let Some((topic, payload)) = after.push { - if socket - .write_all(&publish(topic, payload, v5)) - .await - .is_err() - { + let packet = publish(topic, payload, v5); + if client_max.is_some_and(|max| packet.len() > max as usize) { + seen.lock().unwrap().withheld += 1; + } else if socket.write_all(&packet).await.is_err() { return; } } diff --git a/aimdb-mqtt-connector/tests/tokio_broker.rs b/aimdb-mqtt-connector/tests/tokio_broker.rs index 0a452a7d..3acbfccb 100644 --- a/aimdb-mqtt-connector/tests/tokio_broker.rs +++ b/aimdb-mqtt-connector/tests/tokio_broker.rs @@ -277,3 +277,75 @@ async fn two_connectors_in_one_process_keep_their_own_client_ids() { ids.sort(); assert_eq!(ids, vec!["first-node", "second-node"]); } + +/// Connects a client subscribed to `sensors/temperature` against a broker that +/// pushes `payload_len` bytes after every SUBACK, as it would a retained +/// message, and returns what the broker saw after `wait` plus the length the +/// record last received. +async fn with_retained_push(payload_len: usize, wait: Duration) -> (Seen, Option) { + use aimdb_core::buffer::BufferCfg; + use aimdb_core::AimDbBuilder; + use aimdb_mqtt_connector::MqttConnector; + use aimdb_tokio_adapter::net::TokioNet; + use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; + + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let seen = Arc::new(Mutex::new(Seen::default())); + + let connector = MqttConnector::new(format!("mqtt://127.0.0.1:{port}")) + .transport(TokioNet::tcp()) + .with_client_id("max-packet-size"); + let mut builder = AimDbBuilder::new() + .runtime(Arc::new(TokioAdapter)) + .with_connector(connector); + builder.configure::("temperature", |reg| { + reg.buffer(BufferCfg::SingleLatest) + .link_from("mqtt://sensors/temperature") + .with_deserializer(|_ctx, data: &[u8]| Ok::(data.len() as u64)) + .finish(); + }); + let (db, runner) = builder.build().await.expect("build db"); + let mut reader = db.subscribe::("temperature").expect("subscribe"); + + let payload = vec![b'x'; payload_len]; + let broker = fake_broker( + listener, + seen.clone(), + 0, + Some(("sensors/temperature", payload.as_slice())), + ); + let mut received = None; + let observe = async { + loop { + received = Some(reader.recv().await.expect("record open")); + } + }; + tokio::select! { + _ = runner.run() => panic!("the session loop returned"), + _ = broker => panic!("the broker returned"), + _ = observe => unreachable!(), + _ = tokio::time::sleep(wait) => {} + } + let seen = std::mem::take(&mut *seen.lock().unwrap()); + (seen, received) +} + +/// The CONNECT advertises the largest packet the session always receives, so +/// a broker withholds a larger retained message instead of sending one that +/// would end every session it is replayed into. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_retained_message_over_the_maximum_packet_size_is_withheld() { + let wait = Duration::from_secs(3); + + let (seen, received) = with_retained_push(3000, wait).await; + assert_eq!(seen.max_packet_sizes, [3328]); + assert_eq!((seen.connects, seen.withheld), (1, 0)); + assert_eq!(received, Some(3000), "a message within the limit arrives"); + + let (seen, received) = with_retained_push(4000, wait).await; + assert_eq!(seen.max_packet_sizes, [3328]); + assert_eq!(seen.connects, 1, "no reconnect loop"); + assert_eq!(seen.withheld, 1); + assert_eq!(received, None); +}