-
-
Notifications
You must be signed in to change notification settings - Fork 6
Broker Guide
Complete guide to running and configuring the MQTT broker.
use mqtt5::broker::{BrokerConfig, MqttBroker};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Simple: bind to single address
let mut broker = MqttBroker::bind("0.0.0.0:1883").await?;
// Or with full configuration
let config = BrokerConfig::default()
.with_bind_address("0.0.0.0:1883".parse()?)
.with_max_clients(10_000);
let mut broker = MqttBroker::with_config(config).await?;
broker.run().await?;
Ok(())
}MqttBroker exposes these async/sync methods:
| Method | Purpose |
|---|---|
MqttBroker::bind(addr).await |
Bind a single address with defaults |
MqttBroker::with_config(config).await |
Build from a BrokerConfig
|
MqttBroker::with_config_file(config, path).await |
Build with SIGHUP hot-reload wired to a config file |
broker.run().await |
Run the accept loop (takes &mut self) |
broker.shutdown().await |
Gracefully stop |
broker.stats() |
Arc<BrokerStats> counters |
broker.resource_monitor() |
Arc<ResourceMonitor> |
broker.auth_provider() |
The built Arc<dyn AuthProvider> (wrap it with CompositeAuthProvider) |
broker.with_auth_provider(provider) |
Replace the auth provider (consumes/returns self) |
broker.router() |
Arc<MessageRouter> |
broker.local_addr() / broker.ws_local_addr()
|
Bound socket addresses |
# Basic broker
mqttv5 broker --allow-anonymous
# Custom port (repeatable; -H is the short form)
mqttv5 broker --host 0.0.0.0:1884 --allow-anonymous
# With configuration file (enables SIGHUP hot-reload)
mqttv5 broker --config broker.json
# Generate sample config (subcommand; defaults to stdout)
mqttv5 broker generate-config > broker.json
mqttv5 broker generate-config --format toml --output broker.tomlRun TCP, TLS, WebSocket, and QUIC simultaneously:
┌─────────────────────────────────────────────────────────────────┐
│ MQTT Broker │
├─────────────────────────────────────────────────────────────────┤
│ TCP :1883 │ TLS :8883 │ WS :8080 │ QUIC :14567 │
│ (plaintext) │ (encrypted) │ (browsers) │ (modern) │
└─────────────────────────────────────────────────────────────────┘
use mqtt5::broker::config::{TlsConfig, WebSocketConfig, QuicConfig};
use mqtt5::broker::BrokerConfig;
let config = BrokerConfig::default()
// TCP on port 1883
.with_bind_address("0.0.0.0:1883".parse()?)
// TLS on port 8883
.with_tls(
TlsConfig::new("certs/server.crt".into(), "certs/server.key".into())
.with_bind_address("0.0.0.0:8883".parse()?)
)
// WebSocket on port 8080
.with_websocket(
WebSocketConfig::default()
.with_bind_address("0.0.0.0:8080".parse()?)
.with_path("/mqtt")
)
// QUIC on port 14567 (requires TLS cert/key)
.with_quic(
QuicConfig::new("certs/server.crt".into(), "certs/server.key".into())
.with_bind_address("0.0.0.0:14567".parse()?)
);Related builders: with_websocket_tls(WebSocketConfig) for a TLS WebSocket listener,
and with_cluster_listener(ClusterListenerConfig) for the inter-broker cluster transport.
mqttv5 broker \
--host 0.0.0.0:1883 \
--tls-host 0.0.0.0:8883 \
--tls-cert server.pem \
--tls-key server-key.pem \
--ws-host 0.0.0.0:8080 \
--ws-path /mqtt \
--quic-host 0.0.0.0:14567QUIC requires --tls-cert and --tls-key. WebSocket-over-TLS uses --ws-tls-host.
Field names match BrokerConfig exactly:
{
"bind_addresses": ["0.0.0.0:1883", "[::]:1883"],
"tls_config": {
"cert_file": "certs/server.pem",
"key_file": "certs/server-key.pem",
"ca_file": null,
"require_client_cert": false,
"bind_addresses": ["0.0.0.0:8883"]
},
"websocket_config": {
"bind_addresses": ["0.0.0.0:8080"],
"path": "/mqtt",
"subprotocol": "mqtt",
"use_tls": false
},
"quic_config": {
"cert_file": "certs/server.pem",
"key_file": "certs/server-key.pem",
"ca_file": null,
"require_client_cert": false,
"bind_addresses": ["0.0.0.0:14567"],
"enable_early_data": false
}
}Encrypts traffic but doesn't verify clients:
mqttv5 broker \
--tls-host 0.0.0.0:8883 \
--tls-cert server.pem \
--tls-key server-key.pemRequires client certificates:
mqttv5 broker \
--tls-host 0.0.0.0:8883 \
--tls-cert server.pem \
--tls-key server-key.pem \
--tls-ca-cert ca.pem \
--tls-require-client-certTlsConfig and QuicConfig share the same fields: cert_file, key_file,
ca_file (optional), require_client_cert, and bind_addresses.
- Format: PEM
- Certificate chain: Server cert first, then intermediates
- Private key: Must be unencrypted
- CA file: Can contain multiple CA certificates
The default backend is File with persistence enabled (base dir ./mqtt_storage).
Survives restarts; uses atomic writes (write-temp-then-rename with fsync):
{
"storage_config": {
"backend": "File",
"base_dir": "/var/lib/mqtt",
"cleanup_interval": "1h",
"enable_persistence": true
}
}mqttv5 broker \
--storage-backend file \
--storage-dir /var/lib/mqttFaster, but all retained/queued state is lost on restart:
{
"storage_config": {
"backend": "Memory",
"base_dir": "./mqtt_storage",
"cleanup_interval": "1h",
"enable_persistence": true
}
}mqttv5 broker --storage-backend memory
# or disable persistence entirely:
mqttv5 broker --no-persistenceHow long to keep sessions after disconnect (default 1h):
{
"session_expiry_interval": "1h"
}mqttv5 broker --session-expiry 3600When a client connects with clean_start=false, the broker:
- Restores previous subscriptions (re-authorized against current ACL)
- Delivers queued QoS 1/2 messages
- Resumes in-flight message flows
{
"max_clients": 10000,
"max_packet_size": 268435456,
"topic_alias_maximum": 65535,
"maximum_qos": 2,
"max_subscriptions_per_client": 0,
"max_retained_messages": 0,
"max_retained_message_size": 0
}| Setting | Description | Default |
|---|---|---|
max_clients |
Maximum concurrent connections | 10,000 |
max_packet_size |
Maximum MQTT packet size (bytes) | 256 MB |
topic_alias_maximum |
Max topic aliases per client | 65,535 |
maximum_qos |
Highest QoS level (0, 1, 2) | 2 |
max_subscriptions_per_client |
Subscriptions per client (0 = unlimited) | 0 |
max_retained_messages |
Retained messages (0 = unlimited) | 0 |
max_retained_message_size |
Max retained payload bytes (0 = no cap) | 0 |
client_channel_capacity |
Per-client outbound queue depth | 10,000 |
These are top-level BrokerConfig fields (all 0 = unlimited) and are
hot-reloadable via SIGHUP:
{
"max_message_rate_per_client": 0,
"max_bandwidth_per_client": 0,
"max_outbound_rate_per_client": 0
}| Setting | Type | Description | Default |
|---|---|---|---|
max_message_rate_per_client |
u32 |
Inbound messages/sec (0 = unlimited) | 0 |
max_bandwidth_per_client |
u64 |
Inbound bytes/sec (0 = unlimited) | 0 |
max_outbound_rate_per_client |
u32 |
Outbound messages/sec (0 = unlimited) | 0 |
let config = BrokerConfig::default()
.with_max_outbound_rate_per_client(5_000);Each limit is enforced independently. With the defaults (
0) the broker takes a fast path with no per-message accounting.
Enable or disable MQTT v5 features (all default true):
{
"retain_available": true,
"wildcard_subscription_available": true,
"subscription_identifier_available": true,
"shared_subscription_available": true
}mqttv5 broker --no-retain --no-wildcardsOther advertised-capability fields: server_keep_alive (Option, negotiated
Server Keep Alive), server_receive_maximum (Option), and response_information
(string sent when a client requests it via --response-information).
{
"change_only_delivery_config": {
"enabled": true,
"topic_patterns": ["sensors/#", "status/+"]
},
"echo_suppression_config": {
"enabled": true,
"property_key": "x-origin-client-id"
},
"server_delivery_strategy": "per_topic"
}- Change-only delivery suppresses re-delivery of a topic's payload when it is unchanged from the last value delivered to that subscriber.
-
Echo suppression skips delivering a PUBLISH to a subscriber whose
client_idmatches the configured user-property value (default keyx-origin-client-id); the key is hot-reloadable via SIGHUP. See Authentication & ACL for details. -
server_delivery_strategyselects the QUIC server-side stream mapping:control_only,per_topic(default), orper_publish(CLI:--quic-delivery-strategy control-only|per-topic|per-publish).
A background provider publishes retained statistics to $SYS/# every 10 seconds:
| Topic | Description |
|---|---|
$SYS/broker/version |
Broker version string |
$SYS/broker/implementation |
mqtt-v5 |
$SYS/broker/protocol_version |
5.0 |
$SYS/broker/uptime |
Broker uptime in seconds |
$SYS/broker/clients/connected |
Currently connected clients |
$SYS/broker/clients/total |
Total connections since start |
$SYS/broker/clients/maximum |
Peak simultaneous connections |
$SYS/broker/messages/sent |
Total messages sent |
$SYS/broker/messages/received |
Total messages received |
$SYS/broker/messages/publish/sent |
PUBLISH messages sent |
$SYS/broker/messages/publish/received |
PUBLISH messages received |
$SYS/broker/bytes/sent |
Total bytes sent |
$SYS/broker/bytes/received |
Total bytes received |
$SYS/broker/retained/count |
Retained messages in storage |
$SYS/broker/subscriptions/count |
Active subscribed topics |
Subscribe to monitor:
mqttv5 sub -t '$SYS/#' -vbroker.stats() returns an Arc<BrokerStats> whose counters are atomics:
use std::sync::atomic::Ordering;
let stats = broker.stats();
println!("Connected: {}", stats.clients_connected.load(Ordering::Relaxed));
println!("Messages received: {}", stats.messages_received.load(Ordering::Relaxed));
println!("Bytes sent: {}", stats.bytes_sent.load(Ordering::Relaxed));
println!("Uptime: {}s", stats.uptime_seconds());Requires the opentelemetry feature. Enable distributed tracing:
mqttv5 broker \
--otel-endpoint http://localhost:4317 \
--otel-service-name mqtt-broker \
--otel-sampling 1.0Programmatically via BrokerConfig::with_opentelemetry(TelemetryConfig):
use mqtt5::telemetry::TelemetryConfig;
let telemetry = TelemetryConfig::new("mqtt-broker")
.with_endpoint("http://localhost:4317")
.with_sampling_ratio(1.0);
let config = BrokerConfig::default().with_opentelemetry(telemetry);The broker can redirect clients to backend brokers using the MQTT v5
UseAnotherServer (0x9C) reason code, choosing a backend by consistent hashing on
the client ID:
{
"bind_addresses": ["0.0.0.0:1883"],
"load_balancer": {
"backends": ["mqtt://127.0.0.1:1884", "mqtt://127.0.0.1:1885"]
}
}use mqtt5::broker::config::LoadBalancerConfig;
let config = BrokerConfig::default()
.with_load_balancer(LoadBalancerConfig::new(vec![
"mqtt://backend1:1883".into(),
"mqtt://backend2:1883".into(),
]));Clients connecting to the front port receive a redirect and reconnect to the assigned backend (up to 3 redirect hops).
Implement BrokerEventHandler to react to broker events. Note that
on_client_publish returns a PublishAction (return PublishAction::Continue to
deliver normally); the other hooks return ():
use mqtt5::broker::BrokerConfig;
use mqtt5::broker::events::{
BrokerEventHandler, ClientConnectEvent, ClientPublishEvent, PublishAction,
};
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
struct MyHandler;
impl BrokerEventHandler for MyHandler {
fn on_client_connect<'a>(
&'a self,
event: ClientConnectEvent,
) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>> {
Box::pin(async move {
println!("Client connected: {}", event.client_id);
})
}
fn on_client_publish<'a>(
&'a self,
event: ClientPublishEvent,
) -> Pin<Box<dyn Future<Output = PublishAction> + Send + 'a>> {
Box::pin(async move {
println!("Published to {}: {} bytes", event.topic, event.payload.len());
PublishAction::Continue
})
}
}
let config = BrokerConfig::default()
.with_event_handler(Arc::new(MyHandler));Available hooks:
on_client_connecton_client_subscribeon_client_unsubscribe-
on_client_publish(returnsPublishAction) on_client_disconnecton_retained_seton_message_delivered
ClientPublishEvent fields include client_id, user_id, topic, payload,
qos, retain, packet_id, response_topic, and correlation_data.
When the broker is started from a config file (CLI --config, or programmatic
MqttBroker::with_config_file), configuration is reloaded on SIGHUP or when the
config file changes on disk:
use mqtt5::broker::{BrokerConfig, MqttBroker};
let config = BrokerConfig::default();
let mut broker = MqttBroker::with_config_file(config, "broker.json".into()).await?;
// Trigger a reload programmatically (equivalent to SIGHUP):
if let Some(tx) = broker.manual_reload_sender() {
tx.send(()).await?;
}
broker.run().await?;Internally a HotReloadManager (HotReloadManager::new(config, path),
reload_now()) classifies each change into ConfigChangeType categories:
FullReload, AuthConfig, TlsConfig, ResourceLimits, WebSocketConfig,
BridgeConfig, and StorageConfig. Password and ACL files are also reloaded live.
See CLI Reference for all command-line flags and the JSON schema.
Getting Started
Broker Guide
Client Guide
Platform Guides
CLI Reference
Development