-
-
Notifications
You must be signed in to change notification settings - Fork 6
Mock Testing
Unit test your MQTT code without a real broker using the mock client.
The library provides MqttClientTrait for dependency injection and MockMqttClient
for testing. Both MqttClient and MockMqttClient implement MqttClientTrait:
┌─────────────────────────────────────────────────────────────────┐
│ Your Application │
├─────────────────────────────────────────────────────────────────┤
│ │
│ fn process_sensor<T: MqttClientTrait>(client: &T) { ... } │
│ │
│ │ │ │
│ ▼ ▼ │
│ ┌─────────────────┐ ┌─────────────────┐ │
│ │ MqttClient │ │ MockMqttClient │ │
│ │ (Production) │ │ (Testing) │ │
│ └─────────────────┘ └─────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────────┘
MqttClientTrait is an async trait built on impl Future return types, so its
methods use impl Into<String> / impl Into<Vec<u8>> arguments just like the
concrete client.
Write code that accepts any MqttClientTrait implementation:
use mqtt5::MqttClientTrait;
pub struct SensorService<T: MqttClientTrait> {
client: T,
device_id: String,
}
impl<T: MqttClientTrait> SensorService<T> {
pub fn new(client: T, device_id: String) -> Self {
Self { client, device_id }
}
pub async fn connect(&self, broker: &str) -> Result<(), Box<dyn std::error::Error>> {
self.client.connect(broker).await?;
Ok(())
}
pub async fn publish_temperature(&self, temp: f32) -> Result<(), Box<dyn std::error::Error>> {
let topic = format!("sensors/{}/temperature", self.device_id);
let payload = format!("{temp:.1}");
self.client.publish_qos1(&topic, payload.as_bytes()).await?;
Ok(())
}
pub async fn subscribe_commands(&self) -> Result<(), Box<dyn std::error::Error>> {
let topic = format!("commands/{}/#", self.device_id);
self.client.subscribe(&topic, |msg| {
println!("Command: {}", String::from_utf8_lossy(&msg.payload));
}).await?;
Ok(())
}
}MockMqttClient::new(client_id) creates a mock. Configure canned responses with the
set_*_response methods (all async). MockMqttClient is Clone (it shares state
via an Arc), so a clone handed to your code observes the same recorded calls.
use mqtt5::{MockMqttClient, PublishResult};
#[tokio::test]
async fn test_sensor_service() {
// Create mock client
let mock = MockMqttClient::new("test-device");
// Configure expected responses
mock.set_connect_response(Ok(())).await;
mock.set_publish_response(Ok(PublishResult::QoS1Or2 { packet_id: 1 })).await;
// Use the mock in your code
let service = SensorService::new(mock.clone(), "sensor-001".to_string());
// Test connect
service.connect("mqtt://localhost:1883").await.unwrap();
// Test publish
service.publish_temperature(25.5).await.unwrap();
// Verify calls were made
let calls = mock.get_calls().await;
assert_eq!(calls.len(), 2);
}Available response setters: set_connect_response, set_connect_with_options_response,
set_disconnect_response, set_publish_response, set_subscribe_response,
set_unsubscribe_response. You can also force the connected flag directly with
set_connected(bool) (synchronous), and clear the call log with clear_calls().await.
If no response is configured, the mock uses sensible defaults (connect/disconnect
succeed and flip the connected flag; publish returns QoS0 for QoS 0 and
QoS1Or2 with a generated packet id otherwise; subscribe returns a generated
packet id with the requested/AtMostOnce QoS).
get_calls().await returns a Vec<MockCall>. MockCall is re-exported at the crate
root as mqtt5::MockCall. Its variants and their fields:
| Variant | Fields |
|---|---|
Connect |
address: String |
ConnectWithOptions |
address: String, options: Box<ConnectOptions>
|
Disconnect |
— |
Publish |
topic: String, payload: Vec<u8>
|
PublishWithOptions |
topic: String, payload: Vec<u8>, options: PublishOptions
|
Subscribe |
topic: String |
SubscribeWithOptions |
topic: String, options: SubscribeOptions
|
Unsubscribe |
topic: String |
SetQueueOnDisconnect |
enabled: bool |
Note: the plain Publish variant records only topic and payload. QoS is captured
only when the call went through publish_with_options (recorded as
PublishWithOptions, whose options.qos holds the QoS). The convenience helpers
(publish_qos1, publish_retain, subscribe_with_options, subscribe_many, etc.)
route through the *_with_options methods, so they are recorded as the
...WithOptions variants.
use mqtt5::MockCall;
let calls = mock.get_calls().await;
for call in &calls {
match call {
MockCall::Connect { address } => println!("Connected to {address}"),
MockCall::PublishWithOptions { topic, options, .. } => {
println!("Published to {topic} with QoS {:?}", options.qos);
}
MockCall::Subscribe { topic } => println!("Subscribed to {topic}"),
_ => {}
}
}use mqtt5::{MockCall, QoS};
let calls = mock.get_calls().await;
// Check connect was called
assert!(calls.iter().any(|c| matches!(c, MockCall::Connect { .. })));
// Check a publish to a specific topic (via the with-options variant)
assert!(calls.iter().any(|c| matches!(
c,
MockCall::PublishWithOptions { topic, .. } if topic == "sensors/sensor-001/temperature"
)));
// Check QoS level
assert!(calls.iter().any(|c| matches!(
c,
MockCall::PublishWithOptions { options, .. } if options.qos == QoS::AtLeastOnce
)));MqttError::ConnectionRefused carries a ReasonCode:
use mqtt5::MqttError;
use mqtt5::protocol::v5::reason_codes::ReasonCode;
mock.set_connect_response(Err(MqttError::ConnectionRefused(ReasonCode::NotAuthorized))).await;
let result = service.connect("mqtt://localhost:1883").await;
assert!(result.is_err());use mqtt5::MqttError;
use mqtt5::protocol::v5::reason_codes::ReasonCode;
mock.set_publish_response(Err(MqttError::PublishFailed(ReasonCode::NotAuthorized))).await;
let result = service.publish_temperature(25.5).await;
assert!(result.is_err());Register a subscription callback with subscribe, then deliver a message with
simulate_message(topic, payload, qos). It matches against registered topic filters
(supporting + and #) and returns Err if no subscription matches.
use mqtt5::{MqttClientTrait, QoS};
// Subscribe with a callback
mock.subscribe("commands/#", |msg| {
println!("Received: {}", String::from_utf8_lossy(&msg.payload));
}).await?;
// Simulate an incoming message: (topic, payload, qos)
mock.simulate_message("commands/sensor-001/restart", b"now".to_vec(), QoS::AtMostOnce).await?;use mqtt5::{MockMqttClient, MockCall, MqttClientTrait, PublishResult, QoS};
// Your production code
async fn iot_device_loop<T: MqttClientTrait>(
client: &T,
device_id: &str,
) -> Result<(), Box<dyn std::error::Error>> {
// Connect
client.connect("mqtt://localhost:1883").await?;
// Subscribe to commands
let topic = format!("commands/{device_id}/#");
client.subscribe(&topic, |msg| {
println!("Command: {msg:?}");
}).await?;
// Publish telemetry
for i in 0..3 {
let topic = format!("telemetry/{device_id}/data");
let payload = format!(r#"{{"seq":{i}}}"#);
client.publish_qos1(&topic, payload.as_bytes()).await?;
}
Ok(())
}
#[tokio::test]
async fn test_iot_device_loop() {
// Setup
let mock = MockMqttClient::new("test-device");
mock.set_connect_response(Ok(())).await;
mock.set_subscribe_response(Ok((1, QoS::AtMostOnce))).await;
mock.set_publish_response(Ok(PublishResult::QoS1Or2 { packet_id: 1 })).await;
// Execute
let result = iot_device_loop(&mock, "device-001").await;
assert!(result.is_ok());
// Verify
let calls = mock.get_calls().await;
// 1 connect + 1 subscribe + 3 publishes = 5 calls
assert_eq!(calls.len(), 5);
// Verify connect
assert!(matches!(&calls[0], MockCall::Connect { address } if address == "mqtt://localhost:1883"));
// subscribe() routes through subscribe_with_options -> SubscribeWithOptions
assert!(matches!(
&calls[1],
MockCall::SubscribeWithOptions { topic, .. } if topic == "commands/device-001/#"
));
// publish_qos1() routes through publish_with_options -> PublishWithOptions
let publish_count = calls
.iter()
.filter(|c| matches!(c, MockCall::PublishWithOptions { .. }))
.count();
assert_eq!(publish_count, 3);
}-
Use trait bounds - Accept
T: MqttClientTraitinstead of concreteMqttClient - Inject dependencies - Pass the client as a parameter rather than creating it internally
- Verify important calls - Check topics, payloads, and QoS levels
- Test error paths - Simulate failures to verify error handling
-
Reset between tests - Create a new mock (or
clear_calls) for each test case
#[tokio::test]
async fn test_my_function() {
let mock = MockMqttClient::new("test");
// ... test code
}struct TestFixture {
mock: MockMqttClient,
service: SensorService<MockMqttClient>,
}
impl TestFixture {
async fn new() -> Self {
let mock = MockMqttClient::new("test-device");
mock.set_connect_response(Ok(())).await;
mock.set_publish_response(Ok(PublishResult::QoS0)).await;
let service = SensorService::new(mock.clone(), "test".to_string());
Self { mock, service }
}
}
#[tokio::test]
async fn test_with_fixture() {
let fixture = TestFixture::new().await;
fixture.service.publish_temperature(25.0).await.unwrap();
// ...
}Getting Started
Broker Guide
Client Guide
Platform Guides
CLI Reference
Development