summaryrefslogtreecommitdiff
path: root/schist_core/schist_queries/src/pipes.rs
diff options
context:
space:
mode:
Diffstat (limited to 'schist_core/schist_queries/src/pipes.rs')
-rw-r--r--schist_core/schist_queries/src/pipes.rs55
1 files changed, 55 insertions, 0 deletions
diff --git a/schist_core/schist_queries/src/pipes.rs b/schist_core/schist_queries/src/pipes.rs
new file mode 100644
index 0000000..188d193
--- /dev/null
+++ b/schist_core/schist_queries/src/pipes.rs
@@ -0,0 +1,55 @@
+use anyhow::{Context, Result};
+use diesel::{dsl::sum, QueryDsl, RunQueryDsl, SelectableHelper, SqliteConnection};
+use schist_models::pipe::Pipe;
+use schist_schema::schema::pipes::{self as pipes_schema, dsl::pipes as pipes_table};
+
+pub fn delete_all_pipes(connection: &mut SqliteConnection) -> Result<usize> {
+ 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<Vec<Pipe>> {
+ 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<usize> {
+ 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::<Vec<String>>()
+ .join(", ")
+ )
+}
+
+pub fn sum_pipes_flow_per_bucket_id(connection: &mut SqliteConnection) -> Result<Vec<(i32, i64)>> {
+ let sum = pipes_table
+ .group_by(pipes_schema::bucket_id)
+ .select((pipes_schema::bucket_id, sum(pipes_schema::amount)))
+ .load::<(i32, Option<i64>)>(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)
+}