MQTT Client API.
Provides a simple interface for connecting to MQTT brokers.
Example
# Connect
# `await_connect: true` blocks until the session is live, so the calls
# below work inline. Without it, connect/1 returns immediately and these
# would get {:error, :not_connected} — see "Connection is asynchronous".
{:ok, client} = MqttX.Client.connect(
host: "localhost",
port: 1883,
client_id: "my_app",
await_connect: true
)
# Subscribe (returns {:ok, granted_qos_list})
{:ok, _granted} = MqttX.Client.subscribe(client, "sensors/#", qos: 1)
# Publish
:ok = MqttX.Client.publish(client, "sensors/temp", "25.5")
# Disconnect
:ok = MqttX.Client.disconnect(client)Receiving Messages
To receive messages, provide a handler module:
defmodule MyHandler do
def handle_mqtt_event(:message, {topic, payload, _packet}, state) do
# topic is a list of segments, e.g. ["sensors", "room1", "temp"]
IO.puts("Received: " <> inspect({topic, payload}))
state
end
def handle_mqtt_event(:connected, _data, state) do
IO.puts("Connected!")
state
end
def handle_mqtt_event(:disconnected, reason, state) do
IO.puts("Disconnected: " <> inspect(reason))
state
end
def handle_mqtt_event(:publish_error, {_topic, _packet_id, reason_code}, state) do
IO.puts("Broker rejected publish: " <> inspect(reason_code))
state
end
# Catch-all so future event types don't raise
def handle_mqtt_event(_event, _data, state), do: state
end
{:ok, client} = MqttX.Client.connect(
host: "localhost",
client_id: "my_app",
handler: MyHandler,
handler_state: %{}
)
Summary
Functions
Connect to an MQTT broker.
Connect to an MQTT broker with supervision.
Check if the client is connected.
Disconnect from the broker.
List all registered client connections.
Publish a message to a topic.
Set up an MQTT 5.0 Request/Response exchange.
Subscribe to one or more topics.
Unsubscribe from one or more topics.
Look up a client connection by its client_id.
Functions
@spec connect(keyword()) :: GenServer.on_start() | {:error, term()}
Connect to an MQTT broker.
Options
:host- Broker hostname (required):port- Broker port (default: 1883 TCP / 8883 SSL / 8083 WS / 8084 WSS):client_id- Client identifier (required):username- Optional username:password- Optional password:clean_session- Clean session flag (default: true):keepalive- Keepalive interval in seconds (default: 60):protocol_version-3,4(3.1.1) or5(default: 5):handler- Module to receive callbacks:handler_state- Initial state for handler:name- Optional name for the client process
Transport
:transport-:tcp,:ssl,:wsor:wss(default::tcp):ssl_opts- SSL options for:ssl/:wss, merged over a secure baseline (verify: :verify_peeragainst the OS trust store, SNI, HTTPS hostname checking, TLS 1.2/1.3). Pass[verify: :verify_none]to opt out deliberately; it is logged as a warning.:ws_path- WebSocket path for:ws/:wss(default:"/mqtt"):proxy- Tunnel through an HTTPCONNECTproxy, e.g.[host: "proxy.corp", port: 3128, auth: {"user", "pass"}]. Works for all transports;:portdefaults to 3128 and:auth(Basic) is optional. TLS is negotiated with the broker through the tunnel.
Reliability and limits
:retry_interval- QoS 1/2 retry interval in ms (default: 5000):max_inflight- Max pending QoS 1/2 messages (default: 100):max_packet_size- Reject inbound packets declaring more than this (default: 1 MiB;:infinitydisables):session_store- Session store module or{module, opts}:connect_properties- MQTT 5.0 CONNECT properties, e.g.%{session_expiry_interval: 3600, receive_maximum: 20}
Last Will & Testament
:will_topic,:will_payload,:will_qos,:will_retain,:will_properties
Returns
{:ok, pid} on success, {:error, reason} on failure — including
{:error, {:missing_option, :host | :client_id}} when a required option is
absent.
Connection is asynchronous
connect/1 returns as soon as the client process starts — the session is
not yet live. The CONNECT/CONNACK handshake happens in the background so
that a client can be started before the broker is reachable (a device
booting before its network, a broker still starting up) and keep retrying
with backoff.
This means subscribe/3 and publish/4 called immediately after
connect/1 will return {:error, :not_connected}. Wait for readiness in
one of these ways:
# 1. Act on the :connected event (recommended for long-lived clients)
def handle_mqtt_event(:connected, _info, state) do
MqttX.Client.subscribe(self(), "sensors/#", qos: 1)
state
end
# 2. Block until the first attempt resolves (handy in scripts and tests)
{:ok, client} = MqttX.Client.connect(host: "broker", client_id: "c",
await_connect: true)
# session is live here, or {:error, reason} was returnedWith await_connect: true, a failed first attempt returns
{:error, reason} (for example :econnrefused, a TLS failure, or
{:connack_error, code, info}) and the client is stopped rather than left
retrying.
Connect to an MQTT broker with supervision.
Starts the connection under MqttX.Client.Supervisor, providing
automatic restart on crash. The connection is registered in
MqttX.ClientRegistry for lookup by client_id.
Accepts the same options as connect/1.
Example
{:ok, pid} = MqttX.Client.connect_supervised(
host: "localhost",
port: 1883,
client_id: "my_client"
)
Check if the client is connected.
Disconnect from the broker.
Asynchronous: the returned :ok acknowledges the request, not that the
DISCONNECT packet has been sent — the client process sends it and stops
shortly after.
Options (MQTT 5.0)
:reason_code- Disconnect reason code (default: 0x00 Normal):properties- Disconnect properties map, e.g.%{session_expiry_interval: 0}
List all registered client connections.
Returns a list of {client_id, pid, metadata} tuples for all connections
registered in MqttX.ClientRegistry.
Example
MqttX.Client.list()
#=> [{"my_client", #PID<0.123.0>, %{host: "localhost", port: 1883}}]
Publish a message to a topic.
Options
:qos- QoS level 0, 1, or 2 (default: 0):retain- Retain flag (default: false):properties- PUBLISH properties map (MQTT 5.0), e.g.%{message_expiry_interval: 60, response_topic: "replies/1", correlation_data: "req-42", user_properties: [{"k", "v"}]}
:ok means the packet was written to the socket; for QoS 1/2 the
broker's acknowledgment is tracked in the background (rejections surface
via handle_mqtt_event(:publish_error, {topic, packet_id, reason_code}, state)).
Set up an MQTT 5.0 Request/Response exchange.
Subscribes to the response_topic, then publishes a message with
response_topic and correlation_data properties set. Returns the
generated correlation_data so the caller can match incoming responses
in their handler.
This is a setup helper, not a blocking RPC call. Use handle_mqtt_event/3
in your handler to match responses by correlation_data.
Options
:response_topic- Topic to receive the response on (required):qos- QoS level for both request and subscription (default: 0)
Returns
{:ok, correlation_data}- Request sent, use this to match the response{:error, reason}- Subscribe or publish failed
Example
{:ok, correlation_data} = MqttX.Client.request(
client,
"api/users/get",
Jason.encode!(%{id: 123}),
response_topic: "api/responses/" <> client_id
)
# In your handler:
def handle_mqtt_event(:message, {_topic, payload, packet}, state) do
if packet.properties[:correlation_data] == state.pending_correlation do
# This is the response
end
state
end
Subscribe to one or more topics.
Options
:qos- QoS level 0, 1, or 2 (default: 0)
Unsubscribe from one or more topics.
Look up a client connection by its client_id.
Returns {pid, metadata} if found, or nil if not registered.
Example
MqttX.Client.whereis("my_client")
#=> {#PID<0.123.0>, %{host: "localhost", port: 1883}}