use bytes::Bytes;
use mqttier::{Connection, MqttierClient, MqttierOptionsBuilder};
use std::collections::HashMap;
use std::time::Duration;
use stinger_mqtt_trait::message::{MqttMessage, MqttMessageBuilder, QoS};
use stinger_mqtt_trait::Mqtt5PubSub;
use tokio::sync::broadcast;
use tokio::time::sleep;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt()
.with_writer(std::io::stdout)
.with_thread_ids(true) .with_target(true) .with_env_filter(
tracing_subscriber::EnvFilter::from_default_env()
.add_directive(tracing::Level::DEBUG.into()),
)
.init();
let options = MqttierOptionsBuilder::default()
.connection(Connection::TcpLocalhost(1883))
.client_id("ping_pong_client")
.ack_timeout_ms(5000u64)
.keepalive_secs(60u16)
.session_expiry_interval_secs(1200u16)
.publish_queue_size(128u16)
.max_incoming_packet_size(10u32 * 1024)
.max_inflight_messages(100u16)
.build()
.expect("Failed to build MqttierOptions");
let mut client = MqttierClient::new(options)?;
let (pong_tx, mut pong_rx) = broadcast::channel::<MqttMessage>(32);
let subscription_result = client
.subscribe("example/pong".to_string(), 1, pong_tx)
.await?;
println!("Subscribed to 'example/pong': {:?}", subscription_result);
println!("Starting MQTT client...");
client.start().await?;
println!("Spawning pong handler ...");
let pong_handler = tokio::spawn(tracing::info_span!("pong_handler").in_scope(|| async move {
while let Ok(message) = pong_rx.recv().await {
let payload_str = String::from_utf8_lossy(&message.payload);
println!("🏓 Received PONG: {}", payload_str);
}
}));
println!("Spawning ping publisher. Will publish ping message very 3 seconds (10 messages total) ...");
let mut ping_client = client.clone();
let pubtask = tokio::spawn(
tracing::info_span!("ping_publisher").in_scope(|| async move {
for i in 1..=10 {
let ping_message_payload = format!(
"{{\"ping\":{},\"timestamp\":{}}}",
i,
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs()
);
println!("🏓 Sending PING #{}", i);
let ping_message = MqttMessageBuilder::default()
.topic("example/ping")
.payload(Bytes::from(ping_message_payload))
.qos(QoS::AtLeastOnce)
.retain(false)
.user_properties(HashMap::new())
.build()
.unwrap();
match ping_client.publish(ping_message).await {
Ok(result) => println!(" PING #{} acknowledged: {:?}", i, result),
Err(e) => println!(" PING #{} failed: {:?}", i, e),
}
sleep(Duration::from_secs(3)).await;
}
}),
);
pubtask.await?;
sleep(Duration::from_secs(2)).await;
pong_handler.abort();
println!("Example completed!");
#[cfg(feature = "metrics")]
{
let metrics = client.get_metrics();
println!("\n=== Final Metrics ===");
println!("\n📡 Connection:");
println!(" Connection attempts: {}", metrics.connection_attempts);
println!(
" Successful connections: {}",
metrics.successful_connections
);
println!(" Failed connections: {}", metrics.failed_connections);
println!(" Reconnection count: {}", metrics.reconnection_count);
println!(" Is connected: {}", metrics.is_connected);
if let Some(uptime) = metrics.connection_uptime_ms() {
println!(
" Uptime: {} ms ({:.2} seconds)",
uptime,
uptime as f64 / 1000.0
);
}
println!("\n📤 Publishing:");
println!(" Total published: {}", metrics.total_messages_published());
println!(" - QoS 0: {}", metrics.messages_published_qos0);
println!(" - QoS 1: {}", metrics.messages_published_qos1);
println!(" - QoS 2: {}", metrics.messages_published_qos2);
println!(" Failed: {}", metrics.publish_failures);
println!(" Timeouts: {}", metrics.publish_timeouts);
if let Some(rate) = metrics.publish_success_rate() {
println!(" Success rate: {:.2}%", rate * 100.0);
}
println!(" Total bytes sent: {}", metrics.total_bytes_sent);
if let Some(avg_size) = metrics.avg_sent_message_size() {
println!(" Avg message size: {} bytes", avg_size);
}
println!("\n📥 Receiving:");
println!(" Total received: {}", metrics.total_messages_received());
println!(" - QoS 0: {}", metrics.messages_received_qos0);
println!(" - QoS 1: {}", metrics.messages_received_qos1);
println!(" - QoS 2: {}", metrics.messages_received_qos2);
println!(" Total bytes received: {}", metrics.total_bytes_received);
if let Some(avg_size) = metrics.avg_received_message_size() {
println!(" Avg message size: {} bytes", avg_size);
}
println!("\n� Subscriptions:");
println!(" Active subscriptions: {}", metrics.active_subscriptions);
println!(" Subscription requests: {}", metrics.subscription_requests);
println!(" Subscription failures: {}", metrics.subscription_failures);
println!("\n✅ Reliability:");
println!(" PUBACK received: {}", metrics.puback_received);
println!(" PUBCOMP received: {}", metrics.pubcomp_received);
println!("\n⚡ Performance:");
println!(
" Avg publish latency: {} ms",
metrics.avg_publish_latency_us
);
println!(
" Min publish latency: {} us",
metrics.min_publish_latency_us
);
println!(
" Max publish latency: {} us",
metrics.max_publish_latency_us
);
}
Ok(())
}