Skip to content
Draft
Show file tree
Hide file tree
Changes from all 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
4 changes: 2 additions & 2 deletions src/channels.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ pub struct Channels {
pub from_mqtt: broadcast::Sender<mqtt::ChannelData>,
pub to_mqtt: broadcast::Sender<mqtt::ChannelData>,
pub to_influx: broadcast::Sender<influx::ChannelData>,
//pub to_database: broadcast::Sender<database::ChannelData>,
pub to_database: broadcast::Sender<database::ChannelData>,
pub register_cache: broadcast::Sender<register_cache::ChannelData>,
}

Expand All @@ -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(),
}
}
Expand Down
4 changes: 2 additions & 2 deletions src/coordinator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -572,9 +572,9 @@ impl Coordinator {
}

/*
async fn save_input_all(&self, input: Box<lxp::packet::ReadInputAll>) -> 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?");
}
Expand Down
190 changes: 96 additions & 94 deletions src/database.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use sqlx::{any::AnyConnectOptions, AnyPool, ConnectOptions};

#[derive(PartialEq, Clone, Debug)]
pub enum ChannelData {
ReadInputAll(Box<lxp::packet::ReadInputAll>),
InputData(lxp::register_parser::ParsedData),
Shutdown,
}

Expand Down Expand Up @@ -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;
Expand All @@ -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?;
Expand Down
1 change: 1 addition & 0 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
15 changes: 15 additions & 0 deletions src/lxp/register_parser.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand All @@ -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)]
Expand Down
2 changes: 1 addition & 1 deletion src/prelude.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down