Skip to content

Broker Guide

Fabrício Bracht edited this page Jul 3, 2026 · 1 revision

Broker Guide

Complete guide to running and configuring the MQTT broker.


Starting the Broker

Programmatic

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

CLI

# 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.toml

Multi-Transport Configuration

Run TCP, TLS, WebSocket, and QUIC simultaneously:

┌─────────────────────────────────────────────────────────────────┐
│                        MQTT Broker                              │
├─────────────────────────────────────────────────────────────────┤
│  TCP :1883    │  TLS :8883    │  WS :8080     │  QUIC :14567    │
│  (plaintext)  │  (encrypted)  │  (browsers)   │  (modern)       │
└─────────────────────────────────────────────────────────────────┘

Rust Configuration

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.

CLI Configuration

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:14567

QUIC requires --tls-cert and --tls-key. WebSocket-over-TLS uses --ws-tls-host.

JSON Configuration

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
  }
}

TLS Configuration

Server Certificate Only

Encrypts traffic but doesn't verify clients:

mqttv5 broker \
  --tls-host 0.0.0.0:8883 \
  --tls-cert server.pem \
  --tls-key server-key.pem

Mutual TLS (mTLS)

Requires 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-cert

TlsConfig and QuicConfig share the same fields: cert_file, key_file, ca_file (optional), require_client_cert, and bind_addresses.

Certificate Requirements

  • Format: PEM
  • Certificate chain: Server cert first, then intermediates
  • Private key: Must be unencrypted
  • CA file: Can contain multiple CA certificates

Storage Configuration

The default backend is File with persistence enabled (base dir ./mqtt_storage).

File-Based Persistence (default)

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/mqtt

In-Memory

Faster, 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-persistence

Session Management

Session Expiry

How long to keep sessions after disconnect (default 1h):

{
  "session_expiry_interval": "1h"
}
mqttv5 broker --session-expiry 3600

Clean Start

When a client connects with clean_start=false, the broker:

  1. Restores previous subscriptions (re-authorized against current ACL)
  2. Delivers queued QoS 1/2 messages
  3. Resumes in-flight message flows

Resource Limits

Basic Limits

{
  "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

Per-Client Inbound / Outbound Rate Limits

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.


Feature Toggles

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-wildcards

Other 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).


Advanced Delivery Features

{
  "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_id matches the configured user-property value (default key x-origin-client-id); the key is hot-reloadable via SIGHUP. See Authentication & ACL for details.
  • server_delivery_strategy selects the QUIC server-side stream mapping: control_only, per_topic (default), or per_publish (CLI: --quic-delivery-strategy control-only|per-topic|per-publish).

Monitoring

$SYS Topics

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/#' -v

Programmatic Stats Access

broker.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());

OpenTelemetry

Requires the opentelemetry feature. Enable distributed tracing:

mqttv5 broker \
  --otel-endpoint http://localhost:4317 \
  --otel-service-name mqtt-broker \
  --otel-sampling 1.0

Programmatically 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);

Load Balancer / Server Redirect

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).


Event Hooks (Programmatic)

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_connect
  • on_client_subscribe
  • on_client_unsubscribe
  • on_client_publish (returns PublishAction)
  • on_client_disconnect
  • on_retained_set
  • on_message_delivered

ClientPublishEvent fields include client_id, user_id, topic, payload, qos, retain, packet_id, response_topic, and correlation_data.


Hot Reloading

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.


Complete Configuration Reference

See CLI Reference for all command-line flags and the JSON schema.

Clone this wiki locally