Something went wrong. Try again.
[READ-ONLY] Mirror of https://github.com/SantaClaas/embedded-fan-control.
Something went wrong. Try again.
35 kB · 849 lines
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850use 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<NetDriver<'static>>) -> ! { 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<CriticalSectionRawMutex, [bool; SUBSCRIPTIONS]>, 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<const SUBSCRIPTIONS: usize>( packet_identifier: NonZeroU16, acknowledgements: &Mutex<CriticalSectionRawMutex, [bool; SUBSCRIPTIONS]>, 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 listeningasync fn handle_ping_response( ping_response: PingResponse, signal: &Signal<CriticalSectionRawMutex, PingResponse>,) { 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<publish::Publish<'a>, Error = FromError>, const SEND: usize, const SUBSCRIPTIONS: usize,>( reader: &mut impl Read<Error = ReadError>, sender: &channel::Sender<'_, CriticalSectionRawMutex, Result<Send, FromError>, SEND>, acknowledgements: &Mutex<CriticalSectionRawMutex, [bool; SUBSCRIPTIONS]>, acknowledgement_waker: &AtomicWaker, ping_response_signal: &Signal<CriticalSectionRawMutex, PingResponse>, last_received: &Signal<CriticalSectionRawMutex, Instant>,) -> 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<T> From<T> for Message<'_, T>where T: Publish,{ fn from(publish: T) -> Self { Message::Publish(publish) }}
async fn talk<T: Publish>( writer: &Mutex<CriticalSectionRawMutex, TcpWriter<'_>>, outgoing: &Channel<CriticalSectionRawMutex, Message<'_, T>, 8>, last_packet: &Signal<CriticalSectionRawMutex, Instant>,) { 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<T: Publish, const SUBSCRIPTIONS: usize>( packet_identifier: NonZeroU16, acknowledgements: &Mutex<CriticalSectionRawMutex, [bool; SUBSCRIPTIONS]>, outgoing: &Channel<CriticalSectionRawMutex, Message<'_, T>, 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 taskasync fn keep_alive( writer: &Mutex<CriticalSectionRawMutex, TcpWriter<'_>>, last_sent: &Signal<CriticalSectionRawMutex, Instant>, last_received: &Signal<CriticalSectionRawMutex, Instant>, ping_response: &Signal<CriticalSectionRawMutex, PingResponse>,) -> 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 topicsconst 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<NetDriver<'static>> { 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<Stack<NetDriver<'static>>> = StaticCell::new(); static RESOURCES: StaticCell<StackResources<5>> = 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<NetDriver<'static>>,) { 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<D>, mqtt_broker_address: &str,) -> IpAddresswhere 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 thisasync fn handle_publish_send<'receiver, Send: Publish, const SEND: usize>( receiver: channel::Receiver<'receiver, CriticalSectionRawMutex, Send, SEND>, outgoing: &Channel<CriticalSectionRawMutex, Message<'_, Send>, 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<publish::Publish<'a>, Error = FromError>, const RECEIVE: usize, Send: Publish, const SEND: usize, D: Driver + 'static,>( stack: &Stack<D>, sender: channel::Sender<'sender, CriticalSectionRawMutex, Result<Receive, FromError>, 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<CriticalSectionRawMutex, [bool; SUBSCRIPTIONS_LENGTH]> = 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<CriticalSectionRawMutex, PingResponse> = Signal::new();
let outgoing: Channel<CriticalSectionRawMutex, Message<Send>, 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<CriticalSectionRawMutex, Instant> = Signal::new(); let last_received: Signal<CriticalSectionRawMutex, Instant> = 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::<CriticalSectionRawMutex, TcpWriter<'_>>::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); }}