use anyhow::{Context, Result}; use diesel::{dsl::sum, QueryDsl, RunQueryDsl, SelectableHelper, SqliteConnection}; use schist_models::Pipe; use schist_schema::schema::pipes::{self as pipes_schema, dsl::pipes as pipes_table}; pub fn delete_all_pipes(connection: &mut SqliteConnection) -> Result { let num_rows_deleted = diesel::delete(pipes_table) .execute(connection) .with_context(|| "failed to delete all pipes")?; Ok(num_rows_deleted) } pub fn get_all_pipes(connection: &mut SqliteConnection) -> Result> { let all_pipes = pipes_table .select(Pipe::as_select()) .load(connection) .with_context(|| "failed to get all pipes")?; Ok(all_pipes) } pub fn insert_pipes(pipes: &[Pipe], connection: &mut SqliteConnection) -> Result { let num_rows_inserted = diesel::insert_into(pipes_table) .values(pipes) .execute(connection) .with_context(|| insert_err_msg(&pipes))?; Ok(num_rows_inserted) } fn insert_err_msg(pipes: &[Pipe]) -> String { format!( "failed to insert pipes: [{}]", pipes .iter() .map(|tc| tc.id.to_string()) .collect::>() .join(", ") ) } pub fn sum_pipes_flow_per_bucket_id(connection: &mut SqliteConnection) -> Result> { let sum = pipes_table .group_by(pipes_schema::bucket_id) .select((pipes_schema::bucket_id, sum(pipes_schema::amount))) .load::<(i32, Option)>(connection) .map(|vec| { vec.iter() .map(|(bucket_id, sum)| (*bucket_id, sum.unwrap_or(0))) .collect() }) .with_context(|| "failed to sum pipes flow per bucket ID")?; Ok(sum) }