use std::collections::VecDeque;
use std::net::SocketAddr;
use std::sync::Arc;
use anyhow::{Context, Result};
use blvm_protocol::BlockHeader;
use tokio::task::JoinSet;
use tokio::time::{Duration, timeout};
use tracing::{debug, info, warn};
pub(crate) type HeaderRange = (u64, u64, [u8; 32]);
pub(crate) fn checkpoint_header_ranges(
checkpoints: &[(u64, [u8; 32])],
start_height: u64,
end_height: u64,
) -> Vec<HeaderRange> {
checkpoints
.windows(2)
.filter_map(|w| {
let (range_start, start_hash) = w[0];
let (range_end, _) = w[1];
if range_end < start_height {
return None;
}
let actual_start = range_start.max(start_height);
let actual_end = range_end.min(end_height);
if actual_start > actual_end {
return None;
}
Some((actual_start, actual_end, start_hash))
})
.collect()
}
pub(crate) fn simulate_header_range_schedule(
peer_count: usize,
range_count: usize,
) -> Vec<(usize, usize)> {
if peer_count == 0 || range_count == 0 {
return Vec::new();
}
let mut pending: VecDeque<usize> = (0..range_count).collect();
let mut free: VecDeque<usize> = (0..peer_count).collect();
let mut in_flight: VecDeque<(usize, usize)> = VecDeque::new();
let mut assigned = Vec::with_capacity(range_count);
while !pending.is_empty() || !in_flight.is_empty() {
while !free.is_empty() && !pending.is_empty() {
let peer = free.pop_front().unwrap();
let range = pending.pop_front().unwrap();
in_flight.push_back((peer, range));
assigned.push((peer, range));
}
if let Some((peer, _)) = in_flight.pop_front() {
free.push_back(peer);
}
}
assigned
}
use crate::network::NetworkManager;
use crate::network::peer_scoring::PeerScorer;
use crate::network::protocol::{GetHeadersMessage, ProtocolMessage, ProtocolParser};
use crate::node::event_publisher::EventPublisher;
use crate::storage::blockstore::BlockStore;
use crate::storage::hashing::double_sha256;
use blvm_protocol::GENESIS_BLOCK_HASH_INTERNAL;
pub(crate) struct HeaderSyncResult {
pub tip_height: u64,
}
impl HeaderSyncResult {
fn at_height(tip_height: u64) -> Self {
Self { tip_height }
}
}
pub(crate) fn header_batch_rtt_should_rotate(latency_ms: f64) -> bool {
latency_ms > 400.0
}
pub(crate) fn empty_headers_is_ibd_tip(
fetched_tip: u64,
start_height: u64,
end_height: u64,
) -> bool {
if fetched_tip >= end_height {
return true;
}
fetched_tip + 1 > start_height
}
fn header_links_to_parent(
blockstore: &BlockStore,
header: &BlockHeader,
height: u64,
last_hash: &[u8; 32],
) -> Result<bool> {
if height > 0 {
if let Some(parent) = blockstore.get_header_at_height(height - 1)? {
return Ok(blvm_consensus::block::validate_prev_block_hash(
header, &parent,
));
}
}
Ok(header.prev_block_hash == *last_hash)
}
pub(crate) fn ibd_header_locator_heights(last_known_height: u64) -> Vec<u64> {
let mut out = Vec::with_capacity(32);
let mut h = last_known_height;
let mut step = 1u64;
loop {
out.push(h);
if h == 0 || out.len() >= 32 {
break;
}
if out.len() > 10 {
step = step.saturating_mul(2);
}
h = h.saturating_sub(step);
}
if out.last() != Some(&0) {
out.push(0);
}
out
}
fn ibd_header_locator_hashes(
blockstore: &BlockStore,
last_known_height: u64,
last_hash: &[u8; 32],
) -> Result<Vec<[u8; 32]>> {
let mut hashes = Vec::with_capacity(32);
for h in ibd_header_locator_heights(last_known_height) {
if h == last_known_height {
hashes.push(*last_hash);
continue;
}
match blockstore.get_hash_by_height(h)? {
Some(stored) => hashes.push(stored),
None if h == 0 => hashes.push(GENESIS_BLOCK_HASH_INTERNAL),
None => {}
}
}
if hashes.is_empty() {
hashes.push(*last_hash);
}
Ok(hashes)
}
pub(crate) async fn download_header_range(
network: Arc<NetworkManager>,
peer: SocketAddr,
locator_hash: [u8; 32],
start_height: u64,
end_height: u64,
) -> Result<Vec<blvm_protocol::BlockHeader>> {
let mut all_headers = Vec::new();
let mut current_hash = locator_hash;
let mut current_height = start_height;
let mut consecutive_failures = 0;
let mut current_peer = peer;
const MAX_FAILURES: u32 = 3; const TIMEOUT_SECS: u64 = 10;
while current_height <= end_height {
let get_headers = GetHeadersMessage {
version: 70015,
block_locator_hashes: vec![current_hash],
hash_stop: [0; 32],
};
let wire_msg = ProtocolParser::serialize_message(&ProtocolMessage::GetHeaders(get_headers))
.map_err(|e| anyhow::anyhow!("Failed to serialize GetHeaders: {}", e))?;
let headers_rx = network.register_headers_request(current_peer);
if let Err(e) = network.send_to_peer(current_peer, wire_msg).await {
consecutive_failures += 1;
if consecutive_failures >= MAX_FAILURES {
let fresh = network.get_connected_peer_addresses().await;
if let Some(&next) = fresh.iter().find(|&&p| p != current_peer) {
debug!(
"Range {}-{}: switching from {} to {}",
start_height, end_height, current_peer, next
);
current_peer = next;
consecutive_failures = 0;
} else {
return Err(anyhow::anyhow!(
"No peers for range {}-{}: {}",
start_height,
end_height,
e
));
}
}
tokio::time::sleep(Duration::from_millis(200)).await;
continue;
}
match timeout(Duration::from_secs(TIMEOUT_SECS), headers_rx).await {
Ok(Ok(headers)) => {
consecutive_failures = 0;
if headers.is_empty() {
break;
}
for header in headers {
match blvm_protocol::pow::check_proof_of_work(&header) {
Ok(true) => {}
Ok(false) => {
return Err(anyhow::anyhow!(
"Header at height {} failed PoW — refusing to skip (would corrupt height index)",
current_height
));
}
Err(e) => {
return Err(anyhow::anyhow!(
"Header at height {} PoW check error: {e}",
current_height
));
}
}
let mut header_data = [0u8; 80];
header_data[0..4].copy_from_slice(&(header.version as i32).to_le_bytes());
header_data[4..36].copy_from_slice(&header.prev_block_hash);
header_data[36..68].copy_from_slice(&header.merkle_root);
header_data[68..72].copy_from_slice(&(header.timestamp as u32).to_le_bytes());
header_data[72..76].copy_from_slice(&(header.bits as u32).to_le_bytes());
header_data[76..80].copy_from_slice(&(header.nonce as u32).to_le_bytes());
let header_hash = double_sha256(&header_data);
all_headers.push(header);
current_hash = header_hash;
current_height += 1;
if current_height > end_height {
break;
}
}
let max_headers = network.protocol_limits().max_headers_results;
if all_headers.len() % max_headers != 0 {
break;
}
}
Ok(Err(_)) => {
consecutive_failures += 1;
if consecutive_failures >= MAX_FAILURES {
return Err(anyhow::anyhow!("Headers channel closed too many times"));
}
}
Err(_) => {
consecutive_failures += 1;
if consecutive_failures >= MAX_FAILURES {
return Err(anyhow::anyhow!("Timeout waiting for headers from {}", peer));
}
}
}
}
debug!(
"Downloaded {} headers from {} for range {} - {}",
all_headers.len(),
peer,
start_height,
end_height
);
Ok(all_headers)
}
type HeaderRangeJob = (
SocketAddr,
u64,
u64,
[u8; 32],
Result<Vec<blvm_protocol::BlockHeader>>,
);
fn spawn_header_range_jobs(
join_set: &mut JoinSet<HeaderRangeJob>,
pending: &mut VecDeque<HeaderRange>,
free_peers: &mut VecDeque<SocketAddr>,
network_mgr: &Arc<NetworkManager>,
) {
while !free_peers.is_empty() && !pending.is_empty() {
let peer_addr = free_peers.pop_front().expect("non-empty");
let (actual_start, actual_end, locator_hash) = pending.pop_front().expect("non-empty");
let network_clone = Arc::clone(network_mgr);
join_set.spawn(async move {
let result = download_header_range(
network_clone,
peer_addr,
locator_hash,
actual_start,
actual_end,
)
.await;
(peer_addr, actual_start, actual_end, locator_hash, result)
});
}
}
pub(crate) async fn download_headers(
peer_scorer: Arc<PeerScorer>,
start_height: u64,
end_height: u64,
peer_ids: &[String],
blockstore: &BlockStore,
network: Option<Arc<NetworkManager>>,
headers_timeout_secs: u64,
headers_max_failures: u32,
event_publisher: Option<Arc<EventPublisher>>,
) -> Result<HeaderSyncResult> {
let network = match network.as_ref() {
Some(n) => n,
None => {
warn!("NetworkManager not available, skipping header download");
return Ok(HeaderSyncResult::at_height(start_height));
}
};
if peer_ids.is_empty() {
if !super::synthetic_wan::allow_zero_real_peers() {
return Err(anyhow::anyhow!("No peers available for header download"));
}
let mut h = start_height;
while h <= end_height {
match blockstore.get_hash_by_height(h) {
Ok(Some(_)) => h += 1,
_ => break,
}
}
if h > start_height {
let tip = h.saturating_sub(1);
info!(
"IBD header sync: zero-peer local-replay — using {} already-stored headers (tip={})",
tip.saturating_sub(start_height).saturating_add(1),
tip
);
return Ok(HeaderSyncResult::at_height(tip));
}
return Err(anyhow::anyhow!(
"No peers and no stored headers from {} (synthetic WAN / ALLOW_ZERO_PEERS)",
start_height
));
}
let mut peer_addrs: Vec<SocketAddr> = peer_ids
.iter()
.filter_map(|id| id.parse::<SocketAddr>().ok())
.collect();
if peer_addrs.is_empty() {
return Err(anyhow::anyhow!("No valid peer addresses found"));
}
peer_addrs.sort_by(|a, b| {
let a_score = peer_scorer.get_score(a);
let b_score = peer_scorer.get_score(b);
b_score
.partial_cmp(&a_score)
.unwrap_or(std::cmp::Ordering::Equal)
});
info!(
"Using {} peers for sequential header download",
peer_addrs.len()
);
let genesis_hash = GENESIS_BLOCK_HASH_INTERNAL;
let mut current_height: u64;
let mut last_hash: [u8; 32];
if start_height == 0 {
let genesis_header = blvm_protocol::BlockHeader {
version: 1,
prev_block_hash: [0u8; 32],
merkle_root: [
0x3b, 0xa3, 0xed, 0xfd, 0x7a, 0x7b, 0x12, 0xb2, 0x7a, 0xc7, 0x2c, 0x3e, 0x67, 0x76,
0x8f, 0x61, 0x7f, 0xc8, 0x1b, 0xc3, 0x88, 0x8a, 0x51, 0x32, 0x3a, 0x9f, 0xb8, 0xaa,
0x4b, 0x1e, 0x5e, 0x4a,
],
timestamp: 1231006505,
bits: 0x1d00ffff,
nonce: 2083236893,
};
let mut header_data = [0u8; 80];
header_data[0..4].copy_from_slice(&(genesis_header.version as i32).to_le_bytes());
header_data[4..36].copy_from_slice(&genesis_header.prev_block_hash);
header_data[36..68].copy_from_slice(&genesis_header.merkle_root);
header_data[68..72].copy_from_slice(&(genesis_header.timestamp as u32).to_le_bytes());
header_data[72..76].copy_from_slice(&(genesis_header.bits as u32).to_le_bytes());
header_data[76..80].copy_from_slice(&(genesis_header.nonce as u32).to_le_bytes());
let computed_hash = double_sha256(&header_data);
if computed_hash != genesis_hash {
warn!(
"Genesis hash mismatch! Computed: {}, Expected: {}",
hex::encode(computed_hash),
hex::encode(genesis_hash)
);
}
blockstore
.store_header(&genesis_hash, &genesis_header)
.context("Failed to store genesis header")?;
blockstore
.store_height(0, &genesis_hash)
.context("Failed to store genesis height")?;
info!(
"Stored genesis block (height 0, hash: {})",
hex::encode(genesis_hash)
);
current_height = 1;
last_hash = genesis_hash;
} else {
let parent_h = start_height
.checked_sub(1)
.ok_or_else(|| anyhow::anyhow!("header sync: invalid start_height"))?;
last_hash = blockstore.get_hash_by_height(parent_h)?.ok_or_else(|| {
anyhow::anyhow!(
"Cannot resume header sync at height {}: missing parent hash at height {}. \
Sync from genesis or repair height_index (data may be inconsistent).",
start_height,
parent_h
)
})?;
current_height = start_height;
info!(
"Resuming header sync at height {} (GetHeaders locator = parent {})",
start_height,
hex::encode(last_hash)
);
}
let mut consecutive_failures = 0;
let mut current_peer_idx = 0;
let mut last_progress_log = start_height;
let mut last_progress_event = start_height;
let start_time = std::time::Instant::now();
{
let mut skip_count: u64 = 0;
while current_height <= end_height {
match blockstore.get_hash_by_height(current_height) {
Ok(Some(stored_hash)) => {
last_hash = stored_hash;
skip_count += 1;
current_height += 1;
}
_ => break,
}
}
if skip_count > 0 {
info!(
"IBD header sync: skipped {} already-stored headers (now at height {}); \
only fetching new headers from peers",
skip_count, current_height
);
last_progress_log = current_height;
}
}
while current_height <= end_height {
if peer_addrs.is_empty() {
peer_addrs = network.get_connected_peer_addresses().await;
if peer_addrs.is_empty() {
tokio::time::sleep(Duration::from_secs(5)).await;
peer_addrs = network.get_connected_peer_addresses().await;
if peer_addrs.is_empty() {
return Err(anyhow::anyhow!("No peers available"));
}
}
}
let peer_addr = peer_addrs[current_peer_idx % peer_addrs.len()];
let get_headers = GetHeadersMessage {
version: 70015,
block_locator_hashes: ibd_header_locator_hashes(
blockstore,
current_height.saturating_sub(1),
&last_hash,
)?,
hash_stop: [0; 32],
};
let wire_msg =
match ProtocolParser::serialize_message(&ProtocolMessage::GetHeaders(get_headers)) {
Ok(msg) => msg,
Err(e) => {
warn!("Failed to serialize GetHeaders: {}", e);
return Err(anyhow::anyhow!("Serialization failed"));
}
};
let headers_rx = network.register_headers_request(peer_addr);
let request_start = std::time::Instant::now();
if let Err(e) = network.send_to_peer(peer_addr, wire_msg).await {
debug!("Send failed to {}: {}", peer_addr, e);
peer_addrs.retain(|&a| a != peer_addr);
current_peer_idx += 1;
consecutive_failures += 1;
if consecutive_failures >= headers_max_failures {
return Err(anyhow::anyhow!("Too many failures"));
}
continue;
}
debug!(
"Waiting for headers from {} (timeout: {}s)",
peer_addr, headers_timeout_secs
);
match timeout(Duration::from_secs(headers_timeout_secs), headers_rx).await {
Ok(Ok(headers)) => {
let latency_ms = request_start.elapsed().as_secs_f64() * 1000.0;
peer_scorer.record_latency_sample(peer_addr, latency_ms);
debug!(
"Received {} headers from {} ({}ms)",
headers.len(),
peer_addr,
latency_ms as u64
);
if headers.is_empty() {
let fetched_tip = current_height.saturating_sub(1);
if !empty_headers_is_ibd_tip(fetched_tip, start_height, end_height) {
warn!(
"[IBD_HEADER_EMPTY] peer={} fetched_tip={} start={} end={} — not IBD tip; rotate",
peer_addr, fetched_tip, start_height, end_height
);
consecutive_failures += 1;
current_peer_idx += 1;
if let Some(idx) = peer_addrs.iter().position(|&a| a == peer_addr) {
let p = peer_addrs.remove(idx);
peer_addrs.push(p);
}
continue;
}
info!(
"Header sync COMPLETE at height {} (chain tip reached)",
fetched_tip
);
break;
}
debug!(
"Processing {} headers starting at height {}",
headers.len(),
current_height
);
let mut batch_entries: Vec<(blvm_protocol::Hash, BlockHeader, u64)> =
Vec::with_capacity(headers.len());
for header in &headers {
match blvm_protocol::pow::check_proof_of_work(header) {
Ok(true) => {}
Ok(false) => {
return Err(anyhow::anyhow!(
"Header at height {} failed PoW — refusing to skip (would corrupt height index)",
current_height
));
}
Err(e) => {
return Err(anyhow::anyhow!(
"Header at height {} PoW check error: {e}",
current_height
));
}
}
if !header_links_to_parent(blockstore, header, current_height, &last_hash)? {
if let Ok(Some(attach_h)) =
blockstore.get_height_by_hash(&header.prev_block_hash)
{
if attach_h > 0 && attach_h + 1 < current_height {
warn!(
"[IBD_HEADER_REWIND] attach prev height {} → ask {} (was {}) peer={}",
attach_h,
attach_h + 1,
current_height,
peer_addr
);
last_hash = header.prev_block_hash;
current_height = attach_h + 1;
consecutive_failures = 0;
batch_entries.clear();
break;
}
}
if header.prev_block_hash == GENESIS_BLOCK_HASH_INTERNAL
&& current_height > 1
{
let back = 2048u64.min(current_height.saturating_sub(1));
let new_h = current_height - back;
if let Ok(Some(hh)) =
blockstore.get_hash_by_height(new_h.saturating_sub(1))
{
warn!(
"[IBD_HEADER_REWIND] genesis-start at {} — step back to {} peer={}",
current_height, new_h, peer_addr
);
current_height = new_h;
last_hash = hh;
consecutive_failures = 0;
batch_entries.clear();
break;
}
}
warn!(
"[IBD_HEADER_PEER] chain break at height {} from {}: expected prev {} got {} — trying next peer",
current_height,
peer_addr,
hex::encode(last_hash),
hex::encode(header.prev_block_hash)
);
current_peer_idx += 1;
consecutive_failures += 1;
if consecutive_failures >= headers_max_failures {
return Err(anyhow::anyhow!(
"Header chain break at height {}: expected prev {} got {} ({} peers failed)",
current_height,
hex::encode(last_hash),
hex::encode(header.prev_block_hash),
consecutive_failures
));
}
batch_entries.clear();
break;
}
let mut header_data = [0u8; 80];
header_data[0..4].copy_from_slice(&(header.version as i32).to_le_bytes());
header_data[4..36].copy_from_slice(&header.prev_block_hash);
header_data[36..68].copy_from_slice(&header.merkle_root);
header_data[68..72].copy_from_slice(&(header.timestamp as u32).to_le_bytes());
header_data[72..76].copy_from_slice(&(header.bits as u32).to_le_bytes());
header_data[76..80].copy_from_slice(&(header.nonce as u32).to_le_bytes());
let header_hash = double_sha256(&header_data);
batch_entries.push((header_hash, header.clone(), current_height));
last_hash = header_hash;
current_height += 1;
if current_height > end_height {
break;
}
}
let batch_count = batch_entries.len();
if batch_count == 0 {
continue;
}
consecutive_failures = 0;
debug!("Storing {} headers in batch...", batch_count);
let store_start = std::time::Instant::now();
let blockstore_clone = blockstore.clone();
tokio::task::spawn_blocking(move || {
blockstore_clone.store_headers_batch(&batch_entries)
})
.await
.context("Failed to spawn blocking task")?
.context("Failed to store headers batch")?;
debug!(
"Stored {} headers in {:?}",
batch_count,
store_start.elapsed()
);
if header_batch_rtt_should_rotate(latency_ms) {
warn!(
"[IBD_HEADER_SLOW] peer={} rtt_ms={:.0} n={} — rotate",
peer_addr, latency_ms, batch_count
);
current_peer_idx += 1;
}
if current_height > last_progress_log && current_height - last_progress_log >= 20000
{
let elapsed = start_time.elapsed().as_secs_f64();
let synced = current_height - start_height;
let rate = if elapsed > 0.0 {
synced as f64 / elapsed
} else {
0.0
};
let remaining = end_height.saturating_sub(current_height);
let eta = if rate > 0.0 {
remaining as f64 / rate
} else {
f64::INFINITY
};
info!(
"Header sync: {} / {} ({:.1}%) - {:.0} h/s - ETA: {:.0}s",
current_height,
end_height,
(current_height as f64 / end_height as f64) * 100.0,
rate,
eta
);
last_progress_log = current_height;
}
if current_height > last_progress_event
&& (current_height - last_progress_event) >= 5000
{
if let Some(ref ep) = event_publisher {
let progress_percent = if end_height > start_height {
((current_height - start_height) as f64
/ (end_height - start_height + 1) as f64)
* 100.0
} else {
100.0
};
ep.publish_headers_sync_progress(
current_height.saturating_sub(1),
end_height,
progress_percent,
)
.await;
last_progress_event = current_height;
}
}
let max_headers = network.protocol_limits().max_headers_results;
if headers.len() < max_headers {
let total = current_height - start_height;
let elapsed = start_time.elapsed();
let rate = if elapsed.as_secs_f64() > 0.0 {
total as f64 / elapsed.as_secs_f64()
} else {
0.0
};
info!(
"Header sync COMPLETE: {} headers in {:.1}s ({:.0} h/s) - chain tip reached",
total,
elapsed.as_secs_f64(),
rate
);
return Ok(HeaderSyncResult {
tip_height: current_height.saturating_sub(1),
});
}
}
Ok(Err(_)) => {
debug!("Channel closed for request to {}", peer_addr);
consecutive_failures += 1;
current_peer_idx += 1;
}
Err(_) => {
debug!("Timeout waiting for headers from {}", peer_addr);
consecutive_failures += 1;
current_peer_idx += 1;
if let Some(idx) = peer_addrs.iter().position(|&a| a == peer_addr) {
let p = peer_addrs.remove(idx);
peer_addrs.push(p);
}
}
}
if consecutive_failures >= headers_max_failures {
warn!(
"Too many failures ({}), refreshing peers",
consecutive_failures
);
consecutive_failures = 0;
peer_addrs = network.get_connected_peer_addresses().await;
if peer_addrs.is_empty() {
warn!("Header sync: no connected peers — triggering archive peer re-discovery");
let default_ban = Default::default();
let (arc_res, _) = tokio::join!(network.discover_archive_peers_from_dns(), {
let (net_name, port) =
crate::network::protocol::ProtocolParser::dns_seed_network();
network.discover_peers_from_dns(net_name, port, &default_ban)
},);
let _ = arc_res;
let _ = network.connect_peers_from_database(16).await;
for _ in 0..8 {
tokio::time::sleep(Duration::from_secs(2)).await;
peer_addrs = network.get_connected_peer_addresses().await;
if !peer_addrs.is_empty() {
info!(
"Header sync: {} peer(s) reconnected, resuming",
peer_addrs.len()
);
break;
}
}
if peer_addrs.is_empty() {
return Err(anyhow::anyhow!(
"No peers available after 70s wait + re-discovery"
));
}
}
}
}
let total = current_height - start_height;
let elapsed = start_time.elapsed();
let rate = if elapsed.as_secs_f64() > 0.0 {
total as f64 / elapsed.as_secs_f64()
} else {
0.0
};
info!(
"Header sync COMPLETE: {} headers in {:.1}s ({:.0} h/s)",
total,
elapsed.as_secs_f64(),
rate
);
Ok(HeaderSyncResult {
tip_height: current_height.saturating_sub(1),
})
}
#[cfg(test)]
mod n13_tests {
use super::*;
#[test]
fn n13_checkpoint_ranges_not_capped_by_peer_count() {
let cps = [
(0u64, [0u8; 32]),
(100, [1u8; 32]),
(200, [2u8; 32]),
(300, [3u8; 32]),
(400, [4u8; 32]),
];
let ranges = checkpoint_header_ranges(&cps, 0, 400);
assert_eq!(
ranges.len(),
4,
"all windows kept (old bug: take(peers) dropped tail)"
);
assert_eq!(ranges[0].0, 0);
assert_eq!(ranges[3].1, 400);
assert!(ranges.len() > 2);
}
#[test]
fn ibd_header_locator_heights_doubles_after_ten_steps() {
let hs = ibd_header_locator_heights(20);
assert_eq!(hs[0], 20);
assert_eq!(&hs[0..10], &[20, 19, 18, 17, 16, 15, 14, 13, 12, 11]);
assert!(hs.contains(&0), "locator must include genesis");
let high = ibd_header_locator_heights(961_637);
assert_eq!(high[0], 961_637);
assert!(high.len() > 12);
assert_eq!(*high.last().unwrap(), 0);
assert!(high.contains(&961_628));
}
#[test]
fn n13_schedule_covers_all_ranges_with_one_inflight_per_peer() {
let peer_count = 3;
let range_count = 7;
let assigned = simulate_header_range_schedule(peer_count, range_count);
assert_eq!(assigned.len(), range_count);
let mut seen = vec![false; range_count];
for &(peer, range) in &assigned {
assert!(peer < peer_count);
assert!(!seen[range], "range {range} assigned twice");
seen[range] = true;
}
assert!(seen.iter().all(|&s| s));
let first_wave = peer_count.min(range_count);
let mut inflight = first_wave;
let mut free = peer_count - first_wave;
let mut pending = range_count - first_wave;
let mut max_inflight = inflight;
for _ in first_wave..range_count {
inflight -= 1;
free += 1;
inflight += 1;
free -= 1;
pending -= 1;
max_inflight = max_inflight.max(inflight);
}
assert_eq!(pending, 0);
assert!(max_inflight <= peer_count);
assert_eq!(free + inflight, peer_count);
}
#[test]
fn n13_extra_peers_idle_when_fewer_ranges() {
let assigned = simulate_header_range_schedule(8, 3);
assert_eq!(assigned.len(), 3);
let peers_used: std::collections::HashSet<_> = assigned.iter().map(|(p, _)| *p).collect();
assert_eq!(peers_used.len(), 3);
}
#[test]
fn r216_header_slow_rtt_rotates_r215_stays_on_r214() {
assert!(!header_batch_rtt_should_rotate(103.0));
assert!(!header_batch_rtt_should_rotate(400.0));
assert!(header_batch_rtt_should_rotate(401.0));
assert!(header_batch_rtt_should_rotate(1770.0));
}
#[test]
fn dest_bg_empty_headers_at_genesis_is_not_ibd_tip() {
assert!(!empty_headers_is_ibd_tip(0, 1, 963_969));
assert!(
empty_headers_is_ibd_tip(961_638, 1, 963_969),
"progress then empty is BIP130 this-peer tip (dest-bf 961k complete)"
);
assert!(empty_headers_is_ibd_tip(963_969, 1, 963_969));
assert!(!empty_headers_is_ibd_tip(499_999, 500_000, 963_969));
}
}