From dd4cd8cb9cb2d5b6db941e5a7c89b6d7dc6f9f19 Mon Sep 17 00:00:00 2001 From: Claas Date: Sun, 30 Nov 2025 22:17:46 +0100 Subject: [PATCH] Start work on subscriptions --- fan-controller/src/task.rs | 54 ++++++++++++++++++++------------------ 1 file changed, 29 insertions(+), 25 deletions(-) diff --git a/fan-controller/src/task.rs b/fan-controller/src/task.rs index e831889..7a3f328 100644 --- a/fan-controller/src/task.rs +++ b/fan-controller/src/task.rs @@ -30,6 +30,7 @@ use embassy_rp::pio::PioPin; use embassy_sync::blocking_mutex::raw::CriticalSectionRawMutex; use embassy_sync::channel::{self, Channel}; use embassy_sync::mutex::Mutex; +use embassy_sync::pubsub::subscriber::Sub; use embassy_sync::signal::Signal; use embassy_sync::waitqueue::AtomicWaker; use embassy_time::{Duration, Instant, Timer, with_deadline, with_timeout}; @@ -42,9 +43,9 @@ async fn network_task(stack: &'static Stack>) -> ! { stack.run().await } -async fn handle_subscribe_acknowledgement<'f>( +async fn handle_subscribe_acknowledgement<'f, const SUBSCRIPTIONS: usize>( acknowledgement: &'f SubscribeAcknowledgement<'f>, - acknowledgements: &Mutex, + acknowledgements: &Mutex, ) { info!("Received subscribe acknowledgement"); let mut acknowledgements = acknowledgements.lock().await; @@ -61,9 +62,9 @@ async fn handle_subscribe_acknowledgement<'f>( ); } -async fn wait_for_acknowledgement( +async fn wait_for_acknowledgement( packet_identifier: NonZeroU16, - acknowledgements: &Mutex, + acknowledgements: &Mutex, waker: &AtomicWaker, ) { info!("Waiting for subscribe acknowledgements"); @@ -185,11 +186,12 @@ async fn listen< FromError, Send: for<'a> TryFrom, Error = FromError>, const SEND: usize, + const SUBSCRIPTIONS: usize, >( reader: &mut impl Read, sender: &channel::Sender<'_, CriticalSectionRawMutex, Result, SEND>, client_state: &Signal, - acknowledgements: &Mutex, + acknowledgements: &Mutex, ping_response_signal: &Signal, ) { let mut buffer = [0; 1024]; @@ -427,11 +429,11 @@ async fn talk( } } -async fn set_up_subscriptions( +async fn set_up_subscriptions( packet_identifier: NonZeroU16, - acknowledgements: &Mutex, + acknowledgements: &Mutex, outgoing: &Channel, 8>, - subscriptions: &'static [Subscription<'static>], + subscriptions: &'static [Subscription<'static>; SUBSCRIPTIONS], waker: &AtomicWaker, ) { info!("Setting up subscriptions"); @@ -587,27 +589,28 @@ async fn poll_sensors(fans: ModbusOnceLock) { } } +const SUBSCRIBE_OPTIONS: mqtt::packet::subscribe::Options = mqtt::packet::subscribe::Options::new( + QualityOfService::AtMostOnceDelivery, + false, + false, + // mqtt::packet::subscribe::RetainHandling::DoNotSend, + mqtt::packet::subscribe::RetainHandling::SendAtSubscribe, +); + +const SUBSCRIPTIONS_LENGTH: usize = 3; // Subscribe to home assistant topics -const SUBSCRIPTIONS: [Subscription; 2] = [ +const SUBSCRIPTIONS: [Subscription; SUBSCRIPTIONS_LENGTH] = [ Subscription { topic_filter: topic::fan_controller::COMMAND, - options: mqtt::packet::subscribe::Options::new( - QualityOfService::AtMostOnceDelivery, - false, - false, - // mqtt::packet::subscribe::RetainHandling::DoNotSend, - mqtt::packet::subscribe::RetainHandling::SendAtSubscribe, - ), + options: SUBSCRIBE_OPTIONS, }, Subscription { - topic_filter: "fancontroller/speed/percentage", - options: mqtt::packet::subscribe::Options::new( - QualityOfService::AtMostOnceDelivery, - false, - false, - // mqtt::packet::subscribe::RetainHandling::DoNotSend, - mqtt::packet::subscribe::RetainHandling::SendAtSubscribe, - ), + topic_filter: topic::fan_controller::fan_1::COMMAND, + options: SUBSCRIBE_OPTIONS, + }, + Subscription { + topic_filter: topic::fan_controller::fan_2::COMMAND, + options: SUBSCRIBE_OPTIONS, }, ]; @@ -797,7 +800,8 @@ pub(super) async fn mqtt_with_connect< // with futures so this will be refactored when I made it work // Contains the status of the subscribe packets send out. The packet identifier represents the // index in the array - let acknowledgements: Mutex = Mutex::new([false, false]); + let acknowledgements: Mutex = + Mutex::new([false; SUBSCRIPTIONS_LENGTH]); // The waker needs to be woken to complete the subscribe acknowledgement future. // The embassy documentation does not explain when to use [`AtomicWaker`] but I am assuming // it is useful for cases like this where I need to mutate a static. -- 2.51.2