From 9c2fe39ffd534e35cf767b4d7b993ea70d5b1f2a Mon Sep 17 00:00:00 2001 From: Dimitrios Kouros Date: Sat, 15 Aug 2026 17:16:56 +0300 Subject: [PATCH] connection error --- src/main.rs | 40 ++++++++++++++++++++++++++++------------ 1 file changed, 28 insertions(+), 12 deletions(-) diff --git a/src/main.rs b/src/main.rs index eaaa95a..3aa2f40 100644 --- a/src/main.rs +++ b/src/main.rs @@ -15,6 +15,7 @@ const METEO_MUNICH_API_URL: &str = const LATITUDE_MUC: f64 = 48.163142; const LONGITUDE_MUC: f64 = 11.542922; const PUBLISH_COUNT: i32 = 5; +const MQTT_TIMEOUT: Duration = Duration::from_secs(10); #[tokio::main(flavor = "current_thread")] async fn main() { @@ -26,21 +27,36 @@ async fn main() { task::spawn(async move { publish(client.clone(), &weather).await; }); - println!("test{}", 100); - let mut count = 0; - loop { - let event = eventloop.poll().await; - match &event { - Ok(_) => { - if count == PUBLISH_COUNT * 2 { - break; + println!("{}", 100); + + let poll_result = time::timeout(MQTT_TIMEOUT, async { + let mut count = 0; + let mut err_count = 0; + loop { + let event = eventloop.poll().await; + match &event { + Ok(_) => { + if count >= PUBLISH_COUNT * 2 { + break; + } + count += 1; + } + Err(e) => { + println!("MQTT Error = {e:?}"); + err_count += 1; + if err_count >= 5 { + eprintln!("Exiting MQTT loop after 5 consecutive errors."); + break; + } + time::sleep(Duration::from_millis(200)).await; } - count += 1; - } - Err(e) => { - println!("Error = {e:?}"); } } + }) + .await; + + if poll_result.is_err() { + eprintln!("Timed out waiting for MQTT publishing to complete."); } } Err(e) => {