Skip to main content

max / makenotwork

2.6 KB · 103 lines History Blame Raw
1 //! Peer identity and heartbeat storage for the monitoring mesh.
2
3 use super::{Result, SqlitePool};
4 use tracing::instrument;
5
6 /// Store a peer's identity (first-seen UUID). INSERT OR IGNORE, first wins.
7 #[instrument(skip_all)]
8 pub async fn store_peer_identity(
9 pool: &SqlitePool,
10 peer_name: &str,
11 instance_id: &str,
12 ) -> Result<()> {
13 let now = chrono::Utc::now().to_rfc3339();
14 sqlx::query(
15 "INSERT OR IGNORE INTO peer_identities (peer_name, instance_id, first_seen)
16 VALUES (?, ?, ?)",
17 )
18 .bind(peer_name)
19 .bind(instance_id)
20 .bind(&now)
21 .execute(pool)
22 .await?;
23 Ok(())
24 }
25
26 /// Update a peer's identity (UUID changed, e.g. after reinstall).
27 #[instrument(skip_all)]
28 pub async fn update_peer_identity(
29 pool: &SqlitePool,
30 peer_name: &str,
31 instance_id: &str,
32 ) -> Result<()> {
33 let now = chrono::Utc::now().to_rfc3339();
34 sqlx::query(
35 "UPDATE peer_identities SET instance_id = ?, first_seen = ?
36 WHERE peer_name = ?",
37 )
38 .bind(instance_id)
39 .bind(&now)
40 .bind(peer_name)
41 .execute(pool)
42 .await?;
43 Ok(())
44 }
45
46 /// Get a peer's first-seen instance ID.
47 #[instrument(skip_all)]
48 pub async fn get_peer_identity(pool: &SqlitePool, peer_name: &str) -> Result<Option<String>> {
49 let row = sqlx::query_as::<_, (String,)>(
50 "SELECT instance_id FROM peer_identities WHERE peer_name = ?",
51 )
52 .bind(peer_name)
53 .fetch_optional(pool)
54 .await?;
55 Ok(row.map(|r| r.0))
56 }
57
58 #[instrument(skip_all)]
59 pub async fn insert_peer_heartbeat(
60 pool: &SqlitePool,
61 peer_name: &str,
62 status: &str,
63 latency_ms: i64,
64 ) -> Result<i64> {
65 let now = chrono::Utc::now().to_rfc3339();
66 let result = sqlx::query(
67 "INSERT INTO peer_heartbeats (peer_name, status, latency_ms, checked_at)
68 VALUES (?, ?, ?, ?)",
69 )
70 .bind(peer_name)
71 .bind(status)
72 .bind(latency_ms)
73 .bind(&now)
74 .execute(pool)
75 .await?;
76 Ok(result.last_insert_rowid())
77 }
78
79 #[instrument(skip_all)]
80 pub async fn get_peer_heartbeat_history(
81 pool: &SqlitePool,
82 peer_name: &str,
83 limit: i64,
84 ) -> Result<Vec<PeerHeartbeatRow>> {
85 Ok(sqlx::query_as::<_, PeerHeartbeatRow>(
86 "SELECT id, peer_name, status, latency_ms, checked_at
87 FROM peer_heartbeats WHERE peer_name = ? ORDER BY id DESC LIMIT ?",
88 )
89 .bind(peer_name)
90 .bind(limit)
91 .fetch_all(pool)
92 .await?)
93 }
94
95 #[derive(Debug, sqlx::FromRow, serde::Serialize)]
96 pub struct PeerHeartbeatRow {
97 pub id: i64,
98 pub peer_name: String,
99 pub status: String,
100 pub latency_ms: i64,
101 pub checked_at: String,
102 }
103