use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use super::ParallelIBDConfig;
use super::latch_env;
use super::types::ChunkWorkItem;
#[derive(Debug, Clone)]
struct StickyWanTenure {
peer: String,
start_next_needed: u64,
started_at: Instant,
}
#[derive(Debug, Clone)]
struct TipStreamWindow {
streams: u64,
started: Instant,
last_stream: Instant,
}
#[derive(Debug, Clone)]
struct TipTrial {
sticky: String,
challenger: String,
started: Instant,
sticky_streams_at_start: u64,
challenger_streams_at_start: u64,
next_needed_at_start: u64,
}
static ASSIGNER_GW_WAIT_NS: AtomicU64 = AtomicU64::new(0);
static ASSIGNER_GW_HOLD_NS: AtomicU64 = AtomicU64::new(0);
static ASSIGNER_GW_SAMPLES: AtomicU64 = AtomicU64::new(0);
fn assigner_lock_timers_enabled() -> bool {
matches!(
std::env::var("BLVM_IBD_ASSIGNER_LOCK_TIMERS")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
struct AssignerGetWorkTimer {
wait_ns: u64,
held_at: Instant,
}
impl AssignerGetWorkTimer {
fn start(wait_ns: u64) -> Option<Self> {
if !assigner_lock_timers_enabled() {
return None;
}
Some(Self {
wait_ns,
held_at: Instant::now(),
})
}
}
impl Drop for AssignerGetWorkTimer {
fn drop(&mut self) {
let hold_ns = self.held_at.elapsed().as_nanos() as u64;
ASSIGNER_GW_WAIT_NS.fetch_add(self.wait_ns, Ordering::Relaxed);
ASSIGNER_GW_HOLD_NS.fetch_add(hold_ns, Ordering::Relaxed);
let n = ASSIGNER_GW_SAMPLES.fetch_add(1, Ordering::Relaxed) + 1;
if n % 4096 == 0 {
let wait = ASSIGNER_GW_WAIT_NS.load(Ordering::Relaxed);
let hold = ASSIGNER_GW_HOLD_NS.load(Ordering::Relaxed);
tracing::info!(
"[IBD_ASSIGNER_LOCK] samples={} avg_wait_us={:.1} avg_hold_us={:.1}",
n,
(wait as f64 / n as f64) / 1000.0,
(hold as f64 / n as f64) / 1000.0
);
}
}
}
#[derive(Debug, Clone)]
pub struct BlockChunk {
pub start_height: u64,
pub end_height: u64,
pub peer_id: String,
}
pub fn create_chunks(
config: &ParallelIBDConfig,
start_height: u64,
end_height: u64,
peer_ids: &[String],
scored_peers: Option<&[(String, f64)]>,
) -> Vec<BlockChunk> {
let mut chunks = Vec::new();
let mut current_height = start_height;
let num_peers = peer_ids.len().max(1);
let mut chunk_index: usize = 0;
let use_fastest = (config.mode.eq_ignore_ascii_case("earliest") || config.earliest_first)
&& num_peers > 1
&& scored_peers.map(|s| !s.is_empty()).unwrap_or(false);
let fastest_peer = if use_fastest {
scored_peers.and_then(|s| {
s.iter()
.max_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal))
.map(|(p, _)| p.clone())
})
} else {
None
};
if use_fastest && fastest_peer.is_some() {
tracing::info!("IBD: earliest-first — all chunks to fastest peer");
} else {
tracing::info!(
"Round-robin chunk assignment: {} peers, chunk_size={}",
num_peers,
config.chunk_size
);
}
while current_height <= end_height {
let (chunk_sz, is_bootstrap) = if current_height == 0 && start_height == 0 {
let sz = 128.min(end_height.saturating_add(1));
(sz, true)
} else {
(config.chunk_size, false)
};
let chunk_end = (current_height + chunk_sz - 1).min(end_height);
if is_bootstrap {
tracing::info!(
"IBD: bootstrap chunk 0-{} (99 and 100 in same chunk)",
chunk_end
);
}
let peer_id = fastest_peer.clone().unwrap_or_else(|| {
if peer_ids.is_empty() {
String::new()
} else {
peer_ids[chunk_index % num_peers].clone()
}
});
chunks.push(BlockChunk {
start_height: current_height,
end_height: chunk_end,
peer_id,
});
current_height = chunk_end + 1;
chunk_index += 1;
}
chunks
}
pub(crate) struct ChunkAssigner {
chunks: Vec<(u64, u64)>,
workers: Vec<String>,
extra_workers: Mutex<HashSet<String>>,
preferred_peers: Vec<String>,
next_index: AtomicUsize,
retry_queue: Mutex<VecDeque<ChunkWorkItem>>,
validation_height: Arc<std::sync::atomic::AtomicU64>,
bootstrap_complete: AtomicBool,
start_height: u64,
in_flight_per_peer: Mutex<HashMap<String, Vec<(u64, u64)>>>,
work_stealing: bool,
blacklisted_until: Mutex<HashMap<String, Instant>>,
last_stall_requeue: Mutex<Option<(u64, Instant)>>,
peer_scores: Mutex<HashMap<String, f64>>,
confirmed_body_height_at_start: AtomicU64,
wan_body_tip: AtomicU64,
tip_gap_missing: AtomicBool,
tip_bridge_holes: AtomicU64,
preferred_tip_owner: Mutex<Option<String>>,
tip_owner_fail_until: Mutex<HashMap<String, Instant>>,
tip_cover_claims: Mutex<Vec<(String, u64, u64)>>,
tip_owner_open: AtomicBool,
tip_failover_once_h: AtomicU64,
tip_failover_once_at_ms: AtomicU64,
tip_ahead_hole_freeze: AtomicBool,
tip_ahead_hole_clear_since_ms: AtomicU64,
ibd_ready_peers: Mutex<HashSet<String>>,
sticky_wan_tenure: Mutex<Option<StickyWanTenure>>,
last_a6m_rotate_at: Mutex<Option<Instant>>,
tip_progress_samples: Mutex<VecDeque<(Instant, u64)>>,
peer_tip_streams: Mutex<HashMap<String, TipStreamWindow>>,
tip_hole_depth: Mutex<HashMap<String, usize>>,
tip_trial: Mutex<Option<TipTrial>>,
last_tip_trial_at: Mutex<Option<Instant>>,
tip_trial_post_open_at: Mutex<Option<Instant>>,
header_tip: AtomicU64,
ibd_end_height: AtomicU64,
shutdown: AtomicBool,
synth_tip_dedup_block_since_ms: AtomicU64,
}
fn a6m_min_bps() -> f64 {
std::env::var("BLVM_IBD_A6M_MIN_BPS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(6.0)
}
fn a6m_floor_min_bps() -> f64 {
std::env::var("BLVM_IBD_A6M_FLOOR_MIN_BPS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(22.0)
}
fn a6m_floor_open_slot_min_bps() -> f64 {
std::env::var("BLVM_IBD_A6M_FLOOR_OPEN_SLOT_MIN_BPS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(12.0)
}
fn a6m_tenure_secs() -> u64 {
std::env::var("BLVM_IBD_A6M_TENURE_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(300)
}
fn a6m_recent_window_secs() -> u64 {
std::env::var("BLVM_IBD_A6M_RECENT_WINDOW_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(60)
.clamp(30, 300)
}
fn a6m_rotate_cooldown_secs() -> u64 {
std::env::var("BLVM_IBD_A6M_ROTATE_COOLDOWN")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(600)
}
fn a6m_floor_rotate_cooldown_secs() -> u64 {
std::env::var("BLVM_IBD_A6M_FLOOR_ROTATE_COOLDOWN")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(120)
}
fn a6m_max_getdata_ms() -> u64 {
std::env::var("BLVM_IBD_A6M_MAX_GETDATA_MS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(800)
.clamp(200, 10_000)
}
fn a6m_gd_slow_feeder_keep() -> usize {
std::env::var("BLVM_IBD_A6M_GD_SLOW_FEEDER_KEEP")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(4usize)
.clamp(0, 64)
}
fn a6m_gd_slow_tip_bps_keep() -> f64 {
std::env::var("BLVM_IBD_A6M_GD_SLOW_TIP_BPS_KEEP")
.ok()
.and_then(|s| s.parse::<f64>().ok())
.unwrap_or(80.0)
.clamp(0.0, 400.0)
}
fn tip_trial_post_open_settle_secs() -> u64 {
std::env::var("BLVM_IBD_TIP_TRIAL_POST_OPEN_SETTLE_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(8)
.clamp(0, 30)
}
fn a6m_gd_slow_owner_cooldown_secs() -> u64 {
std::env::var("BLVM_IBD_A6M_GD_SLOW_OWNER_COOLDOWN_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(180)
.clamp(30, 600)
}
fn a6m_gd_slow_force_min_tip_bps() -> f64 {
std::env::var("BLVM_IBD_A6M_GD_SLOW_FORCE_MIN_TIP_BPS")
.ok()
.and_then(|s| s.parse::<f64>().ok())
.unwrap_or(20.0)
.clamp(1.0_f64, 200.0_f64)
}
include!("chunk_assigner_parts/impl_assign.rs");
include!("chunk_assigner_parts/impl_tip_hole.rs");
include!("chunk_assigner_parts/impl_flight.rs");
pub(crate) struct ChunkGuard {
chunk: Option<ChunkWorkItem>,
peer_id: Option<String>,
assigner: Arc<ChunkAssigner>,
}
impl ChunkGuard {
pub(crate) fn new(
start: u64,
end: u64,
exclude: Option<String>,
peer_id: String,
assigner: Arc<ChunkAssigner>,
) -> Self {
Self {
chunk: Some((start, end, exclude)),
peer_id: Some(peer_id),
assigner,
}
}
pub(crate) fn disarm(&mut self) {
self.chunk = None;
self.peer_id = None; }
}
impl Drop for ChunkGuard {
fn drop(&mut self) {
if let Some((start, end, exclude)) = self.chunk.take() {
if let Some(peer_id) = self.peer_id.take() {
self.assigner.on_chunk_complete_range(&peer_id, start, end);
}
self.assigner.requeue(start, end, exclude);
} else if let Some(peer_id) = self.peer_id.take() {
self.assigner.on_chunk_complete(&peer_id);
}
}
}
#[cfg(test)]
#[path = "chunk_assigner_tests.rs"]
mod tests;