diff --git a/CHANGELOG.md b/CHANGELOG.md index 1985b6e..9d64911 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,11 +1,20 @@ # Unreleased +This release deprecates the MQTT topics `lxp/*/inputs/1` (and 2, 3, 4, all). Instead, new applications should switch to using `lxp/*/input/v_bat/parsed` (where `v_bat` is the JSON Key from the old messages). This solves the problem of which messages to publish, and when (inputs/all is particularly problematic!). For more, see https://github.com/celsworth/lxp-bridge/discussions/262 + + * Reconnect to inverter after 15 minutes of not receiving any data (#223) * Fix max/min cell temperature/voltage decoding as reported from BMS (#227) * Add more HA entities: max/min cell temp/voltage, more charge powers (#228) * Add ReadInput4 with EG4 18k generator data (#239, @pmccut) * Add ReadInput4 keys to HA discovery (#240, @jgulick48) * Fix min_chg_curr/max_chg_curr decoding in ReadInputAll packet (#242, @presto8) +* Cache Hold/Input registers internally as they're seen for later use (#248) +* Remove publish_individual_input configuration option (now they always are) (#250) +* Convert HA to use individual inputs messages (#253) +* New configuration open publish_inputs_all_trigger (#256) +* Move time register messages into new parser (start+end times for AC Charge etc) (#259) +* Publish start-end time message when time registers change externally (#261) # 0.13.0 - 27th October 2023 diff --git a/Cargo.lock b/Cargo.lock index de65520..ec944b1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1117,8 +1117,6 @@ dependencies = [ "log", "mockito", "net2", - "nom", - "nom-derive", "num_enum", "reqwest", "rinfluxdb", @@ -1247,28 +1245,6 @@ dependencies = [ "minimal-lexical", ] -[[package]] -name = "nom-derive" -version = "0.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1ff943d68b88d0b87a6e0d58615e8fa07f9fd5a1319fa0a72efc1f62275c79a7" -dependencies = [ - "nom", - "nom-derive-impl", - "rustversion", -] - -[[package]] -name = "nom-derive-impl" -version = "0.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cd0b9a93a84b0d3ec3e70e02d332dc33ac6dfac9cde63e17fcb77172dededa62" -dependencies = [ - "proc-macro2", - "quote", - "syn 1.0.109", -] - [[package]] name = "num-bigint" version = "0.4.4" @@ -1863,12 +1839,6 @@ dependencies = [ "base64 0.21.7", ] -[[package]] -name = "rustversion" -version = "1.0.14" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7ffc183a10b4478d04cbbbfc96d0873219d962dd5accaff2ffbd4ceb7df837f4" - [[package]] name = "ryu" version = "1.0.17" diff --git a/Cargo.toml b/Cargo.toml index fb03706..af914bb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,8 +28,6 @@ env_logger = { version = "~0.10", default-features = false, features = [] } futures = "~0.3" log = "~0.4" net2 = "~0.2" -nom = "~7" -nom-derive = "~0.10" num_enum = "~0.5" rumqttc = "~0.20" serde = { version = "~1 ", features = ["derive"] } diff --git a/config.yaml.example b/config.yaml.example index 30228e4..1df3e04 100644 --- a/config.yaml.example +++ b/config.yaml.example @@ -1,4 +1,5 @@ loglevel: info +publish_inputs_all_trigger: 80 inverters: - enabled: true diff --git a/src/channels.rs b/src/channels.rs index 889b928..0816f09 100644 --- a/src/channels.rs +++ b/src/channels.rs @@ -7,9 +7,8 @@ pub struct Channels { pub from_mqtt: broadcast::Sender, pub to_mqtt: broadcast::Sender, pub to_influx: broadcast::Sender, - pub to_database: broadcast::Sender, - pub read_register_cache: broadcast::Sender, - pub to_register_cache: broadcast::Sender, + //pub to_database: broadcast::Sender, + pub register_cache: broadcast::Sender, } impl Default for Channels { @@ -26,9 +25,8 @@ impl Channels { from_mqtt: Self::channel(), to_mqtt: Self::channel(), to_influx: Self::channel(), - to_database: Self::channel(), - read_register_cache: Self::channel(), - to_register_cache: Self::channel(), + //to_database: Self::channel(), + register_cache: Self::channel(), } } diff --git a/src/config.rs b/src/config.rs index 98289b7..826bd36 100644 --- a/src/config.rs +++ b/src/config.rs @@ -18,6 +18,9 @@ pub struct Config { #[serde(default = "Config::default_loglevel")] pub loglevel: String, + + #[serde(default = "Config::default_publish_inputs_all_trigger")] + pub publish_inputs_all_trigger: u16, } // Inverter {{{ @@ -109,8 +112,6 @@ pub struct Mqtt { #[serde(default = "Config::default_mqtt_homeassistant")] pub homeassistant: HomeAssistant, - - pub publish_individual_input: Option, } impl Mqtt { pub fn enabled(&self) -> bool { @@ -140,10 +141,6 @@ impl Mqtt { pub fn homeassistant(&self) -> &HomeAssistant { &self.homeassistant } - - pub fn publish_individual_input(&self) -> bool { - self.publish_individual_input == Some(true) - } } // }}} // Influx {{{ @@ -329,6 +326,9 @@ impl ConfigWrapper { pub fn loglevel(&self) -> String { self.config.borrow().loglevel.to_owned() } + pub fn publish_inputs_all_trigger(&self) -> u16 { + self.config.borrow().publish_inputs_all_trigger + } } impl Config { @@ -364,6 +364,10 @@ impl Config { fn default_loglevel() -> String { "debug".to_string() } + + pub fn default_publish_inputs_all_trigger() -> u16 { + 80 + } } fn de_serial<'de, D>(deserializer: D) -> Result diff --git a/src/coordinator/commands/time_register_ops.rs b/src/coordinator/commands/time_register_ops.rs index 1fbc452..bcce685 100644 --- a/src/coordinator/commands/time_register_ops.rs +++ b/src/coordinator/commands/time_register_ops.rs @@ -5,20 +5,7 @@ use lxp::{ packet::{DeviceFunction, TranslatedData}, }; -use serde::Serialize; - -pub struct ReadTimeRegister { - channels: Channels, - inverter: config::Inverter, - action: Action, -} - -#[derive(Debug, Serialize)] -struct MqttReplyPayload { - start: String, - end: String, -} - +#[derive(Clone)] pub enum Action { AcCharge(u16), AcFirst(u16), @@ -42,20 +29,15 @@ impl Action { ForcedDischarge(1) => Ok(84), ForcedDischarge(2) => Ok(86), ForcedDischarge(3) => Ok(88), - _ => bail!("unsupported command"), + _ => Err(anyhow!("unsupported command")), } } +} - fn mqtt_reply_topic(&self, datalog: Serial) -> String { - use Action::*; - // no need to be defensive about n here, we checked it already in register() - match self { - AcCharge(n) => format!("{}/ac_charge/{}", datalog, n), - AcFirst(n) => format!("{}/ac_first/{}", datalog, n), - ChargePriority(n) => format!("{}/charge_priority/{}", datalog, n), - ForcedDischarge(n) => format!("{}/forced_discharge/{}", datalog, n), - } - } +pub struct ReadTimeRegister { + channels: Channels, + inverter: config::Inverter, + action: Action, } impl ReadTimeRegister { @@ -76,8 +58,6 @@ impl ReadTimeRegister { values: vec![2, 0], }); - let mut receiver = self.channels.from_inverter.subscribe(); - if self .channels .to_inverter @@ -87,28 +67,7 @@ impl ReadTimeRegister { bail!("send(to_inverter) failed - channel closed?"); } - 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?"); - } - - Ok(()) - } else { - bail!("didn't get expected reply from inverter"); - } + Ok(()) } } @@ -140,23 +99,6 @@ 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?"); - } - Ok(()) } diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index 3f1d9c0..98d5cd6 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -9,8 +9,6 @@ pub enum ChannelData { Shutdown, } -pub type InputsStore = std::collections::HashMap; - pub struct Coordinator { config: ConfigWrapper, channels: Channels, @@ -271,11 +269,14 @@ impl Coordinator { commands::time_register_ops::SetTimeRegister::new( self.channels.clone(), inverter.clone(), - action, + action.clone(), values, ) .run() - .await + .await?; + + // after setting, issue a ReadTimeRegister so we can send a new message out + self.read_time_register(inverter, action).await } async fn set_hold(&self, inverter: config::Inverter, register: U, value: u16) -> Result<()> @@ -317,13 +318,10 @@ impl Coordinator { 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?; + self.process_inverter_packet(packet).await?; } Connected(serial) => { if let Err(e) = self.inverter_connected(serial).await { @@ -339,11 +337,7 @@ impl Coordinator { Ok(()) } - async fn process_inverter_packet( - &self, - packet: lxp::packet::Packet, - inputs_store: &mut InputsStore, - ) -> Result<()> { + async fn process_inverter_packet(&self, packet: lxp::packet::Packet) -> Result<()> { debug!("RX: {:?}", packet); if let Packet::TranslatedData(td) = &packet { @@ -355,88 +349,157 @@ 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}; - - let entry = inputs_store - .entry(td.datalog) - .or_insert_with(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 => { + let register_map = td.register_map(); + let parser = lxp::register_parser::Parser::new(register_map.clone()); + let parsed_inputs = parser.parse_inputs()?; + //debug!("{}", serde_json::to_string(&parsed_inputs)?); + + if self.config.mqtt().enabled() { + // individual message publishing, raw and parsed + self.publish_raw_input_messages(td)?; + self.publish_parsed_input_messages(td, &parsed_inputs)?; + + // inputs/1/2/3/4 + if let Some(topic_fragment) = parser.guess_legacy_inputs_topic() { + self.publish_combined_parsed_input_message( + td, + &parsed_inputs, + topic_fragment, + )?; + }; } - 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; - - entry.set_read_input_4(r4); + for (register, value) in register_map { + self.cache_register(register_cache::Register::Input(register), value)?; + } - if let Some(input) = entry.to_input_all() { + // if we've seen the triggering register in config.publish_inputs_all_trigger + // then feed contents of cache_register into a new parser; + // which should get us "all". feed that to mqtt+influx + if parser.contains_register(self.config.publish_inputs_all_trigger()) { + let cache = register_cache::RegisterCache::dump( + &self.channels, + register_cache::AllRegisters::Input, + ) + .await; + let all_parser = lxp::register_parser::Parser::new(cache); + if all_parser.guess_legacy_inputs_topic() == Some("all") { + let all_parsed_inputs = all_parser.parse_inputs()?; if self.config.mqtt().enabled() { - let message = mqtt::Message::for_input_all(&input, datalog)?; - let channel_data = mqtt::ChannelData::Message(message); - if self.channels.to_mqtt.send(channel_data).is_err() { - bail!("send(to_mqtt) failed - channel closed?"); + // inputs/all + self.publish_combined_parsed_input_message( + td, + &all_parsed_inputs, + "all", + )?; + } + if self.config.influx().enabled() { + let channel_data = influx::ChannelData::InputData( + td.datalog(), + all_parsed_inputs.clone(), + ); + if self.channels.to_influx.send(channel_data).is_err() { + bail!("send(to_influx) failed - channel closed?"); } } + }; - self.save_input_all(Box::new(input)).await?; - } + // clear the cache so we start over next time + register_cache::RegisterCache::clear(&self.channels).await; } - 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?"); - } - } - } + DeviceFunction::ReadHold | DeviceFunction::WriteSingle => { + let register_map = td.register_map(); + + let parser = lxp::register_parser::Parser::new(register_map.clone()); + let parsed_holds = parser.parse_holds()?; + + self.publish_raw_hold_messages(td)?; + self.publish_parsed_hold_messages(td, &parsed_holds)?; + + if td.device_function == DeviceFunction::WriteSingle { + let inverter = self.inverter_config_for_datalog(td.datalog)?; + // if register_map contains an interesting register that's + // part of a multi-register setup (like AC Charge times) then + // issue a ReadHold request to get the other parts so register_parser + // can construct an MQTT message to send out with current data + self.maybe_send_read_holds(register_map, inverter).await?; + } - 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?"); - } + /* not used yet + for (register, value) in register_map { + self.cache_register(register_cache::Register::Hold(register), value)?; } + */ } - Err(e) => { - // log error but avoid exiting loop as then we stop handling - // incoming packets. need better error handling here maybe? - error!("{}", e); - } + DeviceFunction::WriteMulti => {} } } Ok(()) } + async fn maybe_send_read_holds( + &self, + register_map: RegisterMap, + inverter: config::Inverter, + ) -> Result<()> { + // ^ is true if one key is present, but false if none or both are + + if register_map.contains_key(&68) ^ register_map.contains_key(&69) { + self.read_hold(inverter.clone(), 84_u16, 2).await?; // ac_charge/1 + } + if register_map.contains_key(&70) ^ register_map.contains_key(&71) { + self.read_hold(inverter.clone(), 70_u16, 2).await?; // ac_charge/2 + } + if register_map.contains_key(&72) ^ register_map.contains_key(&73) { + self.read_hold(inverter.clone(), 72_u16, 2).await?; // ac_charge/3 + } + if register_map.contains_key(&76) ^ register_map.contains_key(&77) { + self.read_hold(inverter.clone(), 76_u16, 2).await?; // charge_priority/1 + } + if register_map.contains_key(&78) ^ register_map.contains_key(&79) { + self.read_hold(inverter.clone(), 78_u16, 2).await?; // charge_priority/2 + } + if register_map.contains_key(&80) ^ register_map.contains_key(&81) { + self.read_hold(inverter.clone(), 80_u16, 2).await?; // charge_priority/3 + } + if register_map.contains_key(&84) ^ register_map.contains_key(&85) { + self.read_hold(inverter.clone(), 84_u16, 2).await?; // forced_discharge/1 + } + if register_map.contains_key(&86) ^ register_map.contains_key(&87) { + self.read_hold(inverter.clone(), 86_u16, 2).await?; // forced_discharge/2 + } + if register_map.contains_key(&88) ^ register_map.contains_key(&89) { + self.read_hold(inverter.clone(), 88_u16, 2).await?; // forced_discharge/3 + } + if register_map.contains_key(&152) ^ register_map.contains_key(&153) { + self.read_hold(inverter.clone(), 152_u16, 2).await?; // ac_first/1 + } + if register_map.contains_key(&154) ^ register_map.contains_key(&155) { + self.read_hold(inverter.clone(), 154_u16, 2).await?; // ac_first/2 + } + if register_map.contains_key(&156) ^ register_map.contains_key(&157) { + self.read_hold(inverter.clone(), 156_u16, 2).await?; // ac_first/3 + } + + Ok(()) + } + + fn inverter_config_for_datalog(&self, datalog: Serial) -> Result { + self.config + .enabled_inverter_with_datalog(datalog) + .ok_or(anyhow!("Unknown inverter connected: {}", datalog)) + } + // 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. 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), - }; + let inverter = self.inverter_config_for_datalog(datalog)?; if !inverter.publish_holdings_on_connect() { return Ok(()); @@ -484,14 +547,94 @@ 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?"); - } + fn publish_message(&self, topic: String, payload: String, retain: bool) -> Result<()> { + let m = mqtt::Message { + topic, + payload, + retain, + }; + let channel_data = mqtt::ChannelData::Message(m); + if self.channels.to_mqtt.send(channel_data).is_err() { + bail!("send(to_mqtt) failed - channel closed?"); + } + Ok(()) + } + + fn publish_raw_hold_messages(&self, td: &lxp::packet::TranslatedData) -> Result<()> { + for (register, value) in td.register_map() { + self.publish_message( + format!("{}/hold/{}", td.datalog, register), + serde_json::to_string(&value)?, + true, + )?; + } + + Ok(()) + } + + fn publish_raw_input_messages(&self, td: &lxp::packet::TranslatedData) -> Result<()> { + for (register, value) in td.register_map() { + self.publish_message( + format!("{}/input/{}", td.datalog, register), + serde_json::to_string(&value)?, + false, + )?; + } + + Ok(()) + } + + fn publish_combined_parsed_input_message( + &self, + td: &lxp::packet::TranslatedData, + parsed_inputs: &lxp::register_parser::ParsedData, + topic_fragment: &str, + ) -> Result<()> { + self.publish_message( + format!("{}/inputs/{}", td.datalog, topic_fragment), + serde_json::to_string(&parsed_inputs)?, + false, + )?; + + Ok(()) + } + + fn publish_parsed_input_messages( + &self, + td: &lxp::packet::TranslatedData, + parsed_inputs: &lxp::register_parser::ParsedData, + ) -> Result<()> { + for (key, parsed_value) in parsed_inputs.clone() { + self.publish_message( + format!("{}/input/{}/parsed", td.datalog, key), + parsed_value.to_string(), + false, + )?; + } + + Ok(()) + } + + fn publish_parsed_hold_messages( + &self, + td: &lxp::packet::TranslatedData, + parsed_inputs: &lxp::register_parser::ParsedData, + ) -> Result<()> { + for (key, parsed_value) in parsed_inputs.clone() { + self.publish_message( + // no "hold" here, it's done in parse_holds if required! + // (some don't, like ac_charge/1) + format!("{}/{}", td.datalog, key), + parsed_value.to_string(), + true, + )?; } + Ok(()) + } + + /* + async fn save_input_all(&self, input: Box) -> Result<()> { if self.config.have_enabled_database() { let channel_data = database::ChannelData::ReadInputAll(input); if self.channels.to_database.send(channel_data).is_err() { @@ -501,21 +644,15 @@ impl Coordinator { Ok(()) } + */ + + fn cache_register(&self, register: register_cache::Register, value: u16) -> Result<()> { + let channel_data = register_cache::ChannelData::RegisterData(register, value); - 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 + if self.channels.register_cache.send(channel_data).is_err() { + bail!("send(to_register_cache) failed - channel closed?"); } + + Ok(()) } } diff --git a/src/home_assistant.rs b/src/home_assistant.rs index 25ec9bf..aa9a875 100644 --- a/src/home_assistant.rs +++ b/src/home_assistant.rs @@ -1,37 +1,7 @@ use crate::prelude::*; use lxp::packet::Register; -use serde::{Serialize, Serializer}; - -// ValueTemplate {{{ -#[derive(Clone, Debug, PartialEq)] -pub enum ValueTemplate { - None, - Default, // "{{ value_json.$key }}" - String(String), -} -impl ValueTemplate { - pub fn from_default(key: &str) -> Self { - Self::String(format!("{{{{ value_json.{} }}}}", key)) - } - pub fn is_none(&self) -> bool { - *self == Self::None - } - pub fn is_default(&self) -> bool { - *self == Self::Default - } -} -impl Serialize for ValueTemplate { - fn serialize(&self, serializer: S) -> Result - where - S: Serializer, - { - match self { - ValueTemplate::String(str) => serializer.serialize_str(str), - _ => unreachable!(), - } - } -} // }}} +use serde::Serialize; #[derive(Clone, Debug, Serialize)] pub struct Availability { @@ -74,8 +44,6 @@ pub struct Entity<'a> { state_class: Option<&'a str>, #[serde(skip_serializing_if = "Option::is_none")] device_class: Option<&'a str>, - #[serde(skip_serializing_if = "ValueTemplate::is_none")] - value_template: ValueTemplate, #[serde(skip_serializing_if = "Option::is_none")] unit_of_measurement: Option<&'a str>, #[serde(skip_serializing_if = "Option::is_none")] @@ -145,14 +113,7 @@ impl Config { state_class: None, unit_of_measurement: None, icon: None, - value_template: ValueTemplate::Default, // "{{ value_json.$key }}" - // TODO: might change this to an enum that defaults to InputsAll but can be replaced - // with a string for a specific topic? - state_topic: &format!( - "{}/{}/inputs/all", - self.mqtt_config.namespace(), - self.inverter.datalog() - ), + state_topic: &String::default(), device: self.device(), availability: self.availability(), }; @@ -206,12 +167,6 @@ impl Config { Entity { key: "status", name: "Status", - state_topic: &format!( - "{}/{}/input/0/parsed", - self.mqtt_config.namespace(), - self.inverter.datalog() - ), - value_template: ValueTemplate::None, ..base.clone() }, Entity { @@ -226,12 +181,6 @@ impl Config { key: "fault_code", name: "Fault Code", entity_category: Some("diagnostic"), - state_topic: &format!( - "{}/{}/input/fault_code/parsed", - self.mqtt_config.namespace(), - self.inverter.datalog() - ), - value_template: ValueTemplate::None, icon: Some("mdi:alert"), ..base.clone() }, @@ -239,12 +188,6 @@ impl Config { key: "warning_code", name: "Warning Code", entity_category: Some("diagnostic"), - state_topic: &format!( - "{}/{}/input/warning_code/parsed", - self.mqtt_config.namespace(), - self.inverter.datalog() - ), - value_template: ValueTemplate::None, icon: Some("mdi:alert-outline"), ..base.clone() }, @@ -612,14 +555,18 @@ impl Config { sensors .map(|sensor| { - // fill in unique_id and value_template (if default) which are derived from key - let mut sensor = Entity { + // fill in unique_id / state_topic which are derived from key + let sensor = Entity { unique_id: &self.unique_id(sensor.key), + state_topic: &format!( + "{}/{}/input/{}/parsed", + self.mqtt_config.namespace(), + self.inverter.datalog(), + sensor.key + ), + ..sensor }; - if sensor.value_template.is_default() { - sensor.value_template = ValueTemplate::from_default(sensor.key); - } mqtt::Message { topic: self.ha_discovery_topic("sensor", sensor.key), diff --git a/src/influx.rs b/src/influx.rs index a06313f..94ded3f 100644 --- a/src/influx.rs +++ b/src/influx.rs @@ -1,13 +1,12 @@ use crate::prelude::*; -use chrono::TimeZone; use rinfluxdb::line_protocol::{r#async::Client, LineBuilder}; static INPUTS_MEASUREMENT: &str = "inputs"; -#[derive(Eq, PartialEq, Clone, Debug)] +#[derive(PartialEq, Clone, Debug)] pub enum ChannelData { - InputData(serde_json::Value), + InputData(Serial, lxp::register_parser::ParsedData), Shutdown, } @@ -52,6 +51,7 @@ impl Influx { } async fn sender(&self, client: Client) -> Result<()> { + use lxp::register_parser::Value; use ChannelData::*; let mut receiver = self.channels.to_influx.subscribe(); @@ -61,36 +61,28 @@ impl Influx { match receiver.recv().await? { Shutdown => break, - InputData(data) => { - for (key, value) in data.as_object().unwrap() { - let key = key.to_string(); - - 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" { - let value = value.as_str().unwrap_or_else(|| { - panic!("cannot represent {value} as str for {key}") - }); - line.insert_tag(key, value) - } else if value.is_f64() { - let value = value.as_f64().unwrap_or_else(|| { - panic!("cannot represent {value} as f64 for {key}") - }); - line.insert_field(key, value) - } else { - // can't be anything other than int - let value = value.as_i64().unwrap_or_else(|| { - panic!("cannot represent {value} as i64 for {key}") - }); - line.insert_field(key, value) - } + InputData(datalog, data) => { + for (key, value) in data { + line = match (key, value) { + // not required, Influx will just use current time which is good enough + //("time", _) => line.set_timestamp(..), + + // strings make no sense in InfluxDb so use "raw" i64 value + // this is for status/fault_code/warning_code + (_, Value::String(raw, _)) => line.insert_field(key, raw), + (_, Value::StringOwned(raw, _)) => line.insert_field(key, raw), + + (_, Value::Integer(v)) => line.insert_field(key, v), + (_, Value::Float(v)) => line.insert_field(key, v), + }; } + line = line.insert_tag("datalog", datalog.to_string()); + let lines = vec![line.build()]; + debug!("{:?}", lines); + 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; diff --git a/src/lib.rs b/src/lib.rs index 5e3f685..a0181eb 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -2,7 +2,6 @@ pub mod channels; pub mod command; pub mod config; pub mod coordinator; -pub mod database; pub mod home_assistant; pub mod influx; pub mod lxp; @@ -57,14 +56,16 @@ pub async fn app() -> Result<()> { .map(|inverter| Inverter::new(config.clone(), &inverter, channels.clone())) .collect(); - let databases = config - .enabled_databases() - .into_iter() - .map(|database| Database::new(database, channels.clone())) - .collect(); + /* + let databases = config + .enabled_databases() + .into_iter() + .map(|database| Database::new(database, channels.clone())) + .collect(); + */ futures::try_join!( - start_databases(databases), + //start_databases(databases), start_inverters(inverters), scheduler.start(), mqtt.start(), @@ -76,6 +77,7 @@ pub async fn app() -> Result<()> { Ok(()) } +/* async fn start_databases(databases: Vec) -> Result<()> { let futures = databases.iter().map(|d| d.start()); @@ -83,6 +85,7 @@ async fn start_databases(databases: Vec) -> Result<()> { Ok(()) } +*/ async fn start_inverters(inverters: Vec) -> Result<()> { let futures = inverters.iter().map(|i| i.start()); diff --git a/src/lxp/mod.rs b/src/lxp/mod.rs index 69efc8b..b3fca3b 100644 --- a/src/lxp/mod.rs +++ b/src/lxp/mod.rs @@ -1,3 +1,4 @@ pub mod inverter; pub mod packet; pub mod packet_decoder; +pub mod register_parser; diff --git a/src/lxp/packet.rs b/src/lxp/packet.rs index be0246f..1144c0f 100644 --- a/src/lxp/packet.rs +++ b/src/lxp/packet.rs @@ -1,607 +1,9 @@ use crate::prelude::*; use enum_dispatch::*; -use nom_derive::{Nom, Parse}; use num_enum::{IntoPrimitive, TryFromPrimitive}; use serde::Serialize; -pub enum ReadInput { - ReadInputAll(Box), - ReadInput1(ReadInput1), - ReadInput2(ReadInput2), - ReadInput3(ReadInput3), - ReadInput4(ReadInput4), -} - -// {{{ ReadInputAll -#[derive(PartialEq, Clone, Debug, Serialize, Nom)] -#[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, - - pub soc: i8, - pub soh: i8, - - pub internal_fault: u16, - - #[nom(Ignore)] - pub p_pv: u16, - pub p_pv_1: u16, - pub p_pv_2: u16, - pub p_pv_3: u16, - #[nom(Ignore)] - pub p_battery: i32, - pub p_charge: u16, - pub p_discharge: u16, - - #[nom(Parse = "Utils::le_u16_div10")] - pub v_ac_r: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_ac_s: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_ac_t: f64, - #[nom(Parse = "Utils::le_u16_div100")] - pub f_ac: f64, - - pub p_inv: u16, - pub p_rec: u16, - - #[nom(SkipBefore(2))] // IinvRMS - #[nom(Parse = "Utils::le_u16_div1000")] - pub pf: f64, - - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_r: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_s: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_t: f64, - #[nom(Parse = "Utils::le_u16_div100")] - pub f_eps: f64, - pub p_eps: u16, - pub s_eps: u16, - #[nom(Ignore)] - pub p_grid: i32, - pub p_to_grid: u16, - pub p_to_user: u16, - - #[nom(Ignore)] - pub e_pv_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_pv_day_1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_pv_day_2: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_pv_day_3: f64, - - #[nom(Parse = "Utils::le_u16_div10")] - pub e_inv_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_rec_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_chg_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_dischg_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_eps_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_to_grid_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_to_user_day: f64, - - #[nom(Parse = "Utils::le_u16_div10")] - pub v_bus_1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_bus_2: f64, - - #[nom(Ignore)] - pub e_pv_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_pv_all_1: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_pv_all_2: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_pv_all_3: f64, - - #[nom(Parse = "Utils::le_u32_div10")] - pub e_inv_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_rec_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_chg_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_dischg_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_eps_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_to_grid_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_to_user_all: f64, - - pub fault_code: u32, - pub warning_code: u32, - - pub t_inner: u16, - pub t_rad_1: u16, - pub t_rad_2: u16, - pub t_bat: u16, - #[nom(SkipBefore(2))] // reserved - radiator 3? - pub runtime: u32, - // 18 bytes of auto_test stuff here I'm not doing yet - #[nom(SkipBefore(18))] // auto_test stuff, TODO.. - #[nom(SkipBefore(2))] // bat_brand, bat_com_type - #[nom(Parse = "Utils::le_u16_div10")] - pub max_chg_curr: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub max_dischg_curr: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub charge_volt_ref: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub dischg_cut_volt: f64, - - pub bat_status_0: u16, - pub bat_status_1: u16, - pub bat_status_2: u16, - pub bat_status_3: u16, - pub bat_status_4: u16, - pub bat_status_5: u16, - pub bat_status_6: u16, - pub bat_status_7: u16, - pub bat_status_8: u16, - pub bat_status_9: u16, - pub bat_status_inv: u16, - - pub bat_count: u16, - pub bat_capacity: u16, - - #[nom(Parse = "Utils::le_u16_div100")] - pub bat_current: f64, - - pub bms_event_1: u16, // FaultCode_BMS - pub bms_event_2: u16, // WarningCode_BMS - - // TODO: probably floats but need non-zero sample data to check. just guessing at the div100. - #[nom(Parse = "Utils::le_u16_div1000")] - pub max_cell_voltage: f64, - #[nom(Parse = "Utils::le_u16_div1000")] - pub min_cell_voltage: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub max_cell_temp: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub min_cell_temp: f64, - - pub bms_fw_update_state: u16, - - pub cycle_count: u16, - - #[nom(Parse = "Utils::le_u16_div10")] - pub vbat_inv: f64, - - // 14 bytes I'm not sure what they are; possibly generator stuff - #[nom(SkipBefore(14))] - // something about half bus voltage - #[nom(SkipBefore(2))] - #[nom(Parse = "Utils::le_u16_div10")] - pub v_gen: f64, - #[nom(Parse = "Utils::le_u16_div100")] - pub f_gen: f64, - pub p_gen: u16, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_gen_day: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_gen_all: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_l1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_l2: f64, - pub p_eps_l1: u16, - pub p_eps_l2: u16, - pub s_eps_l1: u16, - pub s_eps_l2: u16, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_eps_l1_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_eps_l2_day: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_eps_l1_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_eps_l2_all: f64, - - // EPS data; unsure what this is - - // following are for influx capability only - #[nom(Parse = "Utils::current_time_for_nom")] - pub time: UnixTime, - #[nom(Ignore)] - pub datalog: Serial, -} // }}} - -// {{{ 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, - - pub soc: i8, - pub soh: i8, - - pub internal_fault: u16, - - #[nom(Ignore)] - pub p_pv: u16, - pub p_pv_1: u16, - pub p_pv_2: u16, - pub p_pv_3: u16, - #[nom(Ignore)] - pub p_battery: i32, - pub p_charge: u16, - pub p_discharge: u16, - - #[nom(Parse = "Utils::le_u16_div10")] - pub v_ac_r: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_ac_s: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_ac_t: f64, - #[nom(Parse = "Utils::le_u16_div100")] - pub f_ac: f64, - - pub p_inv: u16, - pub p_rec: u16, - - #[nom(SkipBefore(2))] // IinvRMS - #[nom(Parse = "Utils::le_u16_div1000")] - pub pf: f64, - - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_r: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_s: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_t: f64, - #[nom(Parse = "Utils::le_u16_div100")] - pub f_eps: f64, - pub p_eps: u16, - pub s_eps: u16, - #[nom(Ignore)] - pub p_grid: i32, - pub p_to_grid: u16, - pub p_to_user: u16, - - #[nom(Ignore)] - pub e_pv_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_pv_day_1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_pv_day_2: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_pv_day_3: f64, - - #[nom(Parse = "Utils::le_u16_div10")] - pub e_inv_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_rec_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_chg_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_dischg_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_eps_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_to_grid_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_to_user_day: f64, - - #[nom(Parse = "Utils::le_u16_div10")] - pub v_bus_1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_bus_2: f64, - - #[nom(Parse = "Utils::current_time_for_nom")] - pub time: UnixTime, - #[nom(Ignore)] - pub datalog: Serial, -} // }}} - -// {{{ ReadInput2 -#[derive(Clone, Debug, Serialize, Nom)] -#[nom(Debug, LittleEndian)] -pub struct ReadInput2 { - #[nom(Ignore)] - pub e_pv_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_pv_all_1: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_pv_all_2: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_pv_all_3: f64, - - #[nom(Parse = "Utils::le_u32_div10")] - pub e_inv_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_rec_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_chg_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_dischg_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_eps_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_to_grid_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_to_user_all: f64, - - pub fault_code: u32, - pub warning_code: u32, - - pub t_inner: u16, - pub t_rad_1: u16, - pub t_rad_2: u16, - pub t_bat: u16, - - #[nom(SkipBefore(2))] // reserved - pub runtime: u32, - // 18 bytes of auto_test stuff here I'm not doing yet - // - #[nom(Parse = "Utils::current_time_for_nom")] - pub time: UnixTime, - #[nom(Ignore)] - pub datalog: Serial, -} // }}} - -// {{{ ReadInput3 -#[derive(Clone, Debug, Serialize, Nom)] -#[nom(LittleEndian)] -pub struct ReadInput3 { - #[nom(SkipBefore(2))] // bat_brand, bat_com_type - #[nom(Parse = "Utils::le_u16_div100")] - pub max_chg_curr: f64, - #[nom(Parse = "Utils::le_u16_div100")] - pub max_dischg_curr: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub charge_volt_ref: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub dischg_cut_volt: f64, - - pub bat_status_0: u16, - pub bat_status_1: u16, - pub bat_status_2: u16, - pub bat_status_3: u16, - pub bat_status_4: u16, - pub bat_status_5: u16, - pub bat_status_6: u16, - pub bat_status_7: u16, - pub bat_status_8: u16, - pub bat_status_9: u16, - pub bat_status_inv: u16, - - pub bat_count: u16, - pub bat_capacity: u16, - - #[nom(Parse = "Utils::le_u16_div100")] - pub bat_current: f64, - - pub bms_event_1: u16, - pub bms_event_2: u16, - - // TODO: probably floats but need non-zero sample data to check. just guessing at the div100. - #[nom(Parse = "Utils::le_u16_div1000")] - pub max_cell_voltage: f64, - #[nom(Parse = "Utils::le_u16_div1000")] - pub min_cell_voltage: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub max_cell_temp: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub min_cell_temp: f64, - - pub bms_fw_update_state: u16, - - pub cycle_count: u16, - - #[nom(Parse = "Utils::le_u16_div10")] - pub vbat_inv: f64, - - // following are for influx capability only - #[nom(Parse = "Utils::current_time_for_nom")] - pub time: UnixTime, - #[nom(Ignore)] - pub datalog: Serial, -} // }}} - -#[derive(Clone, Debug, Serialize, Nom)] -#[nom(LittleEndian)] -pub struct ReadInput4 { - // something about half bus voltage - #[nom(SkipBefore(2))] - #[nom(Parse = "Utils::le_u16_div10")] - pub v_gen: f64, - #[nom(Parse = "Utils::le_u16_div100")] - pub f_gen: f64, - pub p_gen: u16, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_gen_day: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_gen_all: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_l1: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub v_eps_l2: f64, - pub p_eps_l1: u16, - pub p_eps_l2: u16, - pub s_eps_l1: u16, - pub s_eps_l2: u16, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_eps_l1_day: f64, - #[nom(Parse = "Utils::le_u16_div10")] - pub e_eps_l2_day: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_eps_l1_all: f64, - #[nom(Parse = "Utils::le_u32_div10")] - pub e_eps_l2_all: f64, - // EPS data; unsure what this is - #[nom(Ignore)] - pub datalog: Serial, -} - -// {{{ ReadInputs -#[derive(Default, Clone, Debug)] -pub struct ReadInputs { - read_input_1: Option, - read_input_2: Option, - read_input_3: Option, - read_input_4: Option, -} - -impl ReadInputs { - pub fn set_read_input_1(&mut self, i: ReadInput1) { - self.read_input_1 = Some(i); - } - pub fn set_read_input_2(&mut self, i: ReadInput2) { - self.read_input_2 = Some(i); - } - pub fn set_read_input_3(&mut self, i: ReadInput3) { - self.read_input_3 = Some(i); - } - pub fn set_read_input_4(&mut self, i: ReadInput4) { - self.read_input_4 = Some(i); - } - - pub fn to_input_all(&self) -> Option { - match ( - self.read_input_1.as_ref(), - self.read_input_2.as_ref(), - 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(), - }), - _ => None, - } - } -} // }}} - // {{{ TcpFunction #[derive(Clone, Copy, Debug, Eq, PartialEq, IntoPrimitive, TryFromPrimitive)] #[repr(u8)] @@ -639,7 +41,7 @@ pub enum Register { AcChargeSocLimit = 67, // AC Charge SOC Limit (%) ChargePriorityPowerCmd = 74, // Charge Priority Charge Rate (%) ChargePrioritySocLimit = 75, // Charge Priority SOC Limit (%) - ForcedDischgSocLimit = 83, // Forced Discarge SOC Limit (%) + ForcedDischgSocLimit = 83, // Forced Discharge SOC Limit (%) DischgCutOffSocEod = 105, // Discharge cut-off SOC (%) EpsDischgCutoffSocEod = 125, // EPS Discharge cut-off SOC (%) AcChargeStartSocLimit = 160, // SOC at which AC charging will begin (%) @@ -861,7 +263,7 @@ pub struct TranslatedData { pub values: Vec, // undecoded, since can be u16 or u32s? } impl TranslatedData { - pub fn pairs(&self) -> Vec<(u16, u16)> { + pub fn register_map(&self) -> RegisterMap { self.values .chunks(2) .enumerate() @@ -869,79 +271,6 @@ impl TranslatedData { .collect() } - pub fn read_input(&self) -> Result { - // note len() is of Vec, so not register count - match (self.register, self.values.len()) { - (0, 254) => Ok(ReadInput::ReadInputAll(Box::new(self.read_input_all()?))), - // (127, 254) has been seen but containing all zeroes, not sure what they are - (0, 80) => Ok(ReadInput::ReadInput1(self.read_input1()?)), - (40, 80) => Ok(ReadInput::ReadInput2(self.read_input2()?)), - (80, 80) => Ok(ReadInput::ReadInput3(self.read_input3()?)), - (120, 80) => Ok(ReadInput::ReadInput4(self.read_input4()?)), - (r1, r2) => bail!("unhandled ReadInput register={} len={}", r1, r2), - } - } - - 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.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; - Ok(r) - } - Err(_) => Err(anyhow!("meh")), - } - } - - 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.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) - } - Err(_) => Err(anyhow!("meh")), - } - } - - fn read_input2(&self) -> Result { - match ReadInput2::parse(&self.values) { - Ok((_, mut r)) => { - 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; - Ok(r) - } - Err(_) => Err(anyhow!("meh")), - } - } - - fn read_input3(&self) -> Result { - match ReadInput3::parse(&self.values) { - Ok((_, mut r)) => { - r.datalog = self.datalog; - Ok(r) - } - Err(_) => Err(anyhow!("meh")), - } - } - - fn read_input4(&self) -> Result { - match ReadInput4::parse(&self.values) { - Ok((_, mut r)) => { - r.datalog = self.datalog; - Ok(r) - } - Err(_) => Err(anyhow!("meh")), - } - } - fn decode(input: &[u8]) -> Result { let len = input.len(); if len < 38 { @@ -1055,7 +384,7 @@ impl PacketCommon for TranslatedData { data[14..16].copy_from_slice(&self.register.to_le_bytes()); if self.device_function == DeviceFunction::WriteMulti { - let register_count = self.pairs().len() as u16; + let register_count = self.register_map().len() as u16; data.extend_from_slice(®ister_count.to_le_bytes()); } @@ -1103,7 +432,7 @@ pub struct ReadParam { pub values: Vec, // undecoded, since can be u16 or i32s? } impl ReadParam { - pub fn pairs(&self) -> Vec<(u16, u16)> { + pub fn register_map(&self) -> RegisterMap { self.values .chunks(2) .enumerate() @@ -1199,7 +528,7 @@ pub struct WriteParam { pub values: Vec, // undecoded, since can be u16 or i32s? } impl WriteParam { - pub fn pairs(&self) -> Vec<(u16, u16)> { + pub fn register_map(&self) -> RegisterMap { self.values .chunks(2) .enumerate() diff --git a/src/lxp/register_parser.rs b/src/lxp/register_parser.rs new file mode 100644 index 0000000..a5b46f9 --- /dev/null +++ b/src/lxp/register_parser.rs @@ -0,0 +1,532 @@ +use crate::prelude::*; +use serde::Serialize; + +use serde_json::json; + +pub type ParsedData = HashMap<&'static str, Value>; + +#[derive(PartialEq, Serialize, Debug, Clone)] +#[serde(untagged)] +pub enum Value { + Integer(i64), + Float(f64), + // for strings (status, fault_code, warning_code), + // we also return the raw i64 value for inserting into InfluxDB. + // time registers are strings too but we don't need it there so just use 0 for now + String(i64, &'static str), + StringOwned(i64, String), +} + +impl Value { + pub fn to_string(&self) -> String { + match self { + Self::Integer(i) => i.to_string(), + Self::Float(f) => f.to_string(), + Self::String(_, s) => s.to_string(), + Self::StringOwned(_, s) => s.to_string(), + } + } +} + +#[derive(Debug, Clone)] +pub struct Parser { + registers: RegisterMap, +} + +impl Parser { + pub fn new(registers: RegisterMap) -> Self { + Self { registers } + } + + // this bodge is to support sending inputs/1 inputs/2 etc with the correct keys in. + // we look at which registers are present and make a guess ;) + pub fn guess_legacy_inputs_topic(&self) -> Option<&'static str> { + let has_register_0 = self.contains_register(0); + let has_register_40 = self.contains_register(40); + let has_register_80 = self.contains_register(80); + let has_register_120 = self.contains_register(120); + + if has_register_0 && has_register_40 && has_register_80 { + // TODO: this should cope with 120 being present too? + // but that may need to be configurable (default off) + Some("all") + } else if has_register_0 && !has_register_40 && !has_register_80 && !has_register_120 { + Some("1") + } else if !has_register_0 && has_register_40 && !has_register_80 && !has_register_120 { + Some("2") + } else if !has_register_0 && !has_register_40 && has_register_80 && !has_register_120 { + Some("3") + } else if !has_register_0 && !has_register_40 && !has_register_80 && has_register_120 { + Some("4") + } else { + None + } + } + + pub fn contains_register(&self, register: u16) -> bool { + self.v_for(register).is_ok() + } + + // given a set of raw input registers from td.registers(), decode what we can + // and return a list of Key -> Value + // + // this will Err if you pass it a subset of registers it doesn't expect, for example + // if you pass in registers that contain 0-7, then it will try to work out p_pv because it + // has seen 7. but p_pv requires 7 + 8 + 9, which aren't present. + // + // I think this is fine because generally we get these registers in lumps of 40, ie 0-39 + // etc. It would be unusual to get a single input register. + // + pub fn parse_inputs(&self) -> Result { + let mut ret = HashMap::new(); + + for (r, v) in self.registers.clone() { + let e = match r { + 0 => vec![("status", self.parse_status(v))], + 1 => vec![("v_pv_1", self.parse_f64_1(v, 10))], + 2 => vec![("v_pv_2", self.parse_f64_1(v, 10))], + 3 => vec![("v_pv_3", self.parse_f64_1(v, 10))], + 4 => vec![("v_bat", self.parse_f64_1(v, 10))], + 5 => vec![("soc", self.parse_i64_l(v)), ("soh", self.parse_i64_h(v))], + 6 => vec![], // reserved + 7 => vec![("p_pv_1", self.parse_i64_1(v)), ("p_pv", self.p_pv()?)], + 8 => vec![("p_pv_2", self.parse_i64_1(v))], + 9 => vec![("p_pv_3", self.parse_i64_1(v))], + 10 => vec![ + ("p_charge", self.parse_i64_1(v)), + ("p_battery", self.p_battery()?), // homebrew net power flow field + ], + 11 => vec![("p_discharge", self.parse_i64_1(v))], + 12 => vec![("v_ac_r", self.parse_f64_1(v, 10))], + 13 => vec![("v_ac_s", self.parse_f64_1(v, 10))], + 14 => vec![("v_ac_t", self.parse_f64_1(v, 10))], + 15 => vec![("f_ac", self.parse_f64_1(v, 100))], + 16 => vec![("p_inv", self.parse_i64_1(v))], + 17 => vec![("p_rec", self.parse_i64_1(v))], + 18 => vec![], // IinvRMS, 0.01A + 19 => vec![("pf", self.parse_f64_1(v, 1000))], + 20 => vec![("v_eps_r", self.parse_f64_1(v, 10))], + 21 => vec![("v_eps_s", self.parse_f64_1(v, 10))], + 22 => vec![("v_eps_t", self.parse_f64_1(v, 10))], + 23 => vec![("f_eps", self.parse_f64_1(v, 100))], + 24 => vec![("p_eps", self.parse_i64_1(v))], + 25 => vec![("s_eps", self.parse_i64_1(v))], + 26 => vec![ + ("p_to_grid", self.parse_i64_1(v)), + ("p_grid", self.p_grid()?), + ], + 27 => vec![("p_to_user", self.parse_i64_1(v))], + 28 => vec![("e_pv_day_1", self.parse_f64_1(v, 10))], + 29 => vec![("e_pv_day_2", self.parse_f64_1(v, 10))], + 30 => vec![("e_pv_day_3", self.parse_f64_1(v, 10))], + 31 => vec![("e_inv_day", self.parse_f64_1(v, 10))], + 32 => vec![("e_rec_day", self.parse_f64_1(v, 10))], + 33 => vec![("e_chg_day", self.parse_f64_1(v, 10))], + 34 => vec![("e_dischg_day", self.parse_f64_1(v, 10))], + 35 => vec![("e_eps_day", self.parse_f64_1(v, 10))], + 36 => vec![("e_to_grid_day", self.parse_f64_1(v, 10))], + 37 => vec![("e_to_user_day", self.parse_f64_1(v, 10))], + 38 => vec![("v_bus_1", self.parse_f64_1(v, 10))], + 39 => vec![("v_bus_2", self.parse_f64_1(v, 10))], + + 40 => vec![("e_pv_all_1", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 41 => vec![], // done in 40 + 42 => vec![("e_pv_all_2", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 43 => vec![], // done in 42 + 44 => vec![("e_pv_all_3", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 45 => vec![], // done in 44 + 46 => vec![("e_inv_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 47 => vec![], // done in 46 + 48 => vec![("e_rec_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 49 => vec![], // done in 48 + 50 => vec![("e_chg_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 51 => vec![], // done in 40 + 52 => vec![("e_dischg_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 53 => vec![], // done in 52 + 54 => vec![("e_eps_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 55 => vec![], // done in 54 + 56 => vec![("e_to_grid_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 57 => vec![], // done in 56 + 58 => vec![("e_to_user_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 59 => vec![], // done in 58 + 60 => vec![("fault_code", self.parse_fault(v, self.v_for(r + 1)?))], + 61 => vec![], // done in 60 + 62 => vec![("warning_code", self.parse_warning(v, self.v_for(r + 1)?))], + 63 => vec![], // done in 62 + 64 => vec![("t_inner", self.parse_i64_1(v))], + 65 => vec![("t_rad_1", self.parse_i64_1(v))], + 66 => vec![("t_rad_2", self.parse_i64_1(v))], + 67 => vec![("t_bat", self.parse_i64_1(v))], + 68 => vec![], // reserved + 69 => vec![("runtime", self.parse_i64_2(v, self.v_for(70)?))], + 70 => vec![], // done in 69 + 71..=79 => vec![], // TODO, rest of ReadInput2 + + 80 => vec![], // bat_brand & bat_com_type + 81 => vec![("max_chg_curr", self.parse_f64_1(v, 100))], + 82 => vec![("max_dischg_curr", self.parse_f64_1(v, 100))], + 83 => vec![("charge_volt_ref", self.parse_f64_1(v, 10))], + 84 => vec![("dischg_cut_volt", self.parse_f64_1(v, 10))], + 85..=95 => vec![], // bat_status_*, not yet parsed + 96 => vec![("bat_count", self.parse_i64_1(v))], + 97 => vec![("bat_capacity", self.parse_i64_1(v))], + 98 => vec![("bat_current", self.parse_f64_1(v, 100))], + 99 => vec![], // bms_event_1 + 100 => vec![], // bms_event_2 + 101 => vec![("max_cell_voltage", self.parse_f64_1(v, 1000))], + 102 => vec![("min_cell_voltage", self.parse_f64_1(v, 1000))], + 103 => vec![("max_cell_temp", self.parse_f64_1(v, 10))], + 104 => vec![("min_cell_temp", self.parse_f64_1(v, 10))], + 105 => vec![], // bms_fw_update_state + 106 => vec![("cycle_count", self.parse_i64_1(v))], + 107 => vec![("vbat_inv", self.parse_f64_1(v, 10))], + 108..=119 => vec![], // TODO, rest of ReadInput3 + + 120 => vec![], // half bus voltage? + 121 => vec![("v_gen", self.parse_f64_1(v, 10))], + 122 => vec![("f_gen", self.parse_f64_1(v, 100))], + 123 => vec![("p_gen", self.parse_i64_1(v))], + 124 => vec![("e_gen_day", self.parse_f64_1(v, 10))], + 125 => vec![("e_gen_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 126 => vec![], // done in 125 + 127 => vec![("v_eps_l1", self.parse_f64_1(v, 10))], + 128 => vec![("v_eps_l2", self.parse_f64_1(v, 10))], + 129 => vec![("p_eps_l1", self.parse_i64_1(v))], + 130 => vec![("p_eps_l2", self.parse_i64_1(v))], + 131 => vec![("s_eps_l1", self.parse_i64_1(v))], + 132 => vec![("s_eps_l2", self.parse_i64_1(v))], + 133 => vec![("e_eps_l1_day", self.parse_f64_1(v, 10))], + 134 => vec![("e_eps_l2_day", self.parse_f64_1(v, 10))], + 135 => vec![("e_eps_l1_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 136 => vec![], // done in 135 + 137 => vec![("e_eps_l2_all", self.parse_f64_2(v, self.v_for(r + 1)?, 10))], + 138 => vec![], // done in 137 + + ..=255 => vec![], // ignore everything else for now + + _ => bail!("unhandled input register {}", r), + }; + + ret.extend(e); + } + + // debug!("{:?}", ret); + + Ok(ret) + } + // given a set of raw hold registers from td.registers(), decode what we can + // and return a list of Key -> Value + // + // unlike parse_inputs, this one does not Err if passed unknown registers. We just silently + // return an empty array. This is because Hold registers are far more likely to be sent + // to us individually, ie a single ReadHold command. Also most hold registers aren't parsed + // and we don't care if we're passed some that aren't implemented. + pub fn parse_holds(&self) -> Result { + let mut ret = HashMap::new(); + + for (r, v) in self.registers.clone() { + let e = match r { + // only place single-register matches in here. If any are missing (which + // is much more likely in ReadHold messages), multi-register operations + // will break! + // + // these should have hold/ prefixed if appropriate! + // + 21 => vec![("hold/21/bits", self.parse_21_bits(v)?)], + 110 => vec![("hold/110/bits", self.parse_110_bits(v)?)], + + // ignore any unknown registers, do not return Err + _ => vec![], + }; + ret.extend(e); + } + + /* self.start_end_tuple hides away a lot of logic to check if the registers + * we pass in are present, and if so, return a Vec suitable for putting into our + * returned HashMap. If any are missing, it does nothing (no error) + * */ + ret.extend(self.start_end_tuple("ac_charge/1", [68, 69])); + ret.extend(self.start_end_tuple("ac_charge/2", [70, 71])); + ret.extend(self.start_end_tuple("ac_charge/3", [72, 73])); + + ret.extend(self.start_end_tuple("charge_priority/1", [76, 77])); + ret.extend(self.start_end_tuple("charge_priority/2", [78, 79])); + ret.extend(self.start_end_tuple("charge_priority/3", [80, 81])); + + ret.extend(self.start_end_tuple("forced_discharge/1", [84, 85])); + ret.extend(self.start_end_tuple("forced_discharge/2", [86, 87])); + ret.extend(self.start_end_tuple("forced_discharge/3", [88, 89])); + + ret.extend(self.start_end_tuple("ac_first/1", [152, 153])); + ret.extend(self.start_end_tuple("ac_first/2", [154, 155])); + ret.extend(self.start_end_tuple("ac_first/3", [156, 157])); + + // debug!("{:?}", ret); + + Ok(ret) + } + + fn start_end_tuple( + &self, + key: &'static str, + registers: [u16; 2], + ) -> Vec<(&'static str, Value)> { + if self.all_registers_present(®isters) { + vec![(key, self.start_end(registers[0], registers[1]).unwrap())] + } else { + vec![] + } + } + + fn start_end(&self, r1: u16, r2: u16) -> Result { + let start = self.v_for(r1)?.to_le_bytes(); + let end = self.v_for(r2)?.to_le_bytes(); + + let payload = json!({ + "start": format!("{:02}:{:02}", start[0], start[1]), + "end": format!("{:02}:{:02}", end[0], end[1]), + }); + + Ok(Value::StringOwned( + 0, // raw value unused, holds are not inserted to Influx anyway + serde_json::to_string(&payload)?, + )) + } + + fn parse_21_bits(&self, value: u16) -> Result { + let bits = lxp::packet::Register21Bits::new(value); + + Ok(Value::StringOwned( + value as i64, + serde_json::to_string(&bits)?, + )) + } + + fn parse_110_bits(&self, value: u16) -> Result { + let bits = lxp::packet::Register110Bits::new(value); + + Ok(Value::StringOwned( + value as i64, + serde_json::to_string(&bits)?, + )) + } + + fn p_pv(&self) -> Result { + let p_pv_1 = self.v_for(7)? as i64; + let p_pv_2 = self.v_for(8)? as i64; + let p_pv_3 = self.v_for(9)? as i64; + + Ok(Value::Integer(p_pv_1 + p_pv_2 + p_pv_3)) + } + + fn p_battery(&self) -> Result { + // special case - use p_charge and p_discharge to return a signed net power flow + let p_charge = self.v_for(10)? as i64; + let p_discharge = self.v_for(11)? as i64; + + Ok(Value::Integer(p_charge - p_discharge)) + } + + fn p_grid(&self) -> Result { + // special case - use p_charge and p_discharge to return a signed net power flow + let p_to_grid = self.v_for(26)? as i64; + let p_to_user = self.v_for(27)? as i64; + + Ok(Value::Integer(p_to_user - p_to_grid)) + } + + fn parse_status(&self, value: u16) -> Value { + Value::String(value as i64, StatusString::from_value(value)) + } + + fn parse_fault(&self, v1: u16, v2: u16) -> Value { + let value: i64 = (v1 as i64) | (v2 as i64) << 16; + Value::String(value, FaultCodeString::from_value(value)) + } + + fn parse_warning(&self, v1: u16, v2: u16) -> Value { + let value: i64 = (v1 as i64) | (v2 as i64) << 16; + Value::String(value, WarningCodeString::from_value(value)) + } + + // one register input, using low half, i64 output + fn parse_i64_l(&self, v1: u16) -> Value { + let r: i64 = (v1 & 0xff) as i64; + Value::Integer(r) + } + + // one register input, using high half, i64 output + fn parse_i64_h(&self, v1: u16) -> Value { + let r: i64 = (v1 >> 8) as i64; + Value::Integer(r) + } + + // one register input, i64 output + fn parse_i64_1(&self, v1: u16) -> Value { + let r: i64 = v1 as i64; + Value::Integer(r) + } + + // two register input, i64 output + fn parse_i64_2(&self, v1: u16, v2: u16) -> Value { + let r: i64 = (v1 as i64) | (v2 as i64) << 16; + Value::Integer(r) + } + + // one register input, f64 output + fn parse_f64_1(&self, v1: u16, divider: i64) -> Value { + let r: i64 = v1 as i64; + Value::Float(r as f64 / divider as f64) + } + + // two register input, f64 output + fn parse_f64_2(&self, v1: u16, v2: u16, divider: i64) -> Value { + let r: i64 = (v1 as i64) | (v2 as i64) << 16; + Value::Float(r as f64 / divider as f64) + } + + // get the value for a given register, or Err + fn v_for(&self, register: u16) -> Result { + self.registers + .get(®ister) + .ok_or(anyhow!("no value found for register {}", register)) + .cloned() + } + + // return true if we have a value for ALL the registers requested + fn all_registers_present(&self, registers: &[u16]) -> bool { + registers + .into_iter() + .all(|register| self.registers.contains_key(register)) + } +} + +struct StatusString; +impl StatusString { + pub fn from_value(status: u16) -> &'static str { + match status { + 0x00 => "Standby", + 0x02 => "FW Updating", + 0x04 => "PV On-grid", + 0x08 => "PV Charge", + 0x0C => "PV Charge On-grid", + 0x10 => "Battery On-grid", + 0x11 => "Bypass", + 0x14 => "PV & Battery On-grid", + 0x19 => "PV Charge + Bypass", + 0x20 => "AC Charge", + 0x28 => "PV & AC Charge", + 0x40 => "Battery Off-grid", + 0x80 => "PV Off-grid", + 0xC0 => "PV & Battery Off-grid", + 0x88 => "PV Charge Off-grid", + + _ => "Unknown", + } + } +} + +struct WarningCodeString; +impl WarningCodeString { + pub fn from_value(value: i64) -> &'static str { + if value == 0 { + return "OK"; + } + + (0..=31) + .find(|i| value & (1 << i) > 0) + .map(Self::from_bit) + .unwrap() + } + + fn from_bit(bit: usize) -> &'static str { + match bit { + 0 => "W000: Battery communication failure", + 1 => "W001: AFCI communication failure", + 2 => "W002: AFCI high", + 3 => "W003: Meter communication failure", + 4 => "W004: Both charge and discharge forbidden by battery", + 5 => "W005: Auto test failed", + 6 => "W006: Reserved", + 7 => "W007: LCD communication failure", + 8 => "W008: FW version mismatch", + 9 => "W009: Fan stuck", + 10 => "W010: Reserved", + 11 => "W011: Parallel number out of range", + 12 => "W012: Bat On Mos", + 13 => "W013: Overtemperature (NTC reading is too high)", + 14 => "W014: Reserved", + 15 => "W015: Battery reverse connection", + 16 => "W016: Grid power outage", + 17 => "W017: Grid voltage out of range", + 18 => "W018: Grid frequency out of range", + 19 => "W019: Reserved", + 20 => "W020: PV insulation low", + 21 => "W021: Leakage current high", + 22 => "W022: DCI high", + 23 => "W023: PV short", + 24 => "W024: Reserved", + 25 => "W025: Battery voltage high", + 26 => "W026: Battery voltage low", + 27 => "W027: Battery open circuit", + 28 => "W028: EPS overload", + 29 => "W029: EPS voltage high", + 30 => "W030: Meter reverse connection", + 31 => "W031: DCV high", + + _ => todo!("Unknown Warning"), + } + } +} + +struct FaultCodeString; +impl FaultCodeString { + pub fn from_value(value: i64) -> &'static str { + if value == 0 { + return "OK"; + } + + (0..=31) + .find(|i| value & (1 << i) > 0) + .map(Self::from_bit) + .unwrap() + } + + fn from_bit(bit: usize) -> &'static str { + match bit { + 0 => "E000: Internal communication fault 1", + 1 => "E001: Model fault", + 2 => "E002: BatOnMosFail", + 3 => "E003: CT Fail", + 4 => "E004: Reserved", + 5 => "E005: Reserved", + 6 => "E006: Reserved", + 7 => "E007: Reserved", + 8 => "E008: CAN communication error in parallel system", + 9 => "E009: master lost in parallel system", + 10 => "E010: multiple master units in parallel system", + 11 => "E011: AC input inconsistent in parallel system", + 12 => "E012: UPS short", + 13 => "E013: Reverse current on UPS output", + 14 => "E014: Bus short", + 15 => "E015: Phase error in three phase system", + 16 => "E016: Relay check fault", + 17 => "E017: Internal communication fault 2", + 18 => "E018: Internal communication fault 3", + 19 => "E019: Bus voltage high", + 20 => "E020: EPS connection fault", + 21 => "E021: PV voltage high", + 22 => "E022: Over current protection", + 23 => "E023: Neutral fault", + 24 => "E024: PV short", + 25 => "E025: Radiator temperature over range", + 26 => "E026: Internal fault", + 27 => "E027: Sample inconsistent between Main CPU and redundant CPU", + 28 => "E028: Reserved", + 29 => "E029: Reserved", + 30 => "E030: Reserved", + 31 => "E031: Internal communication fault 4", + _ => todo!("Unknown Fault"), + } + } +} diff --git a/src/mqtt.rs b/src/mqtt.rs index 6f58c23..a00d301 100644 --- a/src/mqtt.rs +++ b/src/mqtt.rs @@ -16,160 +16,6 @@ pub enum TargetInverter { } impl Message { - pub fn for_param(rp: lxp::packet::ReadParam) -> Result> { - let mut r = Vec::new(); - - for (register, value) in rp.pairs() { - r.push(mqtt::Message { - topic: format!("{}/param/{}", rp.datalog, register), - retain: true, - payload: serde_json::to_string(&value)?, - }); - } - - Ok(r) - } - - pub fn for_hold(td: lxp::packet::TranslatedData) -> Result> { - let mut r = Vec::new(); - - for (register, value) in td.pairs() { - r.push(mqtt::Message { - topic: format!("{}/hold/{}", td.datalog, register), - retain: true, - payload: serde_json::to_string(&value)?, - }); - - if register == 21 { - let bits = lxp::packet::Register21Bits::new(value); - r.push(mqtt::Message { - topic: format!("{}/hold/{}/bits", td.datalog, register), - retain: true, - payload: serde_json::to_string(&bits)?, - }); - } - - if register == 110 { - let bits = lxp::packet::Register110Bits::new(value); - r.push(mqtt::Message { - topic: format!("{}/hold/{}/bits", td.datalog, register), - retain: true, - payload: serde_json::to_string(&bits)?, - }); - } - } - - Ok(r) - } - - pub fn for_input_all( - inputs: &lxp::packet::ReadInputAll, - datalog: lxp::inverter::Serial, - ) -> Result { - Ok(mqtt::Message { - topic: format!("{}/inputs/all", datalog), - retain: false, - payload: serde_json::to_string(&inputs)?, - }) - } - - pub fn for_input( - td: lxp::packet::TranslatedData, - publish_individual: bool, - ) -> Result> { - use lxp::packet::ReadInput; - - let mut r = Vec::new(); - - if publish_individual { - let mut fault_code_registers_seen = false; - let mut fault_code = 0; - let mut warning_code_registers_seen = false; - let mut warning_code = 0; - - for (register, value) in td.pairs() { - r.push(mqtt::Message { - topic: format!("{}/input/{}", td.datalog, register), - retain: false, - payload: serde_json::to_string(&value)?, - }); - - if register == 0 { - r.push(mqtt::Message { - topic: format!("{}/input/{}/parsed", td.datalog, register), - retain: false, - payload: lxp::packet::StatusString::from_value(value).to_owned(), - }); - } - - if register == 60 { - fault_code |= value as u32; - fault_code_registers_seen = true; - } - if register == 61 { - fault_code |= (value as u32) << 16; - fault_code_registers_seen = true; - } - - if register == 62 { - warning_code |= value as u32; - warning_code_registers_seen = true; - } - if register == 63 { - warning_code |= (value as u32) << 16; - warning_code_registers_seen = true; - } - } - - if warning_code_registers_seen { - r.push(mqtt::Message { - topic: format!("{}/input/warning_code/parsed", td.datalog), - retain: false, - payload: lxp::packet::WarningCodeString::from_value(warning_code).to_owned(), - }); - } - - if fault_code_registers_seen { - r.push(mqtt::Message { - topic: format!("{}/input/fault_code/parsed", td.datalog), - retain: false, - payload: lxp::packet::FaultCodeString::from_value(fault_code).to_owned(), - }); - } - } - - match td.read_input() { - Ok(ReadInput::ReadInputAll(r_all)) => r.push(mqtt::Message { - topic: format!("{}/inputs/all", td.datalog), - retain: false, - payload: serde_json::to_string(&r_all)?, - }), - Ok(ReadInput::ReadInput1(r1)) => r.push(mqtt::Message { - topic: format!("{}/inputs/1", td.datalog), - retain: false, - payload: serde_json::to_string(&r1)?, - }), - Ok(ReadInput::ReadInput2(r2)) => r.push(mqtt::Message { - topic: format!("{}/inputs/2", td.datalog), - retain: false, - payload: serde_json::to_string(&r2)?, - }), - Ok(ReadInput::ReadInput3(r3)) => r.push(mqtt::Message { - topic: format!("{}/inputs/3", td.datalog), - retain: false, - payload: serde_json::to_string(&r3)?, - }), - Ok(ReadInput::ReadInput4(r4)) => r.push(mqtt::Message { - topic: format!("{}/inputs/4", td.datalog), - retain: false, - payload: serde_json::to_string(&r4)?, - }), - Err(x) => warn!("ignoring {:?}", x), - } - - Ok(r) - } - pub fn to_command(&self, inverter: config::Inverter) -> Result { use Command::*; diff --git a/src/prelude.rs b/src/prelude.rs index ee3d8e4..364fa1a 100644 --- a/src/prelude.rs +++ b/src/prelude.rs @@ -1,11 +1,14 @@ pub use std::{ cell::{Ref, RefCell, RefMut}, + collections::HashMap, convert::{TryFrom, TryInto}, io::Write, rc::Rc, str::FromStr, }; +pub type RegisterMap = HashMap; + pub use { anyhow::{anyhow, bail, Error, Result}, log::{debug, error, info, trace, warn}, @@ -17,7 +20,7 @@ pub use crate::{ command::Command, config::{self, Config, ConfigWrapper}, coordinator::{self, Coordinator}, - database::{self, Database}, + //database::{self, Database}, home_assistant, influx::{self, Influx}, lxp::{ diff --git a/src/register_cache.rs b/src/register_cache.rs index 460133f..bc24f52 100644 --- a/src/register_cache.rs +++ b/src/register_cache.rs @@ -1,87 +1,138 @@ use crate::prelude::*; -// this just needs to be bigger than the max register we'll see -const REGISTER_COUNT: usize = 256; - #[derive(Clone, Debug)] pub enum ChannelData { - ReadRegister(u16, Rc>>), - RegisterData(u16, u16), + ReadRegister(Register, Rc>>), + ReadAllRegisters(AllRegisters, Rc>>), + RegisterData(Register, u16), + ClearInputRegisters, Shutdown, } +#[derive(Clone, Debug)] +pub enum Register { + Hold(u16), + Input(u16), +} + +#[derive(Clone, Debug)] +pub enum AllRegisters { + Hold, + Input, +} + pub struct RegisterCache { channels: Channels, - register_data: Rc>, + hold_register_data: Rc>, + input_register_data: Rc>, } impl RegisterCache { pub fn new(channels: Channels) -> Self { - let register_data = Rc::new(RefCell::new([0; REGISTER_COUNT])); + let hold_register_data = Rc::new(RefCell::new(RegisterMap::with_capacity(256))); + let input_register_data = Rc::new(RefCell::new(RegisterMap::with_capacity(256))); Self { channels, - register_data, + hold_register_data, + input_register_data, } } pub async fn start(&self) -> Result<()> { - futures::try_join!(self.cache_getter(), self.cache_setter())?; + futures::try_join!(self.runner())?; Ok(()) } // external helper method to simplify access to the cache, use like so: - // // RegisterCache::get(&self.channels, 1); - // - pub async fn get(channels: &Channels, register: u16) -> u16 { + pub async fn get(channels: &Channels, register: Register) -> u16 { let (tx, rx) = oneshot::channel(); + let channel_data = ChannelData::ReadRegister(register, Rc::new(RefCell::new(tx))); - let _ = channels.read_register_cache.send(channel_data); + let _ = channels.register_cache.send(channel_data); rx.await .expect("unexpected error reading from register cache") } - async fn cache_getter(&self) -> Result<()> { - let mut receiver = self.channels.read_register_cache.subscribe(); - - info!("register_cache getter starting"); - - while let ChannelData::ReadRegister(register, reply_tx) = receiver.recv().await? { - if register < REGISTER_COUNT as u16 { - let register_data = self.register_data.borrow(); - let value = register_data[register as usize]; - - let reply_tx = Rc::try_unwrap(reply_tx).unwrap(); - let _ = reply_tx.into_inner().send(value); - } else { - warn!( - "cannot cache register {}, increase REGISTER_COUNT!", - register - ); - } - } - - info!("register_cache getter exiting"); + // RegisterCache::dump(&self.channels, register_cache::AllRegisters::Input) + pub async fn dump(channels: &Channels, register_type: AllRegisters) -> RegisterMap { + let (tx, rx) = oneshot::channel(); - Ok(()) + let channel_data = ChannelData::ReadAllRegisters(register_type, Rc::new(RefCell::new(tx))); + let _ = channels.register_cache.send(channel_data); + rx.await + .expect("unexpected error reading from register cache") } - async fn cache_setter(&self) -> Result<()> { - let mut receiver = self.channels.to_register_cache.subscribe(); - - info!("register_cache setter starting"); + // RegisterCache::dump(&self.channels, register_cache::AllRegisters::Input) + pub async fn clear(channels: &Channels) { + let _ = channels + .register_cache + .send(ChannelData::ClearInputRegisters); + } - while let ChannelData::RegisterData(register, value) = receiver.recv().await? { - if register < REGISTER_COUNT as u16 { - let mut register_data = self.register_data.borrow_mut(); - register_data[register as usize] = value; - } else { - warn!( - "cannot cache register {}, increase REGISTER_COUNT!", - register - ); + async fn runner(&self) -> Result<()> { + let mut receiver = self.channels.register_cache.subscribe(); + + info!("register_cache runner starting"); + + loop { + match receiver.recv().await? { + ChannelData::RegisterData(register, value) => { + // debug!("register_cache setting {:?}={}", register, value); + match register { + Register::Hold(r) => { + let mut register_data = self.hold_register_data.borrow_mut(); + register_data.insert(r, value); + } + Register::Input(r) => { + let mut register_data = self.input_register_data.borrow_mut(); + register_data.insert(r, value); + } + }; + } + + ChannelData::ClearInputRegisters => { + // not used yet + let mut register_data = self.input_register_data.borrow_mut(); + register_data.clear(); + } + + ChannelData::ReadAllRegisters(register_type, reply_tx) => match register_type { + AllRegisters::Hold => { + let register_data = self.hold_register_data.borrow().clone(); + let reply_tx = Rc::try_unwrap(reply_tx).unwrap(); + let _ = reply_tx.into_inner().send(register_data); + } + AllRegisters::Input => { + let register_data = self.input_register_data.borrow().clone(); + let reply_tx = Rc::try_unwrap(reply_tx).unwrap(); + let _ = reply_tx.into_inner().send(register_data); + } + }, + + ChannelData::ReadRegister(register, reply_tx) => { + match register { + Register::Hold(r) => { + let register_data = self.hold_register_data.borrow(); + let value = register_data.get(&r).cloned().unwrap_or(0); + + let reply_tx = Rc::try_unwrap(reply_tx).unwrap(); + let _ = reply_tx.into_inner().send(value); + } + Register::Input(r) => { + let register_data = self.input_register_data.borrow(); + let value = register_data.get(&r).cloned().unwrap_or(0); + + let reply_tx = Rc::try_unwrap(reply_tx).unwrap(); + let _ = reply_tx.into_inner().send(value); + } + }; + } + + ChannelData::Shutdown => break, } } diff --git a/src/utils.rs b/src/utils.rs index 7297864..45e2963 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -1,5 +1,3 @@ -use crate::prelude::*; - // 2022-03-04 05:06:07 hardcoded time for tests #[allow(dead_code)] const HARDCODED_TEST_TIME: i64 = 1646370367; @@ -15,28 +13,6 @@ impl Utils { u16::from_le_bytes([array[offset], array[offset + 1]]) } - 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)) - } - - pub fn le_u32_div10(input: &[u8]) -> nom::IResult<&[u8], f64> { - let (input, num) = nom::number::complete::le_u32(input)?; - Ok((input, num as f64 / 10.0)) - } - - pub fn current_time_for_nom(input: &[u8]) -> nom::IResult<&[u8], UnixTime> { - Ok((input, UnixTime::now())) - } - #[cfg(not(feature = "mocks"))] pub fn utc() -> chrono::DateTime { chrono::Utc::now()