diff --git a/CHANGELOG.md b/CHANGELOG.md index dd763b6a..92bb8d88 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 `, `--quic-stream-window ` 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 diff --git a/README.md b/README.md index 756031f3..a3dddb6c 100644 --- a/README.md +++ b/README.md @@ -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`/...): diff --git a/crates/mqtt5-wasm/Cargo.toml b/crates/mqtt5-wasm/Cargo.toml index d646565a..d579b321 100644 --- a/crates/mqtt5-wasm/Cargo.toml +++ b/crates/mqtt5-wasm/Cargo.toml @@ -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 @@ -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", ] } diff --git a/crates/mqtt5/Cargo.toml b/crates/mqtt5/Cargo.toml index 0b2edc7f..50a0b4e7 100644 --- a/crates/mqtt5/Cargo.toml +++ b/crates/mqtt5/Cargo.toml @@ -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 diff --git a/crates/mqtt5/README.md b/crates/mqtt5/README.md index 6115129e..5ffd4fc0 100644 --- a/crates/mqtt5/README.md +++ b/crates/mqtt5/README.md @@ -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 diff --git a/crates/mqtt5/src/broker/config/transport.rs b/crates/mqtt5/src/broker/config/transport.rs index 9861a0f1..0b4eaf59 100644 --- a/crates/mqtt5/src/broker/config/transport.rs +++ b/crates/mqtt5/src/broker/config/transport.rs @@ -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)] @@ -20,6 +20,12 @@ pub struct QuicConfig { pub bind_addresses: Vec, #[serde(default)] pub enable_early_data: bool, + #[serde(default)] + pub max_concurrent_streams: Option, + #[serde(default)] + pub stream_receive_window: Option, + #[serde(default)] + pub disable_segmentation_offload: bool, } impl QuicConfig { @@ -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, } } @@ -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)] @@ -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(), diff --git a/crates/mqtt5/src/broker/quic_acceptor.rs b/crates/mqtt5/src/broker/quic_acceptor.rs index 28a8174e..495f4d15 100644 --- a/crates/mqtt5/src/broker/quic_acceptor.rs +++ b/crates/mqtt5/src/broker/quic_acceptor.rs @@ -28,7 +28,6 @@ use tracing::{debug, error, instrument, trace, warn}; use super::tls_acceptor::TlsAcceptorConfig; -// [RFC9000§7] QUIC transport parameters pub struct QuicAcceptorConfig { pub cert_chain: Vec>, pub private_key: PrivateKeyDer<'static>, @@ -36,10 +35,13 @@ pub struct QuicAcceptorConfig { pub require_client_cert: bool, pub alpn_protocols: Vec>, pub enable_early_data: bool, + pub max_concurrent_streams: Option, + pub stream_receive_window: Option, + pub disable_segmentation_offload: bool, } impl QuicAcceptorConfig { - #[allow(clippy::must_use_candidate)] + #[must_use] pub fn new( cert_chain: Vec>, private_key: PrivateKeyDer<'static>, @@ -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, } } @@ -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 { let crypto_provider = Arc::new(rustls::crypto::ring::default_provider()); @@ -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) @@ -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, @@ -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, @@ -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( @@ -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, @@ -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, 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 { + 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, packet_tx: mpsc::Sender<(Packet, Option)>, @@ -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)>, diff --git a/crates/mqtt5/src/broker/server.rs b/crates/mqtt5/src/broker/server.rs index f563f967..ffb0cf51 100644 --- a/crates/mqtt5/src/broker/server.rs +++ b/crates/mqtt5/src/broker/server.rs @@ -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 { diff --git a/crates/mqtt5/src/transport/quic.rs b/crates/mqtt5/src/transport/quic.rs index cacc8e97..be96df7c 100644 --- a/crates/mqtt5/src/transport/quic.rs +++ b/crates/mqtt5/src/transport/quic.rs @@ -33,7 +33,6 @@ pub struct QuicConfig { pub enable_datagrams: bool, pub datagram_send_buffer_size: usize, pub datagram_receive_buffer_size: usize, - // [MQoQ§4] Flow header toggle pub enable_flow_headers: bool, pub flow_expire_interval: u64, pub flow_flags: FlowFlags, @@ -232,16 +231,20 @@ impl QuicConfig { let mut client_config = ClientConfig::new(Arc::new(quic_crypto)); let mut transport_config = quinn::TransportConfig::default(); - transport_config.max_idle_timeout(Some( - std::time::Duration::from_secs(120) - .try_into() - .expect("valid duration"), - )); + transport_config.max_idle_timeout(Some(quinn::IdleTimeout::from(quinn::VarInt::from_u32( + 120_000, + )))); transport_config.stream_receive_window(262_144u32.into()); transport_config.receive_window(1_048_576u32.into()); transport_config.send_window(1_048_576); + if let Some(max) = self.max_concurrent_streams { + let max = u32::try_from(max).unwrap_or(u32::MAX); + transport_config.max_concurrent_uni_streams(max.into()); + transport_config.max_concurrent_bidi_streams(max.into()); + } + if self.enable_datagrams { transport_config.datagram_send_buffer_size(self.datagram_send_buffer_size); transport_config.datagram_receive_buffer_size(Some(self.datagram_receive_buffer_size)); @@ -367,7 +370,6 @@ impl QuicTransport { } } -// [RFC9000§7] QUIC connection establishment impl Transport for QuicTransport { #[instrument(skip(self), fields(server_name = %self.config.server_name, addr = %self.config.addr))] async fn connect(&mut self) -> Result<()> { @@ -384,7 +386,7 @@ impl Transport for QuicTransport { self.built_client_config = Some(client_config.clone()); - let mut endpoint = Endpoint::client("0.0.0.0:0".parse().unwrap()) + let mut endpoint = Endpoint::client(SocketAddr::from(([0, 0, 0, 0], 0))) .map_err(|e| MqttError::ConnectionError(format!("Failed to create endpoint: {e}")))?; endpoint.set_default_client_config(client_config); @@ -510,7 +512,6 @@ impl Transport for QuicTransport { } } -// [MQTT5§2] Fixed header parsing impl PacketReader for RecvStream { async fn read_packet(&mut self, protocol_version: u8) -> Result { let mut header_buf = BytesMut::with_capacity(5); @@ -601,7 +602,6 @@ fn packet_to_type(packet: &Packet) -> PacketType { } } -// [MQTT5§2] Packet encoding impl PacketWriter for SendStream { async fn write_packet(&mut self, packet: Packet) -> Result<()> { let packet_type = packet_to_type(&packet); diff --git a/crates/mqtt5/tests/broker_quic_integration.rs b/crates/mqtt5/tests/broker_quic_integration.rs index a7408292..7b8d00b2 100644 --- a/crates/mqtt5/tests/broker_quic_integration.rs +++ b/crates/mqtt5/tests/broker_quic_integration.rs @@ -1,7 +1,9 @@ #![cfg(feature = "broker")] #![cfg(feature = "transport-quic")] -use mqtt5::broker::config::{BrokerConfig, QuicConfig, ServerDeliveryStrategy}; +use mqtt5::broker::config::{ + BrokerConfig, QuicConfig, ServerDeliveryStrategy, StorageBackend, StorageConfig, +}; use mqtt5::broker::MqttBroker; use mqtt5::time::Duration; use mqtt5::transport::StreamStrategy; @@ -51,6 +53,7 @@ async fn start_quic_broker_at(bind_addr: SocketAddr) -> (MqttBroker, SocketAddr) let cert_dir = manifest_dir.join("../../test_certs"); let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) .with_bind_address(([127, 0, 0, 1], 0)) .with_quic( QuicConfig::new(cert_dir.join("server.pem"), cert_dir.join("server.key")) @@ -115,6 +118,7 @@ async fn test_broker_quic_creation() { let _ = rustls::crypto::ring::default_provider().install_default(); let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) .with_bind_address(([127, 0, 0, 1], 0)) .with_quic( QuicConfig::new( @@ -142,6 +146,7 @@ async fn test_broker_quic_creation() { #[tokio::test] async fn test_broker_default_quic_port() { let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) .with_bind_address(([127, 0, 0, 1], 1883)) .with_quic(QuicConfig::new( PathBuf::from("../../test_certs/server.pem"), @@ -715,6 +720,7 @@ async fn start_quic_broker_with_early_data() -> (MqttBroker, SocketAddr) { let _ = rustls::crypto::ring::default_provider().install_default(); let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) .with_bind_address(([127, 0, 0, 1], 0)) .with_quic( QuicConfig::new( @@ -1019,6 +1025,7 @@ async fn test_subscribe_on_data_flow_delivers_on_server_stream() { let _ = rustls::crypto::ring::default_provider().install_default(); let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) .with_bind_address(([127, 0, 0, 1], 0)) .with_server_delivery_strategy(ServerDeliveryStrategy::PerTopic) .with_quic( @@ -1154,6 +1161,7 @@ async fn test_qos1_subscriber_over_quic_receives_message() { let _ = rustls::crypto::ring::default_provider().install_default(); let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) .with_bind_address(([127, 0, 0, 1], 0)) .with_server_delivery_strategy(ServerDeliveryStrategy::PerTopic) .with_quic( diff --git a/crates/mqtt5/tests/quic_multistream_integration.rs b/crates/mqtt5/tests/quic_multistream_integration.rs index 612d263c..ec36b511 100644 --- a/crates/mqtt5/tests/quic_multistream_integration.rs +++ b/crates/mqtt5/tests/quic_multistream_integration.rs @@ -1,7 +1,9 @@ #![cfg(feature = "broker")] #![cfg(feature = "transport-quic")] -use mqtt5::broker::config::{BrokerConfig, QuicConfig}; +use mqtt5::broker::config::{ + BrokerConfig, QuicConfig, ServerDeliveryStrategy, StorageBackend, StorageConfig, +}; use mqtt5::broker::MqttBroker; use mqtt5::session::quic_flow::{FlowRegistry, FlowState, FlowType}; use mqtt5::time::Duration; @@ -18,18 +20,21 @@ fn test_client_id(prefix: &str) -> String { format!("{}-{}", prefix, Ulid::new()) } -async fn start_quic_broker() -> (MqttBroker, SocketAddr) { +fn test_quic_config() -> QuicConfig { + QuicConfig::new( + PathBuf::from("../../test_certs/server.pem"), + PathBuf::from("../../test_certs/server.key"), + ) + .with_bind_address("127.0.0.1:0".parse::().unwrap()) +} + +async fn start_quic_broker_with(quic_config: QuicConfig) -> (MqttBroker, SocketAddr) { let _ = rustls::crypto::ring::default_provider().install_default(); let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) .with_bind_address(([127, 0, 0, 1], 0)) - .with_quic( - QuicConfig::new( - PathBuf::from("../../test_certs/server.pem"), - PathBuf::from("../../test_certs/server.key"), - ) - .with_bind_address("127.0.0.1:0".parse::().unwrap()), - ); + .with_quic(quic_config); let broker = MqttBroker::with_config(config).await.unwrap(); let quic_addr = broker @@ -38,6 +43,59 @@ async fn start_quic_broker() -> (MqttBroker, SocketAddr) { (broker, quic_addr) } +async fn start_quic_broker() -> (MqttBroker, SocketAddr) { + start_quic_broker_with(test_quic_config()).await +} + +async fn publish_and_count( + quic_config: QuicConfig, + strategy: StreamStrategy, + messages: u32, + payload_size: usize, +) -> u32 { + let (mut broker, quic_addr) = start_quic_broker_with(quic_config).await; + let broker_handle = tokio::spawn(async move { broker.run().await }); + + let topic = format!("quic-transport-limit/{}/test", Ulid::new()); + let pub_client = MqttClient::new(test_client_id("quic-limit-pub")); + let sub_client = MqttClient::new(test_client_id("quic-limit-sub")); + + pub_client.set_insecure_tls(true).await; + sub_client.set_insecure_tls(true).await; + pub_client.set_quic_stream_strategy(strategy).await; + + let broker_url = format!("quic://{quic_addr}"); + pub_client.connect(&broker_url).await.unwrap(); + sub_client.connect(&broker_url).await.unwrap(); + + let received = Arc::new(AtomicU32::new(0)); + let received_clone = received.clone(); + sub_client + .subscribe(&topic, move |_msg| { + received_clone.fetch_add(1, Ordering::Relaxed); + }) + .await + .unwrap(); + + tokio::time::sleep(Duration::from_millis(100)).await; + + let payload = vec![0u8; payload_size]; + for _ in 0..messages { + pub_client.publish(&topic, payload.clone()).await.unwrap(); + } + + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + while received.load(Ordering::Relaxed) < messages && tokio::time::Instant::now() < deadline { + tokio::time::sleep(Duration::from_millis(20)).await; + } + + let count = received.load(Ordering::Relaxed); + pub_client.disconnect().await.unwrap(); + sub_client.disconnect().await.unwrap(); + broker_handle.abort(); + count +} + #[tokio::test] async fn test_flow_registry_basic_operations() { let mut registry = FlowRegistry::new(100); @@ -270,14 +328,8 @@ async fn test_quic_data_per_publish_with_broker() { let broker_url = format!("quic://{quic_addr}"); - if pub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } - if sub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + pub_client.connect(&broker_url).await.unwrap(); + sub_client.connect(&broker_url).await.unwrap(); let received = Arc::new(AtomicU32::new(0)); let received_clone = received.clone(); @@ -336,14 +388,8 @@ async fn test_quic_multiple_topics_with_flow_isolation() { let broker_url = format!("quic://{quic_addr}"); - if pub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } - if sub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + pub_client.connect(&broker_url).await.unwrap(); + sub_client.connect(&broker_url).await.unwrap(); let topic1_count = Arc::new(AtomicU32::new(0)); let topic2_count = Arc::new(AtomicU32::new(0)); @@ -419,14 +465,8 @@ async fn test_quic_control_only_strategy() { let broker_url = format!("quic://{quic_addr}"); - if pub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } - if sub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + pub_client.connect(&broker_url).await.unwrap(); + sub_client.connect(&broker_url).await.unwrap(); let received = Arc::new(AtomicU32::new(0)); let received_clone = received.clone(); @@ -484,14 +524,8 @@ async fn test_quic_mixed_qos_with_streams() { let broker_url = format!("quic://{quic_addr}"); - if pub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } - if sub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + pub_client.connect(&broker_url).await.unwrap(); + sub_client.connect(&broker_url).await.unwrap(); let received = Arc::new(AtomicU32::new(0)); let received_clone = received.clone(); @@ -536,10 +570,7 @@ async fn test_quic_concurrent_publishers() { let sub_client = MqttClient::new(test_client_id("quic-conc-sub")); sub_client.set_insecure_tls(true).await; - if sub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + sub_client.connect(&broker_url).await.unwrap(); let received = Arc::new(AtomicU32::new(0)); let received_clone = received.clone(); @@ -615,14 +646,8 @@ async fn test_quic_large_payload_per_stream() { let broker_url = format!("quic://{quic_addr}"); - if pub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } - if sub_client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + pub_client.connect(&broker_url).await.unwrap(); + sub_client.connect(&broker_url).await.unwrap(); let received_sizes = Arc::new(tokio::sync::Mutex::new(Vec::new())); let received_clone = received_sizes.clone(); @@ -746,10 +771,7 @@ async fn test_discard_flow_removes_peer_state() { let broker_url = format!("quic://{quic_addr}"); - if client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + client.connect(&broker_url).await.unwrap(); client.publish(&topic, b"establish flow").await.unwrap(); tokio::time::sleep(Duration::from_millis(100)).await; @@ -786,10 +808,7 @@ async fn test_discard_flow_without_flow_headers_returns_error() { .await; let broker_url = format!("quic://{quic_addr}"); - if client.connect(&broker_url).await.is_err() { - broker_handle.abort(); - return; - } + client.connect(&broker_url).await.unwrap(); let result = client.discard_flow(FlowId::client(1)).await; assert!( @@ -800,3 +819,103 @@ async fn test_discard_flow_without_flow_headers_returns_error() { client.disconnect().await.unwrap(); broker_handle.abort(); } + +#[tokio::test] +async fn test_quic_per_publish_under_tight_stream_limit() { + let received = publish_and_count( + test_quic_config().with_max_concurrent_streams(2), + StreamStrategy::DataPerPublish, + 20, + 16, + ) + .await; + assert_eq!( + received, 20, + "per-publish must retire and reuse streams when the peer limit is below the message count" + ); +} + +#[tokio::test] +async fn test_quic_control_only_under_small_stream_window() { + let received = publish_and_count( + test_quic_config().with_stream_receive_window(2048), + StreamStrategy::ControlOnly, + 50, + 512, + ) + .await; + assert_eq!( + received, 50, + "a single stream must keep flowing when its receive window is far below the bytes sent" + ); +} + +#[tokio::test] +async fn test_quic_delivery_with_segmentation_offload_disabled() { + let received = publish_and_count( + test_quic_config().with_disable_segmentation_offload(true), + StreamStrategy::ControlOnly, + 50, + 512, + ) + .await; + assert_eq!( + received, 50, + "disabling segmentation offload must not affect delivery" + ); +} + +#[tokio::test] +async fn test_quic_per_publish_delivery_under_tight_client_stream_limit() { + let _ = rustls::crypto::ring::default_provider().install_default(); + + let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) + .with_bind_address(([127, 0, 0, 1], 0)) + .with_quic(test_quic_config()) + .with_server_delivery_strategy(ServerDeliveryStrategy::PerPublish); + let mut broker = MqttBroker::with_config(config).await.unwrap(); + let quic_addr = broker + .quic_local_addr() + .expect("QUIC endpoint must be bound"); + let broker_handle = tokio::spawn(async move { broker.run().await }); + + let topic = format!("quic-client-limit/{}/test", Ulid::new()); + let pub_client = MqttClient::new(test_client_id("quic-client-limit-pub")); + let sub_client = MqttClient::new(test_client_id("quic-client-limit-sub")); + pub_client.set_insecure_tls(true).await; + sub_client.set_insecure_tls(true).await; + sub_client.set_quic_max_streams(Some(2)).await; + + let broker_url = format!("quic://{quic_addr}"); + pub_client.connect(&broker_url).await.unwrap(); + sub_client.connect(&broker_url).await.unwrap(); + + let received = Arc::new(AtomicU32::new(0)); + let received_clone = received.clone(); + sub_client + .subscribe(&topic, move |_msg| { + received_clone.fetch_add(1, Ordering::Relaxed); + }) + .await + .unwrap(); + + for _ in 0..20 { + pub_client.publish(&topic, vec![0u8; 16]).await.unwrap(); + } + + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + while received.load(Ordering::Relaxed) < 20 && tokio::time::Instant::now() < deadline { + tokio::time::sleep(Duration::from_millis(20)).await; + } + + assert_eq!( + received.load(Ordering::Relaxed), + 20, + "the broker must reuse streams when the client caps concurrent streams below the message count" + ); + + pub_client.disconnect().await.unwrap(); + sub_client.disconnect().await.unwrap(); + broker_handle.abort(); +} diff --git a/crates/mqtt5/tests/quic_stats_sampling.rs b/crates/mqtt5/tests/quic_stats_sampling.rs new file mode 100644 index 00000000..a32318d0 --- /dev/null +++ b/crates/mqtt5/tests/quic_stats_sampling.rs @@ -0,0 +1,94 @@ +#![cfg(feature = "broker")] +#![cfg(feature = "transport-quic")] + +use mqtt5::broker::config::{ + BrokerConfig, QuicConfig as BrokerQuicConfig, StorageBackend, StorageConfig, +}; +use mqtt5::broker::MqttBroker; +use mqtt5::MqttClient; +use std::net::SocketAddr; +use std::path::{Path, PathBuf}; +use std::time::Duration; + +const HEADER: &str = "timestamp_ns,rtt_us,cwnd,lost_packets,congestion_events,sent_packets,stream_data_blocked,data_blocked,streams_blocked_uni"; + +fn stats_files(dir: &Path) -> Vec { + std::fs::read_dir(dir) + .map(|entries| entries.filter_map(|e| e.ok().map(|e| e.path())).collect()) + .unwrap_or_default() +} + +#[tokio::test] +async fn broker_writes_quic_stats_per_connection_when_the_directory_is_set() { + let _ = rustls::crypto::ring::default_provider().install_default(); + let stats_dir = tempfile::tempdir().unwrap(); + std::env::set_var("MQTT5_QUIC_STATS_DIR", stats_dir.path()); + + let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) + .with_bind_address(([127, 0, 0, 1], 0)) + .with_quic( + BrokerQuicConfig::new( + PathBuf::from("../../test_certs/server.pem"), + PathBuf::from("../../test_certs/server.key"), + ) + .with_bind_address("127.0.0.1:0".parse::().unwrap()) + .with_stream_receive_window(2048), + ); + let mut broker = MqttBroker::with_config(config).await.unwrap(); + let addr = broker + .quic_local_addr() + .expect("QUIC endpoint must be bound"); + let broker_handle = tokio::spawn(async move { + let _ = broker.run().await; + }); + + let client = MqttClient::new("quic-stats-sampler"); + client.set_insecure_tls(true).await; + client.connect(&format!("quic://{addr}")).await.unwrap(); + let payload = vec![0u8; 4096]; + for _ in 0..50 { + client + .publish("quic-stats/burst", payload.clone()) + .await + .unwrap(); + } + client.disconnect().await.unwrap(); + + let deadline = tokio::time::Instant::now() + Duration::from_secs(5); + let rows = loop { + let files = stats_files(stats_dir.path()); + assert!(files.len() <= 1, "one file per connection: {files:?}"); + let rows: Vec = files + .first() + .and_then(|path| std::fs::read_to_string(path).ok()) + .map(|text| text.lines().map(str::to_owned).collect()) + .unwrap_or_default(); + let final_row_written = rows + .last() + .is_some_and(|row| row.split(',').nth(6).is_some_and(|v| v != "0")); + if final_row_written || tokio::time::Instant::now() >= deadline { + break rows; + } + tokio::time::sleep(Duration::from_millis(50)).await; + }; + + std::env::remove_var("MQTT5_QUIC_STATS_DIR"); + broker_handle.abort(); + + assert_eq!(rows.first().map(String::as_str), Some(HEADER)); + assert!(rows.len() >= 3, "expected samples, got {rows:?}"); + let columns = HEADER.split(',').count(); + assert!(rows.iter().all(|row| row.split(',').count() == columns)); + let last: Vec = rows + .last() + .unwrap() + .split(',') + .map(|v| v.parse().unwrap()) + .collect(); + assert!(last[5] > 0, "sent_packets"); + assert!( + last[6] > 0, + "the client's STREAM_DATA_BLOCKED must be counted: {last:?}" + ); +} diff --git a/crates/mqtt5/tests/quic_transport_limits.rs b/crates/mqtt5/tests/quic_transport_limits.rs new file mode 100644 index 00000000..c88dbbed --- /dev/null +++ b/crates/mqtt5/tests/quic_transport_limits.rs @@ -0,0 +1,158 @@ +#![cfg(feature = "broker")] +#![cfg(feature = "transport-quic")] + +use mqtt5::broker::config::{ + BrokerConfig, QuicConfig as BrokerQuicConfig, StorageBackend, StorageConfig, +}; +use mqtt5::broker::quic_acceptor::QuicAcceptorConfig; +use mqtt5::broker::MqttBroker; +use mqtt5::transport::{QuicConfig, QuicSplitResult, QuicTransport}; +use mqtt5::Transport; +use std::net::SocketAddr; +use std::path::PathBuf; +use std::time::Duration; + +const CERT: &str = "../../test_certs/server.pem"; +const KEY: &str = "../../test_certs/server.key"; + +fn install_crypto_provider() { + let _ = rustls::crypto::ring::default_provider().install_default(); +} + +async fn start_broker(quic_config: BrokerQuicConfig) -> (SocketAddr, tokio::task::JoinHandle<()>) { + install_crypto_provider(); + let config = BrokerConfig::default() + .with_storage(StorageConfig::default().with_backend(StorageBackend::Memory)) + .with_bind_address(([127, 0, 0, 1], 0)) + .with_quic(quic_config); + let mut broker = MqttBroker::with_config(config).await.unwrap(); + let quic_addr = broker + .quic_local_addr() + .expect("QUIC endpoint must be bound"); + let handle = tokio::spawn(async move { + let _ = broker.run().await; + }); + (quic_addr, handle) +} + +fn broker_quic_config() -> BrokerQuicConfig { + BrokerQuicConfig::new(PathBuf::from(CERT), PathBuf::from(KEY)) + .with_bind_address("127.0.0.1:0".parse::().unwrap()) +} + +async fn connect(config: QuicConfig) -> QuicSplitResult { + let mut transport = QuicTransport::new(config); + transport.connect().await.unwrap(); + transport.into_split().unwrap() +} + +fn client_config(addr: SocketAddr) -> QuicConfig { + QuicConfig::new(addr, "localhost").with_verify_server_cert(false) +} + +async fn open_uni_streams(connection: &quinn::Connection, count: usize) -> Vec { + let mut streams = Vec::with_capacity(count); + for _ in 0..count { + match tokio::time::timeout(Duration::from_millis(200), connection.open_uni()).await { + Ok(Ok(stream)) => streams.push(stream), + Ok(Err(e)) => panic!("opening a uni stream failed: {e}"), + Err(_) => break, + } + } + tokio::time::sleep(Duration::from_millis(100)).await; + streams +} + +async fn burst_on_one_stream(connection: &quinn::Connection, bytes: usize) { + let mut stream = connection.open_uni().await.unwrap(); + let payload = vec![0u8; bytes]; + let _ = tokio::time::timeout(Duration::from_millis(500), stream.write_all(&payload)).await; + tokio::time::sleep(Duration::from_millis(100)).await; +} + +#[tokio::test] +async fn broker_stream_limit_stops_a_client_opening_more_streams() { + let (addr, broker) = start_broker(broker_quic_config().with_max_concurrent_streams(2)).await; + let split = connect(client_config(addr)).await; + + let opened = open_uni_streams(&split.connection, 20).await; + + assert_eq!(opened.len(), 2); + broker.abort(); +} + +#[tokio::test] +async fn default_broker_stream_limit_does_not_block_twenty_streams() { + let (addr, broker) = start_broker(broker_quic_config()).await; + let split = connect(client_config(addr)).await; + + let opened = open_uni_streams(&split.connection, 20).await; + + assert_eq!(opened.len(), 20); + broker.abort(); +} + +#[tokio::test] +async fn broker_stream_window_blocks_a_burst_on_one_stream() { + let (addr, broker) = start_broker(broker_quic_config().with_stream_receive_window(2048)).await; + let split = connect(client_config(addr)).await; + + burst_on_one_stream(&split.connection, 65_536).await; + + assert!(split.connection.stats().frame_tx.stream_data_blocked > 0); + broker.abort(); +} + +#[tokio::test] +async fn default_broker_stream_window_does_not_block_the_same_burst() { + let (addr, broker) = start_broker(broker_quic_config()).await; + let split = connect(client_config(addr)).await; + + burst_on_one_stream(&split.connection, 65_536).await; + + assert_eq!(split.connection.stats().frame_tx.stream_data_blocked, 0); + broker.abort(); +} + +async fn streams_opened_toward_client(client_limit: Option) -> usize { + install_crypto_provider(); + let certs = QuicAcceptorConfig::load_cert_chain_from_file(CERT) + .await + .unwrap(); + let key = QuicAcceptorConfig::load_private_key_from_file(KEY) + .await + .unwrap(); + let server_config = QuicAcceptorConfig::new(certs, key) + .build_server_config() + .unwrap(); + let endpoint = quinn::Endpoint::server(server_config, "127.0.0.1:0".parse().unwrap()).unwrap(); + let addr = endpoint.local_addr().unwrap(); + + let server = tokio::spawn(async move { + let incoming = endpoint.accept().await.unwrap(); + let connection = incoming.await.unwrap(); + let opened = open_uni_streams(&connection, 20).await.len(); + (opened, endpoint) + }); + + let config = match client_limit { + Some(max) => client_config(addr).with_max_concurrent_streams(max), + None => client_config(addr), + }; + let split = connect(config).await; + + let (opened, endpoint) = server.await.unwrap(); + drop(split); + endpoint.close(0u32.into(), b"done"); + opened +} + +#[tokio::test] +async fn client_stream_limit_stops_a_server_opening_more_streams() { + assert_eq!(streams_opened_toward_client(Some(2)).await, 2); +} + +#[tokio::test] +async fn default_client_stream_limit_does_not_block_twenty_streams() { + assert_eq!(streams_opened_toward_client(None).await, 20); +} diff --git a/crates/mqttv5-cli/CLI_USAGE.md b/crates/mqttv5-cli/CLI_USAGE.md index 2cc13828..bb9dd3aa 100644 --- a/crates/mqttv5-cli/CLI_USAGE.md +++ b/crates/mqttv5-cli/CLI_USAGE.md @@ -38,6 +38,10 @@ Every flag on the `broker`, `pub`, and `sub` subcommands can be set via environm | `--non-interactive` | `MQTT5_NON_INTERACTIVE` | | | `--otel-endpoint` | `MQTT5_OTEL_ENDPOINT` | | +### QUIC Statistics + +`MQTT5_QUIC_STATS_DIR` has no flag. When it names a directory, the broker writes one CSV per QUIC connection, `broker_quic_.csv`, with a row every 100 ms and a last row when the connection closes. Columns: `timestamp_ns`, `rtt_us`, `cwnd`, `lost_packets`, `congestion_events`, `sent_packets` (the broker's own path, so they describe the broker-to-client direction), then `stream_data_blocked`, `data_blocked` and `streams_blocked_uni` (blocked frames received from the client). Clients built on quinn 0.11 never send STREAMS_BLOCKED, so `streams_blocked_uni` stays 0 with them. + ### Precedence CLI flag > environment variable > default value. @@ -122,6 +126,9 @@ mqttv5 broker generate-config [--output FILE] [--format json|toml] | `--quic-host ` | QUIC bind address(es), requires TLS cert/key | None | | `--quic-delivery-strategy ` | QUIC server delivery strategy: `control-only`, `per-topic`, `per-publish` | `per-topic` | | `--quic-early-data` | Enable QUIC 0-RTT early data | `false` | +| `--quic-max-streams ` | Maximum concurrent QUIC streams a client may open, per direction | `100` | +| `--quic-stream-window ` | Per-stream QUIC receive window | `262144` | +| `--quic-disable-offload` | Disable UDP segmentation offload (one datagram per send) | `false` | | `--storage-dir ` | Storage directory for persistence | `./mqtt_storage` | | `--storage-backend ` | Storage backend: `memory` or `file` | `file` | | `--no-persistence` | Disable message persistence | `false` | diff --git a/crates/mqttv5-cli/Cargo.toml b/crates/mqttv5-cli/Cargo.toml index 34c1ea91..4edf646d 100644 --- a/crates/mqttv5-cli/Cargo.toml +++ b/crates/mqttv5-cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqttv5-cli" -version = "0.29.0" +version = "0.29.1" edition.workspace = true rust-version.workspace = true authors.workspace = true @@ -22,7 +22,7 @@ opentelemetry = ["mqtt5/opentelemetry"] codec = ["mqtt5/codec-all"] [dependencies] -mqtt5 = { path = "../mqtt5", version = "0.42" } +mqtt5 = { path = "../mqtt5", version = "0.43" } anyhow = "1.0.103" tracing = "0.1" serde = { version = "1.0", features = ["derive"] } diff --git a/crates/mqttv5-cli/src/commands/broker_cmd.rs b/crates/mqttv5-cli/src/commands/broker_cmd.rs index 29d09d3b..914b88a2 100644 --- a/crates/mqttv5-cli/src/commands/broker_cmd.rs +++ b/crates/mqttv5-cli/src/commands/broker_cmd.rs @@ -195,6 +195,18 @@ pub struct RunArgs { #[arg(long, env = "MQTT5_QUIC_EARLY_DATA")] pub quic_early_data: bool, + /// Maximum concurrent QUIC streams a peer may open (default 100) + #[arg(long, env = "MQTT5_QUIC_MAX_STREAMS")] + pub quic_max_streams: Option, + + /// Per-stream QUIC receive window in bytes (default 262144) + #[arg(long, env = "MQTT5_QUIC_STREAM_WINDOW")] + pub quic_stream_window: Option, + + /// Disable QUIC UDP segmentation offload (send one datagram per syscall) + #[arg(long, env = "MQTT5_QUIC_DISABLE_OFFLOAD")] + pub quic_disable_offload: bool, + /// Storage directory for persistent data #[arg(long, default_value = "./mqtt_storage", env = "MQTT5_STORAGE_DIR")] pub storage_dir: PathBuf, @@ -336,17 +348,10 @@ fn build_example_config() -> BrokerConfig { subprotocol: "mqtt".to_string(), use_tls: true, }), - quic_config: Some(QuicConfig { - cert_file: PathBuf::from("/path/to/server.crt"), - key_file: PathBuf::from("/path/to/server.key"), - ca_file: None, - require_client_cert: false, - bind_addresses: vec![ - "0.0.0.0:14567".parse().unwrap(), - "[::]:14567".parse().unwrap(), - ], - enable_early_data: false, - }), + quic_config: Some(QuicConfig::new( + PathBuf::from("/path/to/server.crt"), + PathBuf::from("/path/to/server.key"), + )), cluster_listener_config: None, storage_config: StorageConfig { backend: StorageBackend::File, @@ -1004,6 +1009,15 @@ fn configure_quic(config: &mut BrokerConfig, cmd: &RunArgs) -> Result<()> { if cmd.quic_early_data { quic_config = quic_config.with_early_data(true); } + if let Some(max) = cmd.quic_max_streams { + quic_config = quic_config.with_max_concurrent_streams(max); + } + if let Some(bytes) = cmd.quic_stream_window { + quic_config = quic_config.with_stream_receive_window(bytes); + } + if cmd.quic_disable_offload { + quic_config = quic_config.with_disable_segmentation_offload(true); + } config.quic_config = Some(quic_config); if let Some(strategy) = cmd.quic_delivery_strategy { config.server_delivery_strategy = strategy; diff --git a/specs/tla/qos1-backlog/Qos1Backlog.cfg b/specs/tla/qos1-backlog/Qos1Backlog.cfg new file mode 100644 index 00000000..5214949b --- /dev/null +++ b/specs/tla/qos1-backlog/Qos1Backlog.cfg @@ -0,0 +1,27 @@ +SPECIFICATION Spec + +CONSTANTS + NPubs = 1 + NSubs = 2 + MaxMsgs = 2 + Cap = 1 + RM = 1 + Batch = 1 + MaxTakeovers = 1 + MaxPurges = 1 + NotifyTrigger = TRUE + +INVARIANTS + TypeOK + InvChanBound + InvWindowBound + InvBatchBound + InvNoLoss + InvNoDuplicate + InvBacklogArmed + InvFifo + +PROPERTIES + Live + +CHECK_DEADLOCK FALSE diff --git a/specs/tla/qos1-backlog/Qos1Backlog.tla b/specs/tla/qos1-backlog/Qos1Backlog.tla new file mode 100644 index 00000000..fd98180b --- /dev/null +++ b/specs/tla/qos1-backlog/Qos1Backlog.tla @@ -0,0 +1,349 @@ +---------------------------- MODULE Qos1Backlog ---------------------------- +(***************************************************************************) +(* Broker QoS 1 delivery under overload: bounded per-client delivery *) +(* channel, outbound window, storage queue, Notify permit, and the *) +(* handler-side backlog drain. *) +(* *) +(* There is no separately stored "backlog flag": a client is behind *) +(* (Backlog(s)) exactly when its storage queue is non-empty, and the *) +(* router's emptiness check, the router's append and the handler's limited *) +(* take are all performed under the storage queue's own lock, so each is *) +(* one atomic step here. (Qos1BacklogLockFree shows why a flag that is *) +(* cleared before the take breaks FIFO.) *) +(* *) +(* Every publisher publishes MaxMsgs messages, each routed to every *) +(* subscriber at QoS 1. The router routes a message for subscriber s: *) +(* - behind the backlog (append to storage) while Backlog(s) or while a *) +(* reconnect hand-off is in progress, so FIFO order is kept; *) +(* - else into the channel if it has room; *) +(* - else it waits with a deadline; a timeout (possibly spurious, even *) +(* when the channel has room) queues to storage. *) +(* After every queue to storage the router posts a permit. *) +(* The handler consumes the channel only while its window has room (RM is *) +(* min(client Receive Maximum, broker cap)), and drains storage only when *) +(* the channel is empty, taking at most Min(free window, Batch) from the *) +(* FRONT of storage; acks are read only between batches. Every trigger -- *) +(* (i) the channel running empty, (ii) an ack freeing a slot, (iii) a *) +(* router queue, (iv) a finished batch with storage left, (v) the end of a *) +(* hand-off -- posts the permit; the drain runs only from the Notify arm. *) +(* *) +(* Purge: the expiry sweep may remove any queued message at any time. *) +(* Takeover (two phases): TakeoverBegin -- the new handler registers; the *) +(* old handler's unacked in-flight (send order), unsent batch and channel *) +(* are in hand-off and every router routes behind; the new handler's *) +(* channel arm and drain wait. TakeoverEnd -- the old handler has *) +(* re-queued them at the FRONT of storage; the new handler gets a permit. *) +(* CleanTakeover: a clean-start reconnect discards everything queued, *) +(* held, in the channel and in flight (MQTT-3.1.2-4) in one step. *) +(* *) +(* CONSTANT NotifyTrigger selects whether the router posts a permit after *) +(* queueing (trigger iii): *) +(* FALSE -- a timeout that raced an already-empty channel leaves storage *) +(* non-empty with nothing to wake the drain. Expected: *) +(* InvBacklogArmed and Live VIOLATED. *) +(* TRUE -- the permit wakes the drain. Expected: holds. *) +(***************************************************************************) +EXTENDS Naturals, Sequences, FiniteSets + +CONSTANTS + NPubs, \* publishers + NSubs, \* subscribers, each subscribed to everything at QoS 1 + MaxMsgs, \* messages per publisher + Cap, \* delivery channel capacity (client_channel_capacity) + RM, \* outbound window: min(client Receive Maximum, max_outbound_inflight) + Batch, \* messages written per drain invocation + MaxTakeovers, \* reconnects allowed per subscriber (state-space bound) + MaxPurges, \* expiry purges allowed per subscriber (state-space bound) + NotifyTrigger \* TRUE = the router posts a permit after every queue + +Pubs == 1..NPubs +Subs == 1..NSubs +Msgs == [p : Pubs, n : 1..MaxMsgs] + +VARIABLES + pubSent, \* per publisher: messages published so far + targets, \* per publisher: subscribers the current message is still being routed to + chan, \* per subscriber: bounded delivery channel (FIFO) + storage, \* per subscriber: storage queue (FIFO) + held, \* per subscriber: batch taken from storage by the running drain (FIFO) + draining, \* per subscriber: a drain batch is between take and finish + permit, \* per subscriber: stored Notify permit + drainPending, \* per subscriber: the Notify arm fired and the drain has not run yet + inflight, \* per subscriber: sent, awaiting PUBACK + delivered, \* per subscriber: acked + purged, \* per subscriber: removed from storage by the expiry sweep + discarded, \* per subscriber: dropped by a clean-start reconnect + sent, \* per subscriber: first-send order to the transport (a re-send does not change it) + handoff, \* per subscriber: a reconnect hand-off is in progress + oldChan, \* per subscriber: what the displaced handler still has to re-queue (age order) + takeovers \* per subscriber: reconnects so far + +vars == <> +subVars == <> + +Range(seq) == {seq[i] : i \in DOMAIN seq} +Min(a, b) == IF a < b THEN a ELSE b +Free(s) == RM - Cardinality(inflight[s]) +Backlog(s) == storage[s] /= <<>> +Behind(s) == Backlog(s) \/ handoff[s] +RemoveAt(seq, i) == SubSeq(seq, 1, i - 1) \o SubSeq(seq, i + 1, Len(seq)) +FirstSend(seq, m) == IF m \in Range(seq) THEN seq ELSE Append(seq, m) + +TypeOK == + /\ pubSent \in [Pubs -> 0..MaxMsgs] + /\ targets \in [Pubs -> SUBSET Subs] + /\ chan \in [Subs -> Seq(Msgs)] + /\ storage \in [Subs -> Seq(Msgs)] + /\ held \in [Subs -> Seq(Msgs)] + /\ draining \in [Subs -> BOOLEAN] + /\ permit \in [Subs -> BOOLEAN] + /\ drainPending \in [Subs -> BOOLEAN] + /\ inflight \in [Subs -> SUBSET Msgs] + /\ delivered \in [Subs -> SUBSET Msgs] + /\ purged \in [Subs -> SUBSET Msgs] + /\ discarded \in [Subs -> SUBSET Msgs] + /\ sent \in [Subs -> Seq(Msgs)] + /\ handoff \in [Subs -> BOOLEAN] + /\ oldChan \in [Subs -> Seq(Msgs)] + /\ takeovers \in [Subs -> 0..MaxTakeovers] + +Init == + /\ pubSent = [p \in Pubs |-> 0] + /\ targets = [p \in Pubs |-> {}] + /\ chan = [s \in Subs |-> <<>>] + /\ storage = [s \in Subs |-> <<>>] + /\ held = [s \in Subs |-> <<>>] + /\ draining = [s \in Subs |-> FALSE] + /\ permit = [s \in Subs |-> FALSE] + /\ drainPending = [s \in Subs |-> FALSE] + /\ inflight = [s \in Subs |-> {}] + /\ delivered = [s \in Subs |-> {}] + /\ purged = [s \in Subs |-> {}] + /\ discarded = [s \in Subs |-> {}] + /\ sent = [s \in Subs |-> <<>>] + /\ handoff = [s \in Subs |-> FALSE] + /\ oldChan = [s \in Subs |-> <<>>] + /\ takeovers = [s \in Subs |-> 0] + +Cur(p) == [p |-> p, n |-> pubSent[p]] + +(* Publisher handler: a PUBLISH arrives; PUBACK is withheld until targets = {} *) +PubStart(p) == + /\ targets[p] = {} + /\ pubSent[p] < MaxMsgs + /\ pubSent' = [pubSent EXCEPT ![p] = @ + 1] + /\ targets' = [targets EXCEPT ![p] = Subs] + /\ UNCHANGED subVars + +(* Router: append to storage under the storage lock, then post a permit *) +QueueBehind(p, s) == + /\ storage' = [storage EXCEPT ![s] = Append(@, Cur(p))] + /\ permit' = [permit EXCEPT ![s] = @ \/ NotifyTrigger] + /\ targets' = [targets EXCEPT ![p] = @ \ {s}] + /\ UNCHANGED <> + +(* Router: storage non-empty or hand-off in progress -- go behind without waiting *) +RouteBehind(p, s) == + /\ s \in targets[p] + /\ Behind(s) + /\ QueueBehind(p, s) + +(* Router: not behind, channel has room *) +RouteEnqueue(p, s) == + /\ s \in targets[p] + /\ ~Behind(s) + /\ Len(chan[s]) < Cap + /\ chan' = [chan EXCEPT ![s] = Append(@, Cur(p))] + /\ targets' = [targets EXCEPT ![p] = @ \ {s}] + /\ UNCHANGED <> + +(* Router: not behind, deadline elapsed -- possibly spuriously *) +RouteTimeout(p, s) == + /\ s \in targets[p] + /\ ~Behind(s) + /\ QueueBehind(p, s) + +(* Subscriber handler: channel arm, polled only while the window has room and no hand-off is pending *) +Consume(s) == + /\ chan[s] /= <<>> + /\ ~draining[s] + /\ ~handoff[s] + /\ Free(s) > 0 + /\ LET m == Head(chan[s]) + rest == Tail(chan[s]) + IN /\ chan' = [chan EXCEPT ![s] = rest] + /\ inflight' = [inflight EXCEPT ![s] = @ \cup {m}] + /\ sent' = [sent EXCEPT ![s] = FirstSend(@, m)] + /\ permit' = [permit EXCEPT ![s] = @ \/ (rest = <<>> /\ Backlog(s))] + /\ UNCHANGED <> + +(* Subscriber handler: the Notify arm fires -- the only entry into the drain *) +Notified(s) == + /\ permit[s] + /\ ~draining[s] + /\ ~handoff[s] + /\ permit' = [permit EXCEPT ![s] = FALSE] + /\ drainPending' = [drainPending EXCEPT ![s] = TRUE] + /\ UNCHANGED <> + +(* drain_backlog: storage empty -- nothing to do *) +DrainIdle(s) == + /\ drainPending[s] + /\ ~draining[s] + /\ ~Backlog(s) + /\ drainPending' = [drainPending EXCEPT ![s] = FALSE] + /\ UNCHANGED <> + +(* drain_backlog: the channel is older than storage, or the window is full -- come back later *) +DrainDefer(s) == + /\ drainPending[s] + /\ ~draining[s] + /\ Backlog(s) + /\ chan[s] /= <<>> \/ Free(s) = 0 + /\ drainPending' = [drainPending EXCEPT ![s] = FALSE] + /\ UNCHANGED <> + +(* drain_backlog: take at most Min(free window, Batch) from the front of storage, under the storage lock *) +DrainBegin(s) == + /\ drainPending[s] + /\ ~draining[s] + /\ Backlog(s) + /\ chan[s] = <<>> + /\ Free(s) > 0 + /\ LET k == Min(Len(storage[s]), Min(Free(s), Batch)) + IN /\ held' = [held EXCEPT ![s] = SubSeq(storage[s], 1, k)] + /\ storage' = [storage EXCEPT ![s] = SubSeq(storage[s], k + 1, Len(storage[s]))] + /\ draining' = [draining EXCEPT ![s] = TRUE] + /\ UNCHANGED <> + +(* drain_backlog: send the next held message *) +DrainSend(s) == + /\ draining[s] + /\ held[s] /= <<>> + /\ LET m == Head(held[s]) + IN /\ held' = [held EXCEPT ![s] = Tail(@)] + /\ inflight' = [inflight EXCEPT ![s] = @ \cup {m}] + /\ sent' = [sent EXCEPT ![s] = FirstSend(@, m)] + /\ UNCHANGED <> + +(* drain_backlog: batch done -- return to select; re-arm through the permit if storage is still non-empty *) +DrainFinish(s) == + /\ draining[s] + /\ held[s] = <<>> + /\ draining' = [draining EXCEPT ![s] = FALSE] + /\ drainPending' = [drainPending EXCEPT ![s] = FALSE] + /\ permit' = [permit EXCEPT ![s] = @ \/ Backlog(s)] + /\ UNCHANGED <> + +(* Subscriber PUBACK is read -- only between batches, only on the live connection *) +Ack(s, m) == + /\ m \in inflight[s] + /\ ~draining[s] + /\ ~handoff[s] + /\ inflight' = [inflight EXCEPT ![s] = @ \ {m}] + /\ delivered' = [delivered EXCEPT ![s] = @ \cup {m}] + /\ permit' = [permit EXCEPT ![s] = @ \/ Backlog(s)] + /\ UNCHANGED <> + +(* Expiry sweep removes one queued message; the count is the queue length, so nothing else changes *) +Purge(s) == + /\ Cardinality(purged[s]) < MaxPurges + /\ \E i \in DOMAIN storage[s] : + /\ purged' = [purged EXCEPT ![s] = @ \cup {storage[s][i]}] + /\ storage' = [storage EXCEPT ![s] = RemoveAt(@, i)] + /\ UNCHANGED <> + +(* Reconnect, phase 1: the new handler registers; the old one's custody enters the hand-off *) +TakeoverBegin(s) == + /\ takeovers[s] < MaxTakeovers + /\ ~draining[s] + /\ ~handoff[s] + /\ LET unacked == SelectSeq(sent[s], LAMBDA m : m \in inflight[s]) + IN oldChan' = [oldChan EXCEPT ![s] = unacked \o held[s] \o chan[s]] + /\ inflight' = [inflight EXCEPT ![s] = {}] + /\ held' = [held EXCEPT ![s] = <<>>] + /\ chan' = [chan EXCEPT ![s] = <<>>] + /\ drainPending' = [drainPending EXCEPT ![s] = FALSE] + /\ permit' = [permit EXCEPT ![s] = FALSE] + /\ handoff' = [handoff EXCEPT ![s] = TRUE] + /\ takeovers' = [takeovers EXCEPT ![s] = @ + 1] + /\ UNCHANGED <> + +(* Reconnect, phase 2: the old handler re-queued its custody at the front; the new handler gets a permit *) +TakeoverEnd(s) == + /\ handoff[s] + /\ storage' = [storage EXCEPT ![s] = oldChan[s] \o @] + /\ oldChan' = [oldChan EXCEPT ![s] = <<>>] + /\ handoff' = [handoff EXCEPT ![s] = FALSE] + /\ permit' = [permit EXCEPT ![s] = TRUE] + /\ UNCHANGED <> + +(* Clean-start reconnect: everything queued, held, in the channel and in flight is discarded in one step *) +CleanTakeover(s) == + /\ takeovers[s] < MaxTakeovers + /\ ~draining[s] + /\ ~handoff[s] + /\ discarded' = [discarded EXCEPT ![s] = @ \cup Range(storage[s]) \cup Range(chan[s]) \cup Range(held[s]) \cup inflight[s]] + /\ storage' = [storage EXCEPT ![s] = <<>>] + /\ chan' = [chan EXCEPT ![s] = <<>>] + /\ held' = [held EXCEPT ![s] = <<>>] + /\ inflight' = [inflight EXCEPT ![s] = {}] + /\ drainPending' = [drainPending EXCEPT ![s] = FALSE] + /\ permit' = [permit EXCEPT ![s] = TRUE] + /\ takeovers' = [takeovers EXCEPT ![s] = @ + 1] + /\ UNCHANGED <> + +Next == + \/ \E p \in Pubs : PubStart(p) + \/ \E p \in Pubs, s \in Subs : RouteBehind(p, s) \/ RouteEnqueue(p, s) \/ RouteTimeout(p, s) + \/ \E s \in Subs : + \/ Consume(s) \/ Notified(s) + \/ DrainIdle(s) \/ DrainDefer(s) \/ DrainBegin(s) \/ DrainSend(s) \/ DrainFinish(s) + \/ Purge(s) \/ TakeoverBegin(s) \/ TakeoverEnd(s) \/ CleanTakeover(s) + \/ \E s \in Subs, m \in Msgs : Ack(s, m) + +Fairness == + /\ \A p \in Pubs : WF_vars(PubStart(p)) + /\ \A p \in Pubs, s \in Subs : WF_vars(RouteBehind(p, s)) /\ WF_vars(RouteEnqueue(p, s)) + /\ \A s \in Subs : + /\ WF_vars(Consume(s)) /\ WF_vars(Notified(s)) /\ WF_vars(TakeoverEnd(s)) + /\ WF_vars(DrainIdle(s)) /\ WF_vars(DrainDefer(s)) /\ WF_vars(DrainBegin(s)) + /\ WF_vars(DrainSend(s)) /\ WF_vars(DrainFinish(s)) + /\ \A s \in Subs, m \in Msgs : WF_vars(Ack(s, m)) + +Spec == Init /\ [][Next]_vars /\ Fairness + +(* A message is routed to s once the router has finished with it for s *) +Routed(s) == + {m \in Msgs : + /\ m.n <= pubSent[m.p] + /\ ~(m.n = pubSent[m.p] /\ s \in targets[m.p])} + +Pending(s) == Range(chan[s]) \cup Range(storage[s]) \cup Range(held[s]) \cup Range(oldChan[s]) \cup inflight[s] +Custody(s) == Pending(s) \cup delivered[s] \cup purged[s] \cup discarded[s] + +InvChanBound == \A s \in Subs : Len(chan[s]) <= Cap +InvWindowBound == \A s \in Subs : Cardinality(inflight[s]) <= RM +InvBatchBound == \A s \in Subs : Len(held[s]) <= Batch +InvNoLoss == \A s \in Subs : Routed(s) = Custody(s) +InvNoDuplicate == \A s \in Subs : + /\ Len(chan[s]) = Cardinality(Range(chan[s])) + /\ Len(storage[s]) = Cardinality(Range(storage[s])) + /\ Len(held[s]) = Cardinality(Range(held[s])) + /\ Len(oldChan[s]) = Cardinality(Range(oldChan[s])) + /\ Len(sent[s]) = Cardinality(Range(sent[s])) + /\ Cardinality(Pending(s)) = Len(chan[s]) + Len(storage[s]) + Len(held[s]) + Len(oldChan[s]) + Cardinality(inflight[s]) + /\ (delivered[s] \cup purged[s] \cup discarded[s]) \cap Pending(s) = {} + /\ delivered[s] \cap purged[s] = {} + /\ delivered[s] \cap discarded[s] = {} + /\ purged[s] \cap discarded[s] = {} +(* Something will still wake the drain: a permit, a pending trigger, a channel message (i), an ack to come (ii) or the end of a hand-off (v) *) +InvBacklogArmed == \A s \in Subs : + (Backlog(s) /\ ~draining[s]) => (permit[s] \/ drainPending[s] \/ chan[s] /= <<>> \/ inflight[s] /= {} \/ handoff[s]) +InvFifo == \A s \in Subs : \A i, j \in DOMAIN sent[s] : + (i < j /\ sent[s][i].p = sent[s][j].p) => sent[s][i].n < sent[s][j].n + +AllDelivered == \A s \in Subs : delivered[s] \cup purged[s] \cup discarded[s] = Msgs +Live == <>AllDelivered + +============================================================================= diff --git a/specs/tla/qos1-backlog/Qos1BacklogLockFree.cfg b/specs/tla/qos1-backlog/Qos1BacklogLockFree.cfg new file mode 100644 index 00000000..5980d4fb --- /dev/null +++ b/specs/tla/qos1-backlog/Qos1BacklogLockFree.cfg @@ -0,0 +1,24 @@ +SPECIFICATION Spec + +CONSTANTS + NPubs = 2 + NSubs = 1 + MaxMsgs = 2 + Cap = 1 + RM = 1 + NotifyTrigger = TRUE + ClearFirst = TRUE + +INVARIANTS + TypeOK + InvChanBound + InvWindowBound + InvNoLoss + InvNoDuplicate + InvBacklogArmed + InvFifo + +PROPERTIES + Live + +CHECK_DEADLOCK FALSE diff --git a/specs/tla/qos1-backlog/Qos1BacklogLockFree.tla b/specs/tla/qos1-backlog/Qos1BacklogLockFree.tla new file mode 100644 index 00000000..cf254750 --- /dev/null +++ b/specs/tla/qos1-backlog/Qos1BacklogLockFree.tla @@ -0,0 +1,281 @@ +---------------------------- MODULE Qos1BacklogLockFree ---------------------------- +(***************************************************************************) +(* Broker QoS 1 delivery under overload: bounded per-client delivery *) +(* channel, Receive-Maximum window, storage queue, backlog flag, Notify *) +(* permit, and the handler-side backlog drain. Lock-free flag: the router *) +(* and the handler never share a lock on the hot path. *) +(* *) +(* Every publisher publishes MaxMsgs messages, each routed to every *) +(* subscriber at QoS 1. The router routes a message for subscriber s: *) +(* - behind the backlog (append to storage) while s's flag is set, so *) +(* an episode of overflow keeps FIFO order; *) +(* - else into the channel if it has room; *) +(* - else it waits with a deadline; a timeout (possibly spurious, even *) +(* when the channel has room) queues to storage. *) +(* After EVERY queue to storage the router sets the flag and posts a *) +(* permit, whatever the flag was -- the handler may have cleared it since *) +(* the router read it. *) +(* The handler consumes the channel only while its window has room (the *) +(* Receive-Maximum gate), and drains storage only when the channel is *) +(* empty: it CLEARS the flag first, then TAKES at most the free window *) +(* from the front of storage in one atomic storage operation, re-arming *) +(* the flag if anything remains; the router may append between the clear *) +(* and the take. Triggers: (i) the channel running empty, (ii) an ack *) +(* freeing a slot, (iii) the Notify permit. A drain loops while the flag *) +(* is armed when its batch finishes. *) +(* *) +(* CONSTANT NotifyTrigger selects whether trigger (iii) exists: *) +(* FALSE -- a timeout that raced an already-empty channel leaves storage *) +(* non-empty with nothing to wake the drain. Expected: Live *) +(* VIOLATED. *) +(* TRUE -- the permit wakes the drain. Expected: holds. *) +(* *) +(* CONSTANT ClearFirst selects when the drain clears the flag: *) +(* FALSE -- cleared unconditionally when the batch is done. A message *) +(* the router queued behind the backlog during the batch is *) +(* orphaned. Expected: InvBacklogArmed VIOLATED. *) +(* TRUE -- cleared before the take, re-armed if storage is left over. *) +(* Expected: holds. *) +(***************************************************************************) +EXTENDS Naturals, Sequences, FiniteSets + +CONSTANTS + NPubs, \* publishers + NSubs, \* subscribers, each subscribed to everything at QoS 1 + MaxMsgs, \* messages per publisher + Cap, \* delivery channel capacity (client_channel_capacity) + RM, \* subscriber Receive Maximum (outbound in-flight window) + NotifyTrigger, \* TRUE = drain has the Notify select branch + ClearFirst \* TRUE = flag cleared before the take, re-armed on leftover + +Pubs == 1..NPubs +Subs == 1..NSubs +Msgs == [p : Pubs, n : 1..MaxMsgs] +Phases == {"idle", "take", "send"} + +VARIABLES + pubSent, \* per publisher: messages published so far + targets, \* per publisher: subscribers the current message is still being routed to + chan, \* per subscriber: bounded delivery channel (FIFO) + storage, \* per subscriber: storage queue (FIFO) + held, \* per subscriber: batch taken from storage by the running drain (FIFO) + phase, \* per subscriber: drain phase + backlog, \* per subscriber: the AtomicBool flag + permit, \* per subscriber: stored Notify permit + drainPending, \* per subscriber: a drain trigger has fired and not yet been serviced + inflight, \* per subscriber: sent, awaiting PUBACK + delivered, \* per subscriber: acked + sent \* per subscriber: every message in the order it was written to the transport + +vars == <> +subVars == <> + +Range(seq) == {seq[i] : i \in DOMAIN seq} +Min(a, b) == IF a < b THEN a ELSE b +Free(s) == RM - Cardinality(inflight[s]) + +TypeOK == + /\ pubSent \in [Pubs -> 0..MaxMsgs] + /\ targets \in [Pubs -> SUBSET Subs] + /\ chan \in [Subs -> Seq(Msgs)] + /\ storage \in [Subs -> Seq(Msgs)] + /\ held \in [Subs -> Seq(Msgs)] + /\ phase \in [Subs -> Phases] + /\ backlog \in [Subs -> BOOLEAN] + /\ permit \in [Subs -> BOOLEAN] + /\ drainPending \in [Subs -> BOOLEAN] + /\ inflight \in [Subs -> SUBSET Msgs] + /\ delivered \in [Subs -> SUBSET Msgs] + /\ sent \in [Subs -> Seq(Msgs)] + +Init == + /\ pubSent = [p \in Pubs |-> 0] + /\ targets = [p \in Pubs |-> {}] + /\ chan = [s \in Subs |-> <<>>] + /\ storage = [s \in Subs |-> <<>>] + /\ held = [s \in Subs |-> <<>>] + /\ phase = [s \in Subs |-> "idle"] + /\ backlog = [s \in Subs |-> FALSE] + /\ permit = [s \in Subs |-> FALSE] + /\ drainPending = [s \in Subs |-> FALSE] + /\ inflight = [s \in Subs |-> {}] + /\ delivered = [s \in Subs |-> {}] + /\ sent = [s \in Subs |-> <<>>] + +Cur(p) == [p |-> p, n |-> pubSent[p]] + +(* Publisher handler: a PUBLISH arrives; PUBACK is withheld until targets = {} *) +PubStart(p) == + /\ targets[p] = {} + /\ pubSent[p] < MaxMsgs + /\ pubSent' = [pubSent EXCEPT ![p] = @ + 1] + /\ targets' = [targets EXCEPT ![p] = Subs] + /\ UNCHANGED subVars + +(* Router: queue to storage, then set the flag and post a permit, whatever the flag was *) +QueueBehind(p, s) == + /\ storage' = [storage EXCEPT ![s] = Append(@, Cur(p))] + /\ backlog' = [backlog EXCEPT ![s] = TRUE] + /\ permit' = [permit EXCEPT ![s] = @ \/ NotifyTrigger] + /\ targets' = [targets EXCEPT ![p] = @ \ {s}] + /\ UNCHANGED <> + +(* Router: read the flag set -- go behind the backlog without waiting *) +RouteBehind(p, s) == + /\ s \in targets[p] + /\ backlog[s] + /\ QueueBehind(p, s) + +(* Router: flag clear, channel has room *) +RouteEnqueue(p, s) == + /\ s \in targets[p] + /\ ~backlog[s] + /\ Len(chan[s]) < Cap + /\ chan' = [chan EXCEPT ![s] = Append(@, Cur(p))] + /\ targets' = [targets EXCEPT ![p] = @ \ {s}] + /\ UNCHANGED <> + +(* Router: flag clear, deadline elapsed -- possibly spuriously *) +RouteTimeout(p, s) == + /\ s \in targets[p] + /\ ~backlog[s] + /\ QueueBehind(p, s) + +(* Subscriber handler: channel arm, polled only while the window has room *) +Consume(s) == + /\ chan[s] /= <<>> + /\ phase[s] = "idle" + /\ Free(s) > 0 + /\ LET m == Head(chan[s]) + rest == Tail(chan[s]) + IN /\ chan' = [chan EXCEPT ![s] = rest] + /\ inflight' = [inflight EXCEPT ![s] = @ \cup {m}] + /\ sent' = [sent EXCEPT ![s] = Append(@, m)] + /\ drainPending' = [drainPending EXCEPT ![s] = @ \/ (rest = <<>> /\ backlog[s])] + /\ UNCHANGED <> + +(* Subscriber handler: the Notify select branch fires *) +Notified(s) == + /\ permit[s] + /\ phase[s] = "idle" + /\ permit' = [permit EXCEPT ![s] = FALSE] + /\ drainPending' = [drainPending EXCEPT ![s] = TRUE] + /\ UNCHANGED <> + +(* drain_backlog: flag clear -- nothing to do *) +DrainIdle(s) == + /\ drainPending[s] + /\ phase[s] = "idle" + /\ ~backlog[s] + /\ drainPending' = [drainPending EXCEPT ![s] = FALSE] + /\ UNCHANGED <> + +(* drain_backlog: the channel is older than storage, or the window is full -- come back later *) +DrainDefer(s) == + /\ drainPending[s] + /\ phase[s] = "idle" + /\ backlog[s] + /\ chan[s] /= <<>> \/ Free(s) = 0 + /\ drainPending' = [drainPending EXCEPT ![s] = FALSE] + /\ UNCHANGED <> + +(* drain_backlog: clear the flag before touching storage *) +DrainClear(s) == + /\ drainPending[s] + /\ phase[s] = "idle" + /\ backlog[s] + /\ chan[s] = <<>> + /\ Free(s) > 0 + /\ backlog' = [backlog EXCEPT ![s] = IF ClearFirst THEN FALSE ELSE @] + /\ phase' = [phase EXCEPT ![s] = "take"] + /\ UNCHANGED <> + +(* drain_backlog: atomic take of at most the free window from the front of storage; re-arm on leftover *) +DrainTake(s) == + /\ phase[s] = "take" + /\ LET k == Min(Len(storage[s]), Free(s)) + rest == SubSeq(storage[s], k + 1, Len(storage[s])) + IN /\ held' = [held EXCEPT ![s] = SubSeq(storage[s], 1, k)] + /\ storage' = [storage EXCEPT ![s] = rest] + /\ backlog' = [backlog EXCEPT ![s] = @ \/ (ClearFirst /\ rest /= <<>>)] + /\ phase' = [phase EXCEPT ![s] = "send"] + /\ UNCHANGED <> + +(* drain_backlog: send the next held message *) +DrainSend(s) == + /\ phase[s] = "send" + /\ held[s] /= <<>> + /\ LET m == Head(held[s]) + IN /\ held' = [held EXCEPT ![s] = Tail(@)] + /\ inflight' = [inflight EXCEPT ![s] = @ \cup {m}] + /\ sent' = [sent EXCEPT ![s] = Append(@, m)] + /\ UNCHANGED <> + +(* drain_backlog: batch done -- loop while the flag is armed *) +DrainFinish(s) == + /\ phase[s] = "send" + /\ held[s] = <<>> + /\ phase' = [phase EXCEPT ![s] = "idle"] + /\ backlog' = [backlog EXCEPT ![s] = IF ClearFirst THEN @ ELSE FALSE] + /\ drainPending' = [drainPending EXCEPT ![s] = backlog'[s]] + /\ UNCHANGED <> + +(* Subscriber PUBACK arrives *) +Ack(s, m) == + /\ m \in inflight[s] + /\ inflight' = [inflight EXCEPT ![s] = @ \ {m}] + /\ delivered' = [delivered EXCEPT ![s] = @ \cup {m}] + /\ drainPending' = [drainPending EXCEPT ![s] = @ \/ backlog[s]] + /\ UNCHANGED <> + +Next == + \/ \E p \in Pubs : PubStart(p) + \/ \E p \in Pubs, s \in Subs : RouteBehind(p, s) \/ RouteEnqueue(p, s) \/ RouteTimeout(p, s) + \/ \E s \in Subs : + \/ Consume(s) \/ Notified(s) + \/ DrainIdle(s) \/ DrainDefer(s) \/ DrainClear(s) \/ DrainTake(s) \/ DrainSend(s) \/ DrainFinish(s) + \/ \E s \in Subs, m \in Msgs : Ack(s, m) + +Fairness == + /\ \A p \in Pubs : WF_vars(PubStart(p)) + /\ \A p \in Pubs, s \in Subs : WF_vars(RouteBehind(p, s)) /\ WF_vars(RouteEnqueue(p, s)) + /\ \A s \in Subs : + /\ WF_vars(Consume(s)) /\ WF_vars(Notified(s)) + /\ WF_vars(DrainIdle(s)) /\ WF_vars(DrainDefer(s)) /\ WF_vars(DrainClear(s)) + /\ WF_vars(DrainTake(s)) /\ WF_vars(DrainSend(s)) /\ WF_vars(DrainFinish(s)) + /\ \A s \in Subs, m \in Msgs : WF_vars(Ack(s, m)) + +Spec == Init /\ [][Next]_vars /\ Fairness + +(* A message is routed to s once the router has finished with it for s *) +Routed(s) == + {m \in Msgs : + /\ m.n <= pubSent[m.p] + /\ ~(m.n = pubSent[m.p] /\ s \in targets[m.p])} + +Custody(s) == Range(chan[s]) \cup Range(storage[s]) \cup Range(held[s]) \cup inflight[s] \cup delivered[s] + +InvChanBound == \A s \in Subs : Len(chan[s]) <= Cap +InvWindowBound == \A s \in Subs : Cardinality(inflight[s]) <= RM +InvNoLoss == \A s \in Subs : Routed(s) = Custody(s) +InvNoDuplicate == \A s \in Subs : + /\ Len(chan[s]) = Cardinality(Range(chan[s])) + /\ Len(storage[s]) = Cardinality(Range(storage[s])) + /\ Len(held[s]) = Cardinality(Range(held[s])) + /\ Len(sent[s]) = Cardinality(Range(sent[s])) + /\ Range(chan[s]) \cap Range(storage[s]) = {} + /\ Range(chan[s]) \cap Range(held[s]) = {} + /\ Range(chan[s]) \cap inflight[s] = {} + /\ Range(storage[s]) \cap Range(held[s]) = {} + /\ Range(storage[s]) \cap inflight[s] = {} + /\ Range(held[s]) \cap inflight[s] = {} + /\ delivered[s] \cap (Range(chan[s]) \cup Range(storage[s]) \cup Range(held[s]) \cup inflight[s]) = {} +InvBacklogArmed == \A s \in Subs : + (storage[s] /= <<>> /\ phase[s] = "idle") => (backlog[s] \/ permit[s] \/ drainPending[s]) +InvFifo == \A s \in Subs : \A i, j \in DOMAIN sent[s] : + (i < j /\ sent[s][i].p = sent[s][j].p) => sent[s][i].n < sent[s][j].n + +AllDelivered == \A s \in Subs : delivered[s] = Msgs +Live == <>AllDelivered + +============================================================================= diff --git a/specs/tla/qos1-backlog/Qos1Backlog_nonotify.cfg b/specs/tla/qos1-backlog/Qos1Backlog_nonotify.cfg new file mode 100644 index 00000000..3fb948a4 --- /dev/null +++ b/specs/tla/qos1-backlog/Qos1Backlog_nonotify.cfg @@ -0,0 +1,27 @@ +SPECIFICATION Spec + +CONSTANTS + NPubs = 1 + NSubs = 2 + MaxMsgs = 2 + Cap = 1 + RM = 1 + Batch = 1 + MaxTakeovers = 0 + MaxPurges = 0 + NotifyTrigger = FALSE + +INVARIANTS + TypeOK + InvChanBound + InvWindowBound + InvBatchBound + InvNoLoss + InvNoDuplicate + InvBacklogArmed + InvFifo + +PROPERTIES + Live + +CHECK_DEADLOCK FALSE diff --git a/specs/tla/qos1-backlog/README.md b/specs/tla/qos1-backlog/README.md new file mode 100644 index 00000000..e17d8107 --- /dev/null +++ b/specs/tla/qos1-backlog/README.md @@ -0,0 +1,103 @@ +# QoS 1 backlog drain — TLA+ verification + +Model of the broker's QoS 1 delivery path under overload: bounded per-client +delivery channel, outbound window, storage queue, Notify permit, and the +handler-side `drain_backlog`, including expiry purges and reconnect takeover. + +Observed failure that motivated it: 8 QoS 1 publishers saturating 8 QoS 1 +subscribers; after ~8 s delivery stalls to zero while every client stays +connected, broker RSS 9 MB → 15 GB. + +## Modules + +| Module / cfg | What it models | Expected | +|---|---|---| +| `Qos1Backlog.tla` + `Qos1Backlog.cfg` | The design: a client is behind iff its storage queue is non-empty or a reconnect hand-off is in progress; the router's check, append and the handler's window-limited FIFO take are atomic w.r.t. each other (storage lock); channel consumed only while the window has room; drain only when the channel is empty, at most `Batch` per invocation, acks read only between batches; every trigger posts the Notify permit and the drain runs only from the Notify arm; expiry purge; two-phase takeover (routers route behind during the hand-off, the displaced handler re-queues unacked in-flight + unsent batch + channel at the front, the new handler starts with a permit); clean-start discard | holds | +| `Qos1Backlog_nonotify.cfg` | Same, but the router posts no permit after queueing | `InvBacklogArmed` violated | +| `Qos1BacklogLockFree.tla` + `.cfg` | Rejected variant: a separately stored backlog flag, cleared before the take and re-armed on leftover | `InvFifo` violated | + +### Actions ↔ code + +| Action | Code path | +|---|---| +| `PubStart` | publisher's PUBLISH arrives; PUBACK withheld until `targets = {}` | +| `RouteBehind` | router: `queue.count() > 0 \|\| queue.handoff` → `queue.push` + `queue.notify()` | +| `RouteEnqueue` | router: `try_send` succeeds | +| `RouteTimeout` | router: timed `reserve()` elapsed or closed (possibly spuriously) → push + notify | +| `Consume` | handler `qos1_rx` arm, polled only while `outbound_inflight < W`, ≤ `Batch` per poll, pending during a hand-off | +| `Notified` | `notified()` arm — the only entry into `drain_backlog` | +| `DrainIdle` / `DrainDefer` | `drain_backlog` steps 1–3: storage empty / channel not yet empty or window full | +| `DrainBegin` | `queue.take(min(free, Batch))`, FIFO, one critical section | +| `DrainSend` / `DrainFinish` | send the batch; if storage is left, `queue.notify()` and return to select | +| `Ack` | PUBACK read between batches → `queue.notify()` | +| `Purge` | `cleanup_expired` / expiry-on-take (count is the queue length, nothing else changes) | +| `TakeoverBegin` | `register_client` sets `handoff` before inserting the new entry; new handler's lane and drain pending | +| `TakeoverEnd` | displaced handler's `requeue_front(unacked ++ held ++ channel)`, then `released` and a permit | +| `CleanTakeover` | clean-start discard at the hand-off point | + +Invariants: `TypeOK`, `InvChanBound`, `InvWindowBound`, `InvBatchBound`, +`InvNoLoss` (routed = channel ∪ storage ∪ held ∪ hand-off ∪ inflight ∪ +delivered ∪ purged ∪ discarded), `InvNoDuplicate`, `InvBacklogArmed` +(non-empty storage always has a permit, a pending drain, a channel message to +consume, an ack to come, or a hand-off about to end), `InvFifo` (per +publisher, the first-send order to each subscriber is publish order; a re-send +after a takeover does not change it). Liveness: `Live == <>AllDelivered` +(delivered ∪ purged ∪ discarded) under weak fairness. + +## Results (2026-09-07, v5) + +| Module / cfg | constants | result | +|---|---|---| +| `Qos1Backlog.cfg` (shipped) | 1p×2s×2m Cap1 RM1 Batch1, 1 takeover, 1 purge | ok, 98 422 states, safety + Live | +| `Qos1Backlog` | 2p×1s×2m, same bounds | ok, 17 839 states, safety + Live | +| `Qos1Backlog` | 2p×2s×1m, same bounds | ok, 272 113 states, safety + Live | +| `Qos1Backlog` | 2p×1s×3m Cap1 RM2 Batch1 | ok, 603 547 states, safety + Live | +| `Qos1Backlog` | 1p×2s×3m Cap1 RM3 Batch1, no takeover, no purge | ok, 23 590 states, safety + Live | +| `Qos1Backlog` | 1p×2s×3m Cap2 RM2 Batch2, no takeover, no purge | ok, 23 978 states, safety + Live | +| `Qos1Backlog` | 2p×2s×2m Cap1 RM1 Batch1, no takeover, no purge, safety only | ok, 913 110 states, 336 s | +| `Qos1Backlog_nonotify.cfg` | 1p×2s×2m | **`InvBacklogArmed` violated** in 2 steps | +| `Qos1BacklogLockFree.cfg` | 2p×1s×2m | **`InvFifo` violated** in 17 steps | + +Larger shapes with takeover and purge enabled (2p×2s×2m) exceed a 10-minute +budget and are inconclusive; the bounds `MaxTakeovers` / `MaxPurges` exist to +keep the checked configurations exhaustive. + +### Counterexample 1 — no router permit + +`PubStart → RouteTimeout`: the deadline fires while the subscriber's channel +and window are both empty. The message is queued, but the two remaining +triggers (channel running empty, an ack) can never fire again. + +### Counterexample 2 — stored flag cleared before the take + +Publisher 2's first message is queued behind the backlog. The drain clears the +flag, publisher 2 reads the cleared flag and enqueues its **second** message +into the channel, the window-limited take takes only an older message from +another publisher and re-arms the flag with publisher 2's first message still +in storage; the channel message is then written first: +`sent = << m(1,1), m(2,2), m(1,2) >>`. Hence: no stored flag; the decision +reads the storage queue itself under its lock, and the take is atomic with +that read. + +### Design decisions the model forced + +- The backlog condition is the storage queue's own emptiness (plus the + hand-off flag), read under the storage lock; never a separately stored + boolean. +- Take at most `min(free window, Batch)` from the front of storage, + atomically; never re-queue at the tail; re-queue at the front only on a + hand-off or a transport error. +- Post the Notify permit after every queue to storage and after every batch + that leaves storage non-empty; run the drain only from the Notify arm. +- Consume the channel only while the window has room, so window-full + propagates to channel-full and the publisher's timed wait is real + backpressure. The window is `min(client Receive Maximum, broker cap)`. +- During a reconnect hand-off routers route behind, so nothing newer reaches + the new channel before the displaced handler's custody reaches the front of + storage. +- Acks are read only between batches, so batches must be bounded. + +Verified for these constants only; not a proof for all sizes. Not modelled: +the QoS 0 lane, transport errors, the router lock order, packet ids, +scheduling and time budgets, the per-client queued cap, the on-disk format, +bridge ingress, broker restart.