//! Bounded scheduling for responder-owned state pulls and change-hint fanout. use std::{collections::HashMap, future::Future, net::IpAddr, pin::Pin, sync::Arc, time::Duration}; use futures::{StreamExt as _, stream::FuturesUnordered}; use lanspread_proto::{ChangeHint, PeerId}; use tokio::sync::{Mutex, mpsc, watch}; use tokio_util::sync::CancellationToken; use crate::{ PeerEvent, context::NetworkServiceCtx, network::{send_call_to_play_changed, send_library_changed}, peer_db::PeerRevisionSnapshot, services::{HandshakeCtx, PeerRefreshOutcome, perform_peer_refresh}, }; const SYNC_QUEUE_CAPACITY: usize = 64; const MAX_TRACKED_PEERS: usize = 64; const MAX_CONCURRENT_PULLS: usize = 8; const MAX_CONCURRENT_HINT_SENDS: usize = 8; pub(crate) const PEER_PULL_COALESCE_WINDOW: Duration = Duration::from_secs(5); #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub(crate) enum StateDomain { Library, CallToPlay, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] struct LocalRevisions { library: u64, call_to_play: u64, } #[derive(Clone, Copy, Debug)] struct HintTrigger { domain: StateDomain, hint: ChangeHint, /// IP address the hint arrived from. Inbound connections are anonymous, /// so this is the only evidence tying `hint.claimed_peer_id` to a sender. source_ip: Option, } struct StateSyncInbox { pinned_rx: Mutex>, hint_rx: Mutex>, } /// Non-blocking ingress for untrusted hints and latest-value local revisions. /// /// Pinned Pong mismatches use a separate bounded queue and an async admission /// method so hostile hint saturation cannot displace reconciliation work. #[derive(Clone)] pub(crate) struct StateSyncHandle { local_peer_id: PeerId, pinned_tx: mpsc::Sender, hint_tx: mpsc::Sender, local_revisions_tx: watch::Sender, inbox: Arc, } impl StateSyncHandle { #[must_use] pub(crate) fn new(local_peer_id: PeerId) -> Self { let (pinned_tx, pinned_rx) = mpsc::channel(SYNC_QUEUE_CAPACITY); let (hint_tx, hint_rx) = mpsc::channel(SYNC_QUEUE_CAPACITY); let (local_revisions_tx, _) = watch::channel(LocalRevisions { library: 0, call_to_play: 0, }); Self { local_peer_id, pinned_tx, hint_tx, local_revisions_tx, inbox: Arc::new(StateSyncInbox { pinned_rx: Mutex::new(pinned_rx), hint_rx: Mutex::new(hint_rx), }), } } /// Enqueues a generation-current pinned Pong mismatch without sharing the /// untrusted hint queue. pub(crate) async fn schedule_pinned_pull( &self, peer_id: PeerId, cancellation: &CancellationToken, ) -> eyre::Result<()> { if peer_id == self.local_peer_id { return Ok(()); } tokio::select! { biased; () = cancellation.cancelled() => eyre::bail!("state-sync admission was cancelled"), result = self.pinned_tx.send(peer_id) => { result.map_err(|_| eyre::eyre!("state-sync pinned queue is closed")) } } } /// Treats an inbound message only as a lossy invalidation hint. /// /// `source_ip` is the address of the anonymous connection that delivered /// the hint. It is later compared against the authenticated address of /// the claimed peer so a third party cannot make this node pull from /// arbitrary known peers. pub(crate) fn schedule_hint( &self, domain: StateDomain, hint: ChangeHint, source_ip: Option, ) { if hint.claimed_peer_id == self.local_peer_id { return; } match self.hint_tx.try_send(HintTrigger { domain, hint, source_ip, }) { Ok(()) => {} Err(mpsc::error::TrySendError::Full(_)) => { log::trace!("Coalescing remote state hint because the bounded queue is full"); } Err(mpsc::error::TrySendError::Closed(_)) => { log::debug!("Ignoring remote state hint because state sync is stopped"); } } } pub(crate) fn publish_library_revision(&self, revision: u64) { self.local_revisions_tx.send_if_modified(|current| { if revision <= current.library { return false; } current.library = revision; true }); } pub(crate) fn publish_call_to_play_revision(&self, revision: u64) { self.local_revisions_tx.send_if_modified(|current| { if revision <= current.call_to_play { return false; } current.call_to_play = revision; true }); } } #[derive(Debug)] struct PeerSlot { in_flight: bool, pending: bool, pinned: bool, next_allowed: tokio::time::Instant, } struct PullCompletion { peer_id: PeerId, result: eyre::Result, } type PullFuture = Pin + Send>>; type FanoutFuture = Pin + Send>>; /// Runs the bounded pull scheduler and local hint fanout in one lexical scope. pub(crate) async fn run_state_sync( ctx: NetworkServiceCtx, tx_notify_ui: tokio::sync::mpsc::UnboundedSender, cancellation: CancellationToken, ) -> eyre::Result<()> { let mut pinned_rx = ctx.state_sync.inbox.pinned_rx.lock().await; let mut hint_rx = ctx.state_sync.inbox.hint_rx.lock().await; let mut local_revisions_rx = ctx.state_sync.local_revisions_tx.subscribe(); let mut last_fanout = *local_revisions_rx.borrow_and_update(); let mut fanout: Option = None; let mut fanout_pending = false; // One actor-owned pinned item preserves backpressure when all tracked // slots are temporarily non-replaceable. While occupied, the bounded // pinned channel is deliberately not drained. let mut blocked_pinned = None; let mut slots = HashMap::::new(); let mut pulls = FuturesUnordered::::new(); let work_cancellation = cancellation.child_token(); let handshake_ctx = HandshakeCtx::from_network(&ctx, &tx_notify_ui) .with_cancellation(work_cancellation.clone()); loop { expire_idle_slots(&mut slots); if let Some(peer_id) = blocked_pinned && try_enqueue_peer(&mut slots, peer_id, true) { blocked_pinned = None; } start_ready_pulls(&mut slots, &mut pulls, &handshake_ctx); if fanout.is_none() && fanout_pending { let target = *local_revisions_rx.borrow_and_update(); fanout_pending = false; if target != last_fanout { fanout = Some(Box::pin(fanout_local_hints( ctx.clone(), last_fanout, target, work_cancellation.clone(), ))); last_fanout = target; } } let deadline = next_slot_deadline(&slots); tokio::select! { biased; () = cancellation.cancelled() => break, peer_id = pinned_rx.recv(), if blocked_pinned.is_none() => { let Some(peer_id) = peer_id else { break; }; if !try_enqueue_peer(&mut slots, peer_id, true) { blocked_pinned = Some(peer_id); } } completed = pulls.next(), if !pulls.is_empty() => { if let Some(completed) = completed && let Some(slot) = slots.get_mut(&completed.peer_id) { settle_pull_completion(slot, completed.peer_id, completed.result); } } () = async { if let Some(fanout) = fanout.as_mut() { fanout.await; } }, if fanout.is_some() => { fanout = None; if *local_revisions_rx.borrow() != last_fanout { fanout_pending = true; } } changed = local_revisions_rx.changed() => { if changed.is_err() { break; } fanout_pending = true; } trigger = hint_rx.recv() => { let Some(trigger) = trigger else { break; }; if hint_requires_pull(&ctx, trigger).await { let _ = try_enqueue_peer( &mut slots, trigger.hint.claimed_peer_id, false, ); } } () = wait_for_deadline(deadline), if deadline.is_some() => {} } } work_cancellation.cancel(); drain_state_sync_children(&mut pulls, fanout).await; Ok(()) } fn settle_pull_completion( slot: &mut PeerSlot, peer_id: PeerId, result: eyre::Result, ) { slot.in_flight = false; match result { Ok(PeerRefreshOutcome::DeferredByCandidate) => { // Candidate ownership is temporary authority, not a completed // refresh. Retain exactly one coalesced retry behind the existing // rate gate. slot.pending = true; slot.pinned = true; } Ok(PeerRefreshOutcome::Completed | PeerRefreshOutcome::Stale) => {} Err(error) => log::warn!("Failed to refresh peer {peer_id}: {error:#}"), } } async fn drain_state_sync_children( pulls: &mut FuturesUnordered, fanout: Option, ) { while let Some(completed) = pulls.next().await { if let Err(error) = completed.result { log::debug!("Peer refresh stopped during state-sync shutdown: {error:#}"); } } if let Some(fanout) = fanout { fanout.await; } } /// Returns whether the trigger was coalesced into a tracked slot. A pinned /// trigger that returns false must remain actor-owned so bounded-channel /// backpressure is preserved until a slot becomes replaceable. fn try_enqueue_peer(slots: &mut HashMap, peer_id: PeerId, pinned: bool) -> bool { if let Some(slot) = slots.get_mut(&peer_id) { slot.pending = true; slot.pinned |= pinned; return true; } if slots.len() >= MAX_TRACKED_PEERS { if !pinned { return false; } let replaceable = slots .iter() .find_map(|(id, slot)| (!slot.in_flight && !slot.pinned).then_some(*id)); let Some(replaceable) = replaceable else { return false; }; slots.remove(&replaceable); } slots.insert( peer_id, PeerSlot { in_flight: false, pending: true, pinned, next_allowed: tokio::time::Instant::now(), }, ); true } fn start_ready_pulls( slots: &mut HashMap, pulls: &mut FuturesUnordered, handshake_ctx: &HandshakeCtx, ) { while pulls.len() < MAX_CONCURRENT_PULLS { let Some(peer_id) = take_ready_peer(slots, tokio::time::Instant::now()) else { break; }; let ctx = handshake_ctx.clone(); pulls.push(Box::pin(async move { let result = perform_refresh_for_peer(ctx, peer_id).await; PullCompletion { peer_id, result } })); } } fn take_ready_peer( slots: &mut HashMap, now: tokio::time::Instant, ) -> Option { let peer_id = slots .iter() .filter(|(_, slot)| slot.pending && !slot.in_flight && slot.next_allowed <= now) .max_by_key(|(_, slot)| slot.pinned) .map(|(peer_id, _)| *peer_id)?; let slot = slots .get_mut(&peer_id) .expect("selected state-sync slot must still exist"); slot.in_flight = true; slot.pending = false; slot.pinned = false; slot.next_allowed = now + PEER_PULL_COALESCE_WINDOW; Some(peer_id) } async fn perform_refresh_for_peer( ctx: HandshakeCtx, peer_id: PeerId, ) -> eyre::Result { let snapshot = ctx.peer_liveness_for(peer_id).await; let Some(snapshot) = snapshot else { return Ok(PeerRefreshOutcome::Stale); }; perform_peer_refresh(ctx, snapshot).await } async fn hint_requires_pull(ctx: &NetworkServiceCtx, trigger: HintTrigger) -> bool { let snapshot = ctx .peer_game_db .read() .await .revision_snapshot(&trigger.hint.claimed_peer_id); hint_requires_pull_from_snapshot(ctx.peer_id, trigger, snapshot.as_ref()) } /// Decides whether an untrusted hint justifies a pinned pull. /// /// Requester connections are not authenticated, so the claimed peer ID in a /// hint proves nothing by itself. The hint is honoured only when it arrived /// from the IP address at which the claimed peer was last authenticated; /// anything else is discarded as a probable forgery. A genuine peer whose /// address changed is still picked up by mDNS rediscovery and pinned /// liveness reconciliation, which do not depend on hints. fn hint_requires_pull_from_snapshot( local_peer_id: PeerId, trigger: HintTrigger, snapshot: Option<&PeerRevisionSnapshot>, ) -> bool { if trigger.hint.claimed_peer_id == local_peer_id { return false; } let Some(snapshot) = snapshot else { return false; }; if trigger.source_ip != Some(snapshot.endpoint.addr.ip()) { log::debug!( "Ignoring change hint for {} from {:?}; peer is authenticated at {}", trigger.hint.claimed_peer_id, trigger.source_ip, snapshot.endpoint.addr ); return false; } if snapshot.runtime_session_id != trigger.hint.runtime_session_id { return true; } match trigger.domain { StateDomain::Library => snapshot.library_revision != Some(trigger.hint.revision), StateDomain::CallToPlay => snapshot.call_to_play_revision != Some(trigger.hint.revision), } } fn expire_idle_slots(slots: &mut HashMap) { let now = tokio::time::Instant::now(); slots.retain(|_, slot| slot.in_flight || slot.pending || slot.next_allowed > now); } fn next_slot_deadline(slots: &HashMap) -> Option { let now = tokio::time::Instant::now(); slots .values() .filter(|slot| !slot.in_flight && slot.next_allowed > now) .map(|slot| slot.next_allowed) .min() } async fn wait_for_deadline(deadline: Option) { if let Some(deadline) = deadline { tokio::time::sleep_until(deadline).await; } } async fn fanout_local_hints( ctx: NetworkServiceCtx, previous: LocalRevisions, target: LocalRevisions, cancellation: CancellationToken, ) { let endpoints = ctx.peer_game_db.read().await.peer_endpoints(); let local_peer_id = ctx.peer_id; let runtime_session_id = ctx.runtime_session_id; let send_library = target.library != previous.library; let send_call_to_play = target.call_to_play != previous.call_to_play; let deliveries = endpoints.into_iter().map(|endpoint| { let quic = ctx.quic.clone(); let cancellation = cancellation.clone(); async move { if send_library { let hint = ChangeHint { claimed_peer_id: local_peer_id, runtime_session_id, revision: target.library, }; if let Err(error) = send_library_changed(&quic, &endpoint, hint, &cancellation).await { log::debug!( "Failed to send library hint to {}: {error:#}", endpoint.addr ); } } if send_call_to_play { let hint = ChangeHint { claimed_peer_id: local_peer_id, runtime_session_id, revision: target.call_to_play, }; if let Err(error) = send_call_to_play_changed(&quic, &endpoint, hint, &cancellation).await { log::debug!( "Failed to send Call-to-Play hint to {}: {error:#}", endpoint.addr ); } } } }); drive_bounded_hint_fanout(deliveries).await; } async fn drive_bounded_hint_fanout(deliveries: impl IntoIterator) where F: Future, { let mut deliveries = futures::stream::iter(deliveries).buffer_unordered(MAX_CONCURRENT_HINT_SENDS); while deliveries.next().await.is_some() {} } #[cfg(test)] mod tests { use std::sync::atomic::{AtomicUsize, Ordering}; use lanspread_proto::RuntimeSessionId; use super::*; fn peer(seed: u8) -> PeerId { PeerId::from_bytes([seed; 32]) } fn session(seed: u8) -> RuntimeSessionId { RuntimeSessionId::from_bytes([seed; 16]) } #[test] fn repeated_triggers_coalesce_to_one_pending_follow_up() { let mut slots = HashMap::new(); assert!(try_enqueue_peer(&mut slots, peer(1), false)); let now = tokio::time::Instant::now(); assert_eq!(take_ready_peer(&mut slots, now), Some(peer(1))); for _ in 0..100 { assert!(try_enqueue_peer(&mut slots, peer(1), false)); } assert_eq!(slots.len(), 1); assert!(slots[&peer(1)].pending); slots .get_mut(&peer(1)) .expect("slot should exist") .in_flight = false; assert_eq!( take_ready_peer( &mut slots, now + PEER_PULL_COALESCE_WINDOW - Duration::from_millis(1) ), None ); assert_eq!( take_ready_peer(&mut slots, now + PEER_PULL_COALESCE_WINDOW), Some(peer(1)) ); assert!(!slots[&peer(1)].pending); slots .get_mut(&peer(1)) .expect("slot should exist") .in_flight = false; assert_eq!( take_ready_peer(&mut slots, now + PEER_PULL_COALESCE_WINDOW), None, "completion or failure must not self-retry" ); } #[test] fn candidate_deferral_rearms_exactly_one_rate_limited_pinned_retry() { let peer_id = peer(1); let mut slots = HashMap::new(); assert!(try_enqueue_peer(&mut slots, peer_id, true)); let now = tokio::time::Instant::now(); assert_eq!(take_ready_peer(&mut slots, now), Some(peer_id)); settle_pull_completion( slots.get_mut(&peer_id).expect("slot should remain"), peer_id, Ok(PeerRefreshOutcome::DeferredByCandidate), ); assert!(slots[&peer_id].pending); assert!(slots[&peer_id].pinned); assert_eq!(slots.len(), 1); assert_eq!( take_ready_peer( &mut slots, now + PEER_PULL_COALESCE_WINDOW - Duration::from_millis(1), ), None ); assert_eq!( take_ready_peer(&mut slots, now + PEER_PULL_COALESCE_WINDOW), Some(peer_id) ); settle_pull_completion( slots.get_mut(&peer_id).expect("slot should remain"), peer_id, Ok(PeerRefreshOutcome::Completed), ); assert!(!slots[&peer_id].pending); assert!(!slots[&peer_id].pinned); assert_eq!( take_ready_peer(&mut slots, now + PEER_PULL_COALESCE_WINDOW), None ); } #[test] fn pinned_trigger_displaces_only_idle_untrusted_slot_at_capacity() { let mut slots = HashMap::new(); for seed in 0..u8::try_from(MAX_TRACKED_PEERS).expect("bound fits u8") { assert!(try_enqueue_peer(&mut slots, peer(seed), false)); } assert!(try_enqueue_peer(&mut slots, peer(200), true)); assert_eq!(slots.len(), MAX_TRACKED_PEERS); assert!(slots.contains_key(&peer(200))); assert!(slots[&peer(200)].pinned); assert_eq!( take_ready_peer(&mut slots, tokio::time::Instant::now()), Some(peer(200)), "pinned work must start before saturated untrusted work" ); } #[test] fn pinned_trigger_waits_when_every_slot_is_nonreplaceable() { let mut slots = HashMap::new(); for seed in 0..u8::try_from(MAX_TRACKED_PEERS).expect("bound fits u8") { assert!(try_enqueue_peer(&mut slots, peer(seed), true)); } let blocked = peer(200); assert!(!try_enqueue_peer(&mut slots, blocked, true)); assert_eq!(slots.len(), MAX_TRACKED_PEERS); assert!(!slots.contains_key(&blocked)); let completed = peer(0); let slot = slots .get_mut(&completed) .expect("tracked pinned slot should exist"); slot.pinned = false; slot.pending = false; slot.in_flight = false; assert!(try_enqueue_peer(&mut slots, blocked, true)); assert_eq!(slots.len(), MAX_TRACKED_PEERS); assert!(slots.contains_key(&blocked)); assert!(!slots.contains_key(&completed)); } #[test] fn quiet_cooldown_slots_have_an_expiry_deadline_before_new_hint_admission() { let now = tokio::time::Instant::now(); let mut slots = HashMap::new(); for seed in 0..u8::try_from(MAX_TRACKED_PEERS).expect("bound fits u8") { slots.insert( peer(seed), PeerSlot { in_flight: false, pending: false, pinned: false, next_allowed: now + PEER_PULL_COALESCE_WINDOW, }, ); } assert_eq!( next_slot_deadline(&slots), Some(now + PEER_PULL_COALESCE_WINDOW) ); } #[test] fn scheduler_starts_at_most_eight_of_sixty_four_ready_peers() { let mut slots = HashMap::new(); for seed in 0..u8::try_from(MAX_TRACKED_PEERS).expect("bound fits u8") { assert!(try_enqueue_peer(&mut slots, peer(seed), false)); } let now = tokio::time::Instant::now(); let started = (0..MAX_CONCURRENT_PULLS) .filter_map(|_| take_ready_peer(&mut slots, now)) .collect::>(); assert_eq!(started.len(), MAX_CONCURRENT_PULLS); assert_eq!(slots.values().filter(|slot| slot.in_flight).count(), 8); assert_eq!(slots.values().filter(|slot| slot.pending).count(), 56); assert_eq!( next_slot_deadline(&slots), None, "ready work waiting on the global pull bound must be completion-driven" ); } #[test] fn unknown_and_self_hints_do_not_allocate_slots() { let local = peer(1); let mut slots = HashMap::new(); for claimed_peer_id in [local, peer(2)] { let trigger = HintTrigger { domain: StateDomain::Library, hint: ChangeHint { claimed_peer_id, runtime_session_id: session(1), revision: 99, }, source_ip: Some(IpAddr::from([192, 168, 1, 2])), }; if hint_requires_pull_from_snapshot(local, trigger, None) { let _ = try_enqueue_peer(&mut slots, claimed_peer_id, false); } } assert!(slots.is_empty()); } #[test] fn hints_are_honoured_only_from_the_claimed_peers_authenticated_address() { use lanspread_proto::PeerEndpoint; use crate::peer_db::PeerEndpointGeneration; let local = peer(1); let claimed = peer(2); let peer_addr: std::net::SocketAddr = "192.168.1.20:42424".parse().expect("addr"); let snapshot = PeerRevisionSnapshot { endpoint: PeerEndpoint::new(claimed, peer_addr), generation: PeerEndpointGeneration::for_tests(1), runtime_session_id: session(1), library_revision: Some(3), call_to_play_revision: Some(0), }; let trigger = |source_ip: Option| HintTrigger { domain: StateDomain::Library, hint: ChangeHint { claimed_peer_id: claimed, runtime_session_id: session(1), revision: 4, }, source_ip, }; assert!(hint_requires_pull_from_snapshot( local, trigger(Some(peer_addr.ip())), Some(&snapshot) )); // A third host claiming the peer's identity must not trigger a pull. assert!(!hint_requires_pull_from_snapshot( local, trigger(Some(IpAddr::from([192, 168, 1, 99]))), Some(&snapshot) )); assert!(!hint_requires_pull_from_snapshot( local, trigger(None), Some(&snapshot) )); // The revision check still applies for a matching source. let same_revision = HintTrigger { hint: ChangeHint { revision: 3, ..trigger(None).hint }, ..trigger(Some(peer_addr.ip())) }; assert!(!hint_requires_pull_from_snapshot( local, same_revision, Some(&snapshot) )); } #[tokio::test] async fn local_revision_watch_coalesces_to_latest_value() { let handle = StateSyncHandle::new(peer(1)); let mut rx = handle.local_revisions_tx.subscribe(); handle.publish_library_revision(1); handle.publish_library_revision(2); handle.publish_library_revision(7); rx.changed().await.expect("watch sender should remain live"); assert_eq!(rx.borrow_and_update().library, 7); assert!(!rx.has_changed().expect("watch sender should remain live")); } #[tokio::test] async fn hint_fanout_never_exceeds_eight_simultaneous_sends() { let active = Arc::new(AtomicUsize::new(0)); let maximum = Arc::new(AtomicUsize::new(0)); let deliveries = (0..64).map(|_| { let active = Arc::clone(&active); let maximum = Arc::clone(&maximum); async move { let current = active.fetch_add(1, Ordering::SeqCst) + 1; maximum.fetch_max(current, Ordering::SeqCst); tokio::task::yield_now().await; active.fetch_sub(1, Ordering::SeqCst); } }); drive_bounded_hint_fanout(deliveries).await; assert!(maximum.load(Ordering::SeqCst) <= MAX_CONCURRENT_HINT_SENDS); assert_eq!(active.load(Ordering::SeqCst), 0); } #[tokio::test] async fn cancellation_drains_pull_and_fanout_futures() { let cancellation = CancellationToken::new(); let completed = Arc::new(AtomicUsize::new(0)); let mut pulls = FuturesUnordered::::new(); for seed in 1..=2 { let cancellation = cancellation.clone(); let completed = Arc::clone(&completed); pulls.push(Box::pin(async move { cancellation.cancelled().await; completed.fetch_add(1, Ordering::SeqCst); PullCompletion { peer_id: peer(seed), result: Ok(PeerRefreshOutcome::Completed), } })); } let fanout_cancellation = cancellation.clone(); let fanout_completed = Arc::clone(&completed); let fanout: FanoutFuture = Box::pin(async move { fanout_cancellation.cancelled().await; fanout_completed.fetch_add(1, Ordering::SeqCst); }); cancellation.cancel(); drain_state_sync_children(&mut pulls, Some(fanout)).await; assert_eq!(completed.load(Ordering::SeqCst), 3); } }