Skip to main content

max / makenotwork

25.3 KB · 751 lines History Blame Raw
1 //! Peer mesh, identity, heartbeat monitoring, and mesh state aggregation.
2
3 use serde::Serialize;
4 use std::collections::HashMap;
5 use std::path::Path;
6 use std::sync::Arc;
7 use tokio::sync::RwLock;
8
9 use tracing::instrument;
10
11 use crate::alerts::Alerter;
12 use crate::config::PeerConfig;
13 use crate::error::Result;
14
15 /// Cap on a peer response body (matches the health probe's `MAX_RESPONSE_BYTES`).
16 /// A hostile or MITM'd peer could otherwise return a multi-GB body that
17 /// `resp.json()` buffers entirely in RAM, and the 10s timeout bounds wall-clock,
18 /// not bytes. Holds regardless of transport.
19 const MAX_PEER_BODY_BYTES: u64 = 10 * 1024 * 1024;
20
21 /// Deserialize a peer response as JSON with a hard body-size cap.
22 async fn read_capped_json<T: serde::de::DeserializeOwned>(
23 peer_name: &str,
24 resp: reqwest::Response,
25 ) -> Option<T> {
26 if let Some(len) = resp.content_length()
27 && len > MAX_PEER_BODY_BYTES
28 {
29 tracing::warn!(
30 "{peer_name}: peer response declares {len} bytes (cap {MAX_PEER_BODY_BYTES}); dropping"
31 );
32 return None;
33 }
34 let bytes = resp.bytes().await.ok()?;
35 if bytes.len() as u64 > MAX_PEER_BODY_BYTES {
36 tracing::warn!(
37 "{peer_name}: peer response body {} bytes exceeded cap; dropping",
38 bytes.len()
39 );
40 return None;
41 }
42 serde_json::from_slice(&bytes).ok()
43 }
44
45 /// Build the base URL for a peer's API from its configured address.
46 ///
47 /// Honors a scheme already present in the address (so an operator can set
48 /// `https://host:port` when the hop crosses an untrusted network), and otherwise
49 /// defaults to `http://`: the tailnet deployment where WireGuard already
50 /// encrypts the transport, and pom's own API server binds plain HTTP.
51 pub(crate) fn peer_base_url(address: &str) -> String {
52 if address.contains("://") {
53 address.trim_end_matches('/').to_string()
54 } else {
55 format!("http://{address}")
56 }
57 }
58
59 // Identity
60
61 /// Load or create a persistent instance ID (UUID v4).
62 ///
63 /// Stored in `data_dir`, which is the directory holding `pom.db` — passed in
64 /// rather than re-derived, so the ID cannot land beside a different database
65 /// than the one this process opened.
66 pub fn load_or_create_instance_id(override_id: Option<&str>, data_dir: &Path) -> Result<String> {
67 if let Some(id) = override_id {
68 return Ok(id.to_string());
69 }
70
71 std::fs::create_dir_all(data_dir)?;
72 let id_path = data_dir.join("instance_id");
73
74 if id_path.exists() {
75 let id = std::fs::read_to_string(&id_path)?.trim().to_string();
76 if !id.is_empty() {
77 return Ok(id);
78 }
79 }
80
81 let id = uuid::Uuid::new_v4().to_string();
82 std::fs::write(&id_path, &id)?;
83 Ok(id)
84 }
85
86 // Types
87
88 #[derive(Debug, Clone, Serialize, serde::Deserialize)]
89 pub struct InstanceInfo {
90 pub id: String,
91 pub name: String,
92 pub version: String,
93 pub targets: Vec<String>,
94 pub started_at: String,
95 }
96
97 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
98 #[serde(rename_all = "lowercase")]
99 pub enum PeerStatus {
100 /// Last heartbeat succeeded and UUID matches the known identity.
101 Online,
102 /// Heartbeat failed but consecutive failures are still below `grace_count`.
103 GracePeriod,
104 /// Heartbeat failed and consecutive failures have reached or exceeded `grace_count`.
105 Missing,
106 /// Initial state before the first heartbeat attempt has completed.
107 Unknown,
108 }
109
110 impl std::fmt::Display for PeerStatus {
111 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
112 match self {
113 Self::Online => write!(f, "online"),
114 Self::GracePeriod => write!(f, "grace_period"),
115 Self::Missing => write!(f, "missing"),
116 Self::Unknown => write!(f, "unknown"),
117 }
118 }
119 }
120
121 #[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, serde::Deserialize)]
122 #[serde(rename_all = "lowercase")]
123 pub enum OnMissing {
124 /// Send an email alert via Postmark when the peer is declared missing.
125 Alert,
126 /// Emit a `tracing::warn` log line (default behavior).
127 #[default]
128 Log,
129 /// Suppress all failure output for this peer.
130 Ignore,
131 }
132
133 #[derive(Debug, Clone, Serialize)]
134 pub struct PeerState {
135 pub address: String,
136 pub on_missing: OnMissing,
137 pub grace_count: u32,
138 pub status: PeerStatus,
139 pub info: Option<InstanceInfo>,
140 pub last_seen: Option<String>,
141 pub latency_ms: Option<u64>,
142 pub consecutive_failures: u32,
143 #[serde(skip)]
144 pub known_id: Option<String>,
145 /// Cached status data from the peer's /api/peer/status endpoint.
146 #[serde(skip)]
147 pub status_data: Option<serde_json::Value>,
148 /// Bearer token for authenticating with this peer's API.
149 #[serde(skip)]
150 pub token: Option<String>,
151 }
152
153 // Mesh State
154
155 #[derive(Debug)]
156 pub struct MeshState {
157 pub instance: InstanceInfo,
158 pub peers: HashMap<String, PeerState>,
159 }
160
161 pub type SharedMeshState = Arc<RwLock<MeshState>>;
162
163 pub fn new_mesh_state<S: std::hash::BuildHasher>(
164 instance: InstanceInfo,
165 peer_configs: &HashMap<String, PeerConfig, S>,
166 ) -> SharedMeshState {
167 let mut peers = HashMap::new();
168 for (name, cfg) in peer_configs {
169 peers.insert(
170 name.clone(),
171 PeerState {
172 address: cfg.address.clone(),
173 on_missing: cfg.on_missing,
174 grace_count: cfg.grace_count.unwrap_or(3),
175 status: PeerStatus::Unknown,
176 info: None,
177 last_seen: None,
178 latency_ms: None,
179 consecutive_failures: 0,
180 known_id: None,
181 status_data: None,
182 token: cfg.token.clone(),
183 },
184 );
185 }
186 Arc::new(RwLock::new(MeshState { instance, peers }))
187 }
188
189 // Heartbeat
190
191 #[instrument(skip_all)]
192 pub async fn spawn_heartbeat_tasks(
193 mesh: SharedMeshState,
194 pool: sqlx::SqlitePool,
195 interval_secs: u64,
196 alerter: Option<Alerter>,
197 token: tokio_util::sync::CancellationToken,
198 ) -> Vec<tokio::task::JoinHandle<()>> {
199 let peer_names: Vec<String> = {
200 let mesh_guard = mesh.read().await;
201 mesh_guard.peers.keys().cloned().collect()
202 };
203
204 let mut handles = Vec::new();
205 for peer_name in peer_names {
206 let mesh = Arc::clone(&mesh);
207 let pool = pool.clone();
208 let alerter = alerter.clone();
209 let token = token.clone();
210 handles.push(tokio::spawn(async move {
211 heartbeat_loop(&peer_name, mesh, pool, interval_secs, alerter, token).await;
212 }));
213 }
214 handles
215 }
216
217 #[instrument(skip_all, fields(peer = %peer_name))]
218 async fn heartbeat_loop(
219 peer_name: &str,
220 mesh: SharedMeshState,
221 pool: sqlx::SqlitePool,
222 interval_secs: u64,
223 alerter: Option<Alerter>,
224 cancel: tokio_util::sync::CancellationToken,
225 ) {
226 let mut interval = tokio::time::interval(std::time::Duration::from_secs(interval_secs));
227 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
228 // Skip the first immediate tick, give peers time to start up
229 interval.tick().await;
230
231 let (address, auth_token) = {
232 let state = mesh.read().await;
233 match state.peers.get(peer_name) {
234 Some(p) => (p.address.clone(), p.token.clone()),
235 None => return,
236 }
237 };
238 let base = peer_base_url(&address);
239
240 let client = match crate::tls::https_client_builder()
241 .timeout(std::time::Duration::from_secs(10))
242 .build()
243 {
244 Ok(c) => c,
245 Err(e) => {
246 // reqwest 0.13's builder can fail and `unwrap_or_default()` panics on
247 // the same failure, silently killing this spawned poll task. Log and
248 // exit visibly instead.
249 tracing::error!(peer = %peer_name, error = %e, "peer: client build failed, poll task exiting");
250 return;
251 }
252 };
253
254 loop {
255 tokio::select! {
256 () = cancel.cancelled() => break,
257 _ = interval.tick() => {}
258 }
259
260 let start = std::time::Instant::now();
261 let mut req = client.get(format!("{base}/api/peer/info"));
262 if let Some(ref t) = auth_token {
263 req = req.bearer_auth(t);
264 }
265 let result = req.send().await;
266 let latency_ms = start.elapsed().as_millis() as u64;
267
268 match result.and_then(reqwest::Response::error_for_status) {
269 Ok(response) => {
270 let info: Option<InstanceInfo> = read_capped_json(peer_name, response).await;
271 handle_heartbeat_success(
272 peer_name,
273 &mesh,
274 &pool,
275 info,
276 latency_ms,
277 alerter.as_ref(),
278 )
279 .await;
280 }
281 Err(_e) => {
282 handle_heartbeat_failure(peer_name, &mesh, &pool, latency_ms, alerter.as_ref())
283 .await;
284 }
285 }
286
287 // Also fetch /api/peer/status for mesh aggregation
288 let mut status_req = client.get(format!("{base}/api/peer/status"));
289 if let Some(ref t) = auth_token {
290 status_req = status_req.bearer_auth(t);
291 }
292 let status_result = status_req.send().await;
293 match status_result {
294 Ok(resp) => {
295 if let Some(data) = read_capped_json::<serde_json::Value>(peer_name, resp).await {
296 let mut state = mesh.write().await;
297 if let Some(peer) = state.peers.get_mut(peer_name) {
298 peer.status_data = Some(data);
299 }
300 }
301 }
302 Err(e) => {
303 tracing::debug!("{peer_name}: failed to fetch /api/peer/status: {e}");
304 }
305 }
306 }
307 }
308
309 async fn handle_heartbeat_success(
310 peer_name: &str,
311 mesh: &SharedMeshState,
312 pool: &sqlx::SqlitePool,
313 info: Option<InstanceInfo>,
314 latency_ms: u64,
315 alerter: Option<&Alerter>,
316 ) {
317 let now = chrono::Utc::now().to_rfc3339();
318
319 // Update in-memory state under lock, collect data for DB writes
320 let (new_identity_id, updated_identity_id, recovery_info) = {
321 let mut state = mesh.write().await;
322 let Some(peer) = state.peers.get_mut(peer_name) else {
323 return;
324 };
325
326 let was_missing =
327 peer.status == PeerStatus::Missing || peer.status == PeerStatus::GracePeriod;
328
329 // Check UUID consistency
330 let mut new_identity = None;
331 let mut updated_identity = None;
332 if let Some(ref info) = info {
333 match &peer.known_id {
334 None => {
335 peer.known_id = Some(info.id.clone());
336 new_identity = Some(info.id.clone());
337 tracing::info!("{peer_name}: first contact, id={}", info.id);
338 }
339 Some(known) if known != &info.id => {
340 // UUID changed, likely a reinstall. Log the error and reset
341 // peer state so the new identity is accepted.
342 tracing::error!(
343 "{peer_name}: UUID mismatch! expected={known}, got={}. Resetting peer state (possible reinstall).",
344 info.id
345 );
346 peer.known_id = Some(info.id.clone());
347 updated_identity = Some(info.id.clone());
348 }
349 _ => {}
350 }
351 }
352
353 // Collect recovery data before mutating state
354 let recovery = if was_missing && peer.on_missing == OnMissing::Alert {
355 Some(peer.address.clone())
356 } else {
357 None
358 };
359
360 if was_missing {
361 tracing::info!("{peer_name}: recovered (was {:?})", peer.status);
362 }
363
364 peer.status = PeerStatus::Online;
365 peer.info = info;
366 peer.last_seen = Some(now);
367 peer.latency_ms = Some(latency_ms);
368 peer.consecutive_failures = 0;
369
370 (new_identity, updated_identity, recovery)
371 };
372 // Lock dropped, DB writes and alerts happen without holding mesh lock
373
374 if let Some(id) = new_identity_id
375 && let Err(e) = crate::db::store_peer_identity(pool, peer_name, &id).await
376 {
377 tracing::warn!(peer = %peer_name, error = %e, "peer: identity store failed");
378 }
379 if let Some(id) = updated_identity_id
380 && let Err(e) = crate::db::update_peer_identity(pool, peer_name, &id).await
381 {
382 tracing::warn!(peer = %peer_name, error = %e, "peer: identity update failed");
383 }
384 if let Err(e) =
385 crate::db::insert_peer_heartbeat(pool, peer_name, "online", latency_ms as i64).await
386 {
387 tracing::warn!(peer = %peer_name, error = %e, "peer: heartbeat insert failed");
388 }
389
390 if let (Some(address), Some(alerter)) = (recovery_info, alerter) {
391 alerter.send_peer_recovery(peer_name, &address).await;
392 }
393 }
394
395 async fn handle_heartbeat_failure(
396 peer_name: &str,
397 mesh: &SharedMeshState,
398 pool: &sqlx::SqlitePool,
399 latency_ms: u64,
400 alerter: Option<&Alerter>,
401 ) {
402 // Update in-memory state under lock, collect data for alert after lock drop
403 let (new_status, alert_info) = {
404 let mut state = mesh.write().await;
405 let Some(peer) = state.peers.get_mut(peer_name) else {
406 return;
407 };
408
409 peer.consecutive_failures += 1;
410
411 let new_status = match peer.status {
412 PeerStatus::Online | PeerStatus::Unknown | PeerStatus::GracePeriod => {
413 if peer.consecutive_failures >= peer.grace_count {
414 PeerStatus::Missing
415 } else {
416 PeerStatus::GracePeriod
417 }
418 }
419 PeerStatus::Missing => PeerStatus::Missing,
420 };
421
422 let transitioned_to_missing =
423 new_status == PeerStatus::Missing && peer.status != PeerStatus::Missing;
424
425 // Collect alert data before mutating state
426 let alert_info = if transitioned_to_missing {
427 match peer.on_missing {
428 OnMissing::Alert => {
429 tracing::warn!(
430 "{peer_name}: MISSING after {} consecutive failures (action: alert)",
431 peer.consecutive_failures
432 );
433 Some((peer.address.clone(), peer.consecutive_failures))
434 }
435 OnMissing::Log => {
436 tracing::info!(
437 "{peer_name}: missing after {} consecutive failures",
438 peer.consecutive_failures
439 );
440 None
441 }
442 OnMissing::Ignore => None,
443 }
444 } else {
445 None
446 };
447
448 peer.status = new_status;
449
450 (new_status, alert_info)
451 };
452 // Lock dropped, DB write and alerts happen without holding mesh lock
453
454 let status_str = new_status.to_string();
455 if let Err(e) =
456 crate::db::insert_peer_heartbeat(pool, peer_name, &status_str, latency_ms as i64).await
457 {
458 tracing::warn!(peer = %peer_name, error = %e, "peer: heartbeat insert failed");
459 }
460
461 if let (Some((address, failures)), Some(alerter)) = (alert_info, alerter) {
462 alerter
463 .send_peer_missing(peer_name, &address, failures)
464 .await;
465 }
466 }
467
468 #[cfg(test)]
469 mod tests {
470 use super::*;
471
472 #[test]
473 fn peer_base_url_honors_scheme() {
474 assert_eq!(peer_base_url("host:9100"), "http://host:9100");
475 assert_eq!(peer_base_url("https://host:9100"), "https://host:9100");
476 assert_eq!(peer_base_url("https://host:9100/"), "https://host:9100");
477 }
478
479 #[test]
480 fn override_id_takes_precedence() {
481 let dir = std::env::temp_dir().join(format!("pom-id-override-{}", std::process::id()));
482 let id = load_or_create_instance_id(Some("override-id"), &dir).unwrap();
483 assert_eq!(id, "override-id");
484 // An override answers without touching the disk at all.
485 assert!(!dir.exists());
486 }
487
488 #[test]
489 fn auto_id_is_valid_uuid_and_persists_beside_the_db() {
490 let dir = std::env::temp_dir().join(format!("pom-id-auto-{}", std::process::id()));
491 std::fs::remove_dir_all(&dir).ok();
492 let id = load_or_create_instance_id(None, &dir).unwrap();
493 assert!(uuid::Uuid::parse_str(&id).is_ok());
494 // Written where the caller said, not where XDG_DATA_HOME points, and
495 // stable across calls.
496 assert_eq!(
497 std::fs::read_to_string(dir.join("instance_id")).unwrap(),
498 id
499 );
500 assert_eq!(load_or_create_instance_id(None, &dir).unwrap(), id);
501 std::fs::remove_dir_all(&dir).ok();
502 }
503
504 #[test]
505 fn on_missing_deserialize() {
506 #[derive(serde::Deserialize)]
507 struct Wrapper {
508 #[serde(default)]
509 on_missing: OnMissing,
510 }
511
512 let w: Wrapper = toml::from_str(r#"on_missing = "alert""#).unwrap();
513 assert_eq!(w.on_missing, OnMissing::Alert);
514
515 let w: Wrapper = toml::from_str(r#"on_missing = "log""#).unwrap();
516 assert_eq!(w.on_missing, OnMissing::Log);
517
518 let w: Wrapper = toml::from_str(r#"on_missing = "ignore""#).unwrap();
519 assert_eq!(w.on_missing, OnMissing::Ignore);
520
521 // Default is Log
522 let w: Wrapper = toml::from_str("").unwrap();
523 assert_eq!(w.on_missing, OnMissing::Log);
524 }
525
526 fn test_instance_info() -> InstanceInfo {
527 InstanceInfo {
528 id: "test-id".to_string(),
529 name: "test".to_string(),
530 version: "0.1.0".to_string(),
531 targets: vec![],
532 started_at: "2026-03-10T00:00:00Z".to_string(),
533 }
534 }
535
536 fn test_mesh_with_peer(grace_count: u32) -> SharedMeshState {
537 let mut peer_configs = HashMap::new();
538 peer_configs.insert(
539 "peer1".to_string(),
540 PeerConfig {
541 address: "10.0.0.1:9100".to_string(),
542 on_missing: OnMissing::Alert,
543 grace_count: Some(grace_count),
544 token: None,
545 },
546 );
547 new_mesh_state(test_instance_info(), &peer_configs)
548 }
549
550 #[tokio::test]
551 async fn heartbeat_failure_transitions_through_grace_to_missing() {
552 let pool = crate::db::connect_in_memory().await.unwrap();
553 let mesh = test_mesh_with_peer(3);
554
555 // Start at Unknown
556 assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Unknown);
557
558 // First failure → GracePeriod
559 handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await;
560 assert_eq!(
561 mesh.read().await.peers["peer1"].status,
562 PeerStatus::GracePeriod
563 );
564
565 // Second failure → still GracePeriod
566 handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await;
567 assert_eq!(
568 mesh.read().await.peers["peer1"].status,
569 PeerStatus::GracePeriod
570 );
571
572 // Third failure (= grace_count) → Missing
573 handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await;
574 assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Missing);
575
576 // Fourth failure → stays Missing
577 handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await;
578 assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Missing);
579 }
580
581 #[tokio::test]
582 async fn heartbeat_success_recovers_from_missing() {
583 let pool = crate::db::connect_in_memory().await.unwrap();
584 let mesh = test_mesh_with_peer(1);
585
586 // Drive to Missing
587 handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await;
588 assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Missing);
589
590 // Success → Online
591 let info = InstanceInfo {
592 id: "remote-id".to_string(),
593 name: "remote".to_string(),
594 version: "0.1.0".to_string(),
595 targets: vec![],
596 started_at: "2026-03-10T00:00:00Z".to_string(),
597 };
598 handle_heartbeat_success("peer1", &mesh, &pool, Some(info), 42, None).await;
599
600 let state = mesh.read().await;
601 let peer = &state.peers["peer1"];
602 assert_eq!(peer.status, PeerStatus::Online);
603 assert_eq!(peer.consecutive_failures, 0);
604 assert_eq!(peer.latency_ms, Some(42));
605 assert_eq!(peer.known_id.as_deref(), Some("remote-id"));
606 }
607
608 #[tokio::test]
609 async fn heartbeat_success_detects_uuid_stored_on_first_contact() {
610 let pool = crate::db::connect_in_memory().await.unwrap();
611 let mesh = test_mesh_with_peer(3);
612
613 let info = InstanceInfo {
614 id: "uuid-abc".to_string(),
615 name: "remote".to_string(),
616 version: "0.1.0".to_string(),
617 targets: vec![],
618 started_at: "2026-03-10T00:00:00Z".to_string(),
619 };
620 handle_heartbeat_success("peer1", &mesh, &pool, Some(info), 10, None).await;
621
622 // UUID should be persisted in DB
623 let stored = crate::db::get_peer_identity(&pool, "peer1").await.unwrap();
624 assert_eq!(stored, Some("uuid-abc".to_string()));
625 }
626
627 #[tokio::test]
628 async fn heartbeat_records_to_db() {
629 let pool = crate::db::connect_in_memory().await.unwrap();
630 let mesh = test_mesh_with_peer(3);
631
632 handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await;
633 handle_heartbeat_success("peer1", &mesh, &pool, None, 55, None).await;
634
635 let history = crate::db::get_peer_heartbeat_history(&pool, "peer1", 10)
636 .await
637 .unwrap();
638 assert_eq!(history.len(), 2);
639 // Most recent first
640 assert_eq!(history[0].status, "online");
641 assert_eq!(history[0].latency_ms, 55);
642 assert_eq!(history[1].status, "grace_period");
643 }
644
645 #[tokio::test]
646 async fn heartbeat_uuid_mismatch_resets_state() {
647 let pool = crate::db::connect_in_memory().await.unwrap();
648 let mesh = test_mesh_with_peer(3);
649
650 // First contact with UUID "aaa"
651 let info_a = InstanceInfo {
652 id: "aaa".to_string(),
653 name: "remote".to_string(),
654 version: "0.1.0".to_string(),
655 targets: vec![],
656 started_at: "2026-03-10T00:00:00Z".to_string(),
657 };
658 handle_heartbeat_success("peer1", &mesh, &pool, Some(info_a), 10, None).await;
659 assert_eq!(
660 mesh.read().await.peers["peer1"].known_id.as_deref(),
661 Some("aaa")
662 );
663 assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Online);
664
665 // UUID stored in DB
666 let stored = crate::db::get_peer_identity(&pool, "peer1").await.unwrap();
667 assert_eq!(stored, Some("aaa".to_string()));
668
669 // Reconnect with UUID "bbb" (simulating reinstall)
670 let info_b = InstanceInfo {
671 id: "bbb".to_string(),
672 name: "remote-reinstalled".to_string(),
673 version: "0.2.0".to_string(),
674 targets: vec![],
675 started_at: "2026-03-13T00:00:00Z".to_string(),
676 };
677 handle_heartbeat_success("peer1", &mesh, &pool, Some(info_b), 20, None).await;
678
679 // State should be reset: new UUID accepted, status online, failures cleared
680 let state = mesh.read().await;
681 let peer = &state.peers["peer1"];
682 assert_eq!(peer.known_id.as_deref(), Some("bbb"));
683 assert_eq!(peer.status, PeerStatus::Online);
684 assert_eq!(peer.consecutive_failures, 0);
685
686 // DB identity should be updated to "bbb"
687 let stored = crate::db::get_peer_identity(&pool, "peer1").await.unwrap();
688 assert_eq!(stored, Some("bbb".to_string()));
689 }
690
691 #[tokio::test]
692 async fn heartbeat_same_uuid_proceeds_normally() {
693 let pool = crate::db::connect_in_memory().await.unwrap();
694 let mesh = test_mesh_with_peer(3);
695
696 let info = InstanceInfo {
697 id: "same-uuid".to_string(),
698 name: "remote".to_string(),
699 version: "0.1.0".to_string(),
700 targets: vec![],
701 started_at: "2026-03-10T00:00:00Z".to_string(),
702 };
703 // First contact
704 handle_heartbeat_success("peer1", &mesh, &pool, Some(info.clone()), 10, None).await;
705 assert_eq!(
706 mesh.read().await.peers["peer1"].known_id.as_deref(),
707 Some("same-uuid")
708 );
709
710 // Second heartbeat with same UUID
711 handle_heartbeat_success("peer1", &mesh, &pool, Some(info), 15, None).await;
712 let state = mesh.read().await;
713 let peer = &state.peers["peer1"];
714 assert_eq!(peer.known_id.as_deref(), Some("same-uuid"));
715 assert_eq!(peer.status, PeerStatus::Online);
716 assert_eq!(peer.latency_ms, Some(15));
717 }
718
719 #[test]
720 fn mesh_state_construction() {
721 let info = InstanceInfo {
722 id: "test-id".to_string(),
723 name: "test".to_string(),
724 version: "0.1.0".to_string(),
725 targets: vec!["mnw".to_string()],
726 started_at: "2026-03-10T00:00:00Z".to_string(),
727 };
728
729 let mut peer_configs = HashMap::new();
730 peer_configs.insert(
731 "peer1".to_string(),
732 PeerConfig {
733 address: "10.0.0.1:9100".to_string(),
734 on_missing: OnMissing::Alert,
735 grace_count: Some(5),
736 token: None,
737 },
738 );
739
740 let mesh = new_mesh_state(info, &peer_configs);
741 let state = mesh.blocking_read();
742 assert_eq!(state.instance.id, "test-id");
743 assert_eq!(state.peers.len(), 1);
744 let peer = state.peers.get("peer1").unwrap();
745 assert_eq!(peer.address, "10.0.0.1:9100");
746 assert_eq!(peer.on_missing, OnMissing::Alert);
747 assert_eq!(peer.grace_count, 5);
748 assert_eq!(peer.status, PeerStatus::Unknown);
749 }
750 }
751