MQTT

hopf-mqtt is an MQTT broker and async client (Gumdrop mqtt port): MQTT 3.1.1 full semantics plus a useful 5.0 core, with an optional MQTT-over-WebSocket bridge sharing the same broker state as the TCP listener.

Versions and scope

Area Support
MQTT 3.1.1 Full semantics — CONNECT/CONNACK, PUBLISH QoS 0/1/2, SUBSCRIBE/SUBACK, UNSUBSCRIBE/UNSUBACK, PINGREQ/PINGRESP, DISCONNECT, retained messages, Will messages, topic wildcards (+ / #), session takeover
MQTT 5.0 (useful core) Wire properties, unified reason codes, Receive Maximum outbound flow control, subscription options (No Local, Retain As Published, Retain Handling), Session Expiry Interval (orphan-and-resume), shared subscriptions ($share/...), topic aliases, Will Delay / Message Expiry enforcement, enhanced AUTH, offline QoS ≥ 1 queues, in-process QoS retransmission
Still limited QoS retry / inflight state does not survive broker process restarts (file-backed store covers offline queues only)

Architecture

TCP (1883/8883)  ──┐
                    ├── MqttFrameParser (push, streaming PUBLISH) ── BrokerState (Arc)
WebSocket (mqtt) ───┘                                                   │
                                                                          └── ConnHandle::send (fan-out)

BrokerState is shared behind Arc by every connection regardless of transport or reactor. Delivery to a subscriber never touches that subscriber's Endpoint directly from another thread — it goes through hopf_core::ConnHandle::send, which hops to the owning reactor. PUBLISH payloads are streamed through the codec (MqttFrameHandler::publish_data chunks), never buffered in full.

Server quick start

use std::sync::Arc;
use hopf_core::{Runtime, RuntimeConfig};
use hopf_mqtt::server::broker::BrokerState;
use hopf_mqtt::server::{MqttConfig, MqttService};

let rt = Runtime::start(RuntimeConfig::default())?;
let broker = Arc::new(BrokerState::new());
let config = MqttConfig::new("0.0.0.0:1883".parse()?, broker);
let svc = MqttService::new(config);
let addr = svc.start(&rt)?;

Example binary: cargo run -p mqtt -- 127.0.0.1:1883.

MqttConfig

Field Default Notes
broker — Shared Arc<BrokerState> — reuse the same one across a TCP listener and a WebSocket bridge to share topics/sessions
credentials None Arc<dyn CredentialStore> for CONNECT username/password; None accepts any CONNECT
max_packet_size 1 MiB Cap on a packet's Remaining Length (with_max_packet_size)
max_publish_payload 1 MiB PUBLISH payload cap before fan-out/spool (with_max_publish_payload); intentional internet default
connect_timeout 10s How long to wait for CONNECT before closing an idle new TCP connection

Handler SPI

CONNECT, PUBLISH, and SUBSCRIBE are staged (Gumdrop parity). Default handlers accept all traffic, or require CONNECT username/password to match a CredentialStore when one is configured. Enhanced AUTH (MQTT 5.0) is driven when CONNECT carries Authentication Method and credentials are configured.

pub trait ConnectHandler: Send {
    fn authorize(&mut self, packet: &ConnectPacket, meta: &MqttConnectionMetadata) -> ConnectDecision;
}
pub trait PublishHandler: Send { /* authorize(client_id, topic, qos, retain, meta) */ }
pub trait SubscribeHandler: Send { /* authorize(client_id, filter, meta) */ }

pub trait MqttHandlerFactory: Send + Sync {
    fn create(&self) -> Box<dyn ConnectHandler>;
    fn create_publish(&self) -> Box<dyn PublishHandler>;
    fn create_subscribe(&self) -> Box<dyn SubscribeHandler>;
}

MqttConnectionMetadata carries peer/local/tls, optional client_id, and (when OTel traces are enabled) a W3C traceparent. Swap in custom policy with MqttService::with_handler_factory. Offline / inflight message storage is pluggable via BrokerState::with_store (InMemoryMessageStore default, FileBackedMessageStore optional).

Client

Async, non-blocking client on the hopf-core Runtime / ProtocolHandler SPI, with DNS resolution via hopf-dns. Unlike the staged POP3/IMAP client drivers, MQTT has only one phase once CONNACK arrives, so MqttClientDriver is a flat set of callbacks over MqttClientControl (publish/subscribe/unsubscribe/disconnect) — QoS 1/2 acknowledgment of incoming messages and the PUBREC→PUBREL leg of the client's own QoS 2 publishes are handled automatically by the endpoint.

use std::sync::Arc;
use hopf_core::{Runtime, RuntimeConfig};
use hopf_mqtt::client::MqttClient;

let rt = Arc::new(Runtime::start(RuntimeConfig::default())?);
MqttClient::new("broker.example.com", 1883, "my-client-id")
    .connect(&rt, Arc::new(driver_factory))?;

Client builder options

Builder method Default Description
MqttClient::new(host, port, client_id) MQTT 5.0 DNS-resolved
MqttClient::from_addr(addr, client_id) — Pre-resolved, skips DNS
.v3_1_1() off (5.0) CONNECT with MQTT 3.1.1
.clean_start(bool) true Clean Session / Clean Start
.keep_alive(Duration) 60s Duration::ZERO disables it
.session_expiry(secs) 0 MQTT 5.0 only
.receive_maximum(n) unset (unlimited) MQTT 5.0 only
.credentials(user, pass) none CONNECT username/password
.will(Will) none Published by the broker on an unclean disconnect
.implicit_tls(connector, name) off MQTTS, typically port 8883
.timeouts(MqttClientTimeouts) dns=5s, connect=30s, connack=30s, pingresp=10s Per-phase deadlines

Example binary: cargo run -p mqtt-pub -- 127.0.0.1 1883 demo/topic "hello" (also subscribes and prints anything published back to it).

WebSocket bridge

Feature websocket (umbrella crate feature mqtt-ws). MqttWsFactory / MqttWsHandler drive the same MqttFrameParser and BrokerState as the TCP listener — a WS-connected client and a TCP-connected client publish and subscribe to each other transparently, sharing one Arc<BrokerState> across both listeners (e.g. via hopf_core::Composition).

use hopf_mqtt::server::DefaultMqttHandlerFactory;
use hopf_mqtt::server::ws::MqttWsFactory;
use hopf_websocket::{WebSocketConfig, WebSocketFactory};

let ws_config = Arc::new(MqttConfig::new(addr, Arc::clone(&broker)));
let events = MqttWsFactory::new(ws_config, Arc::new(DefaultMqttHandlerFactory::new(None)));
let ws = WebSocketFactory::new(events, WebSocketConfig::default());
// wrap `ws` as Arc<dyn ServerHandlerFactory> under CleartextHttpEndpoint / AlpnHttpEndpoint,
// same as any other hopf-websocket factory — see websocket.html

Broker fan-out (BrokerState::publish / deliver_retained) delivers asynchronously, from whichever connection published the message, via ConnHandle::send. A plain ConnHandle writes straight to the raw transport, which would land a bare MQTT packet on the wire unframed and corrupt the WebSocket stream — so the WS bridge registers a hopf_websocket::framed_ws_conn_handle-wrapped handle with the broker instead, which pipes every delivery through WS binary framing first. This is exactly why the WS bridge needed hopf-websocket to expose a per-connection ConnHandle at all — see the WebSocket event handler SPI.

Timers (CONNECT timeout, keepalive, Session Expiry, Will Delay, QoS retransmission) use ConnHandle::schedule_timer on the framed handle — the same semantics as the TCP path.

Examples

cargo run -p mqtt -- 127.0.0.1:1883
cargo run -p mqtt-pub -- 127.0.0.1 1883 demo/topic "hello from hopf"
Package Demonstrates Knobs
examples/mqtt Broker on plain TCP arg1 = listen addr
examples/mqtt-pub Async client: publish once, subscribe and print anything received args = host, port, topic, payload

No environment variables.

Metrics

Process-local MqttServerMetrics (in hopf-mqtt): connections, auth_ok / auth_fail, publishes, and subscribes via MqttService::metrics().

OTLP/JSONL export (when MqttService::with_telemetry(&pipeline) is used — also MqttWsFactory::with_telemetry for the WebSocket bridge): hopf_otel::MqttServerMetrics emits mqtt.server.connections, mqtt.server.active_connections, mqtt.server.auth, mqtt.server.subscribes, mqtt.server.publishes (qos/outcome), mqtt.server.publish.duration (ms), and mqtt.server.publish.size.

With traces enabled on the pipeline, MqttConnectionMetadata::traceparent carries the active W3C traceparent for the connection or current client PUBLISH. Handler SPI implementers can pass it to outbound HTTP with hopf_otel::with_traceparent.

Limitations