From 598438ae2444524aa14dd83b9a429c064159ee1d Mon Sep 17 00:00:00 2001 From: ruv Date: Sat, 22 Aug 2026 14:58:43 -0400 Subject: [PATCH] test(mqtt): keep live subscriber connected --- .../wifi-densepose-bfld/tests/mosquitto_integration.rs | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/v2/crates/wifi-densepose-bfld/tests/mosquitto_integration.rs b/v2/crates/wifi-densepose-bfld/tests/mosquitto_integration.rs index 70db05da..7056a2d4 100644 --- a/v2/crates/wifi-densepose-bfld/tests/mosquitto_integration.rs +++ b/v2/crates/wifi-densepose-bfld/tests/mosquitto_integration.rs @@ -23,9 +23,7 @@ use std::thread; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use rumqttc::{Client, Event, Incoming, MqttOptions, Packet, QoS}; -use wifi_densepose_bfld::{ - publish_event, BfldEvent, PrivacyClass, RumqttPublisher, -}; +use wifi_densepose_bfld::{publish_event, BfldEvent, PrivacyClass, RumqttPublisher}; const SUBSCRIBE_TIMEOUT: Duration = Duration::from_secs(5); const RECEIVE_TIMEOUT: Duration = Duration::from_secs(10); @@ -79,6 +77,11 @@ fn spawn_subscriber( let (incoming_tx, incoming_rx) = channel(); let (suback_tx, suback_rx) = channel(); thread::spawn(move || { + // rumqttc-v4-next stops the connection once every request sender is + // dropped. Keep the subscriber client alive for as long as its pump + // thread runs; otherwise the broker sees a clean disconnect directly + // after SUBACK and no subsequent publications can be delivered. + let _client_guard = client; for notification in connection.iter() { match notification { Ok(Event::Incoming(Packet::SubAck(_))) => {