Skip to main content

aven_core/query/
sync_history.rs

1use anyhow::Result;
2use sqlx::{Row, SqliteConnection};
3
4#[derive(Debug, Clone, Default, PartialEq, Eq)]
5pub struct SyncHistoryStats {
6    pub total_change_rows: i64,
7    pub pending_change_rows: i64,
8    pub synced_change_rows: i64,
9    pub min_server_seq: Option<i64>,
10    pub max_server_seq: Option<i64>,
11    pub payload_bytes: i64,
12}
13
14pub async fn sync_history_stats(conn: &mut SqliteConnection) -> Result<SyncHistoryStats> {
15    let row = sqlx::query(
16        "SELECT
17         COUNT(*) AS total_change_rows,
18         COALESCE(SUM(CASE WHEN server_seq IS NULL THEN 1 ELSE 0 END), 0) AS pending_change_rows,
19         COALESCE(SUM(CASE WHEN server_seq IS NOT NULL THEN 1 ELSE 0 END), 0) AS synced_change_rows,
20         MIN(server_seq) AS min_server_seq,
21         MAX(server_seq) AS max_server_seq,
22         COALESCE(SUM(length(CAST(payload AS BLOB))), 0) AS payload_bytes
23         FROM changes",
24    )
25    .fetch_one(&mut *conn)
26    .await?;
27
28    Ok(SyncHistoryStats {
29        total_change_rows: row.get("total_change_rows"),
30        pending_change_rows: row.get("pending_change_rows"),
31        synced_change_rows: row.get("synced_change_rows"),
32        min_server_seq: row.get("min_server_seq"),
33        max_server_seq: row.get("max_server_seq"),
34        payload_bytes: row.get("payload_bytes"),
35    })
36}