summaryrefslogtreecommitdiff
path: root/schist_core/schist_queries/src/pipes.rs
blob: b471a420bbe1c4021fbf84ed093bdf5da63ff454 (plain)
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)
}