Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
10 changes: 10 additions & 0 deletions aimdb-mqtt-connector/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
69 changes: 68 additions & 1 deletion aimdb-mqtt-connector/src/embedded/session_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u8, RX_CHUNK>;

Expand Down Expand Up @@ -176,10 +183,13 @@ async fn client_loop<D: Delay>(
// 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(),
Expand Down Expand Up @@ -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<u8> {
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::<PACKET_BUFFER_SIZE>::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.
Expand Down
48 changes: 43 additions & 5 deletions aimdb-mqtt-connector/tests/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,10 @@ pub struct Seen {
pub keep_alives: Vec<u16>,
pub subscribes: Vec<Vec<String>>,
pub published: Vec<(String, Vec<u8>)>,
/// The Maximum Packet Size each MQTT 5 CONNECT advertised, when it did.
pub max_packet_sizes: Vec<u32>,
/// Pushes withheld because they exceeded the client's Maximum Packet Size.
pub withheld: usize,
}

impl Seen {
Expand Down Expand Up @@ -127,6 +131,30 @@ fn connect_keep_alive(body: &[u8]) -> Option<u16> {
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<u32> {
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)>)> {
Expand Down Expand Up @@ -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<u32> = None;

loop {
let Some((first, body)) = read_packet(socket, &mut buf).await else {
Expand All @@ -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]
Expand All @@ -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;
}
}
Expand Down
72 changes: 72 additions & 0 deletions aimdb-mqtt-connector/tests/tokio_broker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64>) {
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::<u64>("temperature", |reg| {
reg.buffer(BufferCfg::SingleLatest)
.link_from("mqtt://sensors/temperature")
.with_deserializer(|_ctx, data: &[u8]| Ok::<u64, String>(data.len() as u64))
.finish();
});
let (db, runner) = builder.build().await.expect("build db");
let mut reader = db.subscribe::<u64>("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);
}
Loading