use crate::PingRequest; use crate::mqtt::packet::connect::Connect; use crate::mqtt::packet::disconnect::Disconnect; use crate::mqtt::packet::ping_response::PingResponse; use crate::mqtt::packet::publish; use crate::mqtt::packet::subscribe::{Subscribe, Subscription}; use crate::mqtt::packet::subscribe_acknowledgement::SubscribeAcknowledgement; use crate::mqtt::task::send; use crate::mqtt::{self}; use crate::mqtt::{TryDecode, non_zero_u16}; use crate::{configuration, gain_control}; use ::mqtt::QualityOfService; use core::future::poll_fn; use core::num::NonZeroU16; use core::task::Poll; use cyw43::{Control, NetDriver}; use defmt::{Format, error, info, unwrap, warn}; use embassy_executor::Spawner; use embassy_futures::join::join3; use embassy_futures::select::{Either, Either3, select, select3}; use embassy_net::dns::{DnsQueryType, DnsSocket}; use embassy_net::driver::Driver; use embassy_net::tcp::{TcpSocket, TcpWriter}; use embassy_net::{Config, IpAddress, IpEndpoint, Stack, StackResources}; use embassy_rp::clocks::RoscRng; use embassy_rp::peripherals::{DMA_CH0, PIN_23, PIN_25, PIO0}; 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::signal::Signal; use embassy_sync::waitqueue::AtomicWaker; use embassy_time::{Duration, Instant, TimeoutError, Timer, with_deadline, with_timeout}; use embedded_io_async::Read; use rand::RngCore; use static_cell::StaticCell; #[embassy_executor::task] async fn network_task(stack: &'static Stack>) -> ! { info!("Started network task"); stack.run().await } async fn handle_subscribe_acknowledgement<'f, const SUBSCRIPTIONS: usize>( acknowledgement: &'f SubscribeAcknowledgement<'f>, //TODO rework this to be a channel that sends the packet identifier as acknowledgement acknowledgements: &Mutex, waker: &AtomicWaker, ) { info!("[Subscription] Received subscribe acknowledgement. Waiting for acknowledgement lock"); { let mut acknowledgements = acknowledgements.lock().await; info!("[Subscription] Locked ACKNOWLEDGEMENTS"); // Validate server sends a valid packet identifier or we get bamboozled and panic let Some(value) = acknowledgements.get_mut(acknowledgement.packet_identifier as usize) else { warn!( "[Subscription] Received subscribe acknowledgement for out of bounds packet identifier" ); return; }; // Record the acknowledgement. This is what [`wait_for_acknowledgement`] polls for, so // without it the subscription setup waits for its full timeout even though the broker // answered *value = true; info!( "[Subscription] Received acknowledgement {} Reason codes: {:#04x}", acknowledgement.packet_identifier, acknowledgement.reason_codes ); } // Wake after the lock is released so the waiting future's try_lock succeeds when it polls waker.wake(); } async fn wait_for_acknowledgement( packet_identifier: NonZeroU16, acknowledgements: &Mutex, waker: &AtomicWaker, ) { info!("[Subscription] Waiting for subscribe acknowledgements"); //TODO if this function gets called multiple times it might never be woken up because there // is only one waker and the other call of this function will lock the mutex. To solve this // we could use structs from [embassy-sync::waitqueue] and/or a blocking mutex to remove the // try_lock which is used because the lock function is async and we can not easily await here poll_fn(|context| { // Register before looking at the acknowledgements, never after. The other way round // misses an acknowledgement recorded between the check and the registration: the wake // happens while there is no waker to wake and this future then sleeps until the timeout. // Registering on every poll is required because waking takes the waker back out. waker.register(context.waker()); match acknowledgements.try_lock() { // Checking after registering also covers an acknowledgement that was already // recorded before this future was polled for the first time Ok(guard) => { let packet = guard.get(packet_identifier.get() as usize).unwrap(); if *packet { Poll::Ready(()) } else { Poll::Pending } } // Only recording an acknowledgement takes the lock, and that wakes the waker // registered above once it releases the lock, so this gets polled again Err(_error) => Poll::Pending, } }) .await; info!("[Subscription] Subscribe acknowledgement received") } /// Callback handler for pings received when listening async fn handle_ping_response( ping_response: PingResponse, signal: &Signal, ) { info!("Received ping response"); signal.signal(ping_response); } /// Why an MQTT session ended. Every variant means the connection to the broker is gone and has to /// be established again, so the session futures return this instead of signalling a state that /// nobody listens to. #[derive(Format)] enum SessionEnd { /// The broker closed the TCP connection. BrokerClosedConnection, /// The broker sent a disconnect packet, which means it will not accept anything else on this /// connection. BrokerDisconnected, /// Reading from the socket failed. ReadError, /// A packet could not be split into its parts, which means the stream is out of sync and /// nothing that follows can be trusted. MalformedPacket, /// The keep alive ping request could not be sent. PingRequestFailed, /// The broker did not answer the keep alive ping request within [`configuration::MQTT_TIMEOUT`]. PingResponseTimeout, } async fn listen< ReadError: Format, FromError, Send: for<'a> TryFrom, Error = FromError>, const SEND: usize, const SUBSCRIPTIONS: usize, >( reader: &mut impl Read, sender: &channel::Sender<'_, CriticalSectionRawMutex, Result, SEND>, acknowledgements: &Mutex, acknowledgement_waker: &AtomicWaker, ping_response_signal: &Signal, last_received: &Signal, ) -> SessionEnd { let mut buffer = [0; 1024]; loop { info!("[MQTT/listen] Waiting for packet"); let result = reader.read(&mut buffer).await; let bytes_read = match result { // This indicates the TCP connection was closed. (See embassy-net documentation) Ok(0) => { warn!("MQTT broker closed connection"); return SessionEnd::BrokerClosedConnection; } Ok(bytes_read) => bytes_read, Err(error) => { warn!("Error reading from MQTT broker: {:?}", error); return SessionEnd::ReadError; } }; info!("Packet received"); // Anything at all from the broker proves the connection is still there, whatever kind of // packet it turned out to be last_received.signal(Instant::now()); let parts = match mqtt::packet::get_parts(&buffer[..bytes_read]) { Ok(parts) => parts, Err(error) => { warn!("Error reading MQTT packet: {:?}", error); return SessionEnd::MalformedPacket; } }; info!("Handling packet"); match parts.r#type { publish::Publish::TYPE => { info!("Received publish"); let publish = match publish::Publish::try_decode( parts.flags, parts.variable_header_and_payload, ) { Ok(publish) => publish, Err(error) => { error!("Error reading publish: {:?}", error); continue; } }; // MQTT does not concern itself with validity of the payload // We allow for a TryFrom implementation anyways to reduce boilerplate for users // because if we only allow a From implementation users would have to create their own result type // that implements the From trait and stores the error to be passed on. // If users still want to use a From implementation then they can still do this // as any type that implements From, automatically implements TryFrom let result = Send::try_from(publish); info!("[MQTT/listen] Sending out received send packet from publish"); sender.send(result).await; info!("[MQTT/listen] Sent out received send packet from publish") } SubscribeAcknowledgement::TYPE => { let subscribe_acknowledgement = match SubscribeAcknowledgement::read(parts.variable_header_and_payload) { Ok(acknowledgement) => acknowledgement, Err(error) => { error!( "[Subscription] Error reading subscribe acknowledgement: {:?}", error ); continue; } }; info!("[MQTT/listen] Waiting for subscribe acknowledgement handling"); handle_subscribe_acknowledgement( &subscribe_acknowledgement, acknowledgements, acknowledgement_waker, ) .await; info!("[MQTT/listen] Subscribe acknowledgement handled"); } PingResponse::TYPE => { info!("Received ping response"); let ping_response = match PingResponse::try_decode( parts.flags, parts.variable_header_and_payload, ) { Ok(response) => response, // Matching to get compiler error if this changes Err(_) => { defmt::unreachable!( "Ping response is always empty so decode should always succeed if the protocol did not change" ) } }; info!("[MQTT/listen] Waiting for ping response handling"); handle_ping_response(ping_response, ping_response_signal).await; info!("[MQTT/listen] Ping response handled"); } Disconnect::TYPE => { info!("Received disconnect"); let disconnect = Disconnect::try_decode(parts.flags, parts.variable_header_and_payload); info!("Disconnect {:?}", disconnect); // The broker will not accept anything else on this connection, so end the session // and let the reconnect loop close the socket and start over return SessionEnd::BrokerDisconnected; } other => info!("Unsupported packet type {}", other), } } } enum Message<'a, T: Publish> { Subscribe(Subscribe<'a>), Publish(T), } impl From for Message<'_, T> where T: Publish, { fn from(publish: T) -> Self { Message::Publish(publish) } } async fn talk( writer: &Mutex>, outgoing: &Channel, 8>, last_packet: &Signal, ) { loop { info!("[MQTT/talk] Waiting for next message to send out"); let message = outgoing.receive().await; info!("[MQTT/talk] Message received"); info!("[MQTT/talk] Waiting for lock on writer"); let mut writer = writer.lock().await; info!("[MQTT/talk] Writer lock acquired"); match message { Message::Subscribe(subscribe) => { info!("[MQTT/talk] Sending subscribe"); if let Err(error) = send(&mut *writer, subscribe).await { error!("[MQTT/talk] Error sending subscribe: {:?}", error); continue; } info!("[MQTT/talk] Subscribe sent"); } Message::Publish(publish) => { info!( "[MQTT/talk] Sending publish {:?} {:?}", publish.topic(), core::str::from_utf8(publish.payload()).unwrap_or("Payload is not UTF-8") ); if let Err(error) = send(&mut *writer, publish).await { error!("[MQTT/talk] Error sending publish: {:?}", error); continue; } info!("[MQTT/talk] Publish completed successfully"); } } last_packet.signal(Instant::now()); } } async fn set_up_subscriptions( packet_identifier: NonZeroU16, acknowledgements: &Mutex, outgoing: &Channel, 8>, subscriptions: &'static [Subscription<'static>; SUBSCRIPTIONS], waker: &AtomicWaker, ) { info!("[Subscription] Setting up subscriptions"); for subscription in subscriptions { info!( "[Subscription] Subscribing to MQTT topic: {}", subscription.topic_filter ); } let message = Message::Subscribe(Subscribe { subscriptions, //TODO free identifier management packet_identifier, }); info!("[Subscription] Sending out subscription request"); outgoing.send(message).await; info!("[Subscription] Sent subscription request"); info!("[Subscription] Waiting for subscription acknowledgements"); const TIMEOUT: Duration = Duration::from_secs(30); let result = with_timeout( TIMEOUT, wait_for_acknowledgement(packet_identifier, acknowledgements, waker), ) .await; if let Err(error) = result { error!( "[Subscription] Timed out setting up subscriptions. Waiting for subscription acknowledgements timed out after {:?}: {:?}", TIMEOUT, error ); return; } info!("[Subscription] Set up subscriptions complete") } /// Keep alive task async fn keep_alive( writer: &Mutex>, last_sent: &Signal, last_received: &Signal, ping_response: &Signal, ) -> SessionEnd { // Two different things have to happen within the keep alive interval, and sending is only one // of them. The broker drops a client it has not heard from, which sending covers. A broker // that has silently gone away is only noticed by asking it something and getting an answer, // which only receiving covers. Watching what was sent alone is what let a dead connection go // unnoticed: the fans publish every SENSOR_POLL_INTERVAL, which is shorter than the keep // alive, so the deadline was reset forever and the ping that would have exposed the broker's // absence was never sent. // // So the next deadline belongs to whichever direction has been quiet longer. let now = Instant::now(); let mut last_sent_at = last_sent.try_take().unwrap_or(now); let mut last_received_at = last_received.try_take().unwrap_or(now); loop { // The server waits for 1.5 times the keep alive interval, so being off by a bit due to // network, async overhead or the clock not being exactly precise is fine let deadline = last_sent_at.min(last_received_at) + configuration::KEEP_ALIVE; info!("[Keep Alive] Waiting for traffic in either direction or keep alive timeout"); let result = with_deadline(deadline, select(last_sent.wait(), last_received.wait())).await; match result { Err(TimeoutError {}) => { info!("[Keep Alive] Sending keep alive ping request"); let mut writer = writer.lock().await; // Send keep alive ping request if let Err(error) = send(&mut *writer, PingRequest).await { error!( "[Keep Alive] Error sending keep alive ping request: {:?}", error ); return SessionEnd::PingRequestFailed; } // Keep alive is time from when the last packet was sent and not when the ping response // was received. Therefore, we need to reset it here last_sent_at = Instant::now(); // Wait for ping response info!("[Keep Alive] Waiting for ping response"); if let Err(TimeoutError) = with_timeout(configuration::MQTT_TIMEOUT, ping_response.wait()).await { // Assume disconnect from server error!("[Keep Alive] Timeout waiting for ping response. Disconnecting"); return SessionEnd::PingResponseTimeout; } // The answer is the proof the connection is alive that sending alone never gives last_received_at = Instant::now(); let response_duration = last_received_at - last_sent_at; info!( "[Keep Alive] Received and processed ping response in {:?} ({}μs) after sending ping request", response_duration, response_duration.as_micros() ); } Ok(Either::First(sent_at)) => { info!("[Keep Alive] Sent packet. Resetting timer to send keep alive"); last_sent_at = sent_at; } Ok(Either::Second(received_at)) => { info!("[Keep Alive] Heard from the broker. Resetting timer to send keep alive"); last_received_at = received_at; } } } } 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 = 6; // Subscribe to home assistant topics const SUBSCRIPTIONS: [Subscription; SUBSCRIPTIONS_LENGTH] = [ Subscription { topic_filter: topic::fan_controller::COMMAND, options: SUBSCRIBE_OPTIONS, }, Subscription { topic_filter: topic::fan_controller::fan_1::state::COMMAND, options: SUBSCRIBE_OPTIONS, }, Subscription { topic_filter: topic::fan_controller::fan_1::percentage::COMMAND, options: SUBSCRIBE_OPTIONS, }, Subscription { topic_filter: topic::fan_controller::fan_2::state::COMMAND, options: SUBSCRIBE_OPTIONS, }, Subscription { topic_filter: topic::fan_controller::fan_2::percentage::COMMAND, options: SUBSCRIBE_OPTIONS, }, Subscription { topic_filter: topic::fan_controller::bypass::COMMAND, options: SUBSCRIBE_OPTIONS, }, ]; /// Trait must be implemented by types that represent messages that can be published to MQTT. pub(super) trait Publish { const TYPE: u8 = 3; /// Providing the topic string for the publish through a function allows implementing types to be more flexible. /// For example, they can be defined as an enum and match internally to provide the appropriate string for the enum variant. fn topic(&self) -> &str; fn payload(&self) -> &[u8]; /// Whether the broker should keep this message and hand it to whoever subscribes next. /// Almost nothing wants this: state that is published on every change is better re-read than /// remembered. A message about an event nobody was watching for is the exception. fn is_retained(&self) -> bool { false } } pub(super) async fn set_up_network_stack( spawner: Spawner, pwr_pin: PIN_23, cs_pin: PIN_25, pio: PIO0, dma: DMA_CH0, dio: impl PioPin, clk: impl PioPin, ) -> &'static Stack> { info!("[Network] Initializing network stack"); let (net_device, mut control) = gain_control(spawner, pwr_pin, cs_pin, pio, dma, dio, clk).await; info!("[Network] Initializing network stack"); static STACK: StaticCell>> = StaticCell::new(); static RESOURCES: StaticCell> = StaticCell::new(); let configuration = Config::dhcpv4(Default::default()); let mut random = RoscRng; let seed = random.next_u64(); // Initialize network stack let stack = &*STACK.init(Stack::new( net_device, configuration, RESOURCES.init(StackResources::<5>::new()), seed, )); unwrap!(spawner.spawn(network_task(stack))); // Join Wi-Fi network info!("[Network] Joining wifi network"); join_wifi_network(&mut control).await; info!("[Network] Joined wifi network"); // Wait for DHCP info!("[Network] Waiting for DHCP"); while !stack.is_config_up() { info!("[Network] DHCP not yet up"); Timer::after_millis(100).await; } info!("[Network] DHCP is up"); info!("[Network] Waiting for link up"); while !stack.is_link_up() { info!("[Network] Link not yet up"); Timer::after_millis(500).await; } info!("[Network] Link is up"); info!("[Network] Waiting for stack to be up"); stack.wait_config_up().await; info!("[Network] Stack is up"); // Hand the control handle on rather than dropping it, so the link can be rebuilt later unwrap!(spawner.spawn(rejoin_wifi_routine(control, stack))); stack } /// Joins the Wi-Fi network again whenever the link goes down. /// /// Joining once at start up is enough right up until the access point goes away. The mesh access /// point this controller reached the network through is switched off overnight, and nothing ever /// brought the link back: the fans, the button and the LEDs carried on working perfectly while /// the controller sat unreachable for days, transmitting to an access point that was not there. /// /// The MQTT session already reconnects when the broker goes away. This is the same idea one layer /// down, and without it that reconnect loop retries forever over a radio link that can never work /// again. #[embassy_executor::task] async fn rejoin_wifi_routine( mut control: Control<'static>, stack: &'static Stack>, ) { loop { // Nothing to do for as long as the link holds while stack.is_link_up() { Timer::after(configuration::WIFI_LINK_CHECK_INTERVAL).await; } warn!("[Network] Wi-Fi link is down. Joining again"); join_wifi_network(&mut control).await; info!("[Network] Joined wifi network again"); // An address has to be picked up again before anything can use the stack info!("[Network] Waiting for the stack to come back up"); stack.wait_config_up().await; info!("[Network] Stack is up again"); } } async fn join_wifi_network(control: &mut Control<'_>) { loop { info!("[Join Wifi] Attempting to join Wi-Fi network"); match control .join_wpa2(configuration::WIFI_NETWORK, configuration::WIFI_PASSWORD) .await { Ok(_) => break, Err(error) => { info!( "[Join Wifi] Error joining Wi-Fi network with status: {}", error.status ); // An access point that is switched off for the night is not going to answer any // sooner for being asked continuously Timer::after(configuration::WIFI_JOIN_RETRY_DELAY).await; } } } } async fn resolve_mqtt_broker_address<'a, D>( driver: &'a Stack, mqtt_broker_address: &str, ) -> IpAddress where D: Driver + 'static, { let dns_client = DnsSocket::new(driver); loop { //TODO support IPv6 info!("[Resolve MQTT broker address] Waiting for DNS query"); let result = dns_client.query(mqtt_broker_address, DnsQueryType::A).await; info!("[Resolve MQTT broker address] DNS query completed"); let mut addresses = match result { Ok(addresses) => addresses, Err(error) => { info!( "[Resolve MQTT broker address] Error resolving MQTT broker IP address with {}: {:?}", mqtt_broker_address, error, ); // Exponential backoff doesn't seem necessary here // Maybe the current installation of Home Assistant is in the process of // being set up and the entry is not yet available Timer::after_secs(25).await; continue; } }; if addresses.is_empty() { info!( "[Resolve MQTT broker address] No addresses found for Home Assistant MQTT broker" ); Timer::after_millis(500).await; continue; } return addresses.swap_remove(0); } } /// Set in Homeassitant under Settings > People > Users Tab. Not to be confused with the People tab. /// The Users Tab might only be visible in advanced mode as administrator. /// A separate account is recommended for each device. pub(crate) struct MqttBrokerCredentials<'a> { pub(crate) username: &'a str, pub(crate) password: &'a [u8], } pub(super) struct MqttBrokerConfiguration<'a> { pub(super) client_identifier: &'a str, pub(super) address: &'a str, pub(super) credentials: MqttBrokerCredentials<'a>, pub(super) keep_alive_seconds: Duration, } /// This is piping the output of the receiver into the sender. This is in fact so simple that there should be something built-in to replace this async fn handle_publish_send<'receiver, Send: Publish, const SEND: usize>( receiver: channel::Receiver<'receiver, CriticalSectionRawMutex, Send, SEND>, outgoing: &Channel, 8>, ) { loop { info!("[Publish Send] Waiting for message"); let message = receiver.receive().await; info!("[Publish Send] Received message. Sending to MQTT outgoing"); outgoing.send(message.into()).await; info!("[Publish Send] Message sent to MQTT outgoing"); } } /// Keeps the fan controller connected to the MQTT broker for as long as it is powered. /// /// Runs one session at a time and never returns: when a session ends, everything belonging to it /// is dropped, the socket is closed and a new connection is established after a backoff. The fans /// keep running and the button keeps working while there is no connection, so being disconnected /// is a normal state to recover from rather than a reason to give up. pub(super) async fn mqtt_with_connect< 'tcp, 'sender, 'receiver, 'configuration, FromError, Receive: for<'a> TryFrom, Error = FromError>, const RECEIVE: usize, Send: Publish, const SEND: usize, D: Driver + 'static, >( stack: &Stack, sender: channel::Sender<'sender, CriticalSectionRawMutex, Result, RECEIVE>, receiver: channel::Receiver<'receiver, CriticalSectionRawMutex, Send, SEND>, mqtt_broker_configuration: &MqttBrokerConfiguration<'configuration>, ) // -> Client<'tcp, 'sender, 'receiver, Send, SEND, Receive, RECEIVE> { // Now we can use it let mut receive_buffer = [0; 1024]; let mut send_buffer = [0; 1024]; // Every attempt that fails and every session that ends comes back here. The wait between // attempts doubles up to a cap so a broker that stays down is not asked every second, and it // is reset as soon as an MQTT connection was established. let mut backoff = configuration::MQTT_RECONNECT_BACKOFF_INITIAL; loop { // A fresh socket per attempt. Reusing one would mean waiting for the previous connection // to be fully closed first and the buffers are only borrowed for as long as it lives. let mut socket = TcpSocket::new(stack, &mut receive_buffer, &mut send_buffer); 'attempt: { info!("[MQTT/main] Resolving MQTT broker IP address"); // Get home assistant MQTT broker IP address let address = resolve_mqtt_broker_address(stack, mqtt_broker_configuration.address).await; info!("[MQTT/main] MQTT broker IP address resolved"); info!("[MQTT/main] Connecting to MQTT broker through TCP"); let endpoint = IpEndpoint::new(address, configuration::MQTT_BROKER_PORT); // One attempt only. Retrying is what the surrounding loop does, and doing it here as // well would skip the backoff for the case that is most likely to need it. if let Err(error) = socket.connect(endpoint).await { warn!( "[MQTT/main] Error connecting to Home Assistant MQTT broker: {:?}", error ); break 'attempt; } info!("[MQTT/main] Connected to MQTT broker through TCP"); use crate::mqtt::task; //TODO error handling let packet = defmt::unwrap!(Connect::try_from(mqtt_broker_configuration)); // The session borrows the socket through the reader and writer halves. Keeping it in // its own scope ends those borrows so the socket can be closed below. let session_end = { info!("[MQTT/main] Establishing MQTT connection"); let (mut reader, mut writer) = socket.split(); if let Err(error) = task::connect(&mut writer, &mut reader, packet).await { warn!("[MQTT/main] Error connecting to MQTT broker: {:?}", error); break 'attempt; }; info!("[MQTT/main] MQTT connection established"); // The connection works, so the next failure is a new problem and not a broker that // has been unreachable for a while backoff = configuration::MQTT_RECONNECT_BACKOFF_INITIAL; //TODO yes static "global" state is bad, but I am still learning how to use wakers and polling // 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; 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. let waker: AtomicWaker = AtomicWaker::new(); let ping_response: Signal = Signal::new(); let outgoing: Channel, 8> = Channel::new(); // The instants when a packet was last sent and last received, which together // decide when the next keep alive has to be sent let last_sent: Signal = Signal::new(); let last_received: Signal = Signal::new(); // Using a mutex for the writer, so it can be shared between the task that sends messages (for // subscribing and publishing fan speed updates) and the task that sends the keep alive ping let writer = Mutex::>::new(writer); // Future 1 let listen = listen( &mut reader, &sender, &acknowledgements, &waker, &ping_response, &last_received, ); // Future 2 let talk = talk(&writer, &outgoing, &last_sent); // Future 3 let set_up = set_up_subscriptions( non_zero_u16!(1), &acknowledgements, &outgoing, &SUBSCRIPTIONS, &waker, ); // Future 4 let keep_alive = keep_alive(&writer, &last_sent, &last_received, &ping_response); // Future 5 let handle_publish_send = handle_publish_send(receiver, &outgoing); // Only listening and the keep alive notice that the connection is gone. Selecting // on them cancels the other three by dropping them, which is the point: they would // otherwise keep waiting on a dead socket forever. Anything they had picked up but // not yet written is lost with them, which is acceptable for state updates where // only the latest value matters. match select3(listen, keep_alive, join3(talk, set_up, handle_publish_send)).await { Either3::First(session_end) | Either3::Second(session_end) => session_end, Either3::Third(_) => defmt::unreachable!( "Talking and piping publishes out never return, so neither does joining them" ), } }; warn!( "[MQTT/main] Lost connection to MQTT broker: {:?}. Reconnecting", session_end ); } // Send a reset and wait for it to go out so the broker does not keep a half open // connection around, then release the socket before waiting socket.abort(); if let Err(error) = socket.flush().await { warn!("[MQTT/main] Error closing the TCP connection: {:?}", error); } drop(socket); info!("[MQTT/main] Reconnecting to MQTT broker in {:?}", backoff); Timer::after(backoff).await; backoff = Duration::min(backoff * 2, configuration::MQTT_RECONNECT_BACKOFF_MAX); } }