use mqtt5::protocol::v5::reason_codes::ReasonCode;
use mqtt5::{ConnectOptions, MqttClient, QoS, SubscribeOptions};
use std::time::Duration;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt::init();
let broker =
std::env::var("MQTT_BROKER").unwrap_or_else(|_| "mqtt://127.0.0.1:1883".to_string());
let options = ConnectOptions::new("deferred-ack-demo")
.with_deferred_ack(true)
.with_clean_start(false)
.with_session_expiry_interval(3600)
.with_receive_maximum(16);
let client = MqttClient::with_options(options);
client.connect(&broker).await?;
let subscribe_options = SubscribeOptions {
qos: QoS::ExactlyOnce,
..Default::default()
};
client
.subscribe_with_ack("jobs/#", subscribe_options, |publish, token| match process(
&publish.topic_name,
&publish.payload,
) {
Ok(()) => token.ack(),
Err(error) => {
eprintln!("rejecting {}: {error}", publish.topic_name);
token.reject(ReasonCode::UnspecifiedError);
}
})
.await?;
println!("Waiting for jobs on jobs/# (Ctrl-C to exit)...");
tokio::time::sleep(Duration::from_secs(60)).await;
client.disconnect().await?;
Ok(())
}
fn process(topic: &str, payload: &[u8]) -> Result<(), String> {
if payload.is_empty() {
return Err(format!("empty payload on {topic}"));
}
println!("processing {topic} ({} bytes)", payload.len());
Ok(())
}