use std::{
collections::{BTreeMap, HashMap},
sync::Mutex as StdMutex,
time::Instant,
};
use zakura_chain::block;
use super::{
config::{clamp_advertised_blocks, clamp_advertised_inflight, clamp_advertised_response_bytes},
state::EFFECTIVE_BS_OUTBOUND_INFLIGHT_PER_PEER,
BlockSyncStatus, ServicePeerDirection, ZakuraPeerId,
};
#[derive(Clone, Debug)]
pub(super) struct Entry {
pub(super) direction: ServicePeerDirection,
pub(super) servable_low: block::Height,
pub(super) servable_high: block::Height,
pub(super) received_status: bool,
pub(super) max_blocks_per_response: u32,
pub(super) max_inflight_requests: u32,
pub(super) max_response_bytes: u32,
pub(super) outstanding: BTreeMap<block::Height, OutstandingMeta>,
pub(super) slots: SlotDiagnostics,
pub(super) floor_watchdog_avoid: BTreeMap<block::Height, Instant>,
pub(super) generation: u64,
}
impl Entry {
fn new(
direction: ServicePeerDirection,
config: &super::ZakuraBlockSyncConfig,
generation: u64,
) -> Self {
Self {
direction,
servable_low: block::Height::MIN,
servable_high: block::Height::MIN,
received_status: false,
max_blocks_per_response: config.advertised_max_blocks_per_response(),
max_inflight_requests: config.advertised_max_inflight_requests(),
max_response_bytes: config.advertised_max_response_bytes(),
outstanding: BTreeMap::new(),
slots: SlotDiagnostics::default(),
floor_watchdog_avoid: BTreeMap::new(),
generation,
}
}
}
#[derive(Copy, Clone, Debug, Default)]
pub(super) struct SlotDiagnostics {
pub(super) hard_capacity: usize,
pub(super) effective_window: usize,
pub(super) available_slots: usize,
pub(super) outstanding_requests: usize,
pub(super) bbr_rtprop_ms: Option<u64>,
}
#[derive(Copy, Clone, Debug)]
pub(super) struct OutstandingMeta {
pub(super) hash: block::Hash,
pub(super) estimated_bytes: u64,
pub(super) queued_at: Instant,
pub(super) deadline: Instant,
}
#[derive(Clone, Debug)]
pub(super) struct OutstandingClaim {
pub(super) peer: ZakuraPeerId,
pub(super) height: block::Height,
pub(super) meta: OutstandingMeta,
}
#[derive(Debug)]
pub(super) struct PeerRegistry {
peers: StdMutex<HashMap<ZakuraPeerId, Entry>>,
parked_peers: StdMutex<HashMap<ZakuraPeerId, Instant>>,
next_generation: std::sync::atomic::AtomicU64,
}
impl Default for PeerRegistry {
fn default() -> Self {
Self::new()
}
}
impl PeerRegistry {
pub(super) fn new() -> Self {
Self {
peers: StdMutex::new(HashMap::new()),
parked_peers: StdMutex::new(HashMap::new()),
next_generation: std::sync::atomic::AtomicU64::new(1),
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, HashMap<ZakuraPeerId, Entry>> {
self.peers
.lock()
.expect("peer registry mutex is never poisoned")
}
fn lock_parked(&self) -> std::sync::MutexGuard<'_, HashMap<ZakuraPeerId, Instant>> {
self.parked_peers
.lock()
.expect("peer registry parked-peer mutex is never poisoned")
}
pub(super) fn park_peer_until(&self, peer: &ZakuraPeerId, until: Instant) {
self.lock_parked().insert(peer.clone(), until);
}
pub(super) fn is_peer_parked(&self, peer: &ZakuraPeerId, now: Instant) -> bool {
let mut parked_peers = self.lock_parked();
parked_peers.retain(|_, until| *until > now);
parked_peers.get(peer).is_some_and(|until| *until > now)
}
pub(super) fn admit(
&self,
peer: &ZakuraPeerId,
direction: ServicePeerDirection,
config: &super::ZakuraBlockSyncConfig,
) -> u64 {
let generation = self
.next_generation
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let mut peers = self.lock();
peers
.entry(peer.clone())
.and_modify(|entry| {
entry.direction = direction;
entry.outstanding.clear();
entry.floor_watchdog_avoid.clear();
entry.generation = generation;
})
.or_insert_with(|| Entry::new(direction, config, generation));
generation
}
pub(super) fn remove(&self, peer: &ZakuraPeerId) {
self.lock().remove(peer);
}
pub(super) fn upsert_status(
&self,
peer: &ZakuraPeerId,
generation: u64,
status: BlockSyncStatus,
) {
let mut peers = self.lock();
let Some(entry) = peers.get_mut(peer) else {
return;
};
if entry.generation != generation {
return;
}
entry.servable_low = status.servable_low;
entry.servable_high = status.servable_high;
entry.max_blocks_per_response = clamp_advertised_blocks(status.max_blocks_per_response);
entry.max_inflight_requests = clamp_advertised_inflight(status.max_inflight_requests);
entry.max_response_bytes = clamp_advertised_response_bytes(status.max_response_bytes);
entry.received_status = true;
}
pub(super) fn set_outstanding(
&self,
peer: &ZakuraPeerId,
generation: u64,
outstanding: BTreeMap<block::Height, OutstandingMeta>,
) {
let mut peers = self.lock();
if let Some(entry) = peers.get_mut(peer) {
if entry.generation == generation {
entry.outstanding = outstanding;
}
}
}
pub(super) fn clear_outstanding(&self, peer: &ZakuraPeerId, generation: u64) {
let mut peers = self.lock();
if let Some(entry) = peers.get_mut(peer) {
if entry.generation == generation {
entry.outstanding.clear();
}
}
}
pub(super) fn publish_slots(
&self,
peer: &ZakuraPeerId,
generation: u64,
slots: SlotDiagnostics,
) {
let mut peers = self.lock();
if let Some(entry) = peers.get_mut(peer) {
if entry.generation == generation {
entry.slots = slots;
}
}
}
pub(super) fn slot_summary(&self) -> SlotSummary {
let peers = self.lock();
let mut summary = SlotSummary::default();
for entry in peers.values() {
summary.outstanding_requests = summary
.outstanding_requests
.saturating_add(entry.slots.outstanding_requests);
if !entry.received_status {
continue;
}
summary.capacity = summary.capacity.saturating_add(entry.slots.hard_capacity);
summary.effective_window = summary
.effective_window
.saturating_add(entry.slots.effective_window);
summary.available = summary
.available
.saturating_add(entry.slots.available_slots);
if entry.slots.available_slots == 0 {
summary.saturated_peers = summary.saturated_peers.saturating_add(1);
}
}
summary
}
pub(super) fn has_outstanding_request(&self, height: block::Height, hash: block::Hash) -> bool {
let peers = self.lock();
peers.values().any(|entry| {
entry
.outstanding
.get(&height)
.is_some_and(|meta| meta.hash == hash)
})
}
pub(super) fn has_outstanding_height(&self, height: block::Height) -> bool {
let peers = self.lock();
peers
.values()
.any(|entry| entry.outstanding.contains_key(&height))
}
pub(super) fn peer_has_outstanding_height(
&self,
peer: &ZakuraPeerId,
height: block::Height,
) -> bool {
let peers = self.lock();
peers
.get(peer)
.is_some_and(|entry| entry.outstanding.contains_key(&height))
}
pub(super) fn total_unreceived(&self) -> usize {
let peers = self.lock();
peers.values().map(|entry| entry.outstanding.len()).sum()
}
pub(super) fn any_outstanding_at_or_above(&self, at_or_above: block::Height) -> bool {
let peers = self.lock();
peers.values().any(|entry| {
entry
.outstanding
.keys()
.any(|height| *height >= at_or_above)
})
}
pub(super) fn any_outstanding_conflicts_at(
&self,
height: block::Height,
hash: block::Hash,
) -> bool {
let peers = self.lock();
peers.values().any(|entry| {
entry
.outstanding
.get(&height)
.is_some_and(|expected| expected.hash != hash)
})
}
pub(super) fn has_received_status(&self, peer: &ZakuraPeerId) -> bool {
let peers = self.lock();
peers.get(peer).is_some_and(|entry| entry.received_status)
}
pub(super) fn peers_with_status(&self) -> usize {
let peers = self.lock();
peers.values().filter(|entry| entry.received_status).count()
}
pub(super) fn candidate_snapshot(
&self,
) -> Vec<(ZakuraPeerId, bool, block::Height, block::Height)> {
let peers = self.lock();
peers
.iter()
.map(|(peer, entry)| {
(
peer.clone(),
entry.received_status,
entry.servable_low,
entry.servable_high,
)
})
.collect()
}
pub(super) fn direction_status_counts(&self) -> DirectionStatusCounts {
let peers = self.lock();
let mut counts = DirectionStatusCounts::default();
for entry in peers.values() {
match entry.direction {
ServicePeerDirection::Inbound => {
counts.inbound += 1;
if entry.received_status {
counts.inbound_with_status += 1;
}
}
ServicePeerDirection::Outbound => {
counts.outbound += 1;
if entry.received_status {
counts.outbound_with_status += 1;
}
}
}
}
counts
}
pub(super) fn floor_gap_servable(&self, height: block::Height) -> (usize, usize) {
let peers = self.lock();
let mut servable = 0usize;
let mut outstanding = 0usize;
for entry in peers.values() {
if entry.received_status
&& entry.servable_low <= height
&& height <= entry.servable_high
{
servable = servable.saturating_add(1);
}
if entry.outstanding.contains_key(&height) {
outstanding = outstanding.saturating_add(1);
}
}
(servable, outstanding)
}
pub(super) fn earliest_outstanding_deadline_at(
&self,
height: block::Height,
) -> Option<Instant> {
let peers = self.lock();
peers
.values()
.filter_map(|entry| entry.outstanding.get(&height).map(|meta| meta.deadline))
.min()
}
pub(super) fn floor_has_preferred_unsaturated_server(
&self,
height: block::Height,
self_peer: &ZakuraPeerId,
self_rtprop_ms: Option<u64>,
allow_equal_score: bool,
) -> bool {
let self_score = self_rtprop_ms.unwrap_or(u64::MAX);
let peers = self.lock();
peers.iter().any(|(peer, entry)| {
if peer == self_peer || !entry.can_serve_with_room(height) {
return false;
}
let other_score = entry.slots.bbr_rtprop_ms.unwrap_or(u64::MAX);
if allow_equal_score {
other_score <= self_score
} else {
other_score < self_score
}
})
}
pub(super) fn outstanding_claims_at(&self, height: block::Height) -> Vec<OutstandingClaim> {
let peers = self.lock();
peers
.iter()
.filter_map(|(peer, entry)| {
entry.outstanding.get(&height).map(|meta| OutstandingClaim {
peer: peer.clone(),
height,
meta: *meta,
})
})
.collect()
}
pub(super) fn clear_outstanding_height(&self, peer: &ZakuraPeerId, height: block::Height) {
let mut peers = self.lock();
if let Some(entry) = peers.get_mut(peer) {
entry.outstanding.remove(&height);
}
}
pub(super) fn avoid_floor_height_until(
&self,
peer: &ZakuraPeerId,
height: block::Height,
until: Instant,
) {
let mut peers = self.lock();
if let Some(entry) = peers.get_mut(peer) {
entry.floor_watchdog_avoid.insert(height, until);
}
}
pub(super) fn is_floor_height_avoided(
&self,
peer: &ZakuraPeerId,
height: block::Height,
now: Instant,
) -> bool {
let mut peers = self.lock();
let Some(entry) = peers.get_mut(peer) else {
return false;
};
entry.floor_watchdog_avoid.retain(|_, until| *until > now);
entry
.floor_watchdog_avoid
.get(&height)
.is_some_and(|until| *until > now)
}
pub(super) fn next_floor_avoid_deadline(
&self,
peer: &ZakuraPeerId,
now: Instant,
) -> Option<Instant> {
let mut peers = self.lock();
let entry = peers.get_mut(peer)?;
entry.floor_watchdog_avoid.retain(|_, until| *until > now);
entry.floor_watchdog_avoid.values().min().copied()
}
}
impl Entry {
fn can_serve_with_room(&self, height: block::Height) -> bool {
self.received_status
&& self.servable_low <= height
&& height <= self.servable_high
&& self.slots.available_slots > 0
}
}
#[derive(Copy, Clone, Debug, Default)]
pub(super) struct SlotSummary {
pub(super) capacity: usize,
pub(super) effective_window: usize,
pub(super) available: usize,
pub(super) saturated_peers: usize,
pub(super) outstanding_requests: usize,
}
#[derive(Copy, Clone, Debug, Default)]
pub(super) struct DirectionStatusCounts {
pub(super) inbound: usize,
pub(super) outbound: usize,
pub(super) inbound_with_status: usize,
pub(super) outbound_with_status: usize,
}
pub(super) fn hard_outbound_capacity(max_inflight_requests: u32) -> usize {
usize::try_from(max_inflight_requests)
.expect("u32 max inflight requests fits in usize on supported targets")
.min(EFFECTIVE_BS_OUTBOUND_INFLIGHT_PER_PEER)
}
#[cfg(test)]
mod floor_bias_tests {
use super::*;
fn peer(byte: u8) -> ZakuraPeerId {
ZakuraPeerId::new(vec![byte; 32]).expect("32-byte test peer id is valid")
}
fn register_with_rtprop(
reg: &PeerRegistry,
config: &super::super::ZakuraBlockSyncConfig,
peer: &ZakuraPeerId,
low: u32,
high: u32,
available: usize,
bbr_rtprop_ms: Option<u64>,
) {
let generation = reg.admit(peer, ServicePeerDirection::Outbound, config);
reg.upsert_status(
peer,
generation,
BlockSyncStatus {
servable_low: block::Height(low),
servable_high: block::Height(high),
..BlockSyncStatus::default()
},
);
reg.publish_slots(
peer,
generation,
SlotDiagnostics {
available_slots: available,
bbr_rtprop_ms,
..SlotDiagnostics::default()
},
);
}
fn register(
reg: &PeerRegistry,
config: &super::super::ZakuraBlockSyncConfig,
peer: &ZakuraPeerId,
low: u32,
high: u32,
available: usize,
) {
register_with_rtprop(reg, config, peer, low, high, available, None);
}
#[test]
fn bypass_defers_to_an_equal_or_faster_unsaturated_other_server() {
let config = super::super::ZakuraBlockSyncConfig::default();
let reg = PeerRegistry::new();
let (a, b) = (peer(1), peer(2));
register_with_rtprop(®, &config, &a, 0, 1000, 0, Some(50));
register_with_rtprop(®, &config, &b, 0, 1000, 3, Some(50));
assert!(reg.floor_has_preferred_unsaturated_server(block::Height(100), &a, Some(50), true));
assert!(!reg.floor_has_preferred_unsaturated_server(
block::Height(100),
&b,
Some(50),
true
));
}
#[test]
fn normal_path_defers_only_to_a_strictly_faster_server() {
let config = super::super::ZakuraBlockSyncConfig::default();
let reg = PeerRegistry::new();
let (slow, fast) = (peer(1), peer(2));
register_with_rtprop(®, &config, &slow, 0, 1000, 3, Some(120));
register_with_rtprop(®, &config, &fast, 0, 1000, 3, Some(40));
assert!(reg.floor_has_preferred_unsaturated_server(
block::Height(100),
&slow,
Some(120),
false
));
assert!(!reg.floor_has_preferred_unsaturated_server(
block::Height(100),
&fast,
Some(40),
false
));
}
#[test]
fn normal_path_keeps_equal_carriers_eligible() {
let config = super::super::ZakuraBlockSyncConfig::default();
let reg = PeerRegistry::new();
let (a, b) = (peer(1), peer(2));
register_with_rtprop(®, &config, &a, 0, 1000, 3, Some(50));
register_with_rtprop(®, &config, &b, 0, 1000, 3, Some(50));
assert!(!reg.floor_has_preferred_unsaturated_server(
block::Height(100),
&a,
Some(50),
false
));
assert!(!reg.floor_has_preferred_unsaturated_server(
block::Height(100),
&b,
Some(50),
false
));
}
#[test]
fn saturated_fast_peer_does_not_defer_to_slower_unsaturated_peer() {
let config = super::super::ZakuraBlockSyncConfig::default();
let reg = PeerRegistry::new();
let (fast, slow) = (peer(1), peer(2));
register_with_rtprop(®, &config, &fast, 0, 1000, 0, Some(40));
register_with_rtprop(®, &config, &slow, 0, 1000, 3, Some(120));
assert!(!reg.floor_has_preferred_unsaturated_server(
block::Height(100),
&fast,
Some(40),
true
));
}
#[test]
fn bypasses_when_every_server_is_saturated() {
let config = super::super::ZakuraBlockSyncConfig::default();
let reg = PeerRegistry::new();
let (a, b) = (peer(1), peer(2));
register(®, &config, &a, 0, 1000, 0);
register(®, &config, &b, 0, 1000, 0);
assert!(!reg.floor_has_preferred_unsaturated_server(block::Height(100), &a, None, true));
}
#[test]
fn ignores_an_unsaturated_peer_that_cannot_serve_the_floor() {
let config = super::super::ZakuraBlockSyncConfig::default();
let reg = PeerRegistry::new();
let (a, b) = (peer(1), peer(2));
register(®, &config, &a, 0, 1000, 0);
register(®, &config, &b, 500, 1000, 3);
assert!(!reg.floor_has_preferred_unsaturated_server(block::Height(100), &a, None, true));
}
#[test]
fn floor_avoid_deadline_prunes_expired_entries_and_returns_next_wake() {
let config = super::super::ZakuraBlockSyncConfig::default();
let reg = PeerRegistry::new();
let peer = peer(1);
reg.admit(&peer, ServicePeerDirection::Outbound, &config);
let now = Instant::now();
reg.avoid_floor_height_until(
&peer,
block::Height(1),
now - std::time::Duration::from_secs(1),
);
reg.avoid_floor_height_until(
&peer,
block::Height(2),
now + std::time::Duration::from_secs(2),
);
reg.avoid_floor_height_until(
&peer,
block::Height(3),
now + std::time::Duration::from_secs(1),
);
assert_eq!(
reg.next_floor_avoid_deadline(&peer, now),
Some(now + std::time::Duration::from_secs(1)),
);
assert!(!reg.is_floor_height_avoided(&peer, block::Height(1), now));
assert!(reg.is_floor_height_avoided(&peer, block::Height(2), now));
}
#[test]
fn parked_peer_expires_after_cooldown() {
let reg = PeerRegistry::new();
let peer = peer(1);
let now = Instant::now();
reg.park_peer_until(&peer, now + std::time::Duration::from_secs(1));
assert!(reg.is_peer_parked(&peer, now));
assert!(!reg.is_peer_parked(&peer, now + std::time::Duration::from_secs(2)));
}
}