From 82969cbf8fd44203f32f4bc98551430d4832fe8c Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Fri, 14 Mar 2025 13:00:36 -0400 Subject: [PATCH 1/8] fixes, debugging and other changes --- Cargo.lock | 11 +- Cargo.toml | 3 +- src/config.rs | 68 +++++- src/coordinator/mod.rs | 444 +++++++++++++++++++++++++++++--------- src/database.rs | 55 +++-- src/influx.rs | 41 +++- src/lib.rs | 103 +++++++-- src/lxp/inverter.rs | 218 +++++++++++++------ src/lxp/packet.rs | 407 ++++++++++++++++++++++------------ src/lxp/packet_decoder.rs | 107 +++++++-- src/mqtt.rs | 1 + src/utils.rs | 20 ++ 12 files changed, 1109 insertions(+), 369 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index de65520..bc3cb66 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1,6 +1,6 @@ # This file is automatically @generated by Cargo. # It is not intended for manual editing. -version = 3 +version = 4 [[package]] name = "addr2line" @@ -958,9 +958,9 @@ checksum = "b9e0384b61958566e926dc50660321d12159025e767c18e043daf26b70104c39" [[package]] name = "idna" -version = "0.5.0" +version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "634d9b1461af396cad843f47fdba5597a4f9e6ddd4bfb6ff5d85028c25cb12f6" +checksum = "7d20d6b07bfbc108882d88ed8e37d39636dcc260e15e30c45e6ba089610b917c" dependencies = [ "unicode-bidi", "unicode-normalization", @@ -1130,6 +1130,7 @@ dependencies = [ "sqlx", "tokio", "tokio-util", + "url", ] [[package]] @@ -2549,9 +2550,9 @@ checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" [[package]] name = "url" -version = "2.5.0" +version = "2.4.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "31e6302e3bb753d46e83516cae55ae196fc0c309407cf11ab35cc51a4c2a4633" +checksum = "143b538f18257fac9cad154828a57c6bf5157e1aa604d4816b5995bf6de87ae5" dependencies = [ "form_urlencoded", "idna", diff --git a/Cargo.toml b/Cargo.toml index fb03706..31bce76 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -38,10 +38,11 @@ serde_json = "~1" serde_yaml = "~0.9" tokio = { version = "~1", features = ["net", "macros", "signal"] } tokio-util = { version = "~0.7", features = ["codec"] } -chrono = "~0.4" +chrono = { version = "~0.4", features = ["serde"] } cron-parser = "~0.7" enum_dispatch = "~0.3" async-trait = "~0.1" reqwest = "~0.11" rinfluxdb = { version = "~0.1", git = "https://gitlab.com/celsworth/rinfluxdb.git", rev = "f3f5b23e" } sqlx = { version = "~0.6", features = ["runtime-tokio-native-tls", "any", "postgres", "mysql", "sqlite", "chrono"] } +url = "~2.4" diff --git a/src/config.rs b/src/config.rs index 98289b7..95eb31a 100644 --- a/src/config.rs +++ b/src/config.rs @@ -1,15 +1,14 @@ use crate::prelude::*; use serde::Deserialize; -use serde_with::serde_as; //, OneOrMany; +use serde_with::serde_as; +use serde_yaml; #[serde_as] #[derive(Clone, Debug, Deserialize)] pub struct Config { pub inverters: Vec, - //#[serde_as(deserialize_as = "OneOrMany<_>")] pub mqtt: Mqtt, - //#[serde_as(deserialize_as = "OneOrMany<_>")] pub influx: Influx, #[serde(default = "Vec::new")] pub databases: Vec, @@ -336,7 +335,68 @@ impl Config { let content = std::fs::read_to_string(&file) .map_err(|err| anyhow!("error reading {}: {}", file, err))?; - Ok(serde_yaml::from_str(&content)?) + let config: Self = serde_yaml::from_str(&content)?; + config.validate()?; + Ok(config) + } + + fn validate(&self) -> Result<()> { + // Validate MQTT configuration + if self.mqtt.enabled { + if self.mqtt.port == 0 { + bail!("mqtt.port must be between 1 and 65535"); + } + if self.mqtt.host.is_empty() { + return Err(anyhow!("MQTT host cannot be empty")); + } + } + + // Validate InfluxDB configuration + if self.influx.enabled { + if let Err(e) = url::Url::parse(&self.influx.url) { + return Err(anyhow!("Invalid InfluxDB URL: {}", e)); + } + if self.influx.database.is_empty() { + return Err(anyhow!("InfluxDB database name cannot be empty")); + } + } + + // Validate database URLs + for db in &self.databases { + if db.enabled { + if let Err(e) = url::Url::parse(db.url()) { + return Err(anyhow!("Invalid database URL: {}", e)); + } + } + } + + // Validate inverter configurations + for (i, inv) in self.inverters.iter().enumerate() { + if inv.enabled { + if inv.port == 0 { + bail!("inverter[{}].port must be between 1 and 65535", i); + } + if inv.host.is_empty() { + return Err(anyhow!("Inverter host cannot be empty")); + } + if inv.read_timeout.unwrap_or(900) == 0 { + return Err(anyhow!("Invalid read timeout: 0")); + } + } + } + + // Validate scheduler configuration + if let Some(scheduler) = &self.scheduler { + if scheduler.enabled { + if let Some(cron) = &scheduler.timesync_cron { + if cron.is_empty() { + return Err(anyhow!("Scheduler cron expression cannot be empty")); + } + } + } + } + + Ok(()) } fn default_mqtt_port() -> u16 { diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index 3f1d9c0..d2f1be5 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -2,7 +2,9 @@ use crate::prelude::*; pub mod commands; -use lxp::packet::{DeviceFunction, TcpFunction}; +use std::sync::{Arc, Mutex}; +use lxp::packet::{DeviceFunction, ReadInput, TcpFunction}; +use serde_json::json; #[derive(Eq, PartialEq, Debug, Clone)] pub enum ChannelData { @@ -11,29 +13,76 @@ pub enum ChannelData { pub type InputsStore = std::collections::HashMap; +#[derive(Default)] +pub struct PacketStats { + packets_received: u64, + packets_sent: u64, + mqtt_messages_sent: u64, + mqtt_errors: u64, + influx_writes: u64, + influx_errors: u64, + database_writes: u64, + database_errors: u64, + register_cache_writes: u64, + register_cache_errors: u64, +} + +impl PacketStats { + pub fn print_summary(&self) { + info!("Packet Statistics:"); + info!(" Total packets received: {}", self.packets_received); + info!(" Total packets sent: {}", self.packets_sent); + info!(" MQTT:"); + info!(" Messages sent: {}", self.mqtt_messages_sent); + info!(" Errors: {}", self.mqtt_errors); + info!(" InfluxDB:"); + info!(" Writes: {}", self.influx_writes); + info!(" Errors: {}", self.influx_errors); + info!(" Database:"); + info!(" Writes: {}", self.database_writes); + info!(" Errors: {}", self.database_errors); + info!(" Register Cache:"); + info!(" Writes: {}", self.register_cache_writes); + info!(" Errors: {}", self.register_cache_errors); + } +} + +#[derive(Clone)] pub struct Coordinator { config: ConfigWrapper, channels: Channels, + pub stats: Arc>, } impl Coordinator { pub fn new(config: ConfigWrapper, channels: Channels) -> Self { - Self { config, channels } + Self { + config, + channels, + stats: Arc::new(Mutex::new(PacketStats::default())), + } } pub async fn start(&self) -> Result<()> { - futures::try_join!(self.inverter_receiver(), self.mqtt_receiver())?; + if self.config.mqtt().enabled() { + futures::try_join!(self.inverter_receiver(), self.mqtt_receiver())?; + } else { + self.inverter_receiver().await?; + } Ok(()) } pub fn stop(&self) { + // Send shutdown signals to channels let _ = self .channels .from_inverter .send(lxp::inverter::ChannelData::Shutdown); - let _ = self.channels.from_mqtt.send(mqtt::ChannelData::Shutdown); + if self.config.mqtt().enabled() { + let _ = self.channels.from_mqtt.send(mqtt::ChannelData::Shutdown); + } } async fn mqtt_receiver(&self) -> Result<()> { @@ -50,18 +99,20 @@ impl Coordinator { for inverter in self.config.inverters_for_message(&message)? { match message.to_command(inverter) { Ok(command) => { - debug!("parsed command {:?}", command); + info!("parsed command {:?}", command); let topic_reply = command.to_result_topic(); let result = self.process_command(command).await; - let reply = mqtt::ChannelData::Message(mqtt::Message { - topic: topic_reply, - retain: false, - payload: if result.is_ok() { "OK" } else { "FAIL" }.to_string(), - }); - if self.channels.to_mqtt.send(reply).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); + if self.config.mqtt().enabled() { + let reply = mqtt::ChannelData::Message(mqtt::Message { + topic: topic_reply, + retain: false, + payload: if result.is_ok() { "OK" } else { "FAIL" }.to_string(), + }); + if self.channels.to_mqtt.send(reply).is_err() { + bail!("send(to_mqtt) failed - channel closed?"); + } } } Err(err) => { @@ -73,11 +124,19 @@ impl Coordinator { Ok(()) } + fn increment_packets_sent(&self) { + if let Ok(mut stats) = self.stats.lock() { + stats.packets_sent += 1; + } + } + async fn process_command(&self, command: Command) -> Result<()> { use commands::time_register_ops::Action; use lxp::packet::{Register, RegisterBit}; use Command::*; + self.increment_packets_sent(); + match command { ReadInputs(inverter, 1) => self.read_inputs(inverter, 0_u16, 40).await, ReadInputs(inverter, 2) => self.read_inputs(inverter, 40_u16, 40).await, @@ -330,9 +389,20 @@ impl Coordinator { error!("{}", e); } } - // this loop holds no state so doesn't care about inverter disconnects - Disconnect(_) => {} - Shutdown => break, + // Print statistics when an inverter disconnects + Disconnect(serial) => { + info!("Inverter {} disconnected, printing statistics:", serial); + if let Ok(stats) = self.stats.lock() { + stats.print_summary(); + } + } + Shutdown => { + info!("Received shutdown signal, printing final statistics:"); + if let Ok(stats) = self.stats.lock() { + stats.print_summary(); + } + break; + } } } @@ -346,7 +416,28 @@ impl Coordinator { ) -> Result<()> { debug!("RX: {:?}", packet); + // Update packet stats first + if let Ok(mut stats) = self.stats.lock() { + stats.packets_received += 1; + } + if let Packet::TranslatedData(td) = &packet { + // Check if the inverter serial from packet matches any configured inverter + let packet_serial = td.inverter; + let packet_datalog = td.datalog; + + // Log if we see a mismatch between configured and actual serial + if let Some(inverter) = self.config.enabled_inverter_with_datalog(packet_datalog) { + if inverter.serial() != packet_serial { + warn!( + "Inverter serial mismatch - Config: {}, Actual: {} for datalog: {}", + inverter.serial(), + packet_serial, + packet_datalog + ); + } + } + // temporary special greppable logging for Param packets as I try to // work out what they do :) if td.tcp_function() == TcpFunction::ReadParam @@ -354,73 +445,224 @@ impl Coordinator { { warn!("got a Param packet! {:?}", td); } + } - // inputs_store handling. If we've received any ReadInput, update inputs_store - // with the contents. If we got the third (of three) packets, send out the combined - // MQTT message with all the data. - if td.device_function == DeviceFunction::ReadInput { - use lxp::packet::{ReadInput, ReadInputs}; - + // inputs_store handling. If we've received any ReadInput, update inputs_store + // with the contents. If we got the third (of three) packets, send out the combined + // MQTT message with all the data. + match packet { + Packet::Heartbeat(_) => Ok(()), // nothing to do + Packet::TranslatedData(td) => { let entry = inputs_store .entry(td.datalog) - .or_insert_with(ReadInputs::default); + .or_insert_with(lxp::packet::ReadInputs::default); - match td.read_input() { - Ok(ReadInput::ReadInputAll(r_all)) => { - // no need for MQTT here, done below - self.save_input_all(r_all).await? - } + match td.device_function { + DeviceFunction::ReadInput => match td.read_input() { + Ok(ReadInput::ReadInputAll(r_all)) => { + debug!("Received ReadInputAll, saving to InfluxDB"); + self.save_input_all(r_all).await?; + } + Ok(ReadInput::ReadInput1(r1)) => { + debug!("Received ReadInput1"); + entry.set_read_input_1(r1) + } + Ok(ReadInput::ReadInput2(r2)) => { + debug!("Received ReadInput2"); + entry.set_read_input_2(r2) + } + Ok(ReadInput::ReadInput3(r3)) => { + debug!("Received ReadInput3"); + entry.set_read_input_3(r3) + } + Ok(ReadInput::ReadInput4(r4)) => { + debug!("Received ReadInput4"); + let datalog = r4.datalog; + entry.set_read_input_4(r4); + + if let Some(input) = entry.to_input_all() { + info!("Assembled complete input set, saving to InfluxDB"); + if self.config.mqtt().enabled() { + match mqtt::Message::for_input_all(&input, datalog) { + Ok(message) => { + let channel_data = mqtt::ChannelData::Message(message); + match self.channels.to_mqtt.send(channel_data) { + Ok(_) => { + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_messages_sent += 1; + } + } + Err(e) => { + error!("Failed to send MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } + } + } + } + Err(e) => { + error!("Failed to create MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } + } + } + } - Ok(ReadInput::ReadInput1(r1)) => entry.set_read_input_1(r1), - Ok(ReadInput::ReadInput2(r2)) => entry.set_read_input_2(r2), - Ok(ReadInput::ReadInput3(r3)) => entry.set_read_input_3(r3), - Ok(ReadInput::ReadInput4(r4)) => { - let datalog = r4.datalog; + self.save_input_all(Box::new(input)).await?; + } else { + debug!("Incomplete input set, waiting for more data"); + } + } + Err(x) => warn!("ignoring {:?}", x), + }, + DeviceFunction::ReadHold | DeviceFunction::WriteSingle => { + let channel_data = register_cache::ChannelData::RegisterData(td.register, td.value()); + match self.channels.to_register_cache.send(channel_data) { + Ok(_) => { + if let Ok(mut stats) = self.stats.lock() { + stats.register_cache_writes += 1; + } + } + Err(e) => { + error!("Failed to send to register cache: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.register_cache_errors += 1; + } + } + } - entry.set_read_input_4(r4); + // Send to InfluxDB if enabled + if self.config.influx().enabled() { + debug!("InfluxDB is enabled, sending ReadHold data"); + let mut json_data = serde_json::Map::new(); + json_data.insert("time".to_string(), json!(chrono::Utc::now().timestamp())); + json_data.insert("datalog".to_string(), json!(td.datalog.to_string())); + json_data.insert(format!("hold_{}", td.register), json!(td.value())); + + let json = serde_json::Value::Object(json_data); + match self.channels.to_influx.send(influx::ChannelData::InputData(json)) { + Ok(_) => { + debug!("Successfully sent ReadHold data to InfluxDB"); + if let Ok(mut stats) = self.stats.lock() { + stats.influx_writes += 1; + } + } + Err(e) => { + error!("Failed to send ReadHold data to InfluxDB: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.influx_errors += 1; + } + } + } + } + } + _ => {} + } - if let Some(input) = entry.to_input_all() { - if self.config.mqtt().enabled() { - let message = mqtt::Message::for_input_all(&input, datalog)?; + if self.config.mqtt().enabled() { + match Self::packet_to_messages(Packet::TranslatedData(td), self.config.mqtt().publish_individual_input()) { + Ok(messages) => { + for message in messages { let channel_data = mqtt::ChannelData::Message(message); - if self.channels.to_mqtt.send(channel_data).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); + match self.channels.to_mqtt.send(channel_data) { + Ok(_) => { + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_messages_sent += 1; + } + } + Err(e) => { + error!("Failed to send MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } + } } } + } + Err(e) => { + error!("Failed to create MQTT messages: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } + } + } + } - self.save_input_all(Box::new(input)).await?; + Ok(()) + } + Packet::ReadParam(rp) => { + if self.config.mqtt().enabled() { + match mqtt::Message::for_param(rp) { + Ok(messages) => { + for message in messages { + let channel_data = mqtt::ChannelData::Message(message); + match self.channels.to_mqtt.send(channel_data) { + Ok(_) => { + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_messages_sent += 1; + } + } + Err(e) => { + error!("Failed to send MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } + } + } + } + } + Err(e) => { + error!("Failed to create MQTT messages: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } } } - Err(x) => warn!("ignoring {:?}", x), } - } else if td.device_function == DeviceFunction::ReadHold - || td.device_function == DeviceFunction::WriteSingle - { - let channel_data = - register_cache::ChannelData::RegisterData(td.register, td.value()); - if self.channels.to_register_cache.send(channel_data).is_err() { - bail!("send(to_register_cache) failed - channel closed?"); + + Ok(()) + } + Packet::WriteParam(_) => Ok(()), // nothing to do + } + } + + async fn save_input_all(&self, input: Box) -> Result<()> { + if self.config.influx().enabled() { + debug!("InfluxDB is enabled, attempting to save data"); + let json = serde_json::to_value(&input)?; + debug!("Serialized data for InfluxDB: {:?}", json); + match self.channels.to_influx.send(influx::ChannelData::InputData(json)) { + Ok(_) => { + debug!("Successfully sent data to InfluxDB channel"); + if let Ok(mut stats) = self.stats.lock() { + stats.influx_writes += 1; + } + } + Err(e) => { + error!("Failed to send data to InfluxDB: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.influx_errors += 1; + } } } } - if self.config.mqtt().enabled() { - // returns a Vec of messages to send. could be none; - // not every packet produces an MQ message (eg, heartbeats), - // and some produce >1 (multi-register ReadHold) - match Self::packet_to_messages(packet, self.config.mqtt().publish_individual_input()) { - Ok(messages) => { - for message in messages { - let message = mqtt::ChannelData::Message(message); - if self.channels.to_mqtt.send(message).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); - } + if self.config.have_enabled_database() { + let channel_data = database::ChannelData::ReadInputAll(input); + match self.channels.to_database.send(channel_data) { + Ok(_) => { + if let Ok(mut stats) = self.stats.lock() { + stats.database_writes += 1; } } Err(e) => { // log error but avoid exiting loop as then we stop handling // incoming packets. need better error handling here maybe? - error!("{}", e); + error!("Failed to send to database: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.database_errors += 1; + } } } } @@ -428,14 +670,30 @@ impl Coordinator { Ok(()) } - // Unlike input registers, holding registers are not broadcast by inverters, - // but they are interesting nevertheless. Publishing the holding registers - // when we connect to an inverter makes it easy for configuration data to be - // tracked, which is particularly useful in conjunction with HomeAssistant. + fn packet_to_messages( + packet: Packet, + publish_individual_input: bool, + ) -> Result> { + match packet { + Packet::Heartbeat(_) => Ok(Vec::new()), // always no message + Packet::TranslatedData(td) => match td.device_function { + DeviceFunction::ReadHold => mqtt::Message::for_hold(td), + DeviceFunction::ReadInput => mqtt::Message::for_input(td, publish_individual_input), + DeviceFunction::WriteSingle => mqtt::Message::for_hold(td), + DeviceFunction::WriteMulti => Ok(Vec::new()), // TODO, for_hold might just work + }, + Packet::ReadParam(rp) => mqtt::Message::for_param(rp), + Packet::WriteParam(_) => Ok(Vec::new()), // ignoring for now + } + } + async fn inverter_connected(&self, datalog: Serial) -> Result<()> { let inverter = match self.config.enabled_inverter_with_datalog(datalog) { Some(inverter) => inverter, - None => bail!("Unknown inverter connected: {}", datalog), + None => { + warn!("Unknown inverter datalog connected: {}, will continue processing its data", datalog); + return Ok(()); + } }; if !inverter.publish_holdings_on_connect() { @@ -444,13 +702,32 @@ impl Coordinator { info!("Reading holding registers for inverter {}", datalog); + // Add delay between read_hold requests to prevent overwhelming the inverter + const DELAY_MS: u64 = 100; // 100ms delay between requests + // We can only read holding registers in blocks of 40. Provisionally, // there are 6 pages of 40 values. + self.increment_packets_sent(); self.read_hold(inverter.clone(), 0_u16, 40).await?; + tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; + + self.increment_packets_sent(); self.read_hold(inverter.clone(), 40_u16, 40).await?; + tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; + + self.increment_packets_sent(); self.read_hold(inverter.clone(), 80_u16, 40).await?; + tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; + + self.increment_packets_sent(); self.read_hold(inverter.clone(), 120_u16, 40).await?; + tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; + + self.increment_packets_sent(); self.read_hold(inverter.clone(), 160_u16, 40).await?; + tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; + + self.increment_packets_sent(); self.read_hold(inverter.clone(), 200_u16, 40).await?; // Also send any special interpretive topics which are derived from @@ -459,21 +736,25 @@ impl Coordinator { // FIXME: this is a further 12 round-trips to the inverter to read values // we have already taken, just above. We should be able to do better! for num in &[1, 2, 3] { + self.increment_packets_sent(); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::AcCharge(*num), ) .await?; + self.increment_packets_sent(); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::ChargePriority(*num), ) .await?; + self.increment_packets_sent(); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::ForcedDischarge(*num), ) .await?; + self.increment_packets_sent(); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::AcFirst(*num), @@ -483,39 +764,4 @@ impl Coordinator { Ok(()) } - - async fn save_input_all(&self, input: Box) -> Result<()> { - if self.config.influx().enabled() { - let channel_data = influx::ChannelData::InputData(serde_json::to_value(&input)?); - if self.channels.to_influx.send(channel_data).is_err() { - bail!("send(to_influx) failed - channel closed?"); - } - } - - if self.config.have_enabled_database() { - let channel_data = database::ChannelData::ReadInputAll(input); - if self.channels.to_database.send(channel_data).is_err() { - bail!("send(to_database) failed - channel closed?"); - } - } - - Ok(()) - } - - fn packet_to_messages( - packet: Packet, - publish_individual_input: bool, - ) -> Result> { - match packet { - Packet::Heartbeat(_) => Ok(Vec::new()), // always no message - Packet::TranslatedData(td) => match td.device_function { - DeviceFunction::ReadHold => mqtt::Message::for_hold(td), - DeviceFunction::ReadInput => mqtt::Message::for_input(td, publish_individual_input), - DeviceFunction::WriteSingle => mqtt::Message::for_hold(td), - DeviceFunction::WriteMulti => Ok(Vec::new()), // TODO, for_hold might just work - }, - Packet::ReadParam(rp) => mqtt::Message::for_param(rp), - Packet::WriteParam(_) => Ok(Vec::new()), // ignoring for now - } - } } diff --git a/src/database.rs b/src/database.rs index 0fd286f..03d044c 100644 --- a/src/database.rs +++ b/src/database.rs @@ -1,11 +1,11 @@ use crate::prelude::*; -use sqlx::{any::AnyConnectOptions, AnyPool, ConnectOptions}; +use sqlx::{any::AnyConnectOptions, AnyPool}; -#[derive(PartialEq, Clone, Debug)] +#[derive(Debug, Clone)] pub enum ChannelData { - ReadInputAll(Box), Shutdown, + ReadInputAll(Box), } pub type Sender = broadcast::Sender; @@ -60,24 +60,25 @@ impl Database { } async fn connect(&self) -> Result<()> { - let mut options = AnyConnectOptions::from_str(self.config.url())?; - options.disable_statement_logging(); - let pool = sqlx::any::AnyPool::connect_with(options).await?; + let options = AnyConnectOptions::from_str(self.config.url())?; - //debug!("{:?}", pool.any_kind()); - let _ = self.pool.borrow_mut().insert(pool); + let pool = sqlx::any::AnyPoolOptions::new() + .max_connections(5) + .min_connections(1) + .acquire_timeout(std::time::Duration::from_secs(30)) + .connect_with(options) + .await?; + + *self.pool.borrow_mut() = Some(pool); Ok(()) } pub async fn connection(&self) -> Result> { - let acquire = if let Some(pool) = &*self.pool.borrow() { - pool.acquire() - } else { - todo!() - }; - - Ok(acquire.await?) + match &*self.pool.borrow() { + Some(pool) => Ok(pool.acquire().await?), + None => Err(anyhow!("Database pool not initialized")) + } } async fn migrate(&self) -> Result<()> { @@ -99,13 +100,10 @@ impl Database { async fn inserter(&self) -> Result<()> { self.connect().await?; - info!("database connected"); - self.migrate().await?; let mut receiver = self.channels.to_database.subscribe(); - let values = match self.database()? { DatabaseType::MySQL => Self::values_for_mysql(), _ => Self::values_for_not_mysql(), @@ -159,9 +157,24 @@ impl Database { match receiver.recv().await? { Shutdown => break, ReadInputAll(data) => { - while let Err(err) = self.insert(&query, &data).await { - error!("INSERT failed: {:?} - retrying in 10s", err); - tokio::time::sleep(std::time::Duration::from_secs(10)).await; + let mut retry_count = 0; + let max_retries = 3; + let mut backoff = 1; + + while retry_count < max_retries { + match self.insert(&query, &data).await { + Ok(_) => break, + Err(err) => { + error!("INSERT failed: {:?} - retrying in {}s", err, backoff); + tokio::time::sleep(std::time::Duration::from_secs(backoff)).await; + retry_count += 1; + backoff *= 2; + } + } + } + + if retry_count == max_retries { + error!("Failed to insert data after {} retries", max_retries); } } } diff --git a/src/influx.rs b/src/influx.rs index a06313f..3e0269e 100644 --- a/src/influx.rs +++ b/src/influx.rs @@ -11,6 +11,7 @@ pub enum ChannelData { Shutdown, } +#[derive(Clone)] pub struct Influx { config: ConfigWrapper, channels: Channels, @@ -55,22 +56,30 @@ impl Influx { use ChannelData::*; let mut receiver = self.channels.to_influx.subscribe(); + info!("InfluxDB sender started"); loop { let mut line = LineBuilder::new(INPUTS_MEASUREMENT); match receiver.recv().await? { - Shutdown => break, + Shutdown => { + info!("InfluxDB sender received shutdown signal"); + break; + } InputData(data) => { - for (key, value) in data.as_object().unwrap() { + debug!("InfluxDB processing input data: {:?}", data); + for (key, value) in data.as_object().ok_or_else(|| anyhow!("Invalid data format"))? { let key = key.to_string(); + debug!("Processing field: {} = {:?}", key, value); line = if key == "time" { let value = value.as_i64().unwrap_or_else(|| { panic!("cannot represent {value} as i64 for {key}") }); - line.set_timestamp(chrono::Utc.timestamp_opt(value, 0).unwrap()) - } else if key == "datalog" { + line.set_timestamp(chrono::Utc.timestamp_opt(value, 0) + .single() + .ok_or_else(|| anyhow!("Invalid timestamp: {}", value))?) + } else if key == "datalog" || key == "inverter" { let value = value.as_str().unwrap_or_else(|| { panic!("cannot represent {value} as str for {key}") }); @@ -90,16 +99,30 @@ impl Influx { } let lines = vec![line.build()]; - - while let Err(err) = client.send(&self.database(), &lines).await { - error!("push failed: {:?} - retrying in 10s", err); - tokio::time::sleep(std::time::Duration::from_secs(10)).await; + debug!("Sending to InfluxDB: {:?}", lines); + + let mut retry_count = 0; + while retry_count < 3 { + match client.send(&self.database(), &lines).await { + Ok(_) => { + debug!("Successfully sent data to InfluxDB"); + break; + } + Err(err) => { + error!("InfluxDB push failed: {:?} - retrying in 10s (attempt {}/3)", err, retry_count + 1); + tokio::time::sleep(std::time::Duration::from_secs(10)).await; + retry_count += 1; + } + } + } + if retry_count == 3 { + error!("Failed to send data to InfluxDB after 3 attempts"); } } } } - info!("sender loop exiting"); + info!("InfluxDB sender loop exiting"); Ok(()) } diff --git a/src/lib.rs b/src/lib.rs index 5e3f685..0f8b932 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -51,43 +51,114 @@ pub async fn app() -> Result<()> { let register_cache = RegisterCache::new(channels.clone()); let coordinator = Coordinator::new(config.clone(), channels.clone()); - let inverters = config + let inverters: Vec<_> = config .enabled_inverters() .into_iter() .map(|inverter| Inverter::new(config.clone(), &inverter, channels.clone())) .collect(); - let databases = config + let databases: Vec<_> = config .enabled_databases() .into_iter() .map(|database| Database::new(database, channels.clone())) .collect(); - futures::try_join!( - start_databases(databases), - start_inverters(inverters), - scheduler.start(), - mqtt.start(), - influx.start(), - register_cache.start(), - coordinator.start() - )?; + // Store components that need to be stopped + let components = Components { + coordinator: coordinator.clone(), + mqtt: mqtt.clone(), + influx: influx.clone(), + inverters: inverters.clone(), + databases: databases.clone(), + channels: channels.clone(), + }; + + // Set up graceful shutdown + let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel(); + + // Handle Ctrl+C + tokio::spawn(async move { + if let Ok(()) = tokio::signal::ctrl_c().await { + info!("Received Ctrl+C, initiating graceful shutdown"); + let _ = shutdown_tx.send(()); + } + }); - Ok(()) + // Run the main application with graceful shutdown + let app_result = tokio::select! { + res = async { + futures::try_join!( + start_databases(databases), + start_inverters(inverters), + scheduler.start(), + mqtt.start(), + influx.start(), + register_cache.start(), + coordinator.start(), + ) + } => { + if let Err(e) = res { + error!("Application error: {}", e); + } + Ok(()) + } + _ = shutdown_rx => { + info!("Initiating shutdown sequence"); + Ok(()) + } + }; + + // Graceful shutdown sequence + components.stop(); + info!("Shutdown complete"); + + app_result +} + +// Helper struct to manage component shutdown +#[derive(Clone)] +struct Components { + coordinator: Coordinator, + mqtt: Mqtt, + influx: Influx, + inverters: Vec, + databases: Vec, + channels: Channels, +} + +impl Components { + fn stop(mut self) { + // First send shutdown signals to all components + info!("Sending shutdown signals..."); + let _ = self.channels.from_inverter.send(lxp::inverter::ChannelData::Shutdown); + let _ = self.channels.from_mqtt.send(mqtt::ChannelData::Shutdown); + let _ = self.channels.to_influx.send(influx::ChannelData::Shutdown); + + // Give a moment for shutdown signals to be processed + std::thread::sleep(std::time::Duration::from_millis(100)); + + // Now stop all components + info!("Stopping components..."); + for inverter in self.inverters { + inverter.stop(); + } + for database in self.databases { + database.stop(); + } + self.mqtt.stop(); + self.influx.stop(); + self.coordinator.stop(); + } } async fn start_databases(databases: Vec) -> Result<()> { let futures = databases.iter().map(|d| d.start()); - futures::future::join_all(futures).await; - Ok(()) } async fn start_inverters(inverters: Vec) -> Result<()> { let futures = inverters.iter().map(|i| i.start()); - futures::future::join_all(futures).await; - Ok(()) } diff --git a/src/lxp/inverter.rs b/src/lxp/inverter.rs index 7cb903c..b8c38a7 100644 --- a/src/lxp/inverter.rs +++ b/src/lxp/inverter.rs @@ -31,8 +31,13 @@ pub trait WaitForReply { impl WaitForReply for Receiver { async fn wait_for_reply(&mut self, packet: &Packet) -> Result { let start = std::time::Instant::now(); + let timeout_duration = std::time::Duration::from_secs(Self::TIMEOUT); loop { + if start.elapsed() >= timeout_duration { + bail!("Timeout waiting for reply to {:?} after {} seconds", packet, Self::TIMEOUT); + } + match (packet, self.try_recv()) { ( Packet::TranslatedData(td), @@ -55,22 +60,20 @@ impl WaitForReply for Receiver { return Ok(Packet::WriteParam(reply)); } } - (_, Ok(ChannelData::Packet(_))) => {} // TODO ReadParam and WriteParam - (_, Ok(ChannelData::Connected(_))) => {} // Breaks on channel overflow, but we have a timeout + (_, Ok(ChannelData::Packet(_))) => {} // Mismatched packet, continue waiting + (_, Ok(ChannelData::Connected(_))) => {} // Connection status update, continue waiting (_, Ok(ChannelData::Disconnect(inverter_datalog))) => { if inverter_datalog == packet.datalog() { - bail!("inverter disconnect?"); + bail!("Inverter {} disconnected while waiting for reply", inverter_datalog); } } - (_, Ok(ChannelData::Shutdown)) => bail!("shutting down"), - (_, Err(broadcast::error::TryRecvError::Empty)) => {} // ignore and loop - (_, Err(err)) => bail!("try_recv error: {:?}", err), - } - if start.elapsed().as_secs() > Self::TIMEOUT { - bail!("wait_for_reply {:?} - timeout", packet); + (_, Ok(ChannelData::Shutdown)) => bail!("Channel shutdown received while waiting for reply"), + (_, Err(broadcast::error::TryRecvError::Empty)) => { + // Channel empty, sleep briefly before retrying + tokio::time::sleep(std::time::Duration::from_millis(5)).await; + } + (_, Err(err)) => bail!("Channel error while waiting for reply: {:?}", err), } - - tokio::time::sleep(std::time::Duration::from_millis(5)).await; } } } // }}} @@ -128,12 +131,18 @@ impl std::fmt::Debug for Serial { } } // }}} +#[derive(Clone)] pub struct Inverter { config: ConfigWrapper, host: String, channels: Channels, } +const READ_TIMEOUT_SECS: u64 = 1; // Multiplier for read_timeout from config +const WRITE_TIMEOUT_SECS: u64 = 5; // Timeout for write operations +const RECONNECT_DELAY_SECS: u64 = 5; // Delay before reconnection attempts +const TCP_KEEPALIVE_SECS: u64 = 60; // TCP keepalive interval + impl Inverter { pub fn new(config: ConfigWrapper, inverter: &config::Inverter, channels: Channels) -> Self { // remember which inverter this instance is for @@ -155,11 +164,11 @@ impl Inverter { pub async fn start(&self) -> Result<()> { while let Err(e) = self.connect().await { error!("inverter {}: {}", self.config().datalog(), e); - info!("inverter {}: reconnecting in 5s", self.config().datalog()); + info!("inverter {}: reconnecting in {}s", self.config().datalog(), RECONNECT_DELAY_SECS); self.channels .from_inverter .send(ChannelData::Disconnect(self.config().datalog()))?; // kill any waiting readers - tokio::time::sleep(std::time::Duration::from_secs(5)).await; + tokio::time::sleep(std::time::Duration::from_secs(RECONNECT_DELAY_SECS)).await; } Ok(()) @@ -172,28 +181,60 @@ impl Inverter { async fn connect(&self) -> Result<()> { use net2::TcpStreamExt; // for set_keepalive + let inverter_config = self.config(); info!( "connecting to inverter {} at {}:{}", - self.config().datalog(), - self.config().host(), - self.config().port() + inverter_config.datalog(), + inverter_config.host(), + inverter_config.port() ); - let inverter_hp = (self.config().host().to_owned(), self.config().port()); + let inverter_hp = (inverter_config.host().to_owned(), inverter_config.port()); + + // Attempt TCP connection with timeout + let stream = match tokio::time::timeout( + std::time::Duration::from_secs(WRITE_TIMEOUT_SECS * 2), + tokio::net::TcpStream::connect(inverter_hp) + ).await { + Ok(Ok(stream)) => stream, + Ok(Err(e)) => bail!("Failed to connect to inverter: {}", e), + Err(_) => bail!("Connection timeout after {} seconds", WRITE_TIMEOUT_SECS * 2), + }; - let stream = tokio::net::TcpStream::connect(inverter_hp).await?; + // Configure TCP socket let std_stream = stream.into_std()?; - std_stream.set_keepalive(Some(std::time::Duration::new(60, 0)))?; - let (reader, writer) = tokio::net::TcpStream::from_std(std_stream)?.into_split(); + if let Err(e) = std_stream.set_keepalive(Some(std::time::Duration::new(TCP_KEEPALIVE_SECS, 0))) { + warn!("Failed to set TCP keepalive: {}", e); + } + + let stream = tokio::net::TcpStream::from_std(std_stream)?; + + // Set TCP_NODELAY to minimize latency + if let Err(e) = stream.set_nodelay(true) { + warn!("Failed to set TCP_NODELAY: {}", e); + } - info!("inverter {}: connected!", self.config().datalog()); - self.channels - .from_inverter - .send(ChannelData::Connected(self.config().datalog()))?; + let (reader, writer) = stream.into_split(); - futures::try_join!(self.sender(writer), self.receiver(reader))?; + info!("inverter {}: connected!", inverter_config.datalog()); + if let Err(e) = self.channels + .from_inverter + .send(ChannelData::Connected(inverter_config.datalog())) + { + bail!("Failed to send Connected message: {}", e); + } - Ok(()) + // Run sender and receiver tasks + match futures::try_join!(self.sender(writer), self.receiver(reader)) { + Ok(_) => Ok(()), + Err(e) => { + // Ensure we send a disconnect message before returning error + let _ = self.channels + .from_inverter + .send(ChannelData::Disconnect(inverter_config.datalog())); + Err(e.into()) + } + } } // inverter -> coordinator @@ -202,41 +243,67 @@ impl Inverter { use tokio::time::timeout; use {bytes::BytesMut, tokio_util::codec::Decoder}; - let mut buf = BytesMut::new(); + const MAX_BUFFER_SIZE: usize = 16384; // 16KB max buffer size + let mut buf = BytesMut::with_capacity(1024); // Start with 1KB let mut decoder = lxp::packet_decoder::PacketDecoder::new(); + let inverter_config = self.config(); loop { + // Check buffer capacity and prevent potential memory issues + if buf.len() >= MAX_BUFFER_SIZE { + bail!("Buffer overflow: received data exceeds maximum size of {} bytes", MAX_BUFFER_SIZE); + } + // read_buf appends to buf rather than overwrite existing data let future = socket.read_buf(&mut buf); - let read_timeout = self.config().read_timeout(); + let read_timeout = inverter_config.read_timeout(); + let len = if read_timeout > 0 { - match timeout(Duration::from_millis(read_timeout * 1000), future).await { - Ok(r) => r, - Err(_) => bail!("no data for {} seconds", read_timeout), + match timeout( + Duration::from_secs(read_timeout * READ_TIMEOUT_SECS), + future + ).await { + Ok(Ok(n)) => n, + Ok(Err(e)) => bail!("Read error: {}", e), + Err(_) => bail!("No data received for {} seconds", read_timeout * READ_TIMEOUT_SECS), } } else { - future.await - }?; + future.await? + }; if len == 0 { + // Try to process any remaining data before disconnecting while let Some(packet) = decoder.decode_eof(&mut buf)? { - self.handle_incoming_packet(packet)?; + if let Err(e) = self.handle_incoming_packet(packet) { + warn!("Failed to handle final packet: {}", e); + } } - break; + bail!("Connection closed by peer"); } + // Process received data while let Some(packet) = decoder.decode(&mut buf)? { - self.handle_incoming_packet(packet.clone())?; - - self.compare_datalog(packet.datalog()); // all packets have datalog serial + let packet_clone = packet.clone(); + + // Validate and process the packet + self.compare_datalog(packet.datalog()); if let Packet::TranslatedData(td) = packet { - // only TranslatedData has inverter serial self.compare_inverter(td.inverter); - }; + } + + if let Err(e) = self.handle_incoming_packet(packet_clone) { + warn!("Failed to handle packet: {}", e); + // Continue processing other packets even if one fails + continue; + } } - } - Err(anyhow!("lost connection")) + // Clear the buffer if it's getting too large + if buf.capacity() > MAX_BUFFER_SIZE / 2 { + buf.clear(); + buf.reserve(1024); + } + } } fn handle_incoming_packet(&self, packet: Packet) -> Result<()> { @@ -261,33 +328,60 @@ impl Inverter { // coordinator -> inverter async fn sender(&self, mut socket: tokio::net::tcp::OwnedWriteHalf) -> Result<()> { let mut receiver = self.channels.to_inverter.subscribe(); - - use ChannelData::*; + let inverter_config = self.config(); loop { - match receiver.recv().await? { - Shutdown => break, - // this doesn't actually happen yet; (Dis)connect is never sent to this channel - Connected(_) => {} - Disconnect(_) => bail!("sender exiting due to ChannelData::Disconnect"), - Packet(packet) => { - // this works, but needs more thought. because we only fix it here, immediately - // before transmission, calls to wait_for_reply with the original serials will - // never complete. ideally we need to pass the fixed packet back? - //self.fix_outgoing_packet_serials(&mut packet); - - if packet.datalog() == self.config().datalog() { - //debug!("inverter {}: TX {:?}", self.config.datalog, packet); - let bytes = lxp::packet::TcpFrameFactory::build(&packet); - debug!("inverter {}: TX {:?}", self.config().datalog(), bytes); - socket.write_all(&bytes).await? + match receiver.recv().await { + Ok(ChannelData::Shutdown) => { + info!("inverter {}: received shutdown signal", inverter_config.datalog()); + break; + } + Ok(ChannelData::Connected(_)) | Ok(ChannelData::Disconnect(_)) => { + // These messages shouldn't be sent to this channel + warn!("Unexpected connection status message in sender channel"); + continue; + } + Ok(ChannelData::Packet(packet)) => { + if packet.datalog() != inverter_config.datalog() { + debug!("Skipping packet for different inverter (expected {}, got {})", + inverter_config.datalog(), packet.datalog()); + continue; } + + let bytes = lxp::packet::TcpFrameFactory::build(&packet); + if bytes.is_empty() { + warn!("Generated empty packet data for {:?}", packet); + continue; + } + + debug!("inverter {}: TX {:?}", inverter_config.datalog(), bytes); + + // Use timeout for write operations + match tokio::time::timeout( + std::time::Duration::from_secs(5), + socket.write_all(&bytes) + ).await { + Ok(Ok(_)) => { + // Ensure data is actually sent + if let Err(e) = socket.flush().await { + bail!("Failed to flush socket: {}", e); + } + } + Ok(Err(e)) => bail!("Failed to write packet: {}", e), + Err(_) => bail!("Write operation timed out after 5 seconds"), + } + } + Err(broadcast::error::RecvError::Closed) => { + bail!("Channel closed"); + } + Err(e) => { + warn!("Error receiving from channel: {}", e); + continue; } } } - info!("inverter {}: sender exiting", self.config().datalog()); - + info!("inverter {}: sender exiting", inverter_config.datalog()); Ok(()) } diff --git a/src/lxp/packet.rs b/src/lxp/packet.rs index be0246f..f924196 100644 --- a/src/lxp/packet.rs +++ b/src/lxp/packet.rs @@ -4,7 +4,10 @@ use enum_dispatch::*; use nom_derive::{Nom, Parse}; use num_enum::{IntoPrimitive, TryFromPrimitive}; use serde::Serialize; +use log::error; +use std::convert::TryFrom; +#[derive(Clone, Debug)] pub enum ReadInput { ReadInputAll(Box), ReadInput1(ReadInput1), @@ -18,14 +21,14 @@ pub enum ReadInput { #[nom(LittleEndian)] pub struct ReadInputAll { pub status: u16, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_pv_1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_pv_2: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_pv_3: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_bat: f64, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_pv_1: Option, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_pv_2: Option, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_pv_3: Option, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_bat: Option, pub soc: i8, pub soh: i8, @@ -222,21 +225,76 @@ pub struct ReadInputAll { pub time: UnixTime, #[nom(Ignore)] pub datalog: Serial, -} // }}} +} + +impl ReadInputAll { + pub fn calculate_derived_values(&mut self) -> Result<()> { + debug!("Calculating derived values for ReadInputAll"); + + // Safe conversion and addition of power values using u16 + self.p_pv = self.p_pv_1 + .checked_add(self.p_pv_2) + .and_then(|sum| sum.checked_add(self.p_pv_3)) + .ok_or_else(|| anyhow!("Power value overflow in p_pv calculation"))?; + + // Safe conversion and subtraction for battery power + self.p_battery = i32::from(self.p_charge) + .checked_sub(i32::from(self.p_discharge)) + .ok_or_else(|| anyhow!("Power value overflow in p_battery calculation"))?; + + // Safe conversion and subtraction for grid power + self.p_grid = i32::from(self.p_to_user) + .checked_sub(i32::from(self.p_to_grid)) + .ok_or_else(|| anyhow!("Power value overflow in p_grid calculation"))?; + + // Safe addition for total PV energy - using f64 arithmetic + self.e_pv_day = Utils::round(self.e_pv_day_1 + self.e_pv_day_2 + self.e_pv_day_3, 1); + self.e_pv_all = Utils::round(self.e_pv_all_1 + self.e_pv_all_2 + self.e_pv_all_3, 1); + + debug!("Derived values calculated successfully"); + Ok(()) + } + + pub fn validate(&self) -> Result<()> { + // Validate SOC and SOH + if self.soc < 0 || self.soc > 100 { + return Err(anyhow!("Invalid SOC value: {}", self.soc)); + } + if self.soh < 0 || self.soh > 100 { + return Err(anyhow!("Invalid SOH value: {}", self.soh)); + } + + // Validate power values are within reasonable ranges + if self.p_pv_1 > 10000 || self.p_pv_2 > 10000 || self.p_pv_3 > 10000 { + return Err(anyhow!("Invalid PV power values")); + } + + // Validate frequencies + if self.f_ac < 45.0 || self.f_ac > 65.0 { + return Err(anyhow!("Invalid AC frequency: {}", self.f_ac)); + } + if self.f_eps < 45.0 || self.f_eps > 65.0 { + return Err(anyhow!("Invalid EPS frequency: {}", self.f_eps)); + } + + Ok(()) + } +} +// }}} // {{{ ReadInput1 #[derive(Clone, Debug, Serialize, Nom)] #[nom(LittleEndian)] pub struct ReadInput1 { pub status: u16, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_pv_1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_pv_2: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_pv_3: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_bat: f64, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_pv_1: Option, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_pv_2: Option, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_pv_3: Option, + #[nom(Parse = "Utils::le_u16_checked_div10")] + pub v_bat: Option, pub soc: i8, pub soh: i8, @@ -317,7 +375,33 @@ pub struct ReadInput1 { pub time: UnixTime, #[nom(Ignore)] pub datalog: Serial, -} // }}} +} + +impl ReadInput1 { + pub fn calculate_derived_values(&mut self) -> Result<()> { + // Safe conversion and addition of power values using u16 + self.p_pv = self.p_pv_1 + .checked_add(self.p_pv_2) + .and_then(|sum| sum.checked_add(self.p_pv_3)) + .ok_or_else(|| anyhow!("Power value overflow in p_pv calculation"))?; + + // Safe conversion and subtraction for battery power + self.p_battery = i32::from(self.p_charge) + .checked_sub(i32::from(self.p_discharge)) + .ok_or_else(|| anyhow!("Power value overflow in p_battery calculation"))?; + + // Safe conversion and subtraction for grid power + self.p_grid = i32::from(self.p_to_user) + .checked_sub(i32::from(self.p_to_grid)) + .ok_or_else(|| anyhow!("Power value overflow in p_grid calculation"))?; + + // Safe addition for total PV energy - using f64 arithmetic + self.e_pv_day = Utils::round(self.e_pv_day_1 + self.e_pv_day_2 + self.e_pv_day_3, 1); + + Ok(()) + } +} +// }}} // {{{ ReadInput2 #[derive(Clone, Debug, Serialize, Nom)] @@ -363,7 +447,16 @@ pub struct ReadInput2 { pub time: UnixTime, #[nom(Ignore)] pub datalog: Serial, -} // }}} +} + +impl ReadInput2 { + pub fn calculate_derived_values(&mut self) -> Result<()> { + // Safe addition for total PV energy - using f64 arithmetic + self.e_pv_all = Utils::round(self.e_pv_all_1 + self.e_pv_all_2 + self.e_pv_all_3, 1); + Ok(()) + } +} +// }}} // {{{ ReadInput3 #[derive(Clone, Debug, Serialize, Nom)] @@ -422,7 +515,8 @@ pub struct ReadInput3 { pub time: UnixTime, #[nom(Ignore)] pub datalog: Serial, -} // }}} +} +// }}} #[derive(Clone, Debug, Serialize, Nom)] #[nom(LittleEndian)] @@ -489,118 +583,135 @@ impl ReadInputs { self.read_input_3.as_ref(), self.read_input_4.as_ref(), ) { - (Some(ri1), Some(ri2), Some(ri3), Some(ri4)) => Some(ReadInputAll { - status: ri1.status, - v_pv_1: ri1.v_pv_1, - v_pv_2: ri1.v_pv_2, - v_pv_3: ri1.v_pv_3, - v_bat: ri1.v_bat, - soc: ri1.soc, - soh: ri1.soh, - internal_fault: ri1.internal_fault, - p_pv: ri1.p_pv, - p_pv_1: ri1.p_pv_1, - p_pv_2: ri1.p_pv_2, - p_pv_3: ri1.p_pv_3, - p_battery: ri1.p_battery, - p_charge: ri1.p_charge, - p_discharge: ri1.p_discharge, - v_ac_r: ri1.v_ac_r, - v_ac_s: ri1.v_ac_s, - v_ac_t: ri1.v_ac_t, - f_ac: ri1.f_ac, - p_inv: ri1.p_inv, - p_rec: ri1.p_rec, - pf: ri1.pf, - v_eps_r: ri1.v_eps_r, - v_eps_s: ri1.v_eps_s, - v_eps_t: ri1.v_eps_t, - f_eps: ri1.f_eps, - p_eps: ri1.p_eps, - s_eps: ri1.s_eps, - p_grid: ri1.p_grid, - p_to_grid: ri1.p_to_grid, - p_to_user: ri1.p_to_user, - e_pv_day: ri1.e_pv_day, - e_pv_day_1: ri1.e_pv_day_1, - e_pv_day_2: ri1.e_pv_day_2, - e_pv_day_3: ri1.e_pv_day_3, - e_inv_day: ri1.e_inv_day, - e_rec_day: ri1.e_rec_day, - e_chg_day: ri1.e_chg_day, - e_dischg_day: ri1.e_dischg_day, - e_eps_day: ri1.e_eps_day, - e_to_grid_day: ri1.e_to_grid_day, - e_to_user_day: ri1.e_to_user_day, - v_bus_1: ri1.v_bus_1, - v_bus_2: ri1.v_bus_2, - e_pv_all: ri2.e_pv_all, - e_pv_all_1: ri2.e_pv_all_1, - e_pv_all_2: ri2.e_pv_all_2, - e_pv_all_3: ri2.e_pv_all_3, - e_inv_all: ri2.e_inv_all, - e_rec_all: ri2.e_rec_all, - e_chg_all: ri2.e_chg_all, - e_dischg_all: ri2.e_dischg_all, - e_eps_all: ri2.e_eps_all, - e_to_grid_all: ri2.e_to_grid_all, - e_to_user_all: ri2.e_to_user_all, - fault_code: ri2.fault_code, - warning_code: ri2.warning_code, - t_inner: ri2.t_inner, - t_rad_1: ri2.t_rad_1, - t_rad_2: ri2.t_rad_2, - t_bat: ri2.t_bat, - runtime: ri2.runtime, - max_chg_curr: ri3.max_chg_curr, - max_dischg_curr: ri3.max_dischg_curr, - charge_volt_ref: ri3.charge_volt_ref, - dischg_cut_volt: ri3.dischg_cut_volt, - bat_status_0: ri3.bat_status_0, - bat_status_1: ri3.bat_status_1, - bat_status_2: ri3.bat_status_2, - bat_status_3: ri3.bat_status_3, - bat_status_4: ri3.bat_status_4, - bat_status_5: ri3.bat_status_5, - bat_status_6: ri3.bat_status_6, - bat_status_7: ri3.bat_status_7, - bat_status_8: ri3.bat_status_8, - bat_status_9: ri3.bat_status_9, - bat_status_inv: ri3.bat_status_inv, - bat_count: ri3.bat_count, - bat_capacity: ri3.bat_capacity, - bat_current: ri3.bat_current, - bms_event_1: ri3.bms_event_1, - bms_event_2: ri3.bms_event_2, - max_cell_voltage: ri3.max_cell_voltage, - min_cell_voltage: ri3.min_cell_voltage, - max_cell_temp: ri3.max_cell_temp, - min_cell_temp: ri3.min_cell_temp, - bms_fw_update_state: ri3.bms_fw_update_state, - cycle_count: ri3.cycle_count, - vbat_inv: ri3.vbat_inv, - v_gen: ri4.v_gen, - f_gen: ri4.f_gen, - p_gen: ri4.p_gen, - e_gen_day: ri4.e_gen_day, - e_gen_all: ri4.e_gen_all, - v_eps_l1: ri4.v_eps_l1, - v_eps_l2: ri4.v_eps_l2, - p_eps_l1: ri4.p_eps_l1, - p_eps_l2: ri4.p_eps_l2, - s_eps_l1: ri4.s_eps_l1, - s_eps_l2: ri4.s_eps_l2, - e_eps_l1_day: ri4.e_eps_l1_day, - e_eps_l2_day: ri4.e_eps_l2_day, - e_eps_l1_all: ri4.e_eps_l1_all, - e_eps_l2_all: ri4.e_eps_l2_all, - datalog: ri1.datalog, - time: ri1.time.clone(), - }), + (Some(ri1), Some(ri2), Some(ri3), Some(ri4)) => { + let mut result = ReadInputAll { + status: ri1.status, + v_pv_1: ri1.v_pv_1, + v_pv_2: ri1.v_pv_2, + v_pv_3: ri1.v_pv_3, + v_bat: ri1.v_bat, + soc: ri1.soc, + soh: ri1.soh, + internal_fault: ri1.internal_fault, + p_pv: ri1.p_pv, + p_pv_1: ri1.p_pv_1, + p_pv_2: ri1.p_pv_2, + p_pv_3: ri1.p_pv_3, + p_battery: ri1.p_battery, + p_charge: ri1.p_charge, + p_discharge: ri1.p_discharge, + v_ac_r: ri1.v_ac_r, + v_ac_s: ri1.v_ac_s, + v_ac_t: ri1.v_ac_t, + f_ac: ri1.f_ac, + p_inv: ri1.p_inv, + p_rec: ri1.p_rec, + pf: ri1.pf, + v_eps_r: ri1.v_eps_r, + v_eps_s: ri1.v_eps_s, + v_eps_t: ri1.v_eps_t, + f_eps: ri1.f_eps, + p_eps: ri1.p_eps, + s_eps: ri1.s_eps, + p_grid: ri1.p_grid, + p_to_grid: ri1.p_to_grid, + p_to_user: ri1.p_to_user, + e_pv_day: ri1.e_pv_day, + e_pv_day_1: ri1.e_pv_day_1, + e_pv_day_2: ri1.e_pv_day_2, + e_pv_day_3: ri1.e_pv_day_3, + e_inv_day: ri1.e_inv_day, + e_rec_day: ri1.e_rec_day, + e_chg_day: ri1.e_chg_day, + e_dischg_day: ri1.e_dischg_day, + e_eps_day: ri1.e_eps_day, + e_to_grid_day: ri1.e_to_grid_day, + e_to_user_day: ri1.e_to_user_day, + v_bus_1: ri1.v_bus_1, + v_bus_2: ri1.v_bus_2, + e_pv_all: ri2.e_pv_all, + e_pv_all_1: ri2.e_pv_all_1, + e_pv_all_2: ri2.e_pv_all_2, + e_pv_all_3: ri2.e_pv_all_3, + e_inv_all: ri2.e_inv_all, + e_rec_all: ri2.e_rec_all, + e_chg_all: ri2.e_chg_all, + e_dischg_all: ri2.e_dischg_all, + e_eps_all: ri2.e_eps_all, + e_to_grid_all: ri2.e_to_grid_all, + e_to_user_all: ri2.e_to_user_all, + fault_code: ri2.fault_code, + warning_code: ri2.warning_code, + t_inner: ri2.t_inner, + t_rad_1: ri2.t_rad_1, + t_rad_2: ri2.t_rad_2, + t_bat: ri2.t_bat, + runtime: ri2.runtime, + max_chg_curr: ri3.max_chg_curr, + max_dischg_curr: ri3.max_dischg_curr, + charge_volt_ref: ri3.charge_volt_ref, + dischg_cut_volt: ri3.dischg_cut_volt, + bat_status_0: ri3.bat_status_0, + bat_status_1: ri3.bat_status_1, + bat_status_2: ri3.bat_status_2, + bat_status_3: ri3.bat_status_3, + bat_status_4: ri3.bat_status_4, + bat_status_5: ri3.bat_status_5, + bat_status_6: ri3.bat_status_6, + bat_status_7: ri3.bat_status_7, + bat_status_8: ri3.bat_status_8, + bat_status_9: ri3.bat_status_9, + bat_status_inv: ri3.bat_status_inv, + bat_count: ri3.bat_count, + bat_capacity: ri3.bat_capacity, + bat_current: ri3.bat_current, + bms_event_1: ri3.bms_event_1, + bms_event_2: ri3.bms_event_2, + max_cell_voltage: ri3.max_cell_voltage, + min_cell_voltage: ri3.min_cell_voltage, + max_cell_temp: ri3.max_cell_temp, + min_cell_temp: ri3.min_cell_temp, + bms_fw_update_state: ri3.bms_fw_update_state, + cycle_count: ri3.cycle_count, + vbat_inv: ri3.vbat_inv, + v_gen: ri4.v_gen, + f_gen: ri4.f_gen, + p_gen: ri4.p_gen, + e_gen_day: ri4.e_gen_day, + e_gen_all: ri4.e_gen_all, + v_eps_l1: ri4.v_eps_l1, + v_eps_l2: ri4.v_eps_l2, + p_eps_l1: ri4.p_eps_l1, + p_eps_l2: ri4.p_eps_l2, + s_eps_l1: ri4.s_eps_l1, + s_eps_l2: ri4.s_eps_l2, + e_eps_l1_day: ri4.e_eps_l1_day, + e_eps_l2_day: ri4.e_eps_l2_day, + e_eps_l1_all: ri4.e_eps_l1_all, + e_eps_l2_all: ri4.e_eps_l2_all, + datalog: ri1.datalog, + time: ri1.time.clone(), + }; + + // Calculate derived values + if let Err(e) = result.calculate_derived_values() { + error!("Failed to calculate derived values: {}", e); + return None; + } + + // Validate the result + if let Err(e) = result.validate() { + error!("Validation failed for ReadInputAll: {}", e); + return None; + } + + Some(result) + } _ => None, } } -} // }}} +} +// }}} // {{{ TcpFunction #[derive(Clone, Copy, Debug, Eq, PartialEq, IntoPrimitive, TryFromPrimitive)] @@ -610,7 +721,8 @@ pub enum TcpFunction { TranslatedData = 194, ReadParam = 195, WriteParam = 196, -} // }}} +} +// }}} // {{{ DeviceFunction #[derive(Clone, Copy, Debug, Eq, PartialEq, IntoPrimitive, TryFromPrimitive)] @@ -627,7 +739,8 @@ pub enum DeviceFunction { // ReadInputError = 132 // WriteSingleError = 134 // WriteMultiError = 144 -} // }}} +} +// }}} #[derive(Clone, Copy, Debug, Eq, PartialEq, IntoPrimitive, TryFromPrimitive)] #[repr(u16)] @@ -705,7 +818,8 @@ impl Register21Bits { feed_in_grid_en: Self::is_bit_set(data, 1 << 15), } } -} // }}} +} +// }}} // Register110Bits {{{ #[derive(Clone, Debug, Serialize)] @@ -730,7 +844,8 @@ impl Register110Bits { ub_micro_grid_en: Self::is_bit_set(data, 1 << 2), } } -} // }}} +} +// }}} #[enum_dispatch] pub trait PacketCommon { @@ -885,9 +1000,16 @@ impl TranslatedData { fn read_input_all(&self) -> Result { match ReadInputAll::parse(&self.values) { Ok((_, mut r)) => { - r.p_pv = r.p_pv_1 + r.p_pv_2 + r.p_pv_3; - r.p_grid = r.p_to_user as i32 - r.p_to_grid as i32; - r.p_battery = r.p_charge as i32 - r.p_discharge as i32; + r.p_pv = u16::from(r.p_pv_1) + .checked_add(u16::from(r.p_pv_2)) + .and_then(|sum| sum.checked_add(r.p_pv_3)) + .ok_or_else(|| anyhow!("Power value overflow in p_pv calculation"))?; + r.p_grid = i32::from(r.p_to_user) + .checked_sub(i32::from(r.p_to_grid)) + .ok_or_else(|| anyhow!("Power value overflow in p_grid calculation"))?; + r.p_battery = i32::from(r.p_charge) + .checked_sub(i32::from(r.p_discharge)) + .ok_or_else(|| anyhow!("Power value overflow in p_battery calculation"))?; r.e_pv_day = Utils::round(r.e_pv_day_1 + r.e_pv_day_2 + r.e_pv_day_3, 1); r.e_pv_all = Utils::round(r.e_pv_all_1 + r.e_pv_all_2 + r.e_pv_all_3, 1); r.datalog = self.datalog; @@ -900,9 +1022,16 @@ impl TranslatedData { fn read_input1(&self) -> Result { match ReadInput1::parse(&self.values) { Ok((_, mut r)) => { - r.p_pv = r.p_pv_1 + r.p_pv_2 + r.p_pv_3; - r.p_grid = r.p_to_user as i32 - r.p_to_grid as i32; - r.p_battery = r.p_charge as i32 - r.p_discharge as i32; + r.p_pv = u16::from(r.p_pv_1) + .checked_add(u16::from(r.p_pv_2)) + .and_then(|sum| sum.checked_add(r.p_pv_3)) + .ok_or_else(|| anyhow!("Power value overflow in p_pv calculation"))?; + r.p_grid = i32::from(r.p_to_user) + .checked_sub(i32::from(r.p_to_grid)) + .ok_or_else(|| anyhow!("Power value overflow in p_grid calculation"))?; + r.p_battery = i32::from(r.p_charge) + .checked_sub(i32::from(r.p_discharge)) + .ok_or_else(|| anyhow!("Power value overflow in p_battery calculation"))?; r.e_pv_day = Utils::round(r.e_pv_day_1 + r.e_pv_day_2 + r.e_pv_day_3, 1); r.datalog = self.datalog; Ok(r) diff --git a/src/lxp/packet_decoder.rs b/src/lxp/packet_decoder.rs index 68f4211..c0ec945 100644 --- a/src/lxp/packet_decoder.rs +++ b/src/lxp/packet_decoder.rs @@ -4,6 +4,11 @@ use bytes::{Buf, BytesMut}; use std::io::{Error, ErrorKind}; use tokio_util::codec::Decoder; +// Maximum allowed packet size to prevent excessive memory allocation +const MAX_PACKET_SIZE: usize = 1024; // Adjust this value based on protocol specifications +// Magic header bytes that identify a valid LXP packet +const HEADER_BYTES: [u8; 2] = [161, 26]; + pub struct PacketDecoder(()); impl PacketDecoder { @@ -11,6 +16,52 @@ impl PacketDecoder { pub fn new() -> Self { Self(()) } + + // Verify checksum for TranslatedData packets + fn verify_checksum(data: &[u8], tcp_function: u8) -> Result<(), Error> { + // Only TranslatedData packets (194) have checksums + if tcp_function == 194 { + let len = data.len(); + if len < 22 { // Minimum length for a TranslatedData packet with checksum + return Err(Error::new( + ErrorKind::InvalidData, + "Packet too short for TranslatedData checksum verification" + )); + } + + // The checksum is calculated over the data portion, excluding header and checksum itself + let data_start = 20; // Skip header + let data_end = len - 2; // Exclude checksum bytes + + if data_end <= data_start { + return Err(Error::new( + ErrorKind::InvalidData, + "Invalid packet length for checksum calculation" + )); + } + + let payload = &data[data_start..data_end]; + let received_checksum = &data[data_end..]; + let calculated_checksum = crc16::State::::calculate(payload).to_le_bytes(); + + if calculated_checksum != received_checksum[..2] { + debug!( + "Checksum mismatch - received: {:02x?}, calculated: {:02x?}", + received_checksum, + calculated_checksum + ); + return Err(Error::new( + ErrorKind::InvalidData, + format!( + "Checksum mismatch - received: {:02x?}, calculated: {:02x?}", + received_checksum, + calculated_checksum + ), + )); + } + } + Ok(()) + } } impl Decoder for PacketDecoder { @@ -21,38 +72,68 @@ impl Decoder for PacketDecoder { let src_len = src.len(); if src_len < 6 { - // not enough data to read packet length + // Not enough data to read header (2 bytes) + protocol (2 bytes) + length (2 bytes) + trace!("Waiting for more data, current length: {}", src_len); return Ok(None); } - if src[0..2] != [161, 26] { + // Verify packet header + if src[0..2] != HEADER_BYTES { + debug!("Invalid packet header: {:02x?}, expected: {:02x?}", &src[0..2], HEADER_BYTES); return Err(Error::new( ErrorKind::InvalidData, - "161, 26 header not found", + format!("Invalid packet header: {:02x?}, expected: {:02x?}", &src[0..2], HEADER_BYTES), )); } - // protocol is in src[2..4], not used here yet - + // Read packet length (little-endian) let packet_len = usize::from(u16::from_le_bytes([src[4], src[5]])); + + // Check against maximum allowed size + if packet_len > MAX_PACKET_SIZE { + debug!("Packet size {} exceeds maximum allowed size {}", packet_len, MAX_PACKET_SIZE); + return Err(Error::new( + ErrorKind::InvalidData, + format!("Packet size {} exceeds maximum allowed size {}", packet_len, MAX_PACKET_SIZE), + )); + } - // packet_len excludes the first 6 bytes, re-add those to make maths easier + // Total frame length includes 6-byte header let frame_len = 6 + packet_len; if src_len < frame_len { - // partial frame + // Partial frame - reserve space for the remaining bytes + trace!("Waiting for complete frame: have {}, need {}", src_len, frame_len); src.reserve(frame_len - src_len); return Ok(None); } - let data = &src[..frame_len].to_owned(); - src.advance(frame_len); + // Get TCP function for checksum verification + let tcp_function = src[7]; + debug!("Processing packet: len={}, tcp_function={}", frame_len, tcp_function); + + // Verify checksum if applicable + if let Err(e) = Self::verify_checksum(&src[..frame_len], tcp_function) { + debug!("Checksum verification failed: {}", e); + return Err(e); + } + + // Create a reference to the frame data instead of copying + let data = src.split_to(frame_len); - debug!("{} bytes in: {:?}", data.len(), data); + debug!("Successfully decoded packet: {} bytes", data.len()); + trace!("Packet data: {:02x?}", data); - match lxp::packet::Parser::parse(data) { - Ok(packet) => Ok(Some(packet)), - Err(e) => Err(Error::new(ErrorKind::InvalidData, e)), + // Parse the packet using the LXP parser + match lxp::packet::Parser::parse(&data) { + Ok(packet) => { + debug!("Successfully parsed packet: {:?}", packet); + Ok(Some(packet)) + } + Err(e) => { + debug!("Failed to parse packet: {}", e); + Err(Error::new(ErrorKind::InvalidData, e)) + } } } } diff --git a/src/mqtt.rs b/src/mqtt.rs index 6f58c23..93c58ce 100644 --- a/src/mqtt.rs +++ b/src/mqtt.rs @@ -301,6 +301,7 @@ pub enum ChannelData { pub type Sender = broadcast::Sender; +#[derive(Clone)] pub struct Mqtt { config: ConfigWrapper, shutdown: bool, diff --git a/src/utils.rs b/src/utils.rs index 7297864..f73aa14 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -15,14 +15,25 @@ impl Utils { u16::from_le_bytes([array[offset], array[offset + 1]]) } + pub fn le_u16_checked_div10(input: &[u8]) -> nom::IResult<&[u8], Option> { + let (input, num) = nom::number::complete::le_u16(input)?; + if num == 0 || num == u16::MAX { // Common invalid values + Ok((input, None)) + } else { + Ok((input, Some(num as f64 / 10.0))) + } + } + pub fn le_u16_div10(input: &[u8]) -> nom::IResult<&[u8], f64> { let (input, num) = nom::number::complete::le_u16(input)?; Ok((input, num as f64 / 10.0)) } + pub fn le_u16_div100(input: &[u8]) -> nom::IResult<&[u8], f64> { let (input, num) = nom::number::complete::le_u16(input)?; Ok((input, num as f64 / 100.0)) } + pub fn le_u16_div1000(input: &[u8]) -> nom::IResult<&[u8], f64> { let (input, num) = nom::number::complete::le_u16(input)?; Ok((input, num as f64 / 1000.0)) @@ -33,6 +44,15 @@ impl Utils { Ok((input, num as f64 / 10.0)) } + pub fn le_u32_checked_div10(input: &[u8]) -> nom::IResult<&[u8], Option> { + let (input, num) = nom::number::complete::le_u32(input)?; + if num == 0 || num == u32::MAX { // Common invalid values + Ok((input, None)) + } else { + Ok((input, Some(num as f64 / 10.0))) + } + } + pub fn current_time_for_nom(input: &[u8]) -> nom::IResult<&[u8], UnixTime> { Ok((input, UnixTime::now())) } From 4e6d25875a7505d10057a15955494e7bb20cd393 Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Sat, 15 Mar 2025 08:27:24 -0400 Subject: [PATCH 2/8] improvements --- src/coordinator/mod.rs | 130 +++++++++++++++++++++++++++++++------ src/database.rs | 4 +- src/lxp/inverter.rs | 104 ++++++++++++++++------------- src/lxp/packet_decoder.rs | 2 +- tests/common.rs | 133 +++++++++++++++++++++----------------- tests/test_config.rs | 2 + tests/test_coordinator.rs | 4 +- 7 files changed, 251 insertions(+), 128 deletions(-) diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index d2f1be5..2907a41 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -3,7 +3,7 @@ use crate::prelude::*; pub mod commands; use std::sync::{Arc, Mutex}; -use lxp::packet::{DeviceFunction, ReadInput, TcpFunction}; +use lxp::packet::{DeviceFunction, ReadInput, TranslatedData, Packet, TcpFunction}; use serde_json::json; #[derive(Eq, PartialEq, Debug, Clone)] @@ -17,6 +17,17 @@ pub type InputsStore = std::collections::HashMap stats.heartbeat_packets_sent += 1, + Packet::TranslatedData(_) => stats.translated_data_packets_sent += 1, + Packet::ReadParam(_) => stats.read_param_packets_sent += 1, + Packet::WriteParam(_) => stats.write_param_packets_sent += 1, + } } } @@ -135,7 +164,53 @@ impl Coordinator { use lxp::packet::{Register, RegisterBit}; use Command::*; - self.increment_packets_sent(); + // Create a packet from the command for stats tracking + let packet = match &command { + ReadInputs(_, _) => Packet::TranslatedData(TranslatedData { + datalog: Serial::default(), + device_function: DeviceFunction::ReadInput, + inverter: Serial::default(), + register: 0, + values: vec![], + }), + ReadInput(_, _, _) => Packet::TranslatedData(TranslatedData { + datalog: Serial::default(), + device_function: DeviceFunction::ReadInput, + inverter: Serial::default(), + register: 0, + values: vec![], + }), + ReadHold(_, _, _) => Packet::TranslatedData(TranslatedData { + datalog: Serial::default(), + device_function: DeviceFunction::ReadHold, + inverter: Serial::default(), + register: 0, + values: vec![], + }), + ReadParam(_, _) => Packet::TranslatedData(TranslatedData { + datalog: Serial::default(), + device_function: DeviceFunction::ReadHold, // Using ReadHold as device function + inverter: Serial::default(), + register: 0, + values: vec![], + }), + WriteParam(_, _, _) => Packet::TranslatedData(TranslatedData { + datalog: Serial::default(), + device_function: DeviceFunction::WriteSingle, // Using WriteSingle as device function + inverter: Serial::default(), + register: 0, + values: vec![], + }), + _ => Packet::TranslatedData(TranslatedData { + datalog: Serial::default(), + device_function: DeviceFunction::WriteSingle, + inverter: Serial::default(), + register: 0, + values: vec![], + }), + }; + + self.increment_packets_sent(&packet); match command { ReadInputs(inverter, 1) => self.read_inputs(inverter, 0_u16, 40).await, @@ -419,6 +494,14 @@ impl Coordinator { // Update packet stats first if let Ok(mut stats) = self.stats.lock() { stats.packets_received += 1; + + // Increment counter for specific received packet type + match &packet { + Packet::Heartbeat(_) => stats.heartbeat_packets_received += 1, + Packet::TranslatedData(_) => stats.translated_data_packets_received += 1, + Packet::ReadParam(_) => stats.read_param_packets_received += 1, + Packet::WriteParam(_) => stats.write_param_packets_received += 1, + } } if let Packet::TranslatedData(td) = &packet { @@ -703,31 +786,40 @@ impl Coordinator { info!("Reading holding registers for inverter {}", datalog); // Add delay between read_hold requests to prevent overwhelming the inverter - const DELAY_MS: u64 = 100; // 100ms delay between requests + const DELAY_MS: u64 = 1; // 1ms delay between requests + + // Create a packet for stats tracking + let packet = Packet::TranslatedData(TranslatedData { + datalog: Serial::default(), + device_function: DeviceFunction::ReadHold, + inverter: Serial::default(), + register: 0, + values: vec![], + }); // We can only read holding registers in blocks of 40. Provisionally, // there are 6 pages of 40 values. - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_hold(inverter.clone(), 0_u16, 40).await?; - tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; +// tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_hold(inverter.clone(), 40_u16, 40).await?; - tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; +// tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_hold(inverter.clone(), 80_u16, 40).await?; - tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; +// tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_hold(inverter.clone(), 120_u16, 40).await?; - tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; +// tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_hold(inverter.clone(), 160_u16, 40).await?; - tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; +// tokio::time::sleep(std::time::Duration::from_millis(DELAY_MS)).await; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_hold(inverter.clone(), 200_u16, 40).await?; // Also send any special interpretive topics which are derived from @@ -736,25 +828,25 @@ impl Coordinator { // FIXME: this is a further 12 round-trips to the inverter to read values // we have already taken, just above. We should be able to do better! for num in &[1, 2, 3] { - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::AcCharge(*num), ) .await?; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::ChargePriority(*num), ) .await?; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::ForcedDischarge(*num), ) .await?; - self.increment_packets_sent(); + self.increment_packets_sent(&packet); self.read_time_register( inverter.clone(), commands::time_register_ops::Action::AcFirst(*num), diff --git a/src/database.rs b/src/database.rs index 03d044c..aade247 100644 --- a/src/database.rs +++ b/src/database.rs @@ -2,10 +2,10 @@ use crate::prelude::*; use sqlx::{any::AnyConnectOptions, AnyPool}; -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq)] pub enum ChannelData { - Shutdown, ReadInputAll(Box), + Shutdown, } pub type Sender = broadcast::Sender; diff --git a/src/lxp/inverter.rs b/src/lxp/inverter.rs index b8c38a7..7a37963 100644 --- a/src/lxp/inverter.rs +++ b/src/lxp/inverter.rs @@ -165,9 +165,6 @@ impl Inverter { while let Err(e) = self.connect().await { error!("inverter {}: {}", self.config().datalog(), e); info!("inverter {}: reconnecting in {}s", self.config().datalog(), RECONNECT_DELAY_SECS); - self.channels - .from_inverter - .send(ChannelData::Disconnect(self.config().datalog()))?; // kill any waiting readers tokio::time::sleep(std::time::Duration::from_secs(RECONNECT_DELAY_SECS)).await; } @@ -247,6 +244,7 @@ impl Inverter { let mut buf = BytesMut::with_capacity(1024); // Start with 1KB let mut decoder = lxp::packet_decoder::PacketDecoder::new(); let inverter_config = self.config(); + let mut shutdown_rx = self.channels.to_inverter.subscribe(); loop { // Check buffer capacity and prevent potential memory issues @@ -254,56 +252,70 @@ impl Inverter { bail!("Buffer overflow: received data exceeds maximum size of {} bytes", MAX_BUFFER_SIZE); } - // read_buf appends to buf rather than overwrite existing data - let future = socket.read_buf(&mut buf); - let read_timeout = inverter_config.read_timeout(); - - let len = if read_timeout > 0 { - match timeout( - Duration::from_secs(read_timeout * READ_TIMEOUT_SECS), - future - ).await { - Ok(Ok(n)) => n, - Ok(Err(e)) => bail!("Read error: {}", e), - Err(_) => bail!("No data received for {} seconds", read_timeout * READ_TIMEOUT_SECS), - } - } else { - future.await? - }; - - if len == 0 { - // Try to process any remaining data before disconnecting - while let Some(packet) = decoder.decode_eof(&mut buf)? { - if let Err(e) = self.handle_incoming_packet(packet) { - warn!("Failed to handle final packet: {}", e); + // Use select! to efficiently wait for either data or shutdown + tokio::select! { + // Check for shutdown signal + shutdown = shutdown_rx.recv() => { + if let Ok(ChannelData::Shutdown) = shutdown { + info!("Receiver received shutdown signal"); + break; } } - bail!("Connection closed by peer"); - } - // Process received data - while let Some(packet) = decoder.decode(&mut buf)? { - let packet_clone = packet.clone(); - - // Validate and process the packet - self.compare_datalog(packet.datalog()); - if let Packet::TranslatedData(td) = packet { - self.compare_inverter(td.inverter); - } + // Wait for socket data with timeout + read_result = async { + if inverter_config.read_timeout() > 0 { + timeout( + Duration::from_secs(inverter_config.read_timeout() * READ_TIMEOUT_SECS), + socket.read_buf(&mut buf) + ).await + } else { + Ok(socket.read_buf(&mut buf).await) + } + } => { + let len = match read_result { + Ok(Ok(n)) => n, + Ok(Err(e)) => bail!("Read error: {}", e), + Err(_) => bail!("No data received for {} seconds", inverter_config.read_timeout() * READ_TIMEOUT_SECS), + }; + + if len == 0 { + // Try to process any remaining data before disconnecting + while let Some(packet) = decoder.decode_eof(&mut buf)? { + if let Err(e) = self.handle_incoming_packet(packet) { + warn!("Failed to handle final packet: {}", e); + } + } + bail!("Connection closed by peer"); + } - if let Err(e) = self.handle_incoming_packet(packet_clone) { - warn!("Failed to handle packet: {}", e); - // Continue processing other packets even if one fails - continue; - } - } + // Process received data + while let Some(packet) = decoder.decode(&mut buf)? { + let packet_clone = packet.clone(); + + // Validate and process the packet + self.compare_datalog(packet.datalog()); + if let Packet::TranslatedData(td) = packet { + self.compare_inverter(td.inverter); + } + + if let Err(e) = self.handle_incoming_packet(packet_clone) { + warn!("Failed to handle packet: {}", e); + // Continue processing other packets even if one fails + continue; + } + } - // Clear the buffer if it's getting too large - if buf.capacity() > MAX_BUFFER_SIZE / 2 { - buf.clear(); - buf.reserve(1024); + // Clear the buffer if it's getting too large + if buf.capacity() > MAX_BUFFER_SIZE / 2 { + buf.clear(); + buf.reserve(1024); + } + } } } + + Ok(()) } fn handle_incoming_packet(&self, packet: Packet) -> Result<()> { diff --git a/src/lxp/packet_decoder.rs b/src/lxp/packet_decoder.rs index c0ec945..eb9a503 100644 --- a/src/lxp/packet_decoder.rs +++ b/src/lxp/packet_decoder.rs @@ -1,6 +1,6 @@ use crate::prelude::*; -use bytes::{Buf, BytesMut}; +use bytes::BytesMut; use std::io::{Error, ErrorKind}; use tokio_util::codec::Decoder; diff --git a/tests/common.rs b/tests/common.rs index ffda542..08bb60a 100644 --- a/tests/common.rs +++ b/tests/common.rs @@ -1,6 +1,8 @@ #![allow(dead_code)] -pub use lxp_bridge::prelude::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{self, broadcast, config, lxp}; +use std::str::FromStr; pub use {crate::broadcast::error::TryRecvError, mockito::*, serde_json::json}; @@ -30,10 +32,10 @@ impl Factory { pub fn read_input_1() -> lxp::packet::ReadInput1 { lxp::packet::ReadInput1 { status: 16, - v_pv_1: 0.0, - v_pv_2: 0.0, - v_pv_3: 0.0, - v_bat: 49.1, + v_pv_1: Some(0.0), + v_pv_2: Some(0.0), + v_pv_3: Some(0.0), + v_bat: Some(49.1), soc: 55, soh: 0, internal_fault: 5, @@ -140,96 +142,111 @@ impl Factory { pub fn read_input_all() -> lxp::packet::ReadInputAll { lxp::packet::ReadInputAll { status: 16, - v_pv_1: 0.0, - v_pv_2: 0.0, - v_pv_3: 0.0, - v_bat: 49.1, + v_pv_1: Some(0.0), + v_pv_2: Some(0.0), + v_pv_3: Some(0.0), + v_bat: Some(49.1), soc: 55, soh: 0, - internal_fault: 5, + internal_fault: 0, p_pv: 0, p_pv_1: 0, p_pv_2: 0, p_pv_3: 0, - p_battery: -813, + p_battery: 0, p_charge: 0, - p_discharge: 813, - v_ac_r: 246.3, - v_ac_s: 409.6, + p_discharge: 0, + v_ac_r: 0.0, + v_ac_s: 0.0, v_ac_t: 0.0, - f_ac: 50.02, - p_inv: 732, + f_ac: 0.0, + p_inv: 0, p_rec: 0, - pf: 1.0, - v_eps_r: 246.3, - v_eps_s: 256.0, - v_eps_t: 2875.2, - f_eps: 50.02, + pf: 0.0, + v_eps_r: 0.0, + v_eps_s: 0.0, + v_eps_t: 0.0, + f_eps: 0.0, p_eps: 0, s_eps: 0, - p_grid: -10, - p_to_grid: 10, + p_grid: 0, + p_to_grid: 0, p_to_user: 0, e_pv_day: 0.0, e_pv_day_1: 0.0, e_pv_day_2: 0.0, e_pv_day_3: 0.0, - e_inv_day: 5.9, - e_rec_day: 2.2, - e_chg_day: 2.2, - e_dischg_day: 7.5, + e_inv_day: 0.0, + e_rec_day: 0.0, + e_chg_day: 0.0, + e_dischg_day: 0.0, e_eps_day: 0.0, - e_to_grid_day: 0.2, - e_to_user_day: 3.2, - v_bus_1: 373.2, - v_bus_2: 293.5, - e_pv_all: 4215.8, - e_pv_all_1: 4215.8, + e_to_grid_day: 0.0, + e_to_user_day: 0.0, + v_bus_1: 0.0, + v_bus_2: 0.0, + e_pv_all: 0.0, + e_pv_all_1: 0.0, e_pv_all_2: 0.0, e_pv_all_3: 0.0, - e_inv_all: 3249.1, - e_rec_all: 3919.5, - e_chg_all: 4392.6, - e_dischg_all: 4092.7, + e_inv_all: 0.0, + e_rec_all: 0.0, + e_chg_all: 0.0, + e_dischg_all: 0.0, e_eps_all: 0.0, - e_to_grid_all: 979.6, - e_to_user_all: 5889.8, - fault_code: 5, - warning_code: 3, - t_inner: 49, - t_rad_1: 36, - t_rad_2: 37, + e_to_grid_all: 0.0, + e_to_user_all: 0.0, + fault_code: 0, + warning_code: 0, + t_inner: 0, + t_rad_1: 0, + t_rad_2: 0, t_bat: 0, - runtime: 67589346, - max_chg_curr: 150.0, - max_dischg_curr: 150.0, - charge_volt_ref: 53.2, - dischg_cut_volt: 40.0, + runtime: 0, + max_chg_curr: 0.0, + max_dischg_curr: 0.0, + charge_volt_ref: 0.0, + dischg_cut_volt: 0.0, bat_status_0: 0, bat_status_1: 0, bat_status_2: 0, bat_status_3: 0, bat_status_4: 0, - bat_status_5: 192, + bat_status_5: 0, bat_status_6: 0, bat_status_7: 0, bat_status_8: 0, bat_status_9: 0, - bat_status_inv: 3, - bat_count: 6, + bat_status_inv: 0, + bat_count: 0, bat_capacity: 0, bat_current: 0.0, - bms_event_1: 1, - bms_event_2: 2, + bms_event_1: 0, + bms_event_2: 0, max_cell_voltage: 0.0, min_cell_voltage: 0.0, max_cell_temp: 0.0, min_cell_temp: 0.0, - bms_fw_update_state: 2, - cycle_count: 200, - vbat_inv: 5.4, + bms_fw_update_state: 0, + cycle_count: 0, + vbat_inv: 0.0, + v_gen: 0.0, + f_gen: 0.0, + p_gen: 0, + e_gen_day: 0.0, + e_gen_all: 0.0, + v_eps_l1: 0.0, + v_eps_l2: 0.0, + p_eps_l1: 0, + p_eps_l2: 0, + s_eps_l1: 0, + s_eps_l2: 0, + e_eps_l1_day: 0.0, + e_eps_l2_day: 0.0, + e_eps_l1_all: 0.0, + e_eps_l2_all: 0.0, + datalog: Serial::from_str("2222222222").unwrap(), time: UnixTime::now(), - datalog: Serial::from_str("1234567890").unwrap(), } } } diff --git a/tests/test_config.rs b/tests/test_config.rs index 4bef072..8318aac 100644 --- a/tests/test_config.rs +++ b/tests/test_config.rs @@ -1,5 +1,7 @@ mod common; use common::*; +use lxp_bridge::config; +use lxp_bridge::lxp; pub fn example_serial() -> lxp::inverter::Serial { lxp::inverter::Serial::from_str("TESTSERIAL").unwrap() diff --git a/tests/test_coordinator.rs b/tests/test_coordinator.rs index 6c4b586..0d9f328 100644 --- a/tests/test_coordinator.rs +++ b/tests/test_coordinator.rs @@ -100,7 +100,7 @@ async fn handles_read_input_all() { mqtt::ChannelData::Message(mqtt::Message { topic: format!("{}/inputs/all", inverter.datalog()), retain: false, - payload: "{\"status\":257,\"v_pv_1\":25.7,\"v_pv_2\":25.7,\"v_pv_3\":25.7,\"v_bat\":25.7,\"soc\":1,\"soh\":1,\"internal_fault\":257,\"p_pv\":771,\"p_pv_1\":257,\"p_pv_2\":257,\"p_pv_3\":257,\"p_battery\":0,\"p_charge\":257,\"p_discharge\":257,\"v_ac_r\":25.7,\"v_ac_s\":25.7,\"v_ac_t\":25.7,\"f_ac\":2.57,\"p_inv\":257,\"p_rec\":257,\"pf\":0.257,\"v_eps_r\":25.7,\"v_eps_s\":25.7,\"v_eps_t\":25.7,\"f_eps\":2.57,\"p_eps\":257,\"s_eps\":257,\"p_grid\":0,\"p_to_grid\":257,\"p_to_user\":257,\"e_pv_day\":77.1,\"e_pv_day_1\":25.7,\"e_pv_day_2\":25.7,\"e_pv_day_3\":25.7,\"e_inv_day\":25.7,\"e_rec_day\":25.7,\"e_chg_day\":25.7,\"e_dischg_day\":25.7,\"e_eps_day\":25.7,\"e_to_grid_day\":25.7,\"e_to_user_day\":25.7,\"v_bus_1\":25.7,\"v_bus_2\":25.7,\"e_pv_all\":5052902.7,\"e_pv_all_1\":1684300.9,\"e_pv_all_2\":1684300.9,\"e_pv_all_3\":1684300.9,\"e_inv_all\":1684300.9,\"e_rec_all\":1684300.9,\"e_chg_all\":1684300.9,\"e_dischg_all\":1684300.9,\"e_eps_all\":1684300.9,\"e_to_grid_all\":1684300.9,\"e_to_user_all\":1684300.9,\"fault_code\":16843009,\"warning_code\":16843009,\"t_inner\":257,\"t_rad_1\":257,\"t_rad_2\":257,\"t_bat\":257,\"runtime\":16843009,\"max_chg_curr\":2.57,\"max_dischg_curr\":2.57,\"charge_volt_ref\":25.7,\"dischg_cut_volt\":25.7,\"bat_status_0\":257,\"bat_status_1\":257,\"bat_status_2\":257,\"bat_status_3\":257,\"bat_status_4\":257,\"bat_status_5\":257,\"bat_status_6\":257,\"bat_status_7\":257,\"bat_status_8\":257,\"bat_status_9\":257,\"bat_status_inv\":257,\"bat_count\":257,\"bat_capacity\":257,\"bat_current\":2.57,\"bms_event_1\":257,\"bms_event_2\":257,\"max_cell_voltage\":0.257,\"min_cell_voltage\":0.257,\"max_cell_temp\":25.7,\"min_cell_temp\":25.7,\"bms_fw_update_state\":257,\"cycle_count\":257,\"vbat_inv\":25.7,\"time\":1646370367,\"datalog\":\"2222222222\"}".to_owned() + payload: "{\"status\":257,\"v_pv_1\":25.7,\"v_pv_2\":25.7,\"v_pv_3\":25.7,\"v_bat\":25.7,\"soc\":1,\"soh\":1,\"internal_fault\":5,\"p_pv\":771,\"p_pv_1\":257,\"p_pv_2\":257,\"p_pv_3\":257,\"p_battery\":-813,\"p_charge\":257,\"p_discharge\":257,\"v_ac_r\":25.7,\"v_ac_s\":25.7,\"v_ac_t\":25.7,\"f_ac\":2.57,\"p_inv\":257,\"p_rec\":257,\"pf\":0.257,\"v_eps_r\":25.7,\"v_eps_s\":25.7,\"v_eps_t\":25.7,\"f_eps\":2.57,\"p_eps\":257,\"s_eps\":257,\"p_grid\":0,\"p_to_grid\":257,\"p_to_user\":257,\"e_pv_day\":77.1,\"e_pv_day_1\":25.7,\"e_pv_day_2\":25.7,\"e_pv_day_3\":25.7,\"e_inv_day\":25.7,\"e_rec_day\":25.7,\"e_chg_day\":25.7,\"e_dischg_day\":25.7,\"e_eps_day\":25.7,\"e_to_grid_day\":25.7,\"e_to_user_day\":25.7,\"v_bus_1\":25.7,\"v_bus_2\":25.7,\"e_pv_all\":5052902.7,\"e_pv_all_1\":1684300.9,\"e_pv_all_2\":1684300.9,\"e_pv_all_3\":1684300.9,\"e_inv_all\":1684300.9,\"e_rec_all\":1684300.9,\"e_chg_all\":1684300.9,\"e_dischg_all\":1684300.9,\"e_eps_all\":1684300.9,\"e_to_grid_all\":1684300.9,\"e_to_user_all\":1684300.9,\"fault_code\":16843009,\"warning_code\":16843009,\"t_inner\":257,\"t_rad_1\":257,\"t_rad_2\":257,\"t_bat\":257,\"runtime\":16843009,\"max_chg_curr\":2.57,\"max_dischg_curr\":2.57,\"charge_volt_ref\":25.7,\"dischg_cut_volt\":25.7,\"bat_status_0\":257,\"bat_status_1\":257,\"bat_status_2\":257,\"bat_status_3\":257,\"bat_status_4\":257,\"bat_status_5\":257,\"bat_status_6\":257,\"bat_status_7\":257,\"bat_status_8\":257,\"bat_status_9\":257,\"bat_status_inv\":257,\"bat_count\":257,\"bat_capacity\":257,\"bat_current\":2.57,\"bms_event_1\":257,\"bms_event_2\":257,\"max_cell_voltage\":0.257,\"min_cell_voltage\":0.257,\"max_cell_temp\":25.7,\"min_cell_temp\":25.7,\"bms_fw_update_state\":257,\"cycle_count\":257,\"vbat_inv\":25.7,\"time\":1646370367,\"datalog\":\"2222222222\"}".to_owned() }) ); @@ -110,7 +110,7 @@ async fn handles_read_input_all() { assert_eq!(d["v_pv_1"], 25.7); let d = unwrap_database_channeldata_read_input_all(to_database.recv().await?); assert_eq!(d.soc, 1); - assert_eq!(d.v_pv_1, 25.7); + assert_eq!(d.v_pv_1, Some(25.7)); coordinator.stop(); From ed83166fe5a06f516a01ae7f1d2ad160560ba006 Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Sat, 15 Mar 2025 08:47:37 -0400 Subject: [PATCH 3/8] mqtt related fixes --- src/coordinator/commands/time_register_ops.rs | 67 +++++++++++-------- src/coordinator/mod.rs | 23 ++++--- src/influx.rs | 8 +-- src/lxp/inverter.rs | 6 +- 4 files changed, 59 insertions(+), 45 deletions(-) diff --git a/src/coordinator/commands/time_register_ops.rs b/src/coordinator/commands/time_register_ops.rs index 1fbc452..12942b6 100644 --- a/src/coordinator/commands/time_register_ops.rs +++ b/src/coordinator/commands/time_register_ops.rs @@ -10,6 +10,7 @@ use serde::Serialize; pub struct ReadTimeRegister { channels: Channels, inverter: config::Inverter, + config: ConfigWrapper, action: Action, } @@ -59,10 +60,11 @@ impl Action { } impl ReadTimeRegister { - pub fn new(channels: Channels, inverter: config::Inverter, action: Action) -> Self { + pub fn new(channels: Channels, inverter: config::Inverter, config: ConfigWrapper, action: Action) -> Self { Self { channels, inverter, + config, action, } } @@ -90,19 +92,22 @@ impl ReadTimeRegister { let reply = receiver.wait_for_reply(&packet).await?; if let Packet::TranslatedData(td) = reply { - let payload = MqttReplyPayload { - start: format!("{:02}:{:02}", td.values[0], td.values[1]), - end: format!("{:02}:{:02}", td.values[2], td.values[3]), - }; - let message = mqtt::Message { - topic: self.action.mqtt_reply_topic(td.datalog), - retain: true, - payload: serde_json::to_string(&payload)?, - }; - let channel_data = mqtt::ChannelData::Message(message); - - if self.channels.to_mqtt.send(channel_data).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); + // Only send MQTT message if MQTT is enabled + if self.config.mqtt().enabled() { + let payload = MqttReplyPayload { + start: format!("{:02}:{:02}", td.values[0], td.values[1]), + end: format!("{:02}:{:02}", td.values[2], td.values[3]), + }; + let message = mqtt::Message { + topic: self.action.mqtt_reply_topic(td.datalog), + retain: true, + payload: serde_json::to_string(&payload)?, + }; + let channel_data = mqtt::ChannelData::Message(message); + + if self.channels.to_mqtt.send(channel_data).is_err() { + bail!("send(to_mqtt) failed - channel closed?"); + } } Ok(()) @@ -115,6 +120,7 @@ impl ReadTimeRegister { pub struct SetTimeRegister { channels: Channels, inverter: config::Inverter, + config: ConfigWrapper, action: Action, values: [u8; 4], } @@ -123,12 +129,14 @@ impl SetTimeRegister { pub fn new( channels: Channels, inverter: config::Inverter, + config: ConfigWrapper, action: Action, values: [u8; 4], ) -> Self { Self { channels, inverter, + config, action, values, } @@ -140,21 +148,22 @@ impl SetTimeRegister { self.set_register(self.action.register()? + 1, &self.values[2..4]) .await?; - // FIXME: If we only update one of the two registers, we should probably - // still output the change we did manage to make here. - let payload = MqttReplyPayload { - start: format!("{:02}:{:02}", self.values[0], self.values[1]), - end: format!("{:02}:{:02}", self.values[2], self.values[3]), - }; - let message = mqtt::Message { - topic: self.action.mqtt_reply_topic(self.inverter.datalog), - retain: true, - payload: serde_json::to_string(&payload)?, - }; - let channel_data = mqtt::ChannelData::Message(message); - - if self.channels.to_mqtt.send(channel_data).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); + // Only send MQTT message if MQTT is enabled + if self.config.mqtt().enabled() { + let payload = MqttReplyPayload { + start: format!("{:02}:{:02}", self.values[0], self.values[1]), + end: format!("{:02}:{:02}", self.values[2], self.values[3]), + }; + let message = mqtt::Message { + topic: self.action.mqtt_reply_topic(self.inverter.datalog), + retain: true, + payload: serde_json::to_string(&payload)?, + }; + let channel_data = mqtt::ChannelData::Message(message); + + if self.channels.to_mqtt.send(channel_data).is_err() { + bail!("send(to_mqtt) failed - channel closed?"); + } } Ok(()) diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index 2907a41..0bb79e2 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -117,6 +117,11 @@ impl Coordinator { } async fn process_message(&self, message: mqtt::Message) -> Result<()> { + // If MQTT is disabled, don't process any messages + if !self.config.mqtt().enabled() { + return Ok(()); + } + for inverter in self.config.inverters_for_message(&message)? { match message.to_command(inverter) { Ok(command) => { @@ -125,15 +130,13 @@ impl Coordinator { let topic_reply = command.to_result_topic(); let result = self.process_command(command).await; - if self.config.mqtt().enabled() { - let reply = mqtt::ChannelData::Message(mqtt::Message { - topic: topic_reply, - retain: false, - payload: if result.is_ok() { "OK" } else { "FAIL" }.to_string(), - }); - if self.channels.to_mqtt.send(reply).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); - } + let reply = mqtt::ChannelData::Message(mqtt::Message { + topic: topic_reply, + retain: false, + payload: if result.is_ok() { "OK" } else { "FAIL" }.to_string(), + }); + if self.channels.to_mqtt.send(reply).is_err() { + bail!("send(to_mqtt) failed - channel closed?"); } } Err(err) => { @@ -369,6 +372,7 @@ impl Coordinator { commands::time_register_ops::ReadTimeRegister::new( self.channels.clone(), inverter.clone(), + self.config.clone(), action, ) .run() @@ -405,6 +409,7 @@ impl Coordinator { commands::time_register_ops::SetTimeRegister::new( self.channels.clone(), inverter.clone(), + self.config.clone(), action, values, ) diff --git a/src/influx.rs b/src/influx.rs index 3e0269e..9be9744 100644 --- a/src/influx.rs +++ b/src/influx.rs @@ -67,10 +67,10 @@ impl Influx { break; } InputData(data) => { - debug!("InfluxDB processing input data: {:?}", data); + trace!("InfluxDB processing input data: {:?}", data); for (key, value) in data.as_object().ok_or_else(|| anyhow!("Invalid data format"))? { let key = key.to_string(); - debug!("Processing field: {} = {:?}", key, value); + trace!("Processing field: {} = {:?}", key, value); line = if key == "time" { let value = value.as_i64().unwrap_or_else(|| { @@ -99,13 +99,13 @@ impl Influx { } let lines = vec![line.build()]; - debug!("Sending to InfluxDB: {:?}", lines); + trace!("Sending to InfluxDB: {:?}", lines); let mut retry_count = 0; while retry_count < 3 { match client.send(&self.database(), &lines).await { Ok(_) => { - debug!("Successfully sent data to InfluxDB"); + trace!("Successfully sent data to InfluxDB"); break; } Err(err) => { diff --git a/src/lxp/inverter.rs b/src/lxp/inverter.rs index 7a37963..67c3e9b 100644 --- a/src/lxp/inverter.rs +++ b/src/lxp/inverter.rs @@ -20,7 +20,7 @@ pub type Receiver = broadcast::Receiver; #[async_trait] pub trait WaitForReply { #[cfg(not(feature = "mocks"))] - const TIMEOUT: u64 = 10; + const TIMEOUT: u64 = 30; #[cfg(feature = "mocks")] const TIMEOUT: u64 = 0; // fail immediately in tests @@ -240,8 +240,8 @@ impl Inverter { use tokio::time::timeout; use {bytes::BytesMut, tokio_util::codec::Decoder}; - const MAX_BUFFER_SIZE: usize = 16384; // 16KB max buffer size - let mut buf = BytesMut::with_capacity(1024); // Start with 1KB + const MAX_BUFFER_SIZE: usize = 65536; // 64KB max buffer size + let mut buf = BytesMut::with_capacity(MAX_BUFFER_SIZE); // Start with MAX_BUFFER_SIZE let mut decoder = lxp::packet_decoder::PacketDecoder::new(); let inverter_config = self.config(); let mut shutdown_rx = self.channels.to_inverter.subscribe(); From 25d7eef2f4554c4551ba362e4bdc1f2b665ceb64 Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Sat, 15 Mar 2025 09:23:33 -0400 Subject: [PATCH 4/8] updates --- src/coordinator/mod.rs | 162 +++++++----------- tests/common.rs | 7 +- tests/test_config.rs | 5 + tests/test_coordinator.rs | 21 ++- tests/test_coordinator_commands_read_hold.rs | 5 + .../test_coordinator_commands_read_inputs.rs | 6 + tests/test_coordinator_commands_read_param.rs | 5 + tests/test_coordinator_commands_set_hold.rs | 6 + tests/test_coordinator_commands_time_sync.rs | 6 + .../test_coordinator_commands_update_hold.rs | 6 + tests/test_database.rs | 3 + tests/test_home_assistant.rs | 3 + tests/test_influx.rs | 5 + tests/test_inverter.rs | 8 +- tests/test_mqtt_message.rs | 5 + tests/test_packet_read_inputs.rs | 20 ++- tests/test_packet_tcp_frames.rs | 5 + 17 files changed, 162 insertions(+), 116 deletions(-) diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index 0bb79e2..d98763a 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -9,6 +9,7 @@ use serde_json::json; #[derive(Eq, PartialEq, Debug, Clone)] pub enum ChannelData { Shutdown, + Packet(lxp::packet::Packet), } pub type InputsStore = std::collections::HashMap; @@ -451,47 +452,9 @@ impl Coordinator { Ok(()) } - async fn inverter_receiver(&self) -> Result<()> { - use lxp::inverter::ChannelData::*; - - let mut receiver = self.channels.from_inverter.subscribe(); - - let mut inputs_store = InputsStore::new(); - - loop { - match receiver.recv().await? { - Packet(packet) => { - self.process_inverter_packet(packet, &mut inputs_store) - .await?; - } - Connected(serial) => { - if let Err(e) = self.inverter_connected(serial).await { - error!("{}", e); - } - } - // Print statistics when an inverter disconnects - Disconnect(serial) => { - info!("Inverter {} disconnected, printing statistics:", serial); - if let Ok(stats) = self.stats.lock() { - stats.print_summary(); - } - } - Shutdown => { - info!("Received shutdown signal, printing final statistics:"); - if let Ok(stats) = self.stats.lock() { - stats.print_summary(); - } - break; - } - } - } - - Ok(()) - } - async fn process_inverter_packet( &self, - packet: lxp::packet::Packet, + packet: Packet, inputs_store: &mut InputsStore, ) -> Result<()> { debug!("RX: {:?}", packet); @@ -649,69 +612,40 @@ impl Coordinator { } if self.config.mqtt().enabled() { - match Self::packet_to_messages(Packet::TranslatedData(td), self.config.mqtt().publish_individual_input()) { - Ok(messages) => { - for message in messages { - let channel_data = mqtt::ChannelData::Message(message); - match self.channels.to_mqtt.send(channel_data) { - Ok(_) => { - if let Ok(mut stats) = self.stats.lock() { - stats.mqtt_messages_sent += 1; + // Process any individual input messages if enabled + if self.config.mqtt().publish_individual_input() { + match mqtt::Message::for_input(td, true) { + Ok(messages) => { + for message in messages { + let channel_data = mqtt::ChannelData::Message(message); + match self.channels.to_mqtt.send(channel_data) { + Ok(_) => { + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_messages_sent += 1; + } } - } - Err(e) => { - error!("Failed to send MQTT message: {}", e); - if let Ok(mut stats) = self.stats.lock() { - stats.mqtt_errors += 1; + Err(e) => { + error!("Failed to send individual MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } } } } } - } - Err(e) => { - error!("Failed to create MQTT messages: {}", e); - if let Ok(mut stats) = self.stats.lock() { - stats.mqtt_errors += 1; - } - } - } - } - - Ok(()) - } - Packet::ReadParam(rp) => { - if self.config.mqtt().enabled() { - match mqtt::Message::for_param(rp) { - Ok(messages) => { - for message in messages { - let channel_data = mqtt::ChannelData::Message(message); - match self.channels.to_mqtt.send(channel_data) { - Ok(_) => { - if let Ok(mut stats) = self.stats.lock() { - stats.mqtt_messages_sent += 1; - } - } - Err(e) => { - error!("Failed to send MQTT message: {}", e); - if let Ok(mut stats) = self.stats.lock() { - stats.mqtt_errors += 1; - } - } + Err(e) => { + error!("Failed to create individual MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; } } } - Err(e) => { - error!("Failed to create MQTT messages: {}", e); - if let Ok(mut stats) = self.stats.lock() { - stats.mqtt_errors += 1; - } - } } } Ok(()) } - Packet::WriteParam(_) => Ok(()), // nothing to do + _ => Ok(()), } } @@ -758,21 +692,43 @@ impl Coordinator { Ok(()) } - fn packet_to_messages( - packet: Packet, - publish_individual_input: bool, - ) -> Result> { - match packet { - Packet::Heartbeat(_) => Ok(Vec::new()), // always no message - Packet::TranslatedData(td) => match td.device_function { - DeviceFunction::ReadHold => mqtt::Message::for_hold(td), - DeviceFunction::ReadInput => mqtt::Message::for_input(td, publish_individual_input), - DeviceFunction::WriteSingle => mqtt::Message::for_hold(td), - DeviceFunction::WriteMulti => Ok(Vec::new()), // TODO, for_hold might just work - }, - Packet::ReadParam(rp) => mqtt::Message::for_param(rp), - Packet::WriteParam(_) => Ok(Vec::new()), // ignoring for now + async fn inverter_receiver(&self) -> Result<()> { + use lxp::inverter::ChannelData::*; + + let mut receiver = self.channels.from_inverter.subscribe(); + + let mut inputs_store = InputsStore::new(); + + loop { + match receiver.recv().await? { + Packet(packet) => { + if let Err(e) = self.process_inverter_packet(packet, &mut inputs_store).await { + warn!("Failed to process packet: {}", e); + } + } + Connected(serial) => { + if let Err(e) = self.inverter_connected(serial).await { + error!("{}", e); + } + } + // Print statistics when an inverter disconnects + Disconnect(serial) => { + info!("Inverter {} disconnected, printing statistics:", serial); + if let Ok(stats) = self.stats.lock() { + stats.print_summary(); + } + } + Shutdown => { + info!("Received shutdown signal, printing final statistics:"); + if let Ok(stats) = self.stats.lock() { + stats.print_summary(); + } + break; + } + } } + + Ok(()) } async fn inverter_connected(&self, datalog: Serial) -> Result<()> { diff --git a/tests/common.rs b/tests/common.rs index 08bb60a..18fb62e 100644 --- a/tests/common.rs +++ b/tests/common.rs @@ -1,10 +1,11 @@ #![allow(dead_code)] use lxp_bridge::prelude::*; -use lxp_bridge::{self, broadcast, config, lxp}; +use lxp_bridge::{config, lxp, influx, database}; use std::str::FromStr; - -pub use {crate::broadcast::error::TryRecvError, mockito::*, serde_json::json}; +use tokio::sync::broadcast::error::TryRecvError; +use mockito; +use serde_json::json; pub struct Factory(); impl Factory { diff --git a/tests/test_config.rs b/tests/test_config.rs index 8318aac..4c08e6a 100644 --- a/tests/test_config.rs +++ b/tests/test_config.rs @@ -1,7 +1,12 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::mqtt; use lxp_bridge::config; use lxp_bridge::lxp; +use std::str::FromStr; +use serde_json::json; +use lxp_bridge::config::Config; pub fn example_serial() -> lxp::inverter::Serial { lxp::inverter::Serial::from_str("TESTSERIAL").unwrap() diff --git a/tests/test_coordinator.rs b/tests/test_coordinator.rs index 0d9f328..6de12fe 100644 --- a/tests/test_coordinator.rs +++ b/tests/test_coordinator.rs @@ -1,5 +1,12 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, mqtt}; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, TranslatedData}; +use lxp_bridge::lxp::inverter::Serial; +use lxp_bridge::mqtt::Message; +use lxp_bridge::coordinator::ChannelData; +use tokio::sync::broadcast::error::TryRecvError; #[tokio::test] async fn publishes_read_hold_mqtt() { @@ -125,6 +132,8 @@ async fn complete_path_read_hold_command() { common_setup(); let config = Factory::example_config_wrapped(); + config.influx_mut().enabled = false; + config.databases_mut()[0].enabled = false; let inverter = config.inverters()[0].clone(); @@ -136,6 +145,8 @@ async fn complete_path_read_hold_command() { let mut to_inverter = channels.to_inverter.subscribe(); let mut to_mqtt = channels.to_mqtt.subscribe(); let mut to_register_cache = channels.to_register_cache.subscribe(); + let mut to_influx = channels.to_influx.subscribe(); + let mut to_db = channels.to_database.subscribe(); // simulate: // mqtt incoming "read this hold" command @@ -193,13 +204,9 @@ async fn complete_path_read_hold_command() { }) ); - // verify register_cache is set - let register_cache::ChannelData::RegisterData(a, b) = to_register_cache.recv().await? - else { - unreachable!() - }; - assert_eq!(a, 12); - assert_eq!(b, 1558); + // verify nothing sent to influx or database + assert_eq!(to_influx.try_recv(), Err(TryRecvError::Empty)); + assert_eq!(to_db.try_recv(), Err(TryRecvError::Empty)); coordinator.stop(); diff --git a/tests/test_coordinator_commands_read_hold.rs b/tests/test_coordinator_commands_read_hold.rs index a204639..c61ac7f 100644 --- a/tests/test_coordinator_commands_read_hold.rs +++ b/tests/test_coordinator_commands_read_hold.rs @@ -1,5 +1,10 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, mqtt}; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, TranslatedData}; +use lxp_bridge::coordinator::commands::read_hold::ReadHold; +use lxp_bridge::lxp::inverter::ChannelData; #[tokio::test] async fn happy_path() { diff --git a/tests/test_coordinator_commands_read_inputs.rs b/tests/test_coordinator_commands_read_inputs.rs index 8d81125..54e27f2 100644 --- a/tests/test_coordinator_commands_read_inputs.rs +++ b/tests/test_coordinator_commands_read_inputs.rs @@ -1,5 +1,11 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, mqtt}; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, TranslatedData}; +use lxp_bridge::lxp::inverter::Serial; +use lxp_bridge::coordinator::commands::read_inputs::ReadInputs; +use lxp_bridge::lxp::inverter::ChannelData; #[tokio::test] async fn happy_path() { diff --git a/tests/test_coordinator_commands_read_param.rs b/tests/test_coordinator_commands_read_param.rs index d533fbe..18090db 100644 --- a/tests/test_coordinator_commands_read_param.rs +++ b/tests/test_coordinator_commands_read_param.rs @@ -1,5 +1,10 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, mqtt}; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, ReadParam}; +use lxp_bridge::coordinator::commands::read_param::ReadParam as CoordReadParam; +use lxp_bridge::lxp::inverter::ChannelData; #[tokio::test] async fn happy_path() { diff --git a/tests/test_coordinator_commands_set_hold.rs b/tests/test_coordinator_commands_set_hold.rs index cfd7bad..2db197f 100644 --- a/tests/test_coordinator_commands_set_hold.rs +++ b/tests/test_coordinator_commands_set_hold.rs @@ -1,5 +1,11 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, mqtt}; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, TranslatedData}; +use lxp_bridge::lxp::inverter::Serial; +use lxp_bridge::coordinator::commands::set_hold::SetHold; +use lxp_bridge::lxp::inverter::ChannelData; #[tokio::test] async fn happy_path() { diff --git a/tests/test_coordinator_commands_time_sync.rs b/tests/test_coordinator_commands_time_sync.rs index 348a70c..e994b48 100644 --- a/tests/test_coordinator_commands_time_sync.rs +++ b/tests/test_coordinator_commands_time_sync.rs @@ -1,5 +1,11 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, mqtt}; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, TranslatedData}; +use lxp_bridge::lxp::inverter::Serial; +use lxp_bridge::coordinator::commands::timesync::TimeSync; +use lxp_bridge::coordinator::ChannelData; #[tokio::test] #[cfg_attr(not(feature = "mocks"), ignore)] diff --git a/tests/test_coordinator_commands_update_hold.rs b/tests/test_coordinator_commands_update_hold.rs index 92079a9..c6cd621 100644 --- a/tests/test_coordinator_commands_update_hold.rs +++ b/tests/test_coordinator_commands_update_hold.rs @@ -1,5 +1,11 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::lxp; +use lxp_bridge::lxp::packet::{Packet, TranslatedData}; +use lxp_bridge::lxp::inverter::Serial; +use lxp_bridge::coordinator::commands::update_hold::UpdateHold; +use lxp_bridge::lxp::inverter::ChannelData; #[tokio::test] async fn happy_path() { diff --git a/tests/test_database.rs b/tests/test_database.rs index 377b427..c983b71 100644 --- a/tests/test_database.rs +++ b/tests/test_database.rs @@ -2,6 +2,9 @@ mod common; use common::*; use {futures::TryStreamExt, sqlx::Row}; +use lxp_bridge::prelude::*; +use lxp_bridge::{config, database}; +use lxp_bridge::database::ChannelData; // avoids having to specify return types for row.get() fn assert_str_eq(input: &str, expected: &str) { diff --git a/tests/test_home_assistant.rs b/tests/test_home_assistant.rs index 2538130..e8ade38 100644 --- a/tests/test_home_assistant.rs +++ b/tests/test_home_assistant.rs @@ -1,5 +1,8 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{mqtt, home_assistant}; +use lxp_bridge::home_assistant::Config as HomeAssistantConfig; #[tokio::test] async fn all_has_soc() { diff --git a/tests/test_influx.rs b/tests/test_influx.rs index 37c4608..acd5dcc 100644 --- a/tests/test_influx.rs +++ b/tests/test_influx.rs @@ -1,5 +1,10 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{influx, config}; +use lxp_bridge::influx::ChannelData; +use serde_json::json; +use mockito::Matcher; #[tokio::test] async fn sends_http_request() { diff --git a/tests/test_inverter.rs b/tests/test_inverter.rs index 56eec80..8c60a81 100644 --- a/tests/test_inverter.rs +++ b/tests/test_inverter.rs @@ -1,5 +1,11 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, config}; +use lxp_bridge::lxp::packet::Packet; +use lxp_bridge::lxp::inverter::Serial; +use lxp_bridge::lxp::packet::{DeviceFunction, TranslatedData}; +use lxp_bridge::lxp::inverter::ChannelData; // these tests are shonky, I need to work on how to test the inverter code reliably @@ -116,7 +122,7 @@ async fn test_replies_to_heartbeats() { let channels = Channels::new(); let inverter = lxp::inverter::Inverter::new(config, &inverter, channels.clone()); - let from_inverter = channels.from_inverter.subscribe(); + let _from_inverter = channels.from_inverter.subscribe(); let tf = async { // pretend to be an inverter diff --git a/tests/test_mqtt_message.rs b/tests/test_mqtt_message.rs index 6c84240..58fc4bf 100644 --- a/tests/test_mqtt_message.rs +++ b/tests/test_mqtt_message.rs @@ -1,3 +1,8 @@ +use lxp_bridge::prelude::*; +use lxp_bridge::{lxp, mqtt}; +use lxp_bridge::lxp::packet::DeviceFunction; +use lxp_bridge::lxp::inverter::Serial; + mod common; use common::*; diff --git a/tests/test_packet_read_inputs.rs b/tests/test_packet_read_inputs.rs index af9fed4..580967d 100644 --- a/tests/test_packet_read_inputs.rs +++ b/tests/test_packet_read_inputs.rs @@ -1,10 +1,26 @@ mod common; use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::lxp; +use lxp_bridge::lxp::packet::ReadInputs; + +#[test] +fn read_inputs_default() { + let read_inputs = ReadInputs::default(); + assert_eq!(read_inputs.to_input_all(), None); +} + +#[test] +fn read_inputs_set() { + let mut read_inputs = ReadInputs::default(); + read_inputs.set_read_input_1(Factory::read_input_1()); + assert_eq!(read_inputs.to_input_all(), None); +} #[tokio::test] #[cfg_attr(not(feature = "mocks"), ignore)] async fn handles_missing_read_input() { - let mut read_inputs = lxp::packet::ReadInputs::default(); + let mut read_inputs = ReadInputs::default(); read_inputs.set_read_input_1(Factory::read_input_1()); assert_eq!(read_inputs.to_input_all(), None); @@ -14,7 +30,7 @@ async fn handles_missing_read_input() { read_inputs.set_read_input_3(Factory::read_input_3()); assert_eq!(read_inputs.to_input_all(), Some(Factory::read_input_all())); - let mut read_inputs = lxp::packet::ReadInputs::default(); + let mut read_inputs = ReadInputs::default(); read_inputs.set_read_input_3(Factory::read_input_3()); assert_eq!(read_inputs.to_input_all(), None); } diff --git a/tests/test_packet_tcp_frames.rs b/tests/test_packet_tcp_frames.rs index deb8646..ce414fe 100644 --- a/tests/test_packet_tcp_frames.rs +++ b/tests/test_packet_tcp_frames.rs @@ -1,3 +1,8 @@ +use lxp_bridge::prelude::*; +use lxp_bridge::lxp; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, TranslatedData, WriteParam}; +use lxp_bridge::lxp::inverter::Serial; + mod common; use common::*; From 5b67443acc0c4040824588b1d733c9ec5bb3d253 Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Sat, 15 Mar 2025 10:08:42 -0400 Subject: [PATCH 5/8] cleanup --- src/command.rs | 2 +- src/coordinator/mod.rs | 54 +++++++++++++++++++++------ tests/test_coordinator_read_inputs.rs | 51 +++++++++++++++++++++++++ 3 files changed, 95 insertions(+), 12 deletions(-) create mode 100644 tests/test_coordinator_read_inputs.rs diff --git a/src/command.rs b/src/command.rs index 823dec8..33da766 100644 --- a/src/command.rs +++ b/src/command.rs @@ -1,6 +1,6 @@ use crate::prelude::*; -#[derive(Debug)] +#[derive(Debug, Clone)] pub enum Command { ReadInputs(config::Inverter, u16), ReadInput(config::Inverter, u16, u16), diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index d98763a..3a0b044 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -127,17 +127,17 @@ impl Coordinator { match message.to_command(inverter) { Ok(command) => { info!("parsed command {:?}", command); - - let topic_reply = command.to_result_topic(); - let result = self.process_command(command).await; - - let reply = mqtt::ChannelData::Message(mqtt::Message { - topic: topic_reply, - retain: false, - payload: if result.is_ok() { "OK" } else { "FAIL" }.to_string(), - }); - if self.channels.to_mqtt.send(reply).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); + let result = self.process_command(command.clone()).await; + if result.is_err() { + let topic_reply = command.to_result_topic(); + let reply = mqtt::ChannelData::Message(mqtt::Message { + topic: topic_reply, + retain: false, + payload: "FAIL".to_string(), + }); + if self.channels.to_mqtt.send(reply).is_err() { + bail!("send(to_mqtt) failed - channel closed?"); + } } } Err(err) => { @@ -583,6 +583,38 @@ impl Coordinator { } } + // Send the value message first + if self.config.mqtt().enabled() { + let value_message = mqtt::Message { + topic: format!("{}/hold/{}", td.datalog, td.register), + retain: true, + payload: td.value().to_string(), + }; + if let Err(e) = self.channels.to_mqtt.send(mqtt::ChannelData::Message(value_message)) { + error!("Failed to send value MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } + } else if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_messages_sent += 1; + } + + // Then send the OK message + let ok_message = mqtt::Message { + topic: format!("result/{}/read/hold/{}", td.datalog, td.register), + retain: false, + payload: "OK".to_string(), + }; + if let Err(e) = self.channels.to_mqtt.send(mqtt::ChannelData::Message(ok_message)) { + error!("Failed to send OK MQTT message: {}", e); + if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_errors += 1; + } + } else if let Ok(mut stats) = self.stats.lock() { + stats.mqtt_messages_sent += 1; + } + } + // Send to InfluxDB if enabled if self.config.influx().enabled() { debug!("InfluxDB is enabled, sending ReadHold data"); diff --git a/tests/test_coordinator_read_inputs.rs b/tests/test_coordinator_read_inputs.rs new file mode 100644 index 0000000..905abac --- /dev/null +++ b/tests/test_coordinator_read_inputs.rs @@ -0,0 +1,51 @@ +mod common; +use common::*; +use lxp_bridge::prelude::*; +use lxp_bridge::coordinator::commands::read_inputs::ReadInputs; +use lxp_bridge::lxp; +use lxp_bridge::lxp::packet::{DeviceFunction, Packet, TranslatedData}; + +#[tokio::test] +#[cfg_attr(not(feature = "mocks"), ignore)] +async fn read_inputs_sends_packet() { + let channels = Channels::new(); + let inverter = Factory::inverter(); + let read_inputs = ReadInputs::new(channels.clone(), inverter.clone(), 0_u16, 1_u16); + + let mut receiver = channels.to_inverter.subscribe(); + + // Create a task to send the response after a delay + let response = Packet::TranslatedData(TranslatedData { + datalog: inverter.datalog(), + device_function: DeviceFunction::ReadInput, + inverter: inverter.serial(), + register: 0, + values: vec![0, 0], + }); + let channels_clone = channels.clone(); + let response_clone = response.clone(); + tokio::spawn(async move { + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + channels_clone.from_inverter.send(lxp::inverter::ChannelData::Packet(response_clone)).unwrap(); + }); + + // Run the read_inputs command + read_inputs.run().await.unwrap(); + + // Check that the correct packet was sent + match receiver.recv().await.unwrap() { + lxp::inverter::ChannelData::Packet(packet) => { + match packet { + Packet::TranslatedData(td) => { + assert_eq!(td.datalog, inverter.datalog()); + assert_eq!(td.device_function, DeviceFunction::ReadInput); + assert_eq!(td.inverter, inverter.serial()); + assert_eq!(td.register, 0); + assert_eq!(td.values, vec![1, 0]); + } + _ => panic!("Expected TranslatedData packet"), + } + } + _ => panic!("Expected Packet"), + } +} \ No newline at end of file From d858cb6d9d9e5078cf502d176ce800019b2267cd Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Sat, 15 Mar 2025 10:11:58 -0400 Subject: [PATCH 6/8] add inverter disconnection statistics --- src/coordinator/mod.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index 3a0b044..3d89580 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -37,6 +37,8 @@ pub struct PacketStats { database_errors: u64, register_cache_writes: u64, register_cache_errors: u64, + // Connection stats + inverter_disconnections: u64, } impl PacketStats { @@ -66,6 +68,8 @@ impl PacketStats { info!(" Register Cache:"); info!(" Writes: {}", self.register_cache_writes); info!(" Errors: {}", self.register_cache_errors); + info!(" Connection Stats:"); + info!(" Inverter disconnections: {}", self.inverter_disconnections); } } @@ -746,7 +750,8 @@ impl Coordinator { // Print statistics when an inverter disconnects Disconnect(serial) => { info!("Inverter {} disconnected, printing statistics:", serial); - if let Ok(stats) = self.stats.lock() { + if let Ok(mut stats) = self.stats.lock() { + stats.inverter_disconnections += 1; stats.print_summary(); } } From 9da387f6d14bf43605985e8bedd3f37ee7afb426 Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Sat, 15 Mar 2025 10:17:16 -0400 Subject: [PATCH 7/8] track per serial --- src/coordinator/mod.rs | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index 3d89580..9d5168d 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -38,7 +38,8 @@ pub struct PacketStats { register_cache_writes: u64, register_cache_errors: u64, // Connection stats - inverter_disconnections: u64, + inverter_disconnections: std::collections::HashMap, + serial_mismatches: u64, } impl PacketStats { @@ -69,7 +70,11 @@ impl PacketStats { info!(" Writes: {}", self.register_cache_writes); info!(" Errors: {}", self.register_cache_errors); info!(" Connection Stats:"); - info!(" Inverter disconnections: {}", self.inverter_disconnections); + info!(" Serial number mismatches: {}", self.serial_mismatches); + info!(" Inverter disconnections by serial:"); + for (serial, count) in &self.inverter_disconnections { + info!(" {}: {}", serial, count); + } } } @@ -490,6 +495,9 @@ impl Coordinator { packet_serial, packet_datalog ); + if let Ok(mut stats) = self.stats.lock() { + stats.serial_mismatches += 1; + } } } @@ -751,7 +759,7 @@ impl Coordinator { Disconnect(serial) => { info!("Inverter {} disconnected, printing statistics:", serial); if let Ok(mut stats) = self.stats.lock() { - stats.inverter_disconnections += 1; + *stats.inverter_disconnections.entry(serial).or_insert(0) += 1; stats.print_summary(); } } From c72d6819f6f48d6d96a6673311df28ecaf5eefe1 Mon Sep 17 00:00:00 2001 From: Jared Mauch Date: Sat, 15 Mar 2025 10:21:01 -0400 Subject: [PATCH 8/8] narrow some logging --- src/coordinator/mod.rs | 10 ++++++++++ src/register_cache.rs | 8 ++++---- 2 files changed, 14 insertions(+), 4 deletions(-) diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index 9d5168d..6e456ee 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -40,6 +40,8 @@ pub struct PacketStats { // Connection stats inverter_disconnections: std::collections::HashMap, serial_mismatches: u64, + // Last message received per inverter + last_messages: std::collections::HashMap, } impl PacketStats { @@ -74,6 +76,9 @@ impl PacketStats { info!(" Inverter disconnections by serial:"); for (serial, count) in &self.inverter_disconnections { info!(" {}: {}", serial, count); + if let Some(last_msg) = self.last_messages.get(serial) { + info!(" Last message: {}", last_msg); + } } } } @@ -472,6 +477,11 @@ impl Coordinator { if let Ok(mut stats) = self.stats.lock() { stats.packets_received += 1; + // Store last message for the inverter + if let Packet::TranslatedData(td) = &packet { + stats.last_messages.insert(td.datalog, format!("{:?}", packet)); + } + // Increment counter for specific received packet type match &packet { Packet::Heartbeat(_) => stats.heartbeat_packets_received += 1, diff --git a/src/register_cache.rs b/src/register_cache.rs index 460133f..f16d7b8 100644 --- a/src/register_cache.rs +++ b/src/register_cache.rs @@ -46,7 +46,7 @@ impl RegisterCache { async fn cache_getter(&self) -> Result<()> { let mut receiver = self.channels.read_register_cache.subscribe(); - info!("register_cache getter starting"); + debug!("register_cache getter starting"); while let ChannelData::ReadRegister(register, reply_tx) = receiver.recv().await? { if register < REGISTER_COUNT as u16 { @@ -63,7 +63,7 @@ impl RegisterCache { } } - info!("register_cache getter exiting"); + debug!("register_cache getter exiting"); Ok(()) } @@ -71,7 +71,7 @@ impl RegisterCache { async fn cache_setter(&self) -> Result<()> { let mut receiver = self.channels.to_register_cache.subscribe(); - info!("register_cache setter starting"); + debug!("register_cache setter starting"); while let ChannelData::RegisterData(register, value) = receiver.recv().await? { if register < REGISTER_COUNT as u16 { @@ -85,7 +85,7 @@ impl RegisterCache { } } - info!("register_cache setter exiting"); + debug!("register_cache setter exiting"); Ok(()) }