1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
|
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<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)
}
|