feat(peer)!: cut over to authenticated catalog sharing
Replace address-only trust and pushed peer state with installation identities, SPKI-pinned QUIC, candidate-only discovery, and bounded responder-owned protocol-8 pulls. The runtime now owns each network generation and all admitted work through shutdown. Add exact bundled content identities, reproducible manifest publishing, capability-confined downloads, streaming BLAKE3 verification, quarantine and retry, and crash-recoverable download and install transactions. Ship generated fixture catalogs and fail closed when production manifests are absent. The Tauri backend exposes durable sharing policy, redacted identity state, and attempt-keyed transfer snapshots. Frontend consumption follows in the next commit. Repository-wide test certificates and protocol-7 paths are removed. BREAKING CHANGE: peers must use protocol 8 and exact catalog content artifacts; protocol-7 frames and shared-certificate identities are no longer accepted. Test Plan: - `just test` -- passed on the completed stack (708 workspace tests) - `just clippy` -- passed on the completed stack - `just build` -- passed with fixture catalogs on the completed stack - `just catalog-check-production` -- failed closed because the external production manifest corpus is absent - `git diff --cached --check` -- passed
This commit is contained in:
128 files changed
+51759
-10784
No files matched your search
@@ -0,0 +1,727 @@
|
||||
//! Bounded scheduling for responder-owned state pulls and change-hint fanout.
|
||||
|
||||
use std::{collections::HashMap, future::Future, 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,
|
||||
}
|
||||
|
||||
struct StateSyncInbox {
|
||||
pinned_rx: Mutex<mpsc::Receiver<PeerId>>,
|
||||
hint_rx: Mutex<mpsc::Receiver<HintTrigger>>,
|
||||
}
|
||||
|
||||
/// 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<PeerId>,
|
||||
hint_tx: mpsc::Sender<HintTrigger>,
|
||||
local_revisions_tx: watch::Sender<LocalRevisions>,
|
||||
inbox: Arc<StateSyncInbox>,
|
||||
}
|
||||
|
||||
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.
|
||||
pub(crate) fn schedule_hint(&self, domain: StateDomain, hint: ChangeHint) {
|
||||
if hint.claimed_peer_id == self.local_peer_id {
|
||||
return;
|
||||
}
|
||||
match self.hint_tx.try_send(HintTrigger { domain, hint }) {
|
||||
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<PeerRefreshOutcome>,
|
||||
}
|
||||
|
||||
type PullFuture = Pin<Box<dyn Future<Output = PullCompletion> + Send>>;
|
||||
type FanoutFuture = Pin<Box<dyn Future<Output = ()> + 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<PeerEvent>,
|
||||
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<FanoutFuture> = 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::<PeerId, PeerSlot>::new();
|
||||
let mut pulls = FuturesUnordered::<PullFuture>::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<PeerRefreshOutcome>,
|
||||
) {
|
||||
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<PullFuture>,
|
||||
fanout: Option<FanoutFuture>,
|
||||
) {
|
||||
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<PeerId, PeerSlot>, 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<PeerId, PeerSlot>,
|
||||
pulls: &mut FuturesUnordered<PullFuture>,
|
||||
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<PeerId, PeerSlot>,
|
||||
now: tokio::time::Instant,
|
||||
) -> Option<PeerId> {
|
||||
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<PeerRefreshOutcome> {
|
||||
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())
|
||||
}
|
||||
|
||||
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 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<PeerId, PeerSlot>) {
|
||||
let now = tokio::time::Instant::now();
|
||||
slots.retain(|_, slot| slot.in_flight || slot.pending || slot.next_allowed > now);
|
||||
}
|
||||
|
||||
fn next_slot_deadline(slots: &HashMap<PeerId, PeerSlot>) -> Option<tokio::time::Instant> {
|
||||
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<tokio::time::Instant>) {
|
||||
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<F>(deliveries: impl IntoIterator<Item = F>)
|
||||
where
|
||||
F: Future<Output = ()>,
|
||||
{
|
||||
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::<Vec<_>>();
|
||||
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,
|
||||
},
|
||||
};
|
||||
if hint_requires_pull_from_snapshot(local, trigger, None) {
|
||||
let _ = try_enqueue_peer(&mut slots, claimed_peer_id, false);
|
||||
}
|
||||
}
|
||||
assert!(slots.is_empty());
|
||||
}
|
||||
|
||||
#[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::<PullFuture>::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);
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user