From 267ef16478912564000f891be44ff9a0f587782f Mon Sep 17 00:00:00 2001 From: Richard Russell <2265225+rars@users.noreply.github.com> Date: Wed, 1 Jul 2026 15:26:01 +0100 Subject: [PATCH] refactor: add london date column to help perform aggregations in SQLite DB --- .../down.sql | 6 + .../up.sql | 6 + src-tauri/src/data/consumption.rs | 198 +++++++++++------- src-tauri/src/db.rs | 84 +++++++- src-tauri/src/main.rs | 2 + src-tauri/src/schema.rs | 2 + src-tauri/src/utils.rs | 24 ++- 7 files changed, 240 insertions(+), 82 deletions(-) create mode 100644 src-tauri/migrations/2026-07-01-124959-0000_add_london_date/down.sql create mode 100644 src-tauri/migrations/2026-07-01-124959-0000_add_london_date/up.sql diff --git a/src-tauri/migrations/2026-07-01-124959-0000_add_london_date/down.sql b/src-tauri/migrations/2026-07-01-124959-0000_add_london_date/down.sql new file mode 100644 index 0000000..9a9fb85 --- /dev/null +++ b/src-tauri/migrations/2026-07-01-124959-0000_add_london_date/down.sql @@ -0,0 +1,6 @@ +-- This file should undo anything in `up.sql` +DROP INDEX IF EXISTS idx_electricity_consumption_london_date_id; +DROP INDEX IF EXISTS idx_gas_consumption_london_date_id; + +ALTER TABLE electricity_consumption DROP COLUMN london_date_id; +ALTER TABLE gas_consumption DROP COLUMN london_date_id; diff --git a/src-tauri/migrations/2026-07-01-124959-0000_add_london_date/up.sql b/src-tauri/migrations/2026-07-01-124959-0000_add_london_date/up.sql new file mode 100644 index 0000000..0e86de9 --- /dev/null +++ b/src-tauri/migrations/2026-07-01-124959-0000_add_london_date/up.sql @@ -0,0 +1,6 @@ + +ALTER TABLE electricity_consumption ADD COLUMN london_date_id INTEGER; +ALTER TABLE gas_consumption ADD COLUMN london_date_id INTEGER; + +CREATE INDEX idx_electricity_consumption_london_date_id ON electricity_consumption(london_date_id); +CREATE INDEX idx_gas_consumption_london_date_id ON gas_consumption(london_date_id); diff --git a/src-tauri/src/data/consumption.rs b/src-tauri/src/data/consumption.rs index fdb2492..7d4beff 100644 --- a/src-tauri/src/data/consumption.rs +++ b/src-tauri/src/data/consumption.rs @@ -1,8 +1,7 @@ -use std::collections::BTreeMap; use std::sync::{Arc, Mutex, MutexGuard}; -use chrono::{Datelike, NaiveDate, NaiveDateTime}; -use chrono_tz::Europe::London; +use chrono::{NaiveDate, NaiveDateTime}; +use diesel::dsl::sql; use diesel::insert_into; use diesel::SqliteConnection; use diesel::{prelude::*, upsert::excluded}; @@ -11,11 +10,13 @@ use rust_decimal::prelude::ToPrimitive; use rust_decimal::Decimal; use crate::schema::{electricity_consumption, gas_consumption}; -use crate::utils::london_midnight_as_utc; +use crate::utils::london_date_id_to_naive_date; +use crate::utils::{ + london_midnight_as_utc, naive_date_to_london_date_id, utc_timestamp_to_london_date_id, +}; use super::RepositoryError; -const ENERGY_CONSUMPTION_KWH_ERROR_CODE: f64 = 16777.215f64; const ENERGY_CONSUMPTION_WH_ERROR_CODE: i64 = 16777215i64; const KWH_TO_WH_SCALE: Decimal = Decimal::ONE_THOUSAND; @@ -35,6 +36,7 @@ pub struct GasConsumptionValue { struct NewElectricityConsumption { timestamp: NaiveDateTime, energy_consumption_wh: i64, + london_date_id: i32, } #[derive(Insertable)] @@ -42,6 +44,7 @@ struct NewElectricityConsumption { struct NewGasConsumption { timestamp: NaiveDateTime, energy_consumption_wh: i64, + london_date_id: i32, } #[derive(Queryable)] @@ -49,6 +52,7 @@ pub struct ElectricityConsumptionRecord { pub electricity_consumption_id: i32, pub timestamp: NaiveDateTime, pub energy_consumption_wh: i64, + pub london_date_id: Option, } #[derive(Queryable)] @@ -56,61 +60,11 @@ pub struct GasConsumptionRecord { pub gas_consumption_id: i32, pub timestamp: NaiveDateTime, pub energy_consumption_wh: i64, + pub london_date_id: Option, } type RepositoryResult = Result; -fn group_raw_by_london_day( - raw_data_utc: &Vec, - proj_date_time: F, - proj_energy: G, -) -> Vec<(NaiveDate, i64)> -where - F: Fn(&T) -> &NaiveDateTime, - G: Fn(&T) -> i64, -{ - let mut energy_by_day: BTreeMap = BTreeMap::new(); - - for elem in raw_data_utc { - let date_time = proj_date_time(elem); - let energy = proj_energy(elem); - - let london_date = date_time.and_utc().with_timezone(&London).date_naive(); - - *energy_by_day.entry(london_date).or_insert(0) += energy; - } - - energy_by_day.into_iter().collect() -} - -fn group_raw_by_london_month( - raw_data_utc: &Vec, - proj_date_time: F, - proj_energy: G, -) -> Vec<(NaiveDate, i64)> -where - F: Fn(&T) -> &NaiveDateTime, - G: Fn(&T) -> i64, -{ - let mut energy_by_month: BTreeMap = BTreeMap::new(); - - for elem in raw_data_utc { - let date_time = proj_date_time(elem); - let energy = proj_energy(elem); - - let london_first_of_month = date_time - .and_utc() - .with_timezone(&London) - .date_naive() - .with_day(1) - .expect("Every month should have 1st of the month"); - - *energy_by_month.entry(london_first_of_month).or_insert(0) += energy; - } - - energy_by_month.into_iter().collect() -} - pub trait ConsumptionRepository { fn insert(&self, records: Vec) -> RepositoryResult<()>; @@ -156,6 +110,7 @@ impl ConsumptionRepository RepositoryResult> { - let raw_data = self.get_raw(start, end)?; + use crate::schema::electricity_consumption::dsl::*; + + let mut conn = self.get_connection()?; - Ok(group_raw_by_london_day( - &raw_data, - |x| &x.timestamp, - |x| x.energy_consumption_wh, - )) + let start_london_date_id = naive_date_to_london_date_id(&start); + let end_london_date_id = naive_date_to_london_date_id(&end); + + let daily_consumption = electricity_consumption + .filter(london_date_id.is_not_null()) + .filter(london_date_id.ge(start_london_date_id)) + .filter(london_date_id.lt(end_london_date_id)) + .select(( + london_date_id.assume_not_null(), + sql::("COALESCE(SUM(energy_consumption_wh), 0)"), + )) + .group_by(london_date_id) + .order(london_date_id) + .load::<(i32, i64)>(&mut *conn)?; + + Ok(daily_consumption + .iter() + .map(|(date_id, energy)| { + let date = london_date_id_to_naive_date(*date_id); + (date, *energy) + }) + .collect()) } fn get_monthly( @@ -213,13 +187,35 @@ impl ConsumptionRepository RepositoryResult> { - let raw_data = self.get_raw(start, end)?; + use crate::schema::electricity_consumption::dsl::*; + + let mut conn = self.get_connection()?; - Ok(group_raw_by_london_month( - &raw_data, - |x| &x.timestamp, - |x| x.energy_consumption_wh, - )) + let start_london_date_id = naive_date_to_london_date_id(&start); + let end_london_date_id = naive_date_to_london_date_id(&end); + + let london_month_id = + sql::("london_date_id - (london_date_id % 100) + 1"); + + let monthly_consumption = electricity_consumption + .filter(london_date_id.is_not_null()) + .filter(london_date_id.ge(start_london_date_id)) + .filter(london_date_id.lt(end_london_date_id)) + .select(( + london_month_id.clone().assume_not_null(), + sql::("COALESCE(SUM(energy_consumption_wh), 0)"), + )) + .group_by(london_month_id.clone()) + .order(london_month_id) + .load::<(i32, i64)>(&mut *conn)?; + + Ok(monthly_consumption + .iter() + .map(|(date_id, energy)| { + let date = london_date_id_to_naive_date(*date_id); + (date, *energy) + }) + .collect()) } } @@ -250,6 +246,7 @@ impl ConsumptionRepository energy_consumption_wh: (x.value * KWH_TO_WH_SCALE) .to_i64() .expect("Gas consumption to fit in i64"), + london_date_id: utc_timestamp_to_london_date_id(&x.timestamp), }) .collect(); @@ -298,13 +295,32 @@ impl ConsumptionRepository start: NaiveDate, end: NaiveDate, ) -> RepositoryResult> { - let raw_data = self.get_raw(start, end)?; + use crate::schema::gas_consumption::dsl::*; + + let mut conn = self.get_connection()?; - Ok(group_raw_by_london_day( - &raw_data, - |x| &x.timestamp, - |x| x.energy_consumption_wh, - )) + let start_london_date_id = naive_date_to_london_date_id(&start); + let end_london_date_id = naive_date_to_london_date_id(&end); + + let daily_consumption = gas_consumption + .filter(london_date_id.is_not_null()) + .filter(london_date_id.ge(start_london_date_id)) + .filter(london_date_id.lt(end_london_date_id)) + .select(( + london_date_id.assume_not_null(), + sql::("COALESCE(SUM(energy_consumption_wh), 0)"), + )) + .group_by(london_date_id) + .order(london_date_id) + .load::<(i32, i64)>(&mut *conn)?; + + Ok(daily_consumption + .iter() + .map(|(date_id, energy)| { + let date = london_date_id_to_naive_date(*date_id); + (date, *energy) + }) + .collect()) } fn get_monthly( @@ -312,12 +328,34 @@ impl ConsumptionRepository start: NaiveDate, end: NaiveDate, ) -> RepositoryResult> { - let raw_data = self.get_raw(start, end)?; + use crate::schema::gas_consumption::dsl::*; + + let mut conn = self.get_connection()?; - Ok(group_raw_by_london_month( - &raw_data, - |x| &x.timestamp, - |x| x.energy_consumption_wh, - )) + let start_london_date_id = naive_date_to_london_date_id(&start); + let end_london_date_id = naive_date_to_london_date_id(&end); + + let london_month_id = + sql::("london_date_id - (london_date_id % 100) + 1"); + + let monthly_consumption = gas_consumption + .filter(london_date_id.is_not_null()) + .filter(london_date_id.ge(start_london_date_id)) + .filter(london_date_id.lt(end_london_date_id)) + .select(( + london_month_id.clone().assume_not_null(), + sql::("COALESCE(SUM(energy_consumption_wh), 0)"), + )) + .group_by(london_month_id.clone()) + .order(london_month_id) + .load::<(i32, i64)>(&mut *conn)?; + + Ok(monthly_consumption + .iter() + .map(|(date_id, energy)| { + let date = london_date_id_to_naive_date(*date_id); + (date, *energy) + }) + .collect()) } } diff --git a/src-tauri/src/db.rs b/src-tauri/src/db.rs index 6e56214..66073e3 100644 --- a/src-tauri/src/db.rs +++ b/src-tauri/src/db.rs @@ -1,5 +1,15 @@ -use diesel::{sqlite::SqliteConnection, Connection}; +use chrono_tz::Europe::London; +use diesel::prelude::*; +use diesel::{sqlite::SqliteConnection, Connection, ExpressionMethods}; use diesel_migrations::{embed_migrations, EmbeddedMigrations, MigrationHarness}; +use log::info; + +use crate::data::RepositoryError; +use crate::utils::utc_timestamp_to_london_date_id; +use crate::{ + schema::electricity_consumption::{london_date_id, table}, + AppError, +}; const MIGRATIONS: EmbeddedMigrations = embed_migrations!("migrations"); @@ -19,3 +29,75 @@ pub fn establish_connection(database_url: &str) -> SqliteConnection { SqliteConnection::establish(&database_url) .unwrap_or_else(|_| panic!("Error connecting to {}", database_url)) } + +pub fn populate_missing_london_date_ids( + conn: &mut SqliteConnection, +) -> Result<(), RepositoryError> { + const BATCH_SIZE: i64 = 5000; + + { + use crate::schema::electricity_consumption::dsl::*; + + loop { + let records = electricity_consumption + .filter(london_date_id.is_null()) + .limit(BATCH_SIZE) + .load::(conn)?; + + if records.is_empty() { + break; + } + + info!( + "Populating london_date_id for {} electricity consumptionrecords", + records.len() + ); + + conn.transaction::<_, diesel::result::Error, _>(|transaction_conn| { + for row in records { + let date_id = utc_timestamp_to_london_date_id(&row.timestamp); + + diesel::update(electricity_consumption.find(row.electricity_consumption_id)) + .set(london_date_id.eq(date_id)) + .execute(transaction_conn)?; + } + + Ok(()) + })?; + } + } + + { + use crate::schema::gas_consumption::dsl::*; + + loop { + let records = gas_consumption + .filter(london_date_id.is_null()) + .limit(BATCH_SIZE) + .load::(conn)?; + + if records.is_empty() { + break; + } + + info!( + "Populating london_date_id for {} gas consumption records", + records.len() + ); + + conn.transaction::<_, diesel::result::Error, _>(|transaction_conn| { + for row in records { + let date_id = utc_timestamp_to_london_date_id(&row.timestamp); + + diesel::update(gas_consumption.find(row.gas_consumption_id)) + .set(london_date_id.eq(date_id)) + .execute(transaction_conn)?; + } + + Ok(()) + })?; + } + } + + Ok(()) +} diff --git a/src-tauri/src/main.rs b/src-tauri/src/main.rs index 202a4a0..e7d1103 100644 --- a/src-tauri/src/main.rs +++ b/src-tauri/src/main.rs @@ -22,6 +22,7 @@ use commands::glowmarkt::*; use commands::mqtt::*; use commands::profiles::*; +use crate::db::populate_missing_london_date_ids; use crate::mqtt::start_mqtt_listener; use crate::utils::MqttSettings; use crate::utils::{get_mqtt_settings_opt, MqttAppSettings}; @@ -142,6 +143,7 @@ fn main() { let mut connection = db::establish_connection(db_path.to_str().expect("db path needed")); db::run_migrations(&mut connection); + populate_missing_london_date_ids(&mut connection)?; let store = app.store(SETTINGS_FILE)?; diff --git a/src-tauri/src/schema.rs b/src-tauri/src/schema.rs index f21f625..54eeb9f 100644 --- a/src-tauri/src/schema.rs +++ b/src-tauri/src/schema.rs @@ -5,6 +5,7 @@ diesel::table! { electricity_consumption_id -> Integer, timestamp -> Timestamp, energy_consumption_wh -> BigInt, + london_date_id -> Nullable, } } @@ -49,6 +50,7 @@ diesel::table! { gas_consumption_id -> Integer, timestamp -> Timestamp, energy_consumption_wh -> BigInt, + london_date_id -> Nullable, } } diff --git a/src-tauri/src/utils.rs b/src-tauri/src/utils.rs index 43445b7..0bef73c 100644 --- a/src-tauri/src/utils.rs +++ b/src-tauri/src/utils.rs @@ -1,6 +1,6 @@ use std::sync::{Arc, Mutex}; -use chrono::{NaiveDate, NaiveDateTime, TimeZone, Utc}; +use chrono::{Datelike, NaiveDate, NaiveDateTime, TimeZone, Utc}; use chrono_tz::Europe::London; use diesel::SqliteConnection; use keyring_core::Entry; @@ -29,6 +29,28 @@ pub fn london_midnight_as_utc(date: &NaiveDate) -> NaiveDateTime { london_midnight_utc.naive_utc() } +pub fn utc_timestamp_to_london_date_id(timestamp_utc: &NaiveDateTime) -> i32 { + let london_time = timestamp_utc.and_utc().with_timezone(&London); + + london_time + .format("%Y%m%d") + .to_string() + .parse::() + .unwrap() +} + +pub fn naive_date_to_london_date_id(date: &NaiveDate) -> i32 { + (date.year() * 10000) + (date.month() as i32 * 100) + (date.day() as i32) +} + +pub fn london_date_id_to_naive_date(date_id: i32) -> NaiveDate { + let year = date_id / 10000; + let month = ((date_id % 10000) / 100) as u32; + let day = (date_id % 100) as u32; + + NaiveDate::from_ymd_opt(year, month, day).expect("Invalid date_id in the database") +} + pub fn emit_event(app_handle: &AppHandle, event: &str, payload: T) -> Result<(), AppError> where T: Serialize + Clone,