Skip to content
Open
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,4 +1,11 @@
# Unreleased
# 0.99.0 (road-to-1-0 branch)

* Cache Hold/Input registers internally as they're seen for later use (#248)
* Remove publish_individual_input configuration option (now they always are) (#250)



# Unreleased (rolled into 0.99.0)

* Reconnect to inverter after 15 minutes of not receiving any data (#223)
* Fix max/min cell temperature/voltage decoding as reported from BMS (#227)
Expand Down
6 changes: 0 additions & 6 deletions src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,8 +109,6 @@ pub struct Mqtt {

#[serde(default = "Config::default_mqtt_homeassistant")]
pub homeassistant: HomeAssistant,

pub publish_individual_input: Option<bool>,
}
impl Mqtt {
pub fn enabled(&self) -> bool {
Expand Down Expand Up @@ -140,10 +138,6 @@ impl Mqtt {
pub fn homeassistant(&self) -> &HomeAssistant {
&self.homeassistant
}

pub fn publish_individual_input(&self) -> bool {
self.publish_individual_input == Some(true)
}
} // }}}

// Influx {{{
Expand Down
57 changes: 43 additions & 14 deletions src/coordinator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -355,6 +355,36 @@ impl Coordinator {
warn!("got a Param packet! {:?}", td);
}

match td.device_function {
DeviceFunction::WriteMulti => {}
DeviceFunction::ReadInput => {
for (register, value) in td.pairs() {
self.cache_register(register_cache::Register::Input(register), value)?;
}

// temp bodge to get parsed messages on MQTT
if self.config.mqtt().enabled() {
let parser = lxp::register_parser::ParseInputs::new(td.pairs());
for (key, parsed_value) in parser.parse_inputs()? {
let m = mqtt::Message {
topic: format!("{}/input/{}", td.datalog, key),
retain: false,
payload: parsed_value.to_string(),
};
let channel_data = mqtt::ChannelData::Message(m);
if self.channels.to_mqtt.send(channel_data).is_err() {
bail!("send(to_mqtt) failed - channel closed?");
}
}
}
}
DeviceFunction::ReadHold | DeviceFunction::WriteSingle => {
for (register, value) in td.pairs() {
self.cache_register(register_cache::Register::Hold(register), value)?;
}
}
}

// 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.
Expand Down Expand Up @@ -393,22 +423,14 @@ impl Coordinator {
}
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?");
}
}
}

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()) {
match Self::packet_to_messages(packet) {
Ok(messages) => {
for message in messages {
let message = mqtt::ChannelData::Message(message);
Expand Down Expand Up @@ -502,20 +524,27 @@ impl Coordinator {
Ok(())
}

fn packet_to_messages(
packet: Packet,
publish_individual_input: bool,
) -> Result<Vec<mqtt::Message>> {
fn packet_to_messages(packet: Packet) -> Result<Vec<mqtt::Message>> {
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::ReadInput => mqtt::Message::for_input(td),
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
}
}

fn cache_register(&self, register: register_cache::Register, value: u16) -> Result<()> {
let channel_data = register_cache::ChannelData::RegisterData(register, value);

if self.channels.to_register_cache.send(channel_data).is_err() {
bail!("send(to_register_cache) failed - channel closed?");
}

Ok(())
}
}
1 change: 1 addition & 0 deletions src/lxp/mod.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
pub mod inverter;
pub mod packet;
pub mod packet_decoder;
pub mod register_parser;
Loading