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.
Contents
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.htmlBroker 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
- QoS retransmission / in-flight state is process-local — it does not survive broker restarts (file-backed
MqttMessageStorepersists offline queues only).