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
29 changes: 29 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,35 @@ 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).

## [mqtt5 0.43.0] - 2026-09-27

### Breaking

- **`broker::config::QuicConfig` and `broker::quic_acceptor::QuicAcceptorConfig` have three new public fields**: `max_concurrent_streams`, `stream_receive_window` and `disable_segmentation_offload`. Code that builds either struct with a struct literal must set them; `QuicConfig::new` and `QuicAcceptorConfig::new` set them to their defaults.

### Added

- **Broker QUIC flow-control settings.** `max_concurrent_streams` caps the unidirectional and bidirectional streams a client may open at once (QUIC default 100); `stream_receive_window` sets the per-stream receive window (default 262144 bytes); `disable_segmentation_offload` turns off UDP segmentation offload, so each datagram is sent separately. Set them with `with_max_concurrent_streams`, `with_stream_receive_window` and `with_disable_segmentation_offload`.
- **Per-connection QUIC statistics.** When `MQTT5_QUIC_STATS_DIR` names a directory, the broker writes one CSV per QUIC connection, with a row every 100 ms and a last row at close: RTT, congestion window, lost packets, congestion events and sent packets on the broker's path, and the STREAM_DATA_BLOCKED, DATA_BLOCKED and STREAMS_BLOCKED (uni) frames received from the client. Unset, nothing is sampled.

### Fixed

- **The client's QUIC stream limit is now applied.** `QuicConfig::with_max_concurrent_streams` (and `MqttClient::set_quic_max_streams`, and the bridge's `quic_max_streams`) was stored but never passed to the QUIC transport, so the broker could open up to the QUIC default of 100 streams toward the client. It now caps the streams the broker may open.
- **Eight QUIC multistream integration tests no longer pass without checking anything when a client fails to connect.** They returned early, so with missing test certificates they reported success.

## [mqttv5-cli 0.29.1] - 2026-09-27

### Added

- `mqttv5 broker` flags `--quic-max-streams <N>`, `--quic-stream-window <BYTES>` and `--quic-disable-offload` (environment variables `MQTT5_QUIC_MAX_STREAMS`, `MQTT5_QUIC_STREAM_WINDOW`, `MQTT5_QUIC_DISABLE_OFFLOAD`), backed by the mqtt5 0.43.0 broker settings.
- Requires mqtt5 0.43.

## [mqtt5-wasm 2.1.1] - 2026-09-27

### Changed

- Requires mqtt5 0.43. No change to the wasm API.

## [mqtt5 0.42.0] - 2026-09-25

### Breaking
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ A production-ready MQTT v5.0 and v3.1.1 platform that ships a client library, a

```toml
[dependencies]
mqtt5 = "0.36"
mqtt5 = "0.43"
```

For a lean client-only build (no broker, drops `argon2`/`hyper`/`regex`/`toml`/...):
Expand Down
4 changes: 2 additions & 2 deletions crates/mqtt5-wasm/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "mqtt5-wasm"
version = "2.1.0"
version = "2.1.1"
edition.workspace = true
rust-version.workspace = true
authors.workspace = true
Expand Down Expand Up @@ -28,7 +28,7 @@ codec = ["client", "dep:miniz_oxide"]

[dependencies]
mqtt5-protocol = "0.15.2"
mqtt5 = { version = "0.42", optional = true, default-features = false, features = [
mqtt5 = { version = "0.43", optional = true, default-features = false, features = [
"tokio",
] }

Expand Down
2 changes: 1 addition & 1 deletion crates/mqtt5/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "mqtt5"
version = "0.42.0"
version = "0.43.0"
edition.workspace = true
rust-version.workspace = true
authors.workspace = true
Expand Down
2 changes: 1 addition & 1 deletion crates/mqtt5/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ Full-featured MQTT v5.0 and v3.1.1 client and broker for native platforms (Linux

```toml
[dependencies]
mqtt5 = "0.36"
mqtt5 = "0.43"
```

### Client-only builds
Expand Down
37 changes: 32 additions & 5 deletions crates/mqtt5/src/broker/config/transport.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
use serde::{Deserialize, Serialize};
use std::net::SocketAddr;
use std::net::{Ipv4Addr, Ipv6Addr, SocketAddr};
use std::path::PathBuf;

#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
Expand All @@ -20,6 +20,12 @@ pub struct QuicConfig {
pub bind_addresses: Vec<SocketAddr>,
#[serde(default)]
pub enable_early_data: bool,
#[serde(default)]
pub max_concurrent_streams: Option<u32>,
#[serde(default)]
pub stream_receive_window: Option<u32>,
#[serde(default)]
pub disable_segmentation_offload: bool,
}

impl QuicConfig {
Expand All @@ -35,10 +41,13 @@ impl QuicConfig {
ca_file: None,
require_client_cert: false,
bind_addresses: vec![
"0.0.0.0:14567".parse().expect("valid IPv4 address"),
"[::]:14567".parse().expect("valid IPv6 address"),
SocketAddr::from((Ipv4Addr::UNSPECIFIED, 14567)),
SocketAddr::from((Ipv6Addr::UNSPECIFIED, 14567)),
],
enable_early_data: false,
max_concurrent_streams: None,
stream_receive_window: None,
disable_segmentation_offload: false,
}
}

Expand Down Expand Up @@ -77,6 +86,24 @@ impl QuicConfig {
self.enable_early_data = enable;
self
}

#[must_use]
pub fn with_max_concurrent_streams(mut self, max: u32) -> Self {
self.max_concurrent_streams = Some(max);
self
}

#[must_use]
pub fn with_stream_receive_window(mut self, bytes: u32) -> Self {
self.stream_receive_window = Some(bytes);
self
}

#[must_use]
pub fn with_disable_segmentation_offload(mut self, disable: bool) -> Self {
self.disable_segmentation_offload = disable;
self
}
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
Expand All @@ -92,8 +119,8 @@ impl Default for WebSocketConfig {
fn default() -> Self {
Self {
bind_addresses: vec![
"0.0.0.0:8080".parse().unwrap(),
"[::]:8080".parse().unwrap(),
SocketAddr::from((Ipv4Addr::UNSPECIFIED, 8080)),
SocketAddr::from((Ipv6Addr::UNSPECIFIED, 8080)),
],
path: "/mqtt".to_string(),
subprotocol: "mqtt".to_string(),
Expand Down
119 changes: 104 additions & 15 deletions crates/mqtt5/src/broker/quic_acceptor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,18 +28,20 @@ use tracing::{debug, error, instrument, trace, warn};

use super::tls_acceptor::TlsAcceptorConfig;

// [RFC9000§7] QUIC transport parameters
pub struct QuicAcceptorConfig {
pub cert_chain: Vec<CertificateDer<'static>>,
pub private_key: PrivateKeyDer<'static>,
pub client_ca_certs: Option<Vec<CertificateDer<'static>>>,
pub require_client_cert: bool,
pub alpn_protocols: Vec<Vec<u8>>,
pub enable_early_data: bool,
pub max_concurrent_streams: Option<u32>,
pub stream_receive_window: Option<u32>,
pub disable_segmentation_offload: bool,
}

impl QuicAcceptorConfig {
#[allow(clippy::must_use_candidate)]
#[must_use]
pub fn new(
cert_chain: Vec<CertificateDer<'static>>,
private_key: PrivateKeyDer<'static>,
Expand All @@ -51,6 +53,9 @@ impl QuicAcceptorConfig {
require_client_cert: false,
alpn_protocols: vec![b"MQTT-next".to_vec(), b"mqtt".to_vec()],
enable_early_data: false,
max_concurrent_streams: None,
stream_receive_window: None,
disable_segmentation_offload: false,
}
}

Expand Down Expand Up @@ -98,13 +103,28 @@ impl QuicAcceptorConfig {
self
}

#[must_use]
pub fn with_max_concurrent_streams(mut self, max: u32) -> Self {
self.max_concurrent_streams = Some(max);
self
}

#[must_use]
pub fn with_stream_receive_window(mut self, bytes: u32) -> Self {
self.stream_receive_window = Some(bytes);
self
}

#[must_use]
pub fn with_disable_segmentation_offload(mut self, disable: bool) -> Self {
self.disable_segmentation_offload = disable;
self
}

/// Builds a rustls `ServerConfig` from this QUIC acceptor configuration.
///
/// # Errors
/// Returns an error if certificate loading fails or TLS configuration is invalid.
///
/// # Panics
/// Panics if the idle timeout duration conversion fails (should never happen).
pub fn build_server_config(&self) -> Result<ServerConfig> {
let crypto_provider = Arc::new(rustls::crypto::ring::default_provider());

Expand Down Expand Up @@ -165,18 +185,26 @@ impl QuicAcceptorConfig {
let mut server_config = ServerConfig::with_crypto(Arc::new(quic_config));

let mut transport_config = quinn::TransportConfig::default();
transport_config.max_idle_timeout(Some(
std::time::Duration::from_secs(60)
.try_into()
.expect("valid duration"),
));
transport_config.max_idle_timeout(Some(quinn::IdleTimeout::from(quinn::VarInt::from_u32(
60_000,
))));
transport_config.datagram_receive_buffer_size(Some(65536));
transport_config.datagram_send_buffer_size(65536);

transport_config.stream_receive_window(262_144u32.into());
transport_config
.stream_receive_window(self.stream_receive_window.unwrap_or(262_144).into());
transport_config.receive_window(1_048_576u32.into());
transport_config.send_window(1_048_576);

if let Some(max) = self.max_concurrent_streams {
transport_config.max_concurrent_uni_streams(max.into());
transport_config.max_concurrent_bidi_streams(max.into());
}

if self.disable_segmentation_offload {
transport_config.enable_segmentation_offload(false);
}

server_config.transport_config(Arc::new(transport_config));

Ok(server_config)
Expand All @@ -200,7 +228,7 @@ pub struct QuicStreamWrapper {
}

impl QuicStreamWrapper {
#[allow(clippy::must_use_candidate)]
#[must_use]
pub fn new(send: SendStream, recv: RecvStream, peer_addr: SocketAddr) -> Self {
Self {
send,
Expand Down Expand Up @@ -301,7 +329,6 @@ pub async fn accept_quic_stream(
Ok(QuicStreamWrapper::new(send, recv, peer_addr))
}

// [MQoQ§4.1] Flow type detection
fn is_flow_header_byte(b: u8) -> bool {
matches!(
b,
Expand Down Expand Up @@ -502,7 +529,6 @@ pub(super) async fn read_packet_with_buffer(
Packet::decode_from_body(fixed_header.packet_type, &fixed_header, &mut payload_buf)
}

// [MQoQ§4] QUIC connection handling with flow headers
#[allow(clippy::too_many_arguments)]
#[instrument(skip(connection, config, router, auth_provider, storage, stats, resource_monitor, shutdown_rx), fields(peer_addr = %peer_addr))]
pub async fn run_quic_connection_handler(
Expand Down Expand Up @@ -628,6 +654,7 @@ async fn run_quic_handler_inner(

spawn_datagram_reader(connection.clone(), packet_tx.clone(), peer_addr, label);
spawn_bi_accept_loop(connection.clone(), flow_registry.clone(), peer_addr, label);
spawn_quic_stats_sampler(connection.clone(), peer_addr);
spawn_uni_accept_loop(
connection,
packet_tx,
Expand All @@ -638,6 +665,69 @@ async fn run_quic_handler_inner(
);
}

const QUIC_STATS_DIR_ENV: &str = "MQTT5_QUIC_STATS_DIR";
const QUIC_STATS_HEADER: &str = "timestamp_ns,rtt_us,cwnd,lost_packets,congestion_events,sent_packets,stream_data_blocked,data_blocked,streams_blocked_uni\n";
const QUIC_STATS_INTERVAL: Duration = Duration::from_millis(100);

fn spawn_quic_stats_sampler(connection: Arc<Connection>, peer_addr: SocketAddr) {
let Some(dir) = std::env::var_os(QUIC_STATS_DIR_ENV).filter(|dir| !dir.is_empty()) else {
return;
};
tokio::spawn(async move {
match sample_quic_stats(&connection, std::path::Path::new(&dir), peer_addr).await {
Ok(rows) => debug!(%peer_addr, rows, "QUIC stats sampling finished"),
Err(e) => warn!(%peer_addr, error = %e, "QUIC stats sampling stopped"),
}
});
}

async fn sample_quic_stats(
connection: &Connection,
dir: &std::path::Path,
peer_addr: SocketAddr,
) -> std::io::Result<u64> {
use tokio::io::AsyncWriteExt;

tokio::fs::create_dir_all(dir).await?;
let name = peer_addr.to_string().replace([':', '.', '[', ']'], "-");
let mut out = tokio::fs::File::create(dir.join(format!("broker_quic_{name}.csv"))).await?;
out.write_all(QUIC_STATS_HEADER.as_bytes()).await?;

let mut rows = 0u64;
let mut ticker = tokio::time::interval(QUIC_STATS_INTERVAL);
loop {
let closed = tokio::select! {
_ = connection.closed() => true,
_ = ticker.tick() => false,
};
out.write_all(quic_stats_row(connection).as_bytes()).await?;
rows += 1;
if closed {
return Ok(rows);
}
}
}

fn quic_stats_row(connection: &Connection) -> String {
let stats = connection.stats();
let timestamp_ns = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |elapsed| {
u64::try_from(elapsed.as_nanos()).unwrap_or(u64::MAX)
});
let rtt_us = u64::try_from(stats.path.rtt.as_micros()).unwrap_or(u64::MAX);
format!(
"{timestamp_ns},{rtt_us},{},{},{},{},{},{},{}\n",
stats.path.cwnd,
stats.path.lost_packets,
stats.path.congestion_events,
stats.path.sent_packets,
stats.frame_rx.stream_data_blocked,
stats.frame_rx.data_blocked,
stats.frame_rx.streams_blocked_uni,
)
}

fn spawn_datagram_reader(
connection: Arc<Connection>,
packet_tx: mpsc::Sender<(Packet, Option<u64>)>,
Expand Down Expand Up @@ -859,7 +949,6 @@ fn spawn_discard_handler(
});
}

// [MQoQ§5] Data stream processing
fn spawn_data_stream_reader(
mut recv: RecvStream,
packet_tx: mpsc::Sender<(Packet, Option<u64>)>,
Expand Down
12 changes: 12 additions & 0 deletions crates/mqtt5/src/broker/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -592,6 +592,18 @@ impl MqttBroker {
acceptor_config = acceptor_config.with_early_data(true);
}

if let Some(max) = quic_config.max_concurrent_streams {
acceptor_config = acceptor_config.with_max_concurrent_streams(max);
}

if let Some(bytes) = quic_config.stream_receive_window {
acceptor_config = acceptor_config.with_stream_receive_window(bytes);
}

if quic_config.disable_segmentation_offload {
acceptor_config = acceptor_config.with_disable_segmentation_offload(true);
}

let mut endpoints = Vec::new();
let mut failures = Vec::new();
for addr in &quic_config.bind_addresses {
Expand Down
Loading
Loading