From d16e10c1b815178a30558547f5df04daf2731695 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Robert=20J=C3=B6rdens?= Date: Tue, 14 Jul 2026 17:24:00 +0200 Subject: [PATCH] update docs, examples w.r.t. Connection change --- CHANGELOG.md | 2 +- README.md | 39 +++++----- examples/tls_public_broker.rs | 100 ++++++++++++++------------ src/lib.rs | 2 +- src/mqtt_client/mod.rs | 11 ++- src/mqtt_client/session/mod.rs | 23 +++--- src/mqtt_client/session/operations.rs | 8 ++- 7 files changed, 95 insertions(+), 90 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 17675d6..509f88f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,7 +6,7 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased](https://github.com/quartiq/minimq/compare/v0.12.0...HEAD) +## [Unreleased](https://github.com/quartiq/minimq/compare/v0.12.1...HEAD) ## Changed diff --git a/README.md b/README.md index c47da27..8110494 100644 --- a/README.md +++ b/README.md @@ -17,7 +17,7 @@ The main API is [`Session`]. - [`Disconnect`]: graceful disconnect options - [`Io`]: transport boundary for an established byte stream - [`Session`]: the client you drive -- [`InboundPublish`]: output of [`Session::recv()`] +- [`InboundPublish`]: output of [`Connection::recv()`] ## Example @@ -86,11 +86,11 @@ async fn run() { The attached transport must implement [`embedded_io_async::Read`] and [`embedded_io_async::Write`]. Ordinary lack of inbound data must keep the read future pending; if the transport returns -`TimedOut` or `Interrupted`, [`Session::poll()`] treats that as transport failure and disconnects -the session. +`TimedOut` or `Interrupted`, [`Connection::poll()`] treats that as transport failure and +disconnects the connection. -For a TLS connectivity example and for caller-side cooperative driving via external timeouts, see -`examples/tls_public_broker.rs`. +For a TLS MQTT v5 request/reply example that preserves a subscription across reconnects and reuses +the TLS record buffers, see `examples/tls_public_broker.rs`. ## Errors @@ -106,32 +106,33 @@ You provide packet buffers plus an already-established transport, and a loop tha passes that transport into [`Session::connect()`] to establish or resume the broker session. [`Session::connect()`] takes ownership of the provided transport and performs the unbounded MQTT -`CONNECT` / `CONNACK` handshake. Once connected: -- [`Session::recv()`] blocks until the next inbound publish arrives or the session is lost. -- [`Session::poll()`] blocks until any session progress happens and returns `Ok(None)` for +`CONNECT` / `CONNACK` handshake. It returns a [`Connection`] that borrows the session. Once +connected: +- [`Connection::recv()`] blocks until the next inbound publish arrives or the connection is lost. +- [`Connection::poll()`] blocks until any session progress happens and returns `Ok(None)` for internal-only progress such as ACK handling, replay, or keepalive traffic. -The session drops the transport again on graceful disconnect, connection failure, or -transport/protocol loss. +Dropping the connection releases the session for a later reconnect. Call +[`Connection::disconnect()`] first for a graceful MQTT close. - [`ConnectEvent::Connected`] means the broker created a fresh session. Re-establish subscriptions here. - [`ConnectEvent::Reconnected`] means the broker resumed the existing MQTT session. Existing subscriptions and in-flight QoS state were kept. -- [`Session::recv()`] yields one inbound publish. +- [`Connection::recv()`] yields one inbound publish. -If [`Session::recv()`] or [`Session::poll()`] returns [`Error::Disconnected`], the caller decides +If [`Connection::recv()`] or [`Connection::poll()`] returns [`Error::Disconnected`], the caller +discards the handle and decides when to call [`Session::connect()`] with a fresh transport again. -Other transport/protocol errors already tear down the attached transport locally; callers should -handle the error and reconnect rather than retrying `recv()` or `poll()` on the same session -state. +Other transport/protocol errors mark the handle dead; callers should handle the error and reconnect +rather than retrying network operations on that handle. For cooperative driving: -- use [`Session::drive()`] for immediate local progress without waiting for future inbound reads or - future session deadlines -- wrap cancel-safe blocking [`Session::poll()`] or [`Session::recv()`] in an external timeout such - as [`embassy_time::with_timeout()`] or [`embassy_time::with_deadline()`] +- use [`Connection::drive()`] for immediate local progress without waiting for future inbound reads + or future session deadlines +- wrap cancel-safe blocking [`Connection::poll()`] or [`Connection::recv()`] in an external timeout + such as [`embassy_time::with_timeout()`] or [`embassy_time::with_deadline()`] - if you need real wall-clock limits, enforce them in the transport's `read`, `write`, and `flush` futures; using the same budget as minimq's internal MQTT round-trip timeout keeps keepalive and transport liveness aligned diff --git a/examples/tls_public_broker.rs b/examples/tls_public_broker.rs index 1ccb53f..7dd81af 100644 --- a/examples/tls_public_broker.rs +++ b/examples/tls_public_broker.rs @@ -1,11 +1,8 @@ -//! TLS connectivity example for `embedded-tls`. -//! -//! Use this as a TLS transport example. Bounded/cooperative session driving is done by wrapping -//! the cancel-safe blocking `Session::poll()` in an external timeout at the call site. +//! MQTT v5 request/reply echo demonstrating session resumption over `embedded-tls`. use embedded_io_adapters::tokio_1::FromTokio; use embedded_tls::{Aes128GcmSha256, TlsConfig, TlsConnection, TlsContext, UnsecureProvider}; -use minimq::{ConfigBuilder, Publication, QoS, Session, TopicFilter}; +use minimq::{ConfigBuilder, ConnectEvent, Session, TopicFilter}; use std::error::Error as StdError; use std::time::{SystemTime, UNIX_EPOCH}; use tokio::net::TcpStream; @@ -17,17 +14,19 @@ const PASSWORD: &str = "public"; const TLS_READ_RECORD_BUFFER_LEN: usize = 16_640; const TLS_WRITE_RECORD_BUFFER_LEN: usize = 4_096; -async fn connect_tls( +async fn connect_tls<'a>( host: &str, port: u16, -) -> Result, Aes128GcmSha256>, Box> { + read_record: &'a mut [u8], + write_record: &'a mut [u8], +) -> Result, Aes128GcmSha256>, Box> { let config = TlsConfig::new() .with_server_name(host) .enable_rsa_signatures(); let mut tls = TlsConnection::new( FromTokio::new(TcpStream::connect((host, port)).await?), - Box::leak(Box::new([0u8; TLS_READ_RECORD_BUFFER_LEN])), - Box::leak(Box::new([0u8; TLS_WRITE_RECORD_BUFFER_LEN])), + read_record, + write_record, ); let mut provider = UnsecureProvider::new::(rand::rngs::OsRng); tls.open(TlsContext::new(&config, &mut provider)).await?; @@ -43,51 +42,60 @@ fn unique_id(label: &str) -> String { } async fn run() -> Result<(), Box> { - let topic = format!("minimq/examples/tls/{}", unique_id("topic")); - let payload_str = format!("hello over tls {}", unique_id("msg")); - let payload = payload_str.as_bytes(); - - let mut sub_storage = [0u8; 2048]; + let client_id = unique_id("client"); + let topic = format!("minimq/examples/echo/{}", unique_id("request")); + let mut read_record = vec![0; TLS_READ_RECORD_BUFFER_LEN]; + let mut write_record = vec![0; TLS_WRITE_RECORD_BUFFER_LEN]; + let mut storage = [0u8; 2048]; let mut session = Session::new( - ConfigBuilder::from_buffer(&mut sub_storage, 1024)?.auth(USERNAME, PASSWORD.as_bytes())?, + ConfigBuilder::from_buffer(&mut storage, 1024)? + .client_id(&client_id)? + .auth(USERNAME, PASSWORD.as_bytes())? + .session_expiry_interval(60), ); - let mut subscriber = session - .connect(connect_tls(BROKER_HOST, BROKER_PORT).await?) - .await?; - let sub = subscriber - .subscribe(&[TopicFilter::new(&topic)], &[]) - .await?; - while subscriber.is_pending(&sub) { - subscriber.poll().await?; - } - let mut pub_storage = [0u8; 2048]; - let mut session = Session::new( - ConfigBuilder::from_buffer(&mut pub_storage, 1024)?.auth(USERNAME, PASSWORD.as_bytes())?, - ); - let mut publisher = session - .connect(connect_tls(BROKER_HOST, BROKER_PORT).await?) + let tls = connect_tls( + BROKER_HOST, + BROKER_PORT, + &mut read_record, + &mut write_record, + ) + .await?; + let mut connection = session.connect(tls).await?; + let subscribe = connection + .subscribe(&[TopicFilter::new(&topic)], &[]) .await?; - let pub_ = publisher - .publish(Publication::new(&topic, payload).qos(QoS::AtLeastOnce)) - .await? - .unwrap(); - while publisher.is_pending(&pub_) { - publisher.poll().await?; + while connection.is_pending(&subscribe) { + connection.poll().await?; } + println!("echoing requests on {topic}"); loop { - let message = subscriber.recv().await?; - if message.topic() == topic && message.payload() == payload { - println!( - "received topic={} payload={}", - message.topic(), - payload.escape_ascii() - ); - publisher.disconnect().await?; - subscriber.disconnect().await?; - return Ok(()); + let (reply, payload) = { + let request = connection.recv().await?; + // Preserve the MQTT v5 response topic and correlation data past this receive borrow. + let Some(reply) = request.reply_owned::<256, 64>()? else { + continue; + }; + (reply, request.payload().to_vec()) + }; + + connection.disconnect().await?; + drop(connection); + let tls = connect_tls( + BROKER_HOST, + BROKER_PORT, + &mut read_record, + &mut write_record, + ) + .await?; + connection = session.connect(tls).await?; + if connection.connect_event() != ConnectEvent::Reconnected { + return Err("broker did not resume the MQTT session".into()); } + connection + .publish(reply.publication(payload.as_slice())) + .await?; } } diff --git a/src/lib.rs b/src/lib.rs index 9c21141..f044872 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -137,7 +137,7 @@ pub enum ResourceError { InflightExhausted, } -/// Error returned from [`Session::publish`](crate::Session::publish). +/// Error returned from [`Connection::publish`]. /// /// `P` is the payload serialization error and `T` is the transport error. #[derive(Debug, PartialEq, thiserror::Error)] diff --git a/src/mqtt_client/mod.rs b/src/mqtt_client/mod.rs index 36c2946..7b96a3b 100644 --- a/src/mqtt_client/mod.rs +++ b/src/mqtt_client/mod.rs @@ -9,11 +9,11 @@ use crate::{ }; use embedded_io_async::{ErrorType, Read, Write}; -/// Transport trait required by [`Session`](crate::Session). +/// Transport trait required by [`Connection`]. /// /// Ordinary lack of inbound data must leave the read future pending. If `read()` returns -/// `TimedOut` or `Interrupted`, [`Session::poll`](crate::Session::poll) treats that as transport -/// failure and disconnects the session. +/// `TimedOut` or `Interrupted`, [`Connection::poll`] treats that as +/// transport failure and disconnects the connection. pub trait Io: Read + Write + ErrorType {} impl Io for T where T: Read + Write + ErrorType {} @@ -54,9 +54,8 @@ pub(crate) enum OpStatus { Invalidated, } -/// Inbound MQTT `PUBLISH` surfaced by [`Session::recv`](crate::Session::recv) and by -/// [`Session::drive`](crate::Session::drive) / [`Session::poll`](crate::Session::poll) when they -/// return `Some(...)`. +/// Inbound MQTT `PUBLISH` surfaced by [`Connection::recv`] and by +/// [`Connection::drive`] / [`Connection::poll`] when they return `Some(...)`. #[derive(Debug)] pub struct InboundPublish<'a> { topic: &'a str, diff --git a/src/mqtt_client/session/mod.rs b/src/mqtt_client/session/mod.rs index a3ceff9..884eeb7 100644 --- a/src/mqtt_client/session/mod.rs +++ b/src/mqtt_client/session/mod.rs @@ -19,21 +19,12 @@ use state::{RuntimeState, SessionData}; /// One long-lived MQTT client session. /// -/// Drive the session after [`connect`](Self::connect) has taken ownership of a live transport. -/// Use [`recv`](Self::recv) when you want the next inbound publish, [`poll`](Self::poll) when you -/// need to wait for any session progress, and [`drive`](Self::drive) for cooperative immediate -/// progress. Real time bounds come from the transport: stalled reads, writes, or flushes must -/// eventually error if the caller needs hard latency limits. The same session is also used for -/// outbound `publish`, `subscribe`, and `unsubscribe` operations. +/// [`connect`](Self::connect) borrows the session and returns a live [`Connection`] that owns the +/// transport. Drive all network operations through that handle. Dropping it releases the session +/// for a later reconnect while preserving durable MQTT state such as in-flight QoS replay. /// -/// Cancel safety, assuming the transport's I/O futures are cancel-safe: -/// [`drive`](Self::drive), [`poll`](Self::poll), [`recv`](Self::recv), [`disconnect`](Self::disconnect), -/// [`subscribe`](Self::subscribe), [`unsubscribe`](Self::unsubscribe), and -/// [`publish`](Self::publish) for QoS 1/2 preserve local session state across cancellation. -/// Cancelling [`connect`](Self::connect) drops the supplied transport and leaves the session -/// disconnected; the next `connect()` retries from clean transport-local state. QoS 0 -/// [`publish`](Self::publish) is not cancel-safe because it bypasses retained outbound state and -/// writes directly from temporary TX scratch space. +/// Cancelling `connect()` drops the supplied transport and leaves the session available for a +/// clean retry. pub struct Session<'buf> { client_id: String<64>, packet_reader: PacketReader<'buf>, @@ -145,6 +136,10 @@ impl<'buf> Session<'buf> { /// /// Note that dropping or forgetting the handle is an *ungraceful* MQTT close: **no `DISCONNECT` /// packet is sent** (a sync `Drop` cannot perform the async write). +/// +/// Cancellation guarantees are documented on each network operation. In particular, QoS 1/2 +/// publishes preserve their retained session state, while a QoS 0 publish writes directly from +/// temporary TX scratch space and is not cancel-safe. pub struct Connection<'a, 'buf, IO> { pub(super) session: &'a mut Session<'buf>, pub(super) io: IO, diff --git a/src/mqtt_client/session/operations.rs b/src/mqtt_client/session/operations.rs index 28262c3..6bdff6d 100644 --- a/src/mqtt_client/session/operations.rs +++ b/src/mqtt_client/session/operations.rs @@ -51,7 +51,8 @@ impl<'buf, IO: Io> Connection<'_, 'buf, IO> { /// Send a `SUBSCRIBE`. /// - /// Call this after [`connect`](Self::connect). A resumed [`crate::ConnectEvent::Reconnected`] + /// Call this after [`Session::connect`](crate::Session::connect). A resumed + /// [`ConnectEvent::Reconnected`](crate::ConnectEvent::Reconnected) /// already kept broker-side subscriptions. /// Cancel-safe if the underlying transport I/O futures are cancel-safe. pub async fn subscribe( @@ -152,8 +153,9 @@ impl<'buf, IO: Io> Connection<'_, 'buf, IO> { /// QoS 1 and 2 retain the encoded packet in the session TX buffer until broker ack and are /// cancel-safe if the underlying transport I/O futures are cancel-safe. /// They return `Some(Op)` so the caller can check completion with - /// [`Session::is_pending`](Self::is_pending), [`Session::is_complete`](Self::is_complete), - /// or [`Session::is_invalidated`](Self::is_invalidated). + /// [`Connection::is_pending`](Self::is_pending), + /// [`Connection::is_complete`](Self::is_complete), or + /// [`Connection::is_invalidated`](Self::is_invalidated). /// /// QoS 0 bypasses retained outbound state, encodes into temporary TX scratch space, and writes /// directly to the transport. It therefore does not consume replay/in-flight slots and only