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
2 changes: 1 addition & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
39 changes: 20 additions & 19 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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

Expand All @@ -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
Expand Down
100 changes: 54 additions & 46 deletions examples/tls_public_broker.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -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<TlsConnection<'static, FromTokio<TcpStream>, Aes128GcmSha256>, Box<dyn StdError>> {
read_record: &'a mut [u8],
write_record: &'a mut [u8],
) -> Result<TlsConnection<'a, FromTokio<TcpStream>, Aes128GcmSha256>, Box<dyn StdError>> {
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::<Aes128GcmSha256>(rand::rngs::OsRng);
tls.open(TlsContext::new(&config, &mut provider)).await?;
Expand All @@ -43,51 +42,60 @@ fn unique_id(label: &str) -> String {
}

async fn run() -> Result<(), Box<dyn StdError>> {
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?;
}
}

Expand Down
2 changes: 1 addition & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down
11 changes: 5 additions & 6 deletions src/mqtt_client/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<T> Io for T where T: Read + Write + ErrorType {}
Expand Down Expand Up @@ -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,
Expand Down
23 changes: 9 additions & 14 deletions src/mqtt_client/session/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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>,
Expand Down Expand Up @@ -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,
Expand Down
8 changes: 5 additions & 3 deletions src/mqtt_client/session/operations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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
Expand Down
Loading