fn tip_need_from(validation_height: Option<&Arc<AtomicU64>>) -> Option<u64> {
validation_height.map(|vh| vh.load(Ordering::Relaxed).saturating_add(1))
}
fn block_tx_tip_reserve() -> usize {
latch_env!(usize, {
std::env::var("BLVM_IBD_BLOCK_TX_TIP_RESERVE")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(64)
.clamp(8, 512)
})
}
async fn await_block_tx_tip_reserve(
tx: &tokio::sync::mpsc::Sender<(u64, SharedBlock, SharedWitnesses)>,
height: u64,
tip_need: Option<u64>,
) {
if tip_need == Some(height) || tip_need.is_none() {
return;
}
let reserve = block_tx_tip_reserve();
let mut logged = false;
while tx.capacity() <= reserve {
if !logged {
logged = true;
info!(
"[IBD_BLOCK_TX_TIP_RESERVE] wait h={} free={} (reserve={})",
height,
tx.capacity(),
reserve
);
}
tokio::time::sleep(Duration::from_millis(1)).await;
}
}
pub(crate) fn is_snapshot_sourced_peer(peer_id: &str) -> bool {
is_local_disk_peer(peer_id) || super::synthetic_wan::is_synthetic_peer(peer_id)
}
fn download_received_soft_cap() -> usize {
latch_env!(usize, {
std::env::var("BLVM_IBD_DOWNLOAD_RECEIVED_CAP")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(160)
.clamp(64, 512)
})
}
fn download_received_hard_cap() -> usize {
latch_env!(usize, {
std::env::var("BLVM_IBD_DOWNLOAD_RECEIVED_HARD_CAP")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(320)
.clamp(128, 1024)
.max(download_received_soft_cap())
})
}
fn note_received_insert(n: u64) {
super::memory::DOWNLOAD_RECEIVED_BLOCKS.fetch_add(n, Ordering::Relaxed);
}
fn note_received_remove(n: u64) {
let _ = super::memory::DOWNLOAD_RECEIVED_BLOCKS.fetch_update(
Ordering::Relaxed,
Ordering::Relaxed,
|v| Some(v.saturating_sub(n)),
);
}
fn received_put(
received: &mut BTreeMap<u64, (SharedBlock, SharedWitnesses)>,
h: u64,
v: (SharedBlock, SharedWitnesses),
) {
if received.insert(h, v).is_none() {
note_received_insert(1);
}
}
fn received_take(
received: &mut BTreeMap<u64, (SharedBlock, SharedWitnesses)>,
h: u64,
) -> Option<(SharedBlock, SharedWitnesses)> {
received.remove(&h).inspect(|_| note_received_remove(1))
}
fn received_clone(
received: &BTreeMap<u64, (SharedBlock, SharedWitnesses)>,
h: u64,
) -> Option<(SharedBlock, SharedWitnesses)> {
received
.get(&h)
.map(|(b, w)| (Arc::clone(b), Arc::clone(w)))
}
fn received_drain_all(received: &mut BTreeMap<u64, (SharedBlock, SharedWitnesses)>) {
let n = received.len() as u64;
received.clear();
if n > 0 {
note_received_remove(n);
}
}
fn hard_trim_download_received_far_ahead(
received: &mut BTreeMap<u64, (SharedBlock, SharedWitnesses)>,
need: u64,
hard: usize,
) -> u64 {
let mut forced = 0u64;
while received.len() > hard {
let Some((&h, _)) = received.iter().next_back() else {
break;
};
if h <= need {
break;
}
let _ = received_take(received, h);
forced = forced.saturating_add(1);
}
forced
}
fn trim_download_received(
received: &mut BTreeMap<u64, (SharedBlock, SharedWitnesses)>,
blockstore: &BlockStore,
validation_height: Option<&Arc<AtomicU64>>,
protocol_version: ProtocolVersion,
) {
let soft = download_received_soft_cap();
let hard = download_received_hard_cap();
let need = validation_height
.map(|v| v.load(Ordering::Relaxed).saturating_add(1))
.unwrap_or(0);
let mut trimmed = 0u64;
let mut forced = 0u64;
while let Some((&h, _)) = received.iter().next() {
if h >= need {
break;
}
let _ = received_take(received, h);
trimmed = trimmed.saturating_add(1);
}
while received.len() > soft {
let Some((&h, _)) = received.iter().next_back() else {
break;
};
if h <= need {
break;
}
let Some((block, witnesses)) = received_take(received, h) else {
break;
};
let hash = blockstore.get_block_hash(block.as_ref());
let persisted = try_persist_gap_block_for_local_inject(
blockstore,
validation_height,
h,
hash,
block.as_ref(),
witnesses.as_ref(),
protocol_version,
)
.unwrap_or(false);
if persisted {
trimmed = trimmed.saturating_add(1);
drop((block, witnesses));
continue;
}
received_put(received, h, (block, witnesses));
break;
}
let hard_forced = hard_trim_download_received_far_ahead(received, need, hard);
trimmed = trimmed.saturating_add(hard_forced);
forced = forced.saturating_add(hard_forced);
if trimmed > 0 {
super::memory::DOWNLOAD_RECEIVED_TRIM_BLOCKS.fetch_add(trimmed, Ordering::Relaxed);
static LAST_TRIM_LOG_MS: AtomicU64 = AtomicU64::new(0);
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let prev = LAST_TRIM_LOG_MS.load(Ordering::Relaxed);
if now_ms.saturating_sub(prev) >= 5_000
&& LAST_TRIM_LOG_MS
.compare_exchange(prev, now_ms, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
let n = super::memory::DOWNLOAD_RECEIVED_TRIM_BLOCKS.load(Ordering::Relaxed);
warn!(
"[IBD_DOWNLOAD_RECEIVED_TRIM] dropped {} block(s) (forced={}, soft={}, hard={}, received={}, need={}, total_trimmed={})",
trimmed,
forced,
soft,
hard,
received.len(),
need,
n
);
}
}
}
pub(crate) struct DownloadChunkResult {
pub blocks: Vec<(u64, SharedBlock, SharedWitnesses)>,
pub streamed_block_count: usize,
}
impl DownloadChunkResult {
#[inline]
pub fn block_count(&self) -> usize {
if self.blocks.is_empty() {
self.streamed_block_count
} else {
self.blocks.len()
}
}
}
struct BlockDownloadProgress {
last_block_hash: Option<Hash>,
last_progress_time: std::time::Instant,
current_timeout_seconds: u64,
disconnected_peers_count: usize,
}
impl BlockDownloadProgress {
fn new() -> Self {
Self {
last_block_hash: None,
last_progress_time: std::time::Instant::now(),
current_timeout_seconds: 120,
disconnected_peers_count: 0,
}
}
fn record_progress(&mut self, block_hash: Hash) {
self.last_block_hash = Some(block_hash);
self.last_progress_time = std::time::Instant::now();
}
fn reset_timeout(&mut self) {
self.current_timeout_seconds = 120;
self.disconnected_peers_count = 0;
}
}
async fn send_block_getdata_with_retry(
network: Arc<NetworkManager>,
peer_addr: SocketAddr,
wire_msg: Vec<u8>,
height: u64,
) -> Result<()> {
const MAX_ATTEMPTS: u32 = 30;
const BASE_MS: u64 = 100;
const MAX_WAIT_MS: u64 = 5_000;
let mut attempt: u32 = 0;
let mut reconnect_spawned = false;
loop {
match network.send_to_peer(peer_addr, wire_msg.clone()).await {
Ok(()) => return Ok(()),
Err(e) => {
let msg = e.to_string();
let is_gone = msg.contains("not found") || msg.contains("disconnected");
if !reconnect_spawned && is_gone {
reconnect_spawned = true;
NetworkManager::spawn_outbound_reconnect_attempt(
Arc::clone(&network),
peer_addr,
);
}
attempt += 1;
let wait_ms = BASE_MS
.saturating_mul(1u64 << (attempt - 1).min(6))
.min(MAX_WAIT_MS);
if attempt >= MAX_ATTEMPTS {
return Err(e).with_context(|| {
format!(
"Failed to send GetData for block at height {height} after {MAX_ATTEMPTS} attempts"
)
});
}
if attempt <= 3 || attempt % 5 == 0 {
warn!(
"GetData send failed for height {} (attempt {}/{}): {} — retrying in {}ms",
height, attempt, MAX_ATTEMPTS, e, wait_ms
);
}
tokio::time::sleep(Duration::from_millis(wait_ms)).await;
}
}
}
}
pub(crate) fn resume_download_height(
start_height: u64,
end_height: u64,
validated_tip: u64,
) -> Option<u64> {
if validated_tip >= end_height {
return None;
}
let resume = start_height.max(validated_tip.saturating_add(1));
if resume > end_height {
None
} else {
Some(resume)
}
}
pub(crate) fn chunk_outer_deadline_secs(
start_height: u64,
end_height: u64,
resume_from: u64,
per_block_timeout_secs: u64,
) -> u64 {
if start_height == end_height {
return per_block_timeout_secs.saturating_mul(4).clamp(120, 7200);
}
let blocks_remaining = end_height
.saturating_sub(resume_from)
.saturating_add(1)
.max(1);
per_block_timeout_secs
.saturating_mul(blocks_remaining)
.clamp(35, 7200)
}
pub(crate) fn worker_chunk_outer_deadline_secs(
start_height: u64,
end_height: u64,
resume_from: u64,
per_block_timeout_secs: u64,
confirmed_body_height: u64,
) -> u64 {
let base = chunk_outer_deadline_secs(
start_height,
end_height,
resume_from,
per_block_timeout_secs,
);
if confirmed_body_height > 0
&& start_height > confirmed_body_height
&& end_height.saturating_sub(start_height) >= 63
{
let cap = wan_deep_tip_pipe_chunk_deadline_secs(
start_height,
end_height,
confirmed_body_height,
per_block_timeout_secs,
)
.saturating_mul(2)
.clamp(60, 120);
base.min(cap)
} else {
base
}
}
async fn wait_for_peer_connected(
network: &Arc<NetworkManager>,
peer_addr: SocketAddr,
peer_id: &str,
max_wait: Duration,
tip_enter: &Option<Arc<super::chunk_assigner::ChunkAssigner>>,
) -> Result<()> {
{
let evicted = network.ibd_evicted_ips.read().unwrap();
if evicted.contains(&peer_addr.ip()) {
return Err(anyhow::anyhow!(
"Peer {peer_id} evicted (NODE_NETWORK_LIMITED) — chunk needs retry on another peer"
));
}
}
if network.is_peer_connected(peer_addr).await {
return Ok(());
}
NetworkManager::spawn_outbound_reconnect_attempt(Arc::clone(network), peer_addr);
let deadline = tokio::time::Instant::now() + max_wait;
let mut poll = tokio::time::interval(Duration::from_millis(200));
poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
poll.tick().await;
while tokio::time::Instant::now() < deadline {
if tip_enter
.as_ref()
.is_some_and(|a| a.is_peer_blacklisted(peer_id))
{
return Err(anyhow::anyhow!(
"Peer {peer_id} blacklisted during connect wait — chunk needs retry"
));
}
if network.is_peer_connected(peer_addr).await {
return Ok(());
}
poll.tick().await;
}
Err(anyhow::anyhow!(
"Peer {peer_id} not connected after {}s — chunk needs retry",
max_wait.as_secs()
))
}
async fn wait_for_peer_ibd_ready(
network: &Arc<NetworkManager>,
peer_addr: SocketAddr,
peer_id: &str,
max_wait: Duration,
tip_enter: &Option<Arc<super::chunk_assigner::ChunkAssigner>>,
) -> Result<()> {
let deadline = tokio::time::Instant::now() + max_wait;
let mut poll = tokio::time::interval(Duration::from_millis(200));
poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
poll.tick().await;
while tokio::time::Instant::now() < deadline {
if tip_enter
.as_ref()
.is_some_and(|a| a.is_peer_blacklisted(peer_id))
{
return Err(anyhow::anyhow!(
"Peer {peer_id} blacklisted during handshake wait — chunk needs retry"
));
}
if network.peer_ibd_ready(peer_addr).await {
return Ok(());
}
poll.tick().await;
}
Err(anyhow::anyhow!(
"Peer {peer_id} handshake not complete after {}s — chunk needs retry",
max_wait.as_secs()
))
}
fn ibd_getdata_inv_type(blockstore: &BlockStore, height: u64, block_hash: Hash) -> u32 {
let merkle_root = blockstore
.get_header(&block_hash)
.ok()
.flatten()
.map(|h| h.merkle_root)
.unwrap_or([0u8; 32]);
if crate::module::pipeline::try_filter_block_download_policy(height, block_hash, merkle_root) {
MSG_BLOCK
} else {
MSG_WITNESS_BLOCK
}
}
async fn register_and_request_block(
network: Arc<NetworkManager>,
peer_addr: SocketAddr,
peer_id: &str,
block_hash: Hash,
height: u64,
blockstore: &BlockStore,
) -> Result<tokio::sync::oneshot::Receiver<(Block, Vec<Vec<Witness>>, Option<Vec<u8>>)>> {
if !network.is_peer_connected(peer_addr).await {
return Err(anyhow::anyhow!(
"Peer {peer_id} not connected — cannot request block at height {height}"
));
}
let block_rx = network.register_block_request(peer_addr, block_hash);
let inventory = vec![InventoryVector {
inv_type: ibd_getdata_inv_type(blockstore, height, block_hash),
hash: block_hash,
}];
let wire_msg =
ProtocolParser::serialize_message(&ProtocolMessage::GetData(GetDataMessage { inventory }))?;
if let Err(e) =
send_block_getdata_with_retry(Arc::clone(&network), peer_addr, wire_msg, height).await
{
network.cancel_block_request(peer_addr, block_hash);
return Err(e);
}
super::tip_stage::mark_getdata(height);
Ok(block_rx)
}
type PendingBlockResult = (
u64,
[u8; 32],
std::time::Instant,
Result<
Result<(Block, Vec<Vec<Witness>>, Option<Vec<u8>>), tokio::sync::oneshot::error::RecvError>,
tokio::time::error::Elapsed,
>,
Option<tokio::sync::OwnedSemaphorePermit>,
);
type PendingBlockFuture =
std::pin::Pin<Box<dyn std::future::Future<Output = PendingBlockResult> + Send>>;
fn getdata_batch_size() -> usize {
super::policy::getdata_batch()
}
#[allow(clippy::too_many_arguments)]
async fn enqueue_network_block_batch(
heights_and_hashes: Vec<(u64, [u8; 32])>,
permits: Vec<Option<tokio::sync::OwnedSemaphorePermit>>,
network: &Arc<NetworkManager>,
peer_addr: SocketAddr,
peer_id: &str,
validation_tip: u64,
confirmed_body_height: u64,
chunk_default_secs: u64,
in_flight: &mut FuturesUnordered<PendingBlockFuture>,
in_flight_heights: &mut HashSet<u64>,
inflight_deadlines: &mut HashMap<u64, Arc<AtomicU64>>,
first_block_logged: &mut bool,
start_height: u64,
end_height: u64,
blockstore: &BlockStore,
) -> Result<()> {
debug_assert_eq!(heights_and_hashes.len(), permits.len());
if heights_and_hashes.is_empty() {
return Ok(());
}
let hashes: Vec<[u8; 32]> = heights_and_hashes.iter().map(|(_, h)| *h).collect();
let mut rxs = network.register_block_requests_batch(peer_addr, &hashes);
let inventory: Vec<InventoryVector> = heights_and_hashes
.iter()
.map(|&(height, hash)| InventoryVector {
inv_type: ibd_getdata_inv_type(blockstore, height, hash),
hash,
})
.collect();
let wire_msg =
ProtocolParser::serialize_message(&ProtocolMessage::GetData(GetDataMessage { inventory }))?;
let first_height = heights_and_hashes[0].0;
if let Err(e) =
send_block_getdata_with_retry(Arc::clone(network), peer_addr, wire_msg, first_height).await
{
for &hash in &hashes {
network.cancel_block_request(peer_addr, hash);
}
return Err(e);
}
for &(height, _) in &heights_and_hashes {
super::tip_stage::mark_getdata(height);
}
if !*first_block_logged {
info!(
"[IBD] {} chunk {}-{}: batch-requested {} blocks starting at height {} (hash {})",
peer_id,
start_height,
end_height,
heights_and_hashes.len(),
first_height,
hex::encode(hashes[0])
);
*first_block_logged = true;
}
for (((height, block_hash), permit), rx) in heights_and_hashes
.into_iter()
.zip(permits.into_iter())
.zip(rxs.drain(..))
{
let secs = block_gap_timeout_secs(
height,
validation_tip,
confirmed_body_height,
start_height,
end_height,
chunk_default_secs,
);
push_network_inflight(
in_flight,
in_flight_heights,
inflight_deadlines,
height,
block_hash,
rx,
permit,
secs,
);
}
Ok(())
}
fn try_take_blocks_permit(
blocks_sem: &Option<Arc<Semaphore>>,
) -> Result<Option<Option<tokio::sync::OwnedSemaphorePermit>>> {
match blocks_sem {
None => Ok(Some(None)),
Some(sem) => match sem.clone().try_acquire_owned() {
Ok(p) => Ok(Some(Some(p))),
Err(tokio::sync::TryAcquireError::NoPermits) => Ok(None),
Err(tokio::sync::TryAcquireError::Closed) => {
Err(anyhow::anyhow!("blocks semaphore closed"))
}
},
}
}
async fn enqueue_chunk_block(
height: u64,
block_hash: [u8; 32],
network: &Arc<NetworkManager>,
peer_addr: SocketAddr,
peer_id: &str,
blockstore: &BlockStore,
protocol_version: ProtocolVersion,
validation_tip: u64,
confirmed_body_height: u64,
chunk_default_secs: u64,
blocks_sem: &Option<Arc<Semaphore>>,
in_flight: &mut FuturesUnordered<PendingBlockFuture>,
in_flight_heights: &mut HashSet<u64>,
inflight_deadlines: &mut HashMap<u64, Arc<AtomicU64>>,
first_block_logged: &mut bool,
start_height: u64,
end_height: u64,
local_sourced_heights: &mut HashSet<u64>,
) -> Result<()> {
if in_flight_heights.contains(&height) {
return Ok(());
}
let Some(permit) = try_take_blocks_permit(blocks_sem)? else {
return Ok(());
};
if let Some((block, block_witnesses)) =
try_load_local_ibd_block(blockstore, height, block_hash, protocol_version)?
{
if !*first_block_logged {
info!(
"[IBD] {} chunk {}-{}: local block height {} (hash {})",
peer_id,
start_height,
end_height,
height,
hex::encode(block_hash)
);
*first_block_logged = true;
}
let request_start = std::time::Instant::now();
in_flight_heights.insert(height);
local_sourced_heights.insert(height);
in_flight.push(Box::pin(async move {
let r = Ok(Ok((block, block_witnesses, None)));
(height, block_hash, request_start, r, permit)
}));
} else {
if !*first_block_logged {
info!(
"[IBD] {} chunk {}-{}: registered block height {} (hash {})",
peer_id,
start_height,
end_height,
height,
hex::encode(block_hash)
);
*first_block_logged = true;
}
let block_rx = register_and_request_block(
Arc::clone(network),
peer_addr,
peer_id,
block_hash,
height,
blockstore,
)
.await?;
let secs = block_gap_timeout_secs(
height,
validation_tip,
confirmed_body_height,
start_height,
end_height,
chunk_default_secs,
);
push_network_inflight(
in_flight,
in_flight_heights,
inflight_deadlines,
height,
block_hash,
block_rx,
permit,
secs,
);
}
Ok(())
}