use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use super::ParallelIBDConfig;
use super::latch_env;
use super::types::{ChunkWorkItem, NoprogComplete, RetryEntry, RetryTrack};
#[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 {
match std::env::var("BLVM_IBD_ASSIGNER_LOCK_TIMERS")
.ok()
.as_deref()
.map(str::trim)
{
None => true,
Some("0") | Some("false") | Some("off") | Some("no") => false,
_ => true,
}
}
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
);
}
}
}
struct InFlightLockOwner {
site: &'static str,
}
static IN_FLIGHT_LOCK_SITE: Mutex<Option<(&'static str, String, u64, Instant)>> = Mutex::new(None);
fn note_in_flight_holder(site: &'static str, peer: &str, next: u64) -> InFlightLockOwner {
if let Ok(mut g) = IN_FLIGHT_LOCK_SITE.lock() {
*g = Some((site, peer.to_string(), next, Instant::now()));
}
InFlightLockOwner { site }
}
impl Drop for InFlightLockOwner {
fn drop(&mut self) {
if let Ok(mut g) = IN_FLIGHT_LOCK_SITE.lock() {
if g.as_ref().is_some_and(|(s, ..)| *s == self.site) {
*g = None;
}
}
}
}
fn format_in_flight_holder() -> String {
match IN_FLIGHT_LOCK_SITE.lock() {
Ok(g) => match &*g {
Some((site, peer, next, at)) => format!(
"site={site} peer={peer} next={next} held_ms={}",
at.elapsed().as_millis()
),
None => "site=unknown".to_string(),
},
Err(_) => "site=holder_poison".to_string(),
}
}
fn leftover_get_work_phase(
peer_id: &str,
next: u64,
body_tip: u64,
leftover_force: bool,
phase: &'static str,
) {
if super::leftover_trace_watch(leftover_force, next, body_tip) {
leftover_trace_rate(
"get_work_phase",
format!(
"peer={peer_id} next={next} body_tip={body_tip} leftover_force={leftover_force} phase={phase}"
),
);
}
}
pub(crate) fn leftover_trace_rate(step: &'static str, msg: String) {
static LAST_MS: Mutex<Option<HashMap<&'static str, u64>>> = Mutex::new(None);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
if let Ok(mut slot) = LAST_MS.lock() {
let map = slot.get_or_insert_with(HashMap::new);
let prev = map.get(step).copied().unwrap_or(0);
if now.saturating_sub(prev) < 5_000 {
return;
}
map.insert(step, now);
}
tracing::warn!("[IBD_LEFTOVER_TRACE] step={step} {msg}");
}
#[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<RetryEntry>>,
retry_track: Mutex<HashMap<(u64, u64), RetryTrack>>,
window_done: Mutex<std::collections::BTreeMap<u64, Instant>>,
window_started: Mutex<HashMap<(String, u64, u64), Instant>>,
window_fail: Mutex<HashMap<(u64, u64), (String, Instant)>>,
window_strikes: Mutex<HashMap<String, (u32, Instant)>>,
window_evict: Mutex<Vec<String>>,
window_last_evict: Mutex<Option<Instant>>,
window_bench: Mutex<HashMap<String, Instant>>,
window_front_cool: Mutex<HashMap<String, Instant>>,
window_offense: Mutex<HashMap<String, (u32, Instant)>>,
window_tile_ms_ema: AtomicU64,
window_blk_ms_ema: AtomicU64,
window_peer_blk_ms: Mutex<HashMap<String, (u64, u32)>>,
noprog_ok: Mutex<HashMap<(String, u64, u64), NoprogComplete>>,
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,
leftover_force_getdata: AtomicBool,
tip_bridge_holes: AtomicU64,
window_hole_until: AtomicU64,
window_tip_uncovered: AtomicBool,
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>>,
last_stream_keep_hero: Mutex<Option<String>>,
last_gd_slow_pin: Mutex<Option<(String, Instant)>>,
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,
last_mute_drop_h: AtomicU64,
last_mute_drop_at_ms: AtomicU64,
latched_ahead: Mutex<Vec<(String, u64, u64)>>,
priority_zone: Mutex<Vec<(String, u64, u64)>>,
lookahead_stripes: Mutex<Vec<(String, u64, u64)>>,
fat_probe_retitle_done: AtomicBool,
farm_recv_promote_done: AtomicBool,
farm_recv_fat_rearmed: AtomicBool,
farm_recv_mute_streak: AtomicU32,
farm_recv_last_tick: Mutex<Option<Instant>>,
lookahead_have_hold: Mutex<Vec<(u64, u64)>>,
lookahead_streams: Mutex<HashMap<String, TipStreamWindow>>,
}
macro_rules! prod_or_test {
($prod:expr, $test:expr) => {{
#[cfg(not(test))]
{
$prod
}
#[cfg(test)]
{
$test
}
}};
}
fn a6m_min_bps() -> f64 {
prod_or_test!(
6.0,
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 {
prod_or_test!(
22.0,
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 {
prod_or_test!(
12.0,
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 {
prod_or_test!(
300,
std::env::var("BLVM_IBD_A6M_TENURE_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(300)
)
}
fn a6m_recent_window_secs() -> u64 {
prod_or_test!(
60,
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 {
prod_or_test!(
600,
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 {
prod_or_test!(
120,
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 {
prod_or_test!(
800,
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 {
prod_or_test!(
4,
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 {
prod_or_test!(
0.0,
std::env::var("BLVM_IBD_A6M_GD_SLOW_TIP_BPS_KEEP")
.ok()
.and_then(|s| s.parse::<f64>().ok())
.unwrap_or(0.0)
.clamp(0.0, 400.0)
)
}
fn tip_trial_post_open_settle_secs() -> u64 {
prod_or_test!(
8,
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 {
prod_or_test!(
180,
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 {
prod_or_test!(
20.0,
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)
)
}
fn gd_slow_pin_protect_secs() -> u64 {
prod_or_test!(
30,
std::env::var("BLVM_IBD_GD_SLOW_PIN_PROTECT_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(30)
.clamp(8, 90)
)
}
const H_SLOW_MIN_HEIGHT: u64 = 180_000;
const DUMP_WAREHOUSE_REORDER: u64 = 2048;
const H_SLOW_COOLDOWN_SECS: u64 = 120;
pub(crate) fn h_slow_rotate_enabled() -> bool {
prod_or_test!(
false,
matches!(
std::env::var("BLVM_IBD_H_SLOW_ROTATE").ok().as_deref(),
Some("1") | Some("true") | Some("on") | Some("yes")
)
)
}
static LAST_H_SLOW_GD_MS: AtomicU64 = AtomicU64::new(0);
static TIP_RERACE_HEIGHT: AtomicU64 = AtomicU64::new(0);
pub(crate) fn body_h_rtt_should_rotate(
latency_ms: f64,
sticky_bps: f64,
feeder: u64,
next_needed: u64,
already_rotated_this_gd: bool,
) -> bool {
next_needed >= H_SLOW_MIN_HEIGHT
&& feeder == 0
&& latency_ms > 1500.0
&& sticky_bps < 60.0
&& !already_rotated_this_gd
}
pub(crate) fn h_slow_recv_leader_holds(sticky_recv_mbps: f64, top_recv_mbps: f64) -> bool {
sticky_recv_mbps >= top_recv_mbps
}
pub(crate) fn h_slow_recv_successor_beats(sticky_recv_mbps: f64, cand_recv_mbps: f64) -> bool {
cand_recv_mbps > sticky_recv_mbps
}
pub(crate) fn body_dump_warehouse_should_rotate(
sticky_bps: f64,
feeder: u64,
next_needed: u64,
reorder_ahead: u64,
already_rotated_this_gd: bool,
) -> bool {
next_needed < H_SLOW_MIN_HEIGHT
&& feeder == 0
&& sticky_bps < 150.0
&& reorder_ahead >= DUMP_WAREHOUSE_REORDER
&& !already_rotated_this_gd
}
#[cfg(test)]
pub(crate) fn test_reset_h_slow_gd() {
LAST_H_SLOW_GD_MS.store(0, Ordering::Relaxed);
}
include!("chunk_assigner_parts/impl_assign.rs");
include!("chunk_assigner_parts/impl_tip_hole.rs");
include!("chunk_assigner_parts/impl_flight.rs");
include!("chunk_assigner_parts/impl_window.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_reason(start, end, exclude, "guard_drop");
} 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;