aven_core/query/
sync_history.rs1use 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}