| 1 |
|
| 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 |
|
| 16 |
|
| 17 |
|
| 18 |
|
| 19 |
const MAX_PEER_BODY_BYTES: u64 = 10 * 1024 * 1024; |
| 20 |
|
| 21 |
|
| 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 |
|
| 46 |
|
| 47 |
|
| 48 |
|
| 49 |
|
| 50 |
|
| 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 |
|
| 60 |
|
| 61 |
|
| 62 |
|
| 63 |
|
| 64 |
|
| 65 |
|
| 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 |
|
| 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 |
|
| 101 |
Online, |
| 102 |
|
| 103 |
GracePeriod, |
| 104 |
|
| 105 |
Missing, |
| 106 |
|
| 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 |
|
| 125 |
Alert, |
| 126 |
|
| 127 |
#[default] |
| 128 |
Log, |
| 129 |
|
| 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 |
|
| 146 |
#[serde(skip)] |
| 147 |
pub status_data: Option<serde_json::Value>, |
| 148 |
|
| 149 |
#[serde(skip)] |
| 150 |
pub token: Option<String>, |
| 151 |
} |
| 152 |
|
| 153 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 247 |
|
| 248 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 341 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 495 |
|
| 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 |
|
| 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 |
|
| 556 |
assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Unknown); |
| 557 |
|
| 558 |
|
| 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 |
|
| 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 |
|
| 573 |
handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await; |
| 574 |
assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Missing); |
| 575 |
|
| 576 |
|
| 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 |
|
| 587 |
handle_heartbeat_failure("peer1", &mesh, &pool, 0, None).await; |
| 588 |
assert_eq!(mesh.read().await.peers["peer1"].status, PeerStatus::Missing); |
| 589 |
|
| 590 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 666 |
let stored = crate::db::get_peer_identity(&pool, "peer1").await.unwrap(); |
| 667 |
assert_eq!(stored, Some("aaa".to_string())); |
| 668 |
|
| 669 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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 |
|