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
4 changes: 3 additions & 1 deletion minimq/src/mqtt_client/outbound.rs
Original file line number Diff line number Diff line change
Expand Up @@ -289,8 +289,10 @@ impl<'a> Outbound<'a> {
payload: P,
) -> Result<usize, PubError<P::Error, E>> {
let start = self.retained_bytes();
// Adaptive payloads must see the retained budget, including encoder workspace.
let end = self.retained_capacity();
let (offset, packet) =
MqttSerializer::encode_publish_with_offset(&mut self.buf[start..], header, payload)?;
MqttSerializer::encode_publish_with_offset(&mut self.buf[start..end], header, payload)?;
let len = packet.len();
self.pack_encoded(offset, len)
.map_err(|err| PubError::Session(Error::Resource(err)))?;
Expand Down
16 changes: 10 additions & 6 deletions minimq/tests/async_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1554,8 +1554,6 @@ fn admitted_qos1_publish_leaves_space_to_reconnect() {
const TX_LEN: usize = 256;
// This payload fits the transient encode space, but not retained state plus the next CONNECT.
const TOO_LARGE_PAYLOAD_LEN: usize = 240;
// This smaller payload is admitted and must remain replayable after reconnect.
const REPLAYABLE_PAYLOAD_LEN: usize = 180;
// MQTT PUBLISH type (3), QoS 1, with and without the DUP flag.
const PUBLISH_QOS1: u8 = 0b0011_0010;

Expand All @@ -1578,11 +1576,17 @@ fn admitted_qos1_publish_leaves_space_to_reconnect() {
let mut conn = expect_connected(&mut session, &connector);
assert_eq!(
publish_qos1(&mut conn, "x", &[0; TOO_LARGE_PAYLOAD_LEN]),
Err(PubError::Session(Error::Resource(
ResourceError::BufferTooSmall
)))
Err(PubError::Payload(()))
);
publish_qos1_ok(&mut conn, "x", &[0; REPLAYABLE_PAYLOAD_LEN]);
assert!(conn.is_connected());
// QoS0 may still use transient space reserved against retained packets.
block_on(conn.publish(Publication::bytes("x", &[0; TOO_LARGE_PAYLOAD_LEN]))).unwrap();
let publication = Publication::new("x", |buffer: &mut [u8]| {
buffer.fill(0x55);
Ok::<_, ()>(buffer.len())
})
.qos(QoS::AtLeastOnce);
block_on(conn.publish(publication)).unwrap();

let mut replay = first_inspect.tx().last().unwrap().clone();
assert_eq!(replay[0], PUBLISH_QOS1);
Expand Down