diff --git a/fan-controller/src/fan.rs b/fan-controller/src/fan.rs index 2db7b58..4735db9 100644 --- a/fan-controller/src/fan.rs +++ b/fan-controller/src/fan.rs @@ -24,6 +24,17 @@ pub(crate) fn get_configuration() -> uart::Config { } pub(crate) const MAX_SET_POINT: u16 = 64_000; +/// Like a string with length and capacity of 5. Used for sending publish packets to Home Assistant through MQTT +pub(crate) struct SettingStringBuffer { + buffer: [u8; 5], + start_index: usize, +} + +impl SettingStringBuffer { + pub(crate) fn as_bytes(&self) -> &[u8] { + &self.buffer[self.start_index..] + } +} #[derive(Debug, Format)] pub(crate) struct Setting(u16); @@ -44,6 +55,28 @@ impl Setting { const fn get(&self) -> u16 { self.0 } + + pub(crate) const fn to_string_buffer(&self) -> SettingStringBuffer { + // The largest value the set point can assume is 64000 which is 5 characters long + let mut buffer = [0; 5]; + let mut index = 4; + let mut remainder = self.0; + loop { + let digit = remainder % 10; + // Convert digit to ASCII (which is also valid utf-8) + buffer[index] = digit as u8 + b'0'; + remainder /= 10; + if remainder <= 0 { + break; + } + index -= 1; + } + + SettingStringBuffer { + buffer, + start_index: index, + } + } } #[derive(Default, Format, Debug)] @@ -90,7 +123,7 @@ impl<'a, UART: uart::Instance, PIN: Pin> Client<'a, UART, PIN> { pub(crate) async fn set_set_point( &mut self, - Setting(set_point): Setting, + Setting(set_point): &Setting, ) -> Result<(), TimeoutError> { // Send update through UART to MAX845 to modbus fans // Form message to fan 1 @@ -104,7 +137,7 @@ impl<'a, UART: uart::Instance, PIN: Pin> Client<'a, UART, PIN> { holding_registers::REFERENCE_SET_POINT[1], // Value to set (set_point >> 8) as u8, - set_point as u8, + *set_point as u8, // CRC is set later 0, 0, diff --git a/fan-controller/src/main.rs b/fan-controller/src/main.rs index 267fa56..a537cfa 100644 --- a/fan-controller/src/main.rs +++ b/fan-controller/src/main.rs @@ -37,7 +37,6 @@ use embassy_sync::waitqueue::AtomicWaker; use embassy_time::{with_deadline, with_timeout, Duration, Instant, Ticker, TimeoutError, Timer}; use embedded_io_async::{Read, Write}; use embedded_nal_async::{AddrType, Dns, SocketAddr, TcpConnect}; -use mqtt::client::ConnectError; use mqtt::packet::disconnect::Disconnect; use mqtt::TryDecode; use rand::RngCore; @@ -49,7 +48,6 @@ use {defmt_rtt as _, panic_probe as _}; use self::mqtt::packet; use self::mqtt::packet::Packet; use crate::async_callback::AsyncCallback; -use crate::mqtt::client::runner::State; use crate::mqtt::non_zero_u16; use crate::mqtt::packet::connect::Connect; use crate::mqtt::packet::connect_acknowledgement::ConnectReasonCode; @@ -174,6 +172,121 @@ enum MqttError { UnexpectedPacketType(u8), WritePublishError(publish::EncodeError), } +#[embassy_executor::task] +async fn mqtt_task_client( + spawner: Spawner, + pwr_pin: PIN_23, + cs_pin: PIN_25, + pio: PIO0, + dma: DMA_CH0, + dio: impl PioPin, + clk: impl PioPin, +) -> () { + let (net_device, mut control) = + gain_control(spawner, pwr_pin, cs_pin, pio, dma, dio, clk).await; + + 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 + loop { + match control + .join_wpa2(configuration::WIFI_NETWORK, configuration::WIFI_PASSWORD) + .await + { + Ok(_) => break, + Err(error) => info!("Error joining Wi-Fi network with status: {}", error.status), + } + } + + // Wait for DHCP + info!("Waiting for DHCP"); + while !stack.is_config_up() { + Timer::after_millis(100).await; + } + + info!("DHCP is up"); + + info!("Waiting for link up"); + while !stack.is_link_up() { + Timer::after_millis(500).await; + } + + info!("Link is up"); + + info!("Waiting for stack to be up"); + stack.wait_config_up().await; + info!("Stack is up"); + + // Now we can use it + let mut receive_buffer = [0; 1024]; + let mut send_buffer = [0; 1024]; + let mut socket = TcpSocket::new(stack, &mut receive_buffer, &mut send_buffer); + + let dns_client = DnsSocket::new(stack); + + info!("Resolving MQTT broker IP address"); + // Get home assistant MQTT broker IP address + let address = loop { + //TODO support IPv6 + let result = dns_client + .query(configuration::MQTT_BROKER_ADDRESS, DnsQueryType::A) + .await; + + let mut addresses = match result { + Ok(addresses) => addresses, + Err(error) => { + info!( + "Error resolving Home Assistant MQTT broker IP address with {}: {:?}", + configuration::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!("No addresses found for Home Assistant MQTT broker"); + Timer::after_millis(500).await; + continue; + } + + break addresses.swap_remove(0); + }; + info!("MQTT broker IP address resolved"); + + info!("Connecting to MQTT broker through TCP"); + let endpoint = IpEndpoint::new(address, configuration::MQTT_BROKER_PORT); + // Connect + while let Err(error) = socket.connect(endpoint).await { + info!( + "Error connecting to Home Assistant MQTT broker: {:?}", + error + ); + Timer::after_millis(500).await; + } + + use crate::mqtt::client::Client as MqttClient; + + let (reader, writer) = socket.split(); + MqttClient::connect(reader, writer); +} #[embassy_executor::task] async fn mqtt_task( @@ -384,6 +497,7 @@ async fn mqtt_task( } }; + info!("Setting fan to set point {} from publish", set_point); let Ok(setting) = fan::Setting::new(set_point) else { warn!( "Setting fan speed out of bounds. Not accepting new setting: {}", @@ -398,12 +512,37 @@ async fn mqtt_task( return; }; - if let Err(TimeoutError) = fan.set_set_point(setting).await { + if let Err(TimeoutError) = fan.set_set_point(&setting).await { error!("Setting fan speed from publish timed out"); + + // Turn fan state to off in home assistant as we don't know the current state the fan is in + const PACKET: Message = Message::Publish(Publish { + //TODO use shared constant for topic + topic_name: "testfan/on/state", + payload: b"OFF", + }); + + OUTGOING.send(PACKET).await; } + + // let s : str = core::fmt:: + // Max is "64000" which is 5 characters + let mut buffer = [0; 5]; + + // Update the fan state after successful update + // let packet: Message = Message::Publish(Publish { + // //TODO use shared constant for topic + // topic_name: "testfan/speed/percentage_state", + // // payload: set_point.to_string().as_bytes(), + // payload: publish.payload, + // }); + + info!("Fan speed set. Updating homeassistant"); + let packet = Message::PredefinedPublish(PublishPercentageState { setting }); + OUTGOING.send(packet).await; } - other => info!( + other => warn!( "Unexpected topic: {} with payload: {}", other, publish.payload ), @@ -516,9 +655,14 @@ async fn mqtt_task( // Future 1 let listen = listen(&mut reader); + struct PublishPercentageState { + setting: fan::Setting, + } + enum Message<'a> { Subscribe(Subscribe<'a>), Publish(Publish<'a>), + PredefinedPublish(PublishPercentageState), } static OUTGOING: Channel = Channel::new(); @@ -544,6 +688,19 @@ async fn mqtt_task( continue; } } + Message::PredefinedPublish(PublishPercentageState { setting }) => { + let buffer = setting.to_string_buffer(); + let packet = Publish { + topic_name: "testfan/speed/percentage_state", + payload: buffer.as_bytes(), + }; + + info!("Sending predefined publish"); + if let Err(error) = send(&mut *writer, packet).await { + error!("Error sending predefined publish: {:?}", error); + continue; + } + } } LAST_PACKET.signal(Instant::now()); @@ -812,7 +969,7 @@ async fn input_task(pin_18: PIN_18) { }; info!("Button pressed, setting fan speed to: {}", setting); - if let Err(TimeoutError) = fan.set_set_point(setting).await { + if let Err(TimeoutError) = fan.set_set_point(&setting).await { error!("Timeout setting fan speed from button press"); } } diff --git a/fan-controller/src/mqtt/client.rs b/fan-controller/src/mqtt/client.rs index e7da899..39f3e60 100644 --- a/fan-controller/src/mqtt/client.rs +++ b/fan-controller/src/mqtt/client.rs @@ -1,601 +1,28 @@ -use crate::mqtt::packet::{FromPublish, FromSubscribeAcknowledgement, Packet}; +use embassy_net::tcp::{TcpReader, TcpSocket, TcpWriter}; -use super::packet::{connect, connect_acknowledgement::ConnectReasonCode}; -use crate::mqtt::packet::connect::Connect; -use crate::mqtt::ConnectErrorReasonCode; -use crate::non_zero_u16; -use core::marker::PhantomData; -use embassy_net::tcp::{self, TcpReader, TcpSocket, TcpWriter}; -use embassy_time::{with_timeout, Duration, TimeoutError}; - -pub(crate) struct NotConnected; -pub(crate) struct Connected; - -#[derive(Debug)] -pub(crate) enum ConnectError { - WriteConnectError(connect::EncodeError), - TcpWriteError(tcp::Error), - TcpFlushError(tcp::Error), - ReadError(ReadError), - DecodePacketError(crate::mqtt::packet::ReadError), - /// Received invalid packet type. The server had to respond with a Connect Acknowledgement (CONNACK) packet but send another one instead - InvalidPacket(u8), - /// Received a connect acknowledgement with an error code - ErrorCode(ConnectErrorReasonCode), -} - -#[derive(Debug)] -pub(crate) enum ReadError { - TcpReadError(tcp::Error), - Timeout(TimeoutError), -} - -pub(crate) struct MqttClient<'a, S> { - _state: PhantomData, - socket: TcpSocket<'a>, - // No experience how big this buffer should be - send_buffer: [u8; 256], - receive_buffer: [u8; 1024], - timeout: Duration, -} -impl<'a, S> MqttClient<'a, S> { - async fn read(&mut self) -> Result<&[u8], ReadError> { - let read = self.socket.read(&mut self.receive_buffer); - let bytes_read = with_timeout(self.timeout, read) - .await - .map_err(ReadError::Timeout)? - .map_err(ReadError::TcpReadError)?; - - Ok(&self.receive_buffer[..bytes_read]) - } -} - -impl<'a> MqttClient<'a, NotConnected> { - pub(crate) const fn new(socket: TcpSocket<'a>) -> Self { - Self { - socket, - _state: PhantomData, - send_buffer: [0; 256], - receive_buffer: [0; 1024], - timeout: Duration::from_secs(5), - } - } - - /// The server must not send any data before send - pub(crate) async fn connect( - mut self, - username: &str, - password: &[u8], - keep_alive: Duration, - ) -> Result, (Self, ConnectError)> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - // Send MQTT connect packet - let packet = Connect { - client_identifier: "testfan", - username, - password, - keep_alive_seconds: keep_alive.as_secs() as u16, - }; - - let mut offset = 0; - if let Err(error) = packet.encode(&mut self.send_buffer, &mut offset) { - return Err((self, ConnectError::WriteConnectError(error))); - }; - - if let Err(error) = self.socket.write(&self.send_buffer[..offset]).await { - return Err((self, ConnectError::TcpWriteError(error))); - } - - if let Err(error) = self.socket.flush().await { - return Err((self, ConnectError::TcpFlushError(error))); - } - - // (if connect is accepted) The Server MUST send a CONNACK with a 0x00 (Success) Reason Code before sending any - // Packet other than AUTH [MQTT-3.2.0-1]. - // Wait for connect acknowledgement - let result = self.read().await; - let bytes = match result { - Ok(bytes) => bytes, - Err(error) => return Err((self, ConnectError::ReadError(error))), - }; - - let packet = match Packet::::read(bytes) { - Ok(packet) => packet, - Err(error) => return Err((self, ConnectError::DecodePacketError(error))), - }; - - let Packet::ConnectAcknowledgement(acknowledgement) = packet else { - let r#type = packet.get_type(); - return Err((self, ConnectError::InvalidPacket(r#type))); - }; - - if let ConnectReasonCode::ErrorCode(error_code) = acknowledgement.connect_reason_code { - return Err((self, ConnectError::ErrorCode(error_code))); - } - - Ok(MqttClient:: { - socket: self.socket, - _state: PhantomData, - send_buffer: self.send_buffer, - receive_buffer: self.receive_buffer, - timeout: self.timeout, - }) - } -} - -impl<'a> MqttClient<'a, Connected> { - pub(crate) fn split(self: &mut Self) -> (MqttReceiver<'_>, MqttSender<'_>) { - let (reader, writer) = self.socket.split(); - ( - MqttReceiver { - reader, - receive_buffer: self.receive_buffer, - timeout: self.timeout, - }, - MqttSender { - writer, - send_buffer: self.send_buffer, - }, - ) - } -} - -pub(crate) struct MqttSender<'a> { - send_buffer: [u8; 256], - writer: TcpWriter<'a>, -} - -#[derive(Debug)] -pub(crate) enum ReceiveError { - ReadError(ReadError), - DecodePacketError(crate::mqtt::packet::ReadError), -} - -pub(crate) struct MqttReceiver<'a> { - receive_buffer: [u8; 1024], - reader: TcpReader<'a>, - timeout: Duration, +pub(crate) struct ClientBuilder<'socket> { + socket: TcpSocket<'socket>, } -impl<'a> MqttReceiver<'a> { - async fn read(&mut self) -> Result<&[u8], ReadError> { - let read = self.reader.read(&mut self.receive_buffer); - let bytes_read = with_timeout(self.timeout, read) - .await - .map_err(ReadError::Timeout)? - .map_err(ReadError::TcpReadError)?; - - Ok(&self.receive_buffer[..bytes_read]) - } +// impl<'socket> ClientBuilder<'socket> { +// pub(crate) fn new(socket: TcpSocket<'socket>) -> Self { +// Self { socket } +// } - pub(crate) async fn receive(&mut self) -> Result, ReceiveError> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - let result = self.read().await; - let bytes = match result { - Ok(bytes) => bytes, - Err(error) => return Err(ReceiveError::ReadError(error)), - }; +// pub(crate) fn build(self) -> Client<'socket> { +// self +// } +// } - Packet::read(bytes).map_err(ReceiveError::DecodePacketError) - } +pub(crate) struct Client<'socket> { + reader: TcpReader<'socket>, + writer: TcpWriter<'socket>, } -pub(crate) mod runner { - use core::num::NonZero; - - use crate::mqtt::packet::connect::Connect; - use crate::mqtt::packet::connect_acknowledgement::{ConnectAcknowledgement, ConnectReasonCode}; - use crate::mqtt::packet::subscribe::{Subscribe, Subscription}; - use crate::mqtt::packet::subscribe_acknowledgement::SubscribeErrorReasonCode; - use crate::mqtt::packet::{connect, subscribe}; - use crate::mqtt::packet::{FromPublish, FromSubscribeAcknowledgement, Packet}; - use crate::mqtt::ConnectErrorReasonCode; - use crate::mqtt::{self, TryEncode}; - use defmt::{warn, Format}; - use embassy_futures::select::{select, Either}; - use embassy_net::tcp; - use embassy_net::tcp::{TcpReader, TcpSocket, TcpWriter}; - use embassy_sync::blocking_mutex::raw::CriticalSectionRawMutex; - use embassy_sync::channel::{Channel, Receiver, Sender}; - use embassy_time::{with_deadline, with_timeout, Duration, Instant, TimeoutError}; - - mod message { - use core::num::NonZeroU16; - - use crate::mqtt::packet::connect_acknowledgement::ConnectReasonCode; - use crate::mqtt::packet::subscribe::Subscription; - use crate::mqtt::packet::{FromPublish, FromSubscribeAcknowledgement}; - - pub(super) enum Outgoing<'a> { - Connect { - client_identifier: &'a str, - username: &'a str, - password: &'a [u8], - }, - Subscribe { - subscriptions: &'a [Subscription<'a>], - packet_identifier: NonZeroU16, - }, - } - - pub(super) enum Incoming - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - ConnectAcknowledgement(ConnectReasonCode), - SubscribeAcknowledgement(S), - Publish(T), - } - } - - enum Response { - ConnectAcknowledgement { - connect_reason_code: ConnectReasonCode, - }, - } - - pub(crate) struct State<'ch, T, S> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - incoming: Channel, 8>, - outgoing: Channel, 8>, - } - - impl<'ch, T, S> State<'ch, T, S> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - pub(crate) const fn new() -> Self { - Self { - incoming: Channel::new(), - outgoing: Channel::new(), - } - } - } - - #[derive(Clone, Debug, Format)] - enum ReceiveError { - ReadError(ReadError), - DecodePacketError(mqtt::packet::ReadError), - } - - struct MqttRunner<'socket, 'state, T, S> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - tcp_receiver: TcpReader<'socket>, - /// Subscribe to messages to be sent from the client (like an actor handle) - outgoing: Receiver<'state, CriticalSectionRawMutex, message::Outgoing<'socket>, 8>, - incoming: Sender<'state, CriticalSectionRawMutex, message::Incoming, 8>, - receive_buffer: [u8; 1024], - send_buffer: [u8; 256], - timeout: Duration, - tcp_writer: TcpWriter<'socket>, - } - - impl<'socket, 'state, T, S> MqttRunner<'socket, 'state, T, S> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - pub(self) async fn read(&mut self) -> Result<&[u8], ReadError> { - let read = self.tcp_receiver.read(&mut self.receive_buffer); - let bytes_read = with_timeout(self.timeout, read) - .await - .map_err(ReadError::Timeout)? - .map_err(ReadError::TcpReadError)?; - - Ok(&self.receive_buffer[..bytes_read]) - } - - pub(crate) async fn run(&'socket mut self) -> ! { - loop { - let result = select( - self.outgoing.receive(), - self.tcp_receiver.read(&mut self.receive_buffer), - ) - .await; - match result { - Either::First(message) => { - match message { - message::Outgoing::Connect { - client_identifier, - username, - password, - } => { - let packet = Connect { - client_identifier, - username, - password, - //TODO validate seconds < u16::MAX - keep_alive_seconds: self.timeout.as_secs() as u16, - }; - - let mut offset = 0; - if let Err(error) = - packet.encode(&mut self.send_buffer, &mut offset) - { - //TODO handle error - warn!("Error encoding connect packet: {:?}", error); - continue; - }; - - if let Err(error) = - self.tcp_writer.write(&self.send_buffer[..offset]).await - { - //TODO handle error - warn!("Error writing connect packet: {:?}", error); - continue; - } - - if let Err(error) = self.tcp_writer.flush().await { - //TODO handle error - warn!("Error flushing connect packet: {:?}", error); - continue; - } - } - message::Outgoing::Subscribe { - subscriptions, - packet_identifier, - } => { - //TODO packet identifier - let packet = Subscribe { - subscriptions, - packet_identifier, - }; - - let mut offset = 0; - if let Err(error) = - packet.try_encode(&mut self.send_buffer, &mut offset) - { - //TODO handle error - warn!("Error encoding subscribe packet: {:?}", error); - continue; - }; - - if let Err(error) = - self.tcp_writer.write(&self.send_buffer[..offset]).await - { - //TODO handle error - warn!("Error writing subscribe packet: {:?}", error); - continue; - } - - if let Err(error) = self.tcp_writer.flush().await { - //TODO handle error - warn!("Error flushing subscribe packet: {:?}", error); - continue; - } - } - } - } - Either::Second(result) => { - let bytes_read = match result { - Ok(bytes_read) => bytes_read, - Err(error) => { - //TODO handle error - warn!("Error reading from TCP: {:?}", error); - continue; - } - }; - - let result = Packet::::read(&self.receive_buffer[..bytes_read]); - let packet = match result { - Ok(packet) => packet, - Err(error) => { - //TODO handle error - warn!("Error decoding packet: {:?}", error); - continue; - } - }; - - match packet { - Packet::ConnectAcknowledgement(ConnectAcknowledgement { - connect_reason_code, - is_session_present: _, - }) => { - let message = - message::Incoming::ConnectAcknowledgement(connect_reason_code); - self.incoming.send(message).await; - continue; - } - - //TODO - Packet::SubscribeAcknowledgement(acknowledgement) => { - let message = - message::Incoming::SubscribeAcknowledgement(acknowledgement); - self.incoming.send(message).await; - continue; - } - Packet::Publish(_) => {} - } - } - } - } - } - } - - #[derive(Debug, Clone, Format)] - pub(crate) enum ReadError { - TcpReadError(tcp::Error), - Timeout(TimeoutError), - } - - #[derive(Debug)] - pub(crate) enum ConnectError { - WriteConnectError(connect::EncodeError), - TcpWriteError(tcp::Error), - TcpFlushError(tcp::Error), - ReceiveError(ReceiveError), - /// Received a connect acknowledgement with an error code - ErrorCode(ConnectErrorReasonCode), - Timeout(TimeoutError), - } - - pub(crate) enum SubscribeError { - WriteSubscribeError(subscribe::EncodeError), - TcpWriteError(tcp::Error), - TcpFlushError(tcp::Error), - Timeout(TimeoutError), - } - - pub(crate) struct MqttClient<'a, T, S> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - { - timeout: Duration, - /// Publisher for sending messages to the runner to execute - outgoing: Sender<'a, CriticalSectionRawMutex, message::Outgoing<'a>, 8>, - incoming: Receiver<'a, CriticalSectionRawMutex, message::Incoming, 8>, - } - - impl<'state, 'socket, T, S> MqttClient<'state, T, S> - where - T: FromPublish, - S: FromSubscribeAcknowledgement, - 'state: 'socket, - { - pub(crate) fn new( - state: &'state State<'state, T, S>, - mut socket: &'socket mut TcpSocket<'socket>, - ) -> (MqttRunner<'socket, 'state, T, S>, Self) { - let (reader, writer) = socket.split(); - //TODO handle out of subscribers/publishers error - let send_incoming = state.incoming.sender(); - let receive_incoming = state.incoming.receiver(); - let send_outgoing = state.outgoing.sender(); - let receive_outgoing = state.outgoing.receiver(); - - let timeout = Duration::from_secs(5); - - ( - MqttRunner { - incoming: send_incoming, - outgoing: receive_outgoing, - tcp_receiver: reader, - receive_buffer: [0; 1024], - send_buffer: [0; 256], - timeout: timeout.clone(), - tcp_writer: writer, - }, - Self { - timeout, - incoming: receive_incoming, - outgoing: send_outgoing, - }, - ) - } - - pub async fn connect( - &self, - client_identifier: &'state str, - username: &'state str, - password: &'state [u8], - ) -> Result<(), ConnectError> { - let message = message::Outgoing::Connect { - client_identifier, - username, - password, - }; - - self.outgoing.send(message).await; - - // Wait for connect acknowledgement - // Discard all messages before the connect acknowledgement - // The server has to send a connect acknowledgement before sending any other packet - // TCP should ensure the order of packets (afaik), so they should not arrive out of order - // Using a deadline because the loop could be run multiple times - let deadline = Instant::now() + self.timeout; - loop { - let result = with_deadline(deadline, self.incoming.receive()).await; - let message = match result { - Ok(message) => message, - Err(error) => { - //TODO handle error - warn!("Timed out receiving connect acknowledgement: {:?}", error); - return Err(ConnectError::Timeout(error)); - } - }; - match message { - message::Incoming::ConnectAcknowledgement(reason_code) => { - return match reason_code { - ConnectReasonCode::Success => Ok(()), - ConnectReasonCode::ErrorCode(error_code) => { - Err(ConnectError::ErrorCode(error_code)) - } - } - } - _other => continue, - }; - } - } - - pub async fn subscribe( - &self, - subscriptions: &'state [Subscription<'state>; 2], - ) -> Result<[Result<(), SubscribeErrorReasonCode>; 2], SubscribeError> { - //TODO manage packet identifier - let packet_identifier: NonZero = crate::mqtt::non_zero_u16!(42); - let message = message::Outgoing::Subscribe { - subscriptions, - packet_identifier, - }; - self.outgoing.send(message).await; - - let deadline = Instant::now() + self.timeout; - loop { - todo!(); - // let result = with_deadline(deadline, self.incoming.receive()).await; - // let message = match result { - // Ok(message) => message, - // Err(error) => { - // //TODO handle error - // warn!("Timed out receiving subscribe acknowledgement: {:?}", error); - // return Err(SubscribeError::Timeout(error)); - // } - // }; - // let (identifier, reason_codes) = match message { - // message::Incoming::SubscribeAcknowledgement(SubscribeAcknowledgement { - // packet_identifier, - // reason_codes, - // }) => (packet_identifier, reason_codes), - // _other => continue, - // }; - // - // if packet_identifier != identifier { - // continue; - // } +impl<'socket> Client<'socket> { + pub(crate) async fn connect(reader: TcpReader<'socket>, writer: TcpWriter<'socket>) -> Self { + // let (reader, writer) = socket.split(); - //TODO don't assume we have 2 results - // let mut results = [Ok(()), Ok(())]; - // - // // Write codes to results - // for (index, code) in reason_codes.into_iter().enumerate() { - // if index >= results.len() { - // break; - // } - // - // let Some(SubscribeReasonCode::ErrorCode(code)) = code else { - // continue; - // }; - // - // results[index] = Err(code); - // } - // - // return Ok(results); - } - } + Self { reader, writer } } } diff --git a/fan-controller/src/mqtt/packet/publish.rs b/fan-controller/src/mqtt/packet/publish.rs index 4a45bc9..c318dab 100644 --- a/fan-controller/src/mqtt/packet/publish.rs +++ b/fan-controller/src/mqtt/packet/publish.rs @@ -101,7 +101,6 @@ impl TryEncode for Publish<'_> { buffer[*offset] = 0; *offset += 1; - info!("Encode 2"); // Payload // No need to set length as it will be calculated for byte in self.payload { @@ -160,8 +159,6 @@ impl<'a> TryDecode<'a> for Publish<'a> { // Payload //TODO validate there is enough space left in the buffer let payload = &variable_header_and_payload[offset..]; - - debug!("8"); Ok(Publish { topic_name, payload,