//! Pinned liveness checks and generation-conditional peer cleanup. use std::{sync::Arc, time::Duration}; use futures::{StreamExt as _, stream}; use crate::{ PeerEventSender, config::{PEER_PING_IDLE_SECS, PEER_PING_INTERVAL_SECS, peer_stale_timeout}, content_quarantine::ContentQuarantine, context::{NetworkServiceCtx, OperationKind}, network::ping_peer, peer_db::PeerLivenessSnapshot, scoped_blocking::scoped_blocking, services::{HandshakeCtx, remote_state}, }; const MAX_CONCURRENT_PINGS: usize = 8; /// Runs revision-bearing pinned liveness checks. The idle gate deliberately /// uses `last_revision_check`; inbound and content traffic only affect /// `last_seen`, which remains the stale-pruning clock. pub async fn run_ping_service( tx_notify_ui: PeerEventSender, ctx: NetworkServiceCtx, ) -> eyre::Result<()> { log::info!( "Starting ping service ({PEER_PING_INTERVAL_SECS}s interval, \ {}s revision-check idle threshold, {}s stale timeout)", PEER_PING_IDLE_SECS, peer_stale_timeout().as_secs() ); let mut interval = tokio::time::interval(Duration::from_secs(PEER_PING_INTERVAL_SECS)); let remote_ctx = HandshakeCtx::from_network(&ctx, &tx_notify_ui); loop { tokio::select! { biased; () = ctx.shutdown.cancelled() => return Ok(()), _ = interval.tick() => {} } let snapshots = ctx .peer_game_db .read() .await .peer_liveness_snapshot() .into_iter() .filter(revision_check_due) .collect::>(); let mut checks = stream::iter(snapshots.into_iter().map(|snapshot| { let ctx = ctx.clone(); let remote_ctx = remote_ctx.clone(); async move { check_peer_liveness(&ctx, &remote_ctx, snapshot).await } })) .buffer_unordered(MAX_CONCURRENT_PINGS); while checks.next().await.is_some() {} if ctx.shutdown.is_cancelled() { return Ok(()); } prune_stale_peers(&ctx, &remote_ctx).await?; } } fn revision_check_due(snapshot: &PeerLivenessSnapshot) -> bool { snapshot.last_revision_check.elapsed() >= Duration::from_secs(PEER_PING_IDLE_SECS) } async fn check_peer_liveness( ctx: &NetworkServiceCtx, remote_ctx: &HandshakeCtx, snapshot: PeerLivenessSnapshot, ) { match ping_peer(&ctx.quic, &snapshot.endpoint, &ctx.shutdown).await { Ok(revisions) => { match remote_state::observe_pinned_pong(remote_ctx, snapshot, revisions).await { Ok(remote_state::PongCommit::NeedsPull { .. }) => { if let Err(error) = ctx .state_sync .schedule_pinned_pull(snapshot.endpoint.peer_id, &ctx.shutdown) .await { log::debug!( "Could not schedule revision refresh for {}: {error:#}", snapshot.endpoint.peer_id ); } } Ok( remote_state::PongCommit::Current | remote_state::PongCommit::StaleGeneration, ) => {} Err(error) => log::error!( "Failed to apply Pong from {}: {error:#}", snapshot.endpoint.addr ), } } Err(error) => { log::warn!( "Pinned ping to {} failed: {error:#}", snapshot.endpoint.addr ); // A capacity reset or transient transport failure is not topology // authority. Leave both clocks unchanged and require repeated // generation-current failures plus the stale timeout before the // normal pruning path removes state. ctx.peer_game_db .write() .await .record_ping_failure_if_generation(snapshot); } } } async fn prune_stale_peers(ctx: &NetworkServiceCtx, remote_ctx: &HandshakeCtx) -> eyre::Result<()> { let stale = ctx .peer_game_db .read() .await .stale_peer_liveness_snapshots(peer_stale_timeout()); let mut removed_any = false; for snapshot in stale { removed_any |= remote_state::remove_peer_if_generation(remote_ctx, snapshot).await?; } if removed_any { handle_active_downloads_without_peers(ctx).await; } Ok(()) } async fn handle_active_downloads_without_peers(ctx: &NetworkServiceCtx) { let active_ids = ctx .active_operations .read() .await .iter() .filter_map(|(id, kind)| (*kind == OperationKind::Downloading).then_some(id.clone())) .collect::>(); for id in active_ids { if eligible_source_remains(ctx, &id).await { continue; } { // Signalling is one-shot; the download task retains ownership of // operation-map cleanup until every transfer worker has drained. let active_downloads = ctx.active_downloads.read().await; let Some(download) = active_downloads.get(&id) else { continue; }; download.cancel_sources_exhausted(); } } } async fn eligible_source_remains(ctx: &NetworkServiceCtx, game_id: &str) -> bool { let catalog = Arc::clone(&ctx.catalog); let game_id_owned = game_id.to_owned(); let Ok(manifest) = scoped_blocking(move || catalog.manifest(&game_id_owned)) else { return false; }; let content_id = manifest.content_id(); let endpoints = ctx .peer_game_db .read() .await .peer_endpoints_with_content(game_id, content_id); has_nonquarantined_source(&endpoints, content_id, &ctx.content_quarantine) } fn has_nonquarantined_source( endpoints: &[lanspread_proto::PeerEndpoint], content_id: lanspread_db::content_manifest::ContentId, quarantine: &ContentQuarantine, ) -> bool { endpoints .iter() .any(|endpoint| !quarantine.is_quarantined(endpoint, content_id)) } #[cfg(test)] mod tests { use lanspread_db::content_manifest::ContentId; use lanspread_proto::{LibrarySnapshot, PeerEndpoint, PeerId, RuntimeSessionId}; use super::*; use crate::peer_db::PeerGameDB; #[tokio::test(start_paused = true)] async fn dropped_hint_is_revision_checked_within_configured_bound() { let endpoint = PeerEndpoint::new( PeerId::from_bytes([1; 32]), std::net::SocketAddr::from(([127, 0, 0, 1], 12001)), ); let mut interval = tokio::time::interval(Duration::from_secs(PEER_PING_INTERVAL_SECS)); interval.tick().await; tokio::time::advance(Duration::from_millis(1)).await; let mut db = PeerGameDB::new(); let ticket = db .begin_candidate_negotiation(endpoint) .expect("candidate should reserve"); db.commit_authenticated_snapshot( endpoint, ticket, RuntimeSessionId::from_bytes([1; 16]), Some(LibrarySnapshot { revision: 0, games: Vec::new(), }), ) .expect("peer should commit") .expect("ticket should remain current"); let start = tokio::time::Instant::now(); interval.tick().await; let mut snapshot = db .peer_liveness_for(&endpoint.peer_id) .expect("peer should remain"); snapshot.last_seen = tokio::time::Instant::now(); assert!( !revision_check_due(&snapshot), "a tick just before the idle threshold should not ping" ); interval.tick().await; let mut snapshot = db .peer_liveness_for(&endpoint.peer_id) .expect("peer should remain"); // Simulate repeated inbound/content activity. It may refresh liveness, // but must not touch the independently captured revision-check clock. snapshot.last_seen = tokio::time::Instant::now(); assert!(revision_check_due(&snapshot)); assert!( start.elapsed() <= Duration::from_secs(PEER_PING_IDLE_SECS + PEER_PING_INTERVAL_SECS) ); } #[tokio::test(start_paused = true)] async fn one_transient_ping_failure_never_removes_and_success_resets_failure_history() { let endpoint = PeerEndpoint::new( PeerId::from_bytes([9; 32]), std::net::SocketAddr::from(([127, 0, 0, 1], 12009)), ); let mut db = PeerGameDB::new(); let ticket = db .begin_candidate_negotiation(endpoint) .expect("candidate should reserve"); db.commit_authenticated_snapshot( endpoint, ticket, RuntimeSessionId::from_bytes([9; 16]), Some(LibrarySnapshot { revision: 0, games: Vec::new(), }), ) .expect("peer should commit") .expect("ticket should remain current"); let probe = db .peer_liveness_for(&endpoint.peer_id) .expect("peer should exist"); tokio::time::advance(peer_stale_timeout() + Duration::from_secs(1)).await; assert!(db.record_ping_failure_if_generation(probe)); assert!( db.stale_peer_liveness_snapshots(peer_stale_timeout()) .is_empty(), "one transport failure must preserve authenticated state" ); assert!(matches!( db.observe_pong_if_generation( probe, lanspread_proto::PeerRevisions { runtime_session_id: RuntimeSessionId::from_bytes([9; 16]), library_revision: 0, call_to_play_revision: 0, }, ), crate::peer_db::PongObservation::Current | crate::peer_db::PongObservation::RevisionMismatch )); let refreshed = db .peer_liveness_for(&endpoint.peer_id) .expect("successful Pong should preserve peer"); assert_eq!(refreshed.consecutive_ping_failures, 0); } #[test] fn wrong_content_and_quarantined_sources_do_not_keep_download_alive() { let expected = ContentId::from_bytes([7; 32]); let endpoint = PeerEndpoint::new( PeerId::from_bytes([2; 32]), std::net::SocketAddr::from(([127, 0, 0, 1], 12002)), ); let quarantine = ContentQuarantine::default(); // Exact-content filtering happens before this seam, so a wrong-content // peer produces the empty eligible endpoint set. assert!(!has_nonquarantined_source(&[], expected, &quarantine)); assert!(has_nonquarantined_source( &[endpoint], expected, &quarantine )); quarantine.record_integrity_failure(&endpoint, expected); assert!(!has_nonquarantined_source( &[endpoint], expected, &quarantine )); } }