Skip to main content

max / makenotwork

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