From 25e7d605252cce553d0b08035c882641a0d52be9 Mon Sep 17 00:00:00 2001 From: Chris Elsworth Date: Fri, 15 Mar 2024 10:18:20 +0000 Subject: [PATCH] Port database insertion over to new register_parser --- src/channels.rs | 4 +- src/coordinator/mod.rs | 4 +- src/database.rs | 190 +++++++++++++++++++------------------ src/lib.rs | 1 + src/lxp/register_parser.rs | 15 +++ src/prelude.rs | 2 +- 6 files changed, 117 insertions(+), 99 deletions(-) diff --git a/src/channels.rs b/src/channels.rs index 0816f09..6916975 100644 --- a/src/channels.rs +++ b/src/channels.rs @@ -7,7 +7,7 @@ pub struct Channels { pub from_mqtt: broadcast::Sender, pub to_mqtt: broadcast::Sender, pub to_influx: broadcast::Sender, - //pub to_database: broadcast::Sender, + pub to_database: broadcast::Sender, pub register_cache: broadcast::Sender, } @@ -25,7 +25,7 @@ impl Channels { from_mqtt: Self::channel(), to_mqtt: Self::channel(), to_influx: Self::channel(), - //to_database: Self::channel(), + to_database: Self::channel(), register_cache: Self::channel(), } } diff --git a/src/coordinator/mod.rs b/src/coordinator/mod.rs index f285176..0a23ffe 100644 --- a/src/coordinator/mod.rs +++ b/src/coordinator/mod.rs @@ -572,9 +572,9 @@ impl Coordinator { } /* - async fn save_input_all(&self, input: Box) -> Result<()> { + async fn save_input_all(&self, input: &lxp::register_parser::ParsedData) -> Result<()> { if self.config.have_enabled_database() { - let channel_data = database::ChannelData::ReadInputAll(input); + let channel_data = database::ChannelData::InputData(input); if self.channels.to_database.send(channel_data).is_err() { bail!("send(to_database) failed - channel closed?"); } diff --git a/src/database.rs b/src/database.rs index 0fd286f..b088dfc 100644 --- a/src/database.rs +++ b/src/database.rs @@ -4,7 +4,7 @@ use sqlx::{any::AnyConnectOptions, AnyPool, ConnectOptions}; #[derive(PartialEq, Clone, Debug)] pub enum ChannelData { - ReadInputAll(Box), + InputData(lxp::register_parser::ParsedData), Shutdown, } @@ -158,7 +158,7 @@ impl Database { match receiver.recv().await? { Shutdown => break, - ReadInputAll(data) => { + InputData(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; @@ -170,101 +170,103 @@ impl Database { Ok(()) } - async fn insert(&self, query: &str, data: &lxp::packet::ReadInputAll) -> Result<()> { + async fn insert(&self, query: &str, data: &lxp::register_parser::ParsedData) -> Result<()> { let mut conn = self.connection().await?; sqlx::query(query) - .bind(data.status as i32) - .bind(data.v_pv_1) - .bind(data.v_pv_2) - .bind(data.v_pv_3) - .bind(data.v_bat) - .bind(data.soc as i16) - .bind(data.soh as i16) - .bind(data.internal_fault as i32) - .bind(data.p_pv as i32) - .bind(data.p_pv_1 as i32) - .bind(data.p_pv_2 as i32) - .bind(data.p_pv_3 as i32) - .bind(data.p_battery) - .bind(data.p_charge as i32) - .bind(data.p_discharge as i32) - .bind(data.v_ac_r) - .bind(data.v_ac_s) - .bind(data.v_ac_t) - .bind(data.f_ac) - .bind(data.p_inv as i32) - .bind(data.p_rec as i32) - .bind(data.pf) - .bind(data.v_eps_r) - .bind(data.v_eps_s) - .bind(data.v_eps_t) - .bind(data.f_eps) - .bind(data.p_eps as i32) - .bind(data.s_eps as i32) - .bind(data.p_grid) - .bind(data.p_to_grid as i32) - .bind(data.p_to_user as i32) - .bind(data.e_pv_day) - .bind(data.e_pv_day_1) - .bind(data.e_pv_day_2) - .bind(data.e_pv_day_3) - .bind(data.e_inv_day) - .bind(data.e_rec_day) - .bind(data.e_chg_day) - .bind(data.e_dischg_day) - .bind(data.e_eps_day) - .bind(data.e_to_grid_day) - .bind(data.e_to_user_day) - .bind(data.v_bus_1) - .bind(data.v_bus_2) - .bind(data.e_pv_all) - .bind(data.e_pv_all_1) - .bind(data.e_pv_all_2) - .bind(data.e_pv_all_3) - .bind(data.e_inv_all) - .bind(data.e_rec_all) - .bind(data.e_chg_all) - .bind(data.e_dischg_all) - .bind(data.e_eps_all) - .bind(data.e_to_grid_all) - .bind(data.e_to_user_all) - .bind(data.fault_code as i64) - .bind(data.warning_code as i64) - .bind(data.t_inner as i32) - .bind(data.t_rad_1 as i32) - .bind(data.t_rad_2 as i32) - .bind(data.t_bat as i32) - .bind(data.runtime as i32) // TODO - .bind(data.max_chg_curr) - .bind(data.max_dischg_curr) - .bind(data.charge_volt_ref) - .bind(data.dischg_cut_volt) - .bind(data.bat_status_0 as i32) - .bind(data.bat_status_1 as i32) - .bind(data.bat_status_2 as i32) - .bind(data.bat_status_3 as i32) - .bind(data.bat_status_4 as i32) - .bind(data.bat_status_5 as i32) - .bind(data.bat_status_6 as i32) - .bind(data.bat_status_7 as i32) - .bind(data.bat_status_8 as i32) - .bind(data.bat_status_9 as i32) - .bind(data.bat_status_inv as i32) - .bind(data.bat_count as i32) - .bind(data.bat_capacity as i32) - .bind(data.bat_current) - .bind(data.bms_event_1 as i32) - .bind(data.bms_event_2 as i32) - .bind(data.max_cell_voltage) - .bind(data.min_cell_voltage) - .bind(data.max_cell_temp) - .bind(data.min_cell_temp) - .bind(data.bms_fw_update_state as i32) - .bind(data.cycle_count as i32) - .bind(data.vbat_inv) - .bind(data.datalog.to_string()) - .bind(data.time.0) + .bind(data.get("status").unwrap().unwrap_i64()) + .bind(data.get("v_pv_1").unwrap().unwrap_f64()) + .bind(data.get("v_pv_2").unwrap().unwrap_f64()) + .bind(data.get("v_pv_3").unwrap().unwrap_f64()) + .bind(data.get("v_bat").unwrap().unwrap_f64()) + .bind(data.get("soc").unwrap().unwrap_i64()) + .bind(data.get("soh").unwrap().unwrap_i64()) + /* + .bind(data.internal_fault as i32) + .bind(data.p_pv as i32) + .bind(data.p_pv_1 as i32) + .bind(data.p_pv_2 as i32) + .bind(data.p_pv_3 as i32) + .bind(data.p_battery) + .bind(data.p_charge as i32) + .bind(data.p_discharge as i32) + .bind(data.v_ac_r) + .bind(data.v_ac_s) + .bind(data.v_ac_t) + .bind(data.f_ac) + .bind(data.p_inv as i32) + .bind(data.p_rec as i32) + .bind(data.pf) + .bind(data.v_eps_r) + .bind(data.v_eps_s) + .bind(data.v_eps_t) + .bind(data.f_eps) + .bind(data.p_eps as i32) + .bind(data.s_eps as i32) + .bind(data.p_grid) + .bind(data.p_to_grid as i32) + .bind(data.p_to_user as i32) + .bind(data.e_pv_day) + .bind(data.e_pv_day_1) + .bind(data.e_pv_day_2) + .bind(data.e_pv_day_3) + .bind(data.e_inv_day) + .bind(data.e_rec_day) + .bind(data.e_chg_day) + .bind(data.e_dischg_day) + .bind(data.e_eps_day) + .bind(data.e_to_grid_day) + .bind(data.e_to_user_day) + .bind(data.v_bus_1) + .bind(data.v_bus_2) + .bind(data.e_pv_all) + .bind(data.e_pv_all_1) + .bind(data.e_pv_all_2) + .bind(data.e_pv_all_3) + .bind(data.e_inv_all) + .bind(data.e_rec_all) + .bind(data.e_chg_all) + .bind(data.e_dischg_all) + .bind(data.e_eps_all) + .bind(data.e_to_grid_all) + .bind(data.e_to_user_all) + .bind(data.fault_code as i64) + .bind(data.warning_code as i64) + .bind(data.t_inner as i32) + .bind(data.t_rad_1 as i32) + .bind(data.t_rad_2 as i32) + .bind(data.t_bat as i32) + .bind(data.runtime as i32) // TODO + .bind(data.max_chg_curr) + .bind(data.max_dischg_curr) + .bind(data.charge_volt_ref) + .bind(data.dischg_cut_volt) + .bind(data.bat_status_0 as i32) + .bind(data.bat_status_1 as i32) + .bind(data.bat_status_2 as i32) + .bind(data.bat_status_3 as i32) + .bind(data.bat_status_4 as i32) + .bind(data.bat_status_5 as i32) + .bind(data.bat_status_6 as i32) + .bind(data.bat_status_7 as i32) + .bind(data.bat_status_8 as i32) + .bind(data.bat_status_9 as i32) + .bind(data.bat_status_inv as i32) + .bind(data.bat_count as i32) + .bind(data.bat_capacity as i32) + .bind(data.bat_current) + .bind(data.bms_event_1 as i32) + .bind(data.bms_event_2 as i32) + .bind(data.max_cell_voltage) + .bind(data.min_cell_voltage) + .bind(data.max_cell_temp) + .bind(data.min_cell_temp) + .bind(data.bms_fw_update_state as i32) + .bind(data.cycle_count as i32) + .bind(data.vbat_inv) + .bind(data.datalog.to_string()) + .bind(data.time.0) + */ .persistent(true) .fetch_optional(&mut conn) .await?; diff --git a/src/lib.rs b/src/lib.rs index a0181eb..3ef064b 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -2,6 +2,7 @@ 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; diff --git a/src/lxp/register_parser.rs b/src/lxp/register_parser.rs index ce65a9d..7dc3cb8 100644 --- a/src/lxp/register_parser.rs +++ b/src/lxp/register_parser.rs @@ -16,6 +16,7 @@ pub enum Value { } impl Value { + // must be a neater way to do this pub fn to_string(&self) -> String { match self { Self::Integer(i) => i.to_string(), @@ -24,6 +25,20 @@ impl Value { Self::StringOwned(_, s) => s.to_string(), } } + + pub fn unwrap_i64(&self) -> i64 { + match self { + Self::Integer(i) => *i, + _ => panic!("not an Integer"), + } + } + + pub fn unwrap_f64(&self) -> f64 { + match self { + Self::Float(f) => *f, + _ => panic!("not a Float"), + } + } } #[derive(Debug, Serialize)] diff --git a/src/prelude.rs b/src/prelude.rs index 77faee2..303a9ae 100644 --- a/src/prelude.rs +++ b/src/prelude.rs @@ -18,7 +18,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::{