Skip to content

Mock Testing

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

Mock Testing

Unit test your MQTT code without a real broker using the mock client.


Overview

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.


Using the Trait

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

Creating a Mock Client

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


Verifying Calls

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.

Get All Calls

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}"),
        _ => {}
    }
}

Assert Specific Calls

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

Simulating Errors

Connection Failure

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

Publish Rejection

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

Simulating Incoming Messages

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?;

Complete Test Example

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

Best Practices

  1. Use trait bounds - Accept T: MqttClientTrait instead of concrete MqttClient
  2. Inject dependencies - Pass the client as a parameter rather than creating it internally
  3. Verify important calls - Check topics, payloads, and QoS levels
  4. Test error paths - Simulate failures to verify error handling
  5. Reset between tests - Create a new mock (or clear_calls) for each test case

Integration with Test Frameworks

With #[tokio::test]

#[tokio::test]
async fn test_my_function() {
    let mock = MockMqttClient::new("test");
    // ... test code
}

With Test Fixtures

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

Clone this wiki locally