fetch.rs (16475B)
1 use std::{collections::BTreeMap, net::SocketAddr, time::Instant}; 2 3 use anyhow::{Context, Result}; 4 use tokio::{ 5 net::{TcpStream, tcp::OwnedReadHalf}, 6 time::timeout, 7 }; 8 9 use crate::{ 10 app::{ 11 ChainBootstrap, GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION, ProtocolHello, 12 debug_logging_enabled, now_ms, protocol_capabilities, validate_network_genesis, 13 validate_protocol_capabilities, validate_transaction_v2_peer_capability, 14 }, 15 domain::{Block, ChainSnapshot, LaunchProfile, Ledger, verify_vdf}, 16 }; 17 18 use super::{ 19 GossipNetwork, JOIN_RESPONSE_TIMEOUT, LimitedLineReader, MAX_BLOCK_BATCH, 20 MAX_JOIN_RESPONSE_ENVELOPES, PeerStatus, parse_envelope, write_envelope, 21 }; 22 23 pub async fn fetch_snapshot(peer: &str) -> Result<ChainSnapshot> { 24 fetch_snapshot_with_announcement(peer, None, &LaunchProfile::default().profile_id).await 25 } 26 27 pub async fn fetch_peer_height(peer: &str) -> Result<u64> { 28 fetch_peer_status(peer).await.map(|status| status.height) 29 } 30 31 /// Validates a downloaded snapshot without making VDF verification a serial part of replay. 32 /// State-dependent consensus rules are checked first, then independent VDF proofs are checked 33 /// across a bounded number of worker threads before the candidate ledger is returned. 34 pub async fn validate_chain_snapshot(snapshot: ChainSnapshot) -> Result<Ledger> { 35 let validation_time_ms = now_ms(); 36 tokio::task::spawn_blocking(move || { 37 let started = Instant::now(); 38 let ledger = Ledger::from_preverified_snapshot_at(snapshot, validation_time_ms)?; 39 let replay_elapsed = started.elapsed(); 40 41 let vdf_started = Instant::now(); 42 verify_block_vdfs_parallel(ledger.chain().iter().skip(1))?; 43 if debug_logging_enabled() { 44 eprintln!( 45 "initial chain validation: state={:.3}s vdf={:.3}s blocks={}", 46 replay_elapsed.as_secs_f64(), 47 vdf_started.elapsed().as_secs_f64(), 48 ledger.chain().len().saturating_sub(1), 49 ); 50 } 51 Ok(ledger) 52 }) 53 .await 54 .context("chain snapshot validation worker failed")? 55 } 56 57 async fn fetch_peer_status(peer: &str) -> Result<PeerStatus> { 58 let stream = TcpStream::connect(peer) 59 .await 60 .with_context(|| format!("connecting to peer {peer}"))?; 61 let (reader, _writer) = stream.into_split(); 62 let mut reader = LimitedLineReader::new(reader); 63 let line = reader 64 .read_line() 65 .await? 66 .with_context(|| format!("peer {peer} closed before sending its peer status"))?; 67 match parse_envelope(&line)? { 68 GossipEnvelope::Hello(hello) => { 69 if hello.protocol_version != PROTOCOL_VERSION { 70 anyhow::bail!( 71 "unsupported protocol version {}; expected {}", 72 hello.protocol_version, 73 PROTOCOL_VERSION 74 ); 75 } 76 if hello.network_id != NETWORK_ID { 77 anyhow::bail!( 78 "wrong network {}; expected {}", 79 hello.network_id, 80 NETWORK_ID 81 ); 82 } 83 validate_protocol_capabilities(&hello.capabilities)?; 84 validate_transaction_v2_peer_capability(&hello.capabilities, 0, hello.height)?; 85 Ok(PeerStatus::with_time( 86 hello.height, 87 hello.tip_hash, 88 hello.time_ms, 89 )) 90 } 91 GossipEnvelope::PeerStatus { 92 height, 93 tip_hash, 94 time_ms, 95 } => Ok(PeerStatus::from_envelope(height, tip_hash, time_ms)), 96 other => anyhow::bail!("peer {peer} sent {other:?} instead of peer status"), 97 } 98 } 99 100 pub async fn fetch_snapshot_with_announcement( 101 peer: &str, 102 _advertised_addr: Option<SocketAddr>, 103 expected_profile_id: &str, 104 ) -> Result<ChainSnapshot> { 105 let stream = TcpStream::connect(peer) 106 .await 107 .with_context(|| format!("connecting to join peer {peer}"))?; 108 let (reader, mut writer) = stream.into_split(); 109 let mut reader = LimitedLineReader::new(reader); 110 let line = reader 111 .read_line() 112 .await? 113 .with_context(|| format!("join peer {peer} closed before sending its peer status"))?; 114 match parse_envelope(&line)? { 115 GossipEnvelope::Hello(hello) => { 116 if hello.protocol_version != PROTOCOL_VERSION { 117 anyhow::bail!( 118 "unsupported protocol version {}; expected {}", 119 hello.protocol_version, 120 PROTOCOL_VERSION 121 ); 122 } 123 if hello.network_id != NETWORK_ID { 124 anyhow::bail!( 125 "wrong network {}; expected {}", 126 hello.network_id, 127 NETWORK_ID 128 ); 129 } 130 validate_protocol_capabilities(&hello.capabilities)?; 131 validate_transaction_v2_peer_capability(&hello.capabilities, 0, hello.height)?; 132 } 133 GossipEnvelope::PeerStatus { .. } => {} 134 other => anyhow::bail!("join peer {peer} sent {other:?} instead of peer status"), 135 } 136 137 write_envelope(&mut writer, &join_client_hello()).await?; 138 write_envelope(&mut writer, &GossipEnvelope::ChainBootstrapRequest).await?; 139 let bootstrap = read_join_bootstrap_response(peer, &mut reader).await?; 140 validate_bootstrap_genesis(expected_profile_id, &bootstrap)?; 141 142 let mut snapshot = ChainSnapshot { 143 genesis_allocations: bootstrap.genesis_allocations, 144 vdf_rounds: bootstrap.vdf_rounds, 145 launch_profile: bootstrap.launch_profile, 146 blocks: vec![bootstrap.genesis_block], 147 }; 148 while snapshot.blocks.last().map_or(0, |block| block.height) < bootstrap.height { 149 let from_height = snapshot.blocks.last().map_or(0, |block| block.height) + 1; 150 let remaining = bootstrap.height - from_height + 1; 151 write_envelope( 152 &mut writer, 153 &GossipEnvelope::BlockRangeRequest { 154 from_height, 155 limit: remaining.min(MAX_BLOCK_BATCH as u64) as usize, 156 }, 157 ) 158 .await?; 159 let blocks = read_join_blocks_response(peer, &mut reader).await?; 160 if blocks.is_empty() { 161 anyhow::bail!("join peer {peer} returned an empty block page at height {from_height}"); 162 } 163 if blocks[0].height != from_height { 164 anyhow::bail!("join peer {peer} returned a non-contiguous block page"); 165 } 166 snapshot.blocks.extend(blocks); 167 } 168 if snapshot.blocks.last().map(|block| &block.hash) != Some(&bootstrap.tip_hash) { 169 anyhow::bail!("join peer {peer} changed tips while serving block pages"); 170 } 171 172 Ok(snapshot) 173 } 174 175 fn join_client_hello() -> GossipEnvelope { 176 let setup = Ledger::new(BTreeMap::new(), 1); 177 GossipEnvelope::Hello(ProtocolHello { 178 protocol_version: PROTOCOL_VERSION, 179 capabilities: protocol_capabilities(), 180 network_id: NETWORK_ID.to_string(), 181 genesis_hash: setup.genesis_hash().to_string(), 182 listen_addr: None, 183 node_id: None, 184 height: 0, 185 tip_hash: setup.tip_hash().to_string(), 186 time_ms: now_ms(), 187 }) 188 } 189 190 async fn read_join_bootstrap_response( 191 peer: &str, 192 reader: &mut LimitedLineReader<OwnedReadHalf>, 193 ) -> Result<ChainBootstrap> { 194 for _ in 0..MAX_JOIN_RESPONSE_ENVELOPES { 195 let envelope = read_join_envelope(peer, reader, "chain bootstrap").await?; 196 match envelope { 197 GossipEnvelope::ChainBootstrap(bootstrap) => return Ok(bootstrap), 198 envelope if is_join_control_envelope(&envelope) => continue, 199 other => anyhow::bail!("join peer {peer} sent {other:?} instead of chain bootstrap"), 200 } 201 } 202 anyhow::bail!("join peer {peer} sent too many control envelopes while joining") 203 } 204 205 async fn read_join_blocks_response( 206 peer: &str, 207 reader: &mut LimitedLineReader<OwnedReadHalf>, 208 ) -> Result<Vec<Block>> { 209 for _ in 0..MAX_JOIN_RESPONSE_ENVELOPES { 210 let envelope = read_join_envelope(peer, reader, "block page").await?; 211 match envelope { 212 GossipEnvelope::Blocks { blocks } => return Ok(blocks), 213 envelope if is_join_control_envelope(&envelope) => continue, 214 other => anyhow::bail!("join peer {peer} sent {other:?} instead of a block page"), 215 } 216 } 217 anyhow::bail!("join peer {peer} sent too many control envelopes while joining") 218 } 219 220 async fn read_join_envelope( 221 peer: &str, 222 reader: &mut LimitedLineReader<OwnedReadHalf>, 223 expected: &str, 224 ) -> Result<GossipEnvelope> { 225 let line = timeout(JOIN_RESPONSE_TIMEOUT, reader.read_line()) 226 .await 227 .with_context(|| format!("join peer {peer} timed out waiting for {expected}"))?? 228 .with_context(|| format!("join peer {peer} closed before sending {expected}"))?; 229 parse_envelope(&line) 230 } 231 232 fn is_join_control_envelope(envelope: &GossipEnvelope) -> bool { 233 matches!( 234 envelope, 235 GossipEnvelope::Hello(_) 236 | GossipEnvelope::PeerStatus { .. } 237 | GossipEnvelope::PeerList { .. } 238 | GossipEnvelope::PeerVerificationChallenge { .. } 239 | GossipEnvelope::PeerVerificationResponse { .. } 240 | GossipEnvelope::Inventory { .. } 241 ) 242 } 243 244 pub(super) async fn validate_chain_bootstrap( 245 expected_profile_id: &str, 246 bootstrap: ChainBootstrap, 247 now_ms: u64, 248 ) -> Result<Ledger> { 249 validate_bootstrap_genesis(expected_profile_id, &bootstrap)?; 250 let snapshot = ChainSnapshot { 251 genesis_allocations: bootstrap.genesis_allocations, 252 vdf_rounds: bootstrap.vdf_rounds, 253 launch_profile: bootstrap.launch_profile, 254 blocks: vec![bootstrap.genesis_block], 255 }; 256 tokio::task::spawn_blocking(move || Ledger::from_snapshot_at(snapshot, now_ms)) 257 .await 258 .context("chain bootstrap adoption worker failed")? 259 } 260 261 fn validate_bootstrap_genesis(expected_profile_id: &str, bootstrap: &ChainBootstrap) -> Result<()> { 262 if bootstrap.launch_profile.profile_id != expected_profile_id { 263 anyhow::bail!( 264 "chain bootstrap profile {} does not match expected profile {expected_profile_id}", 265 bootstrap.launch_profile.profile_id 266 ); 267 } 268 validate_network_genesis(expected_profile_id, &bootstrap.genesis_block.hash) 269 } 270 271 pub(super) async fn validate_blocks_extension( 272 mut ledger: Ledger, 273 blocks: Vec<Block>, 274 now_ms: u64, 275 on_progress: impl Fn(u64) + Send + 'static, 276 ) -> Result<Ledger> { 277 if blocks.is_empty() { 278 return Ok(ledger); 279 } 280 281 tokio::task::spawn_blocking(move || { 282 let state_started = Instant::now(); 283 let target_height = blocks.last().map(|block| block.height); 284 if blocks[0].prev_hash != ledger.tip_hash() { 285 let mut candidate = ledger.snapshot(); 286 let ancestor = candidate 287 .blocks 288 .iter() 289 .position(|block| block.hash == blocks[0].prev_hash) 290 .ok_or(super::SyncError::BlockPageHasNoCommonAncestor)?; 291 candidate.blocks.truncate(ancestor + 1); 292 candidate.blocks.extend(blocks.iter().cloned()); 293 ledger.extend_from_preverified_snapshot_at(candidate, now_ms)?; 294 let state_elapsed = state_started.elapsed(); 295 let vdf_started = Instant::now(); 296 verify_block_vdfs_parallel(&blocks)?; 297 log_batch_validation_timing(blocks.len(), state_elapsed, vdf_started.elapsed(), true); 298 if let Some(target_height) = target_height { 299 on_progress(target_height); 300 } 301 return Ok(ledger); 302 } 303 for block in blocks.iter().cloned() { 304 ledger.apply_preverified_block_at(block, now_ms)?; 305 } 306 let state_elapsed = state_started.elapsed(); 307 let vdf_started = Instant::now(); 308 verify_block_vdfs_parallel(&blocks)?; 309 log_batch_validation_timing(blocks.len(), state_elapsed, vdf_started.elapsed(), false); 310 for block in &blocks { 311 on_progress(block.height); 312 } 313 Ok(ledger) 314 }) 315 .await 316 .context("block batch extension worker failed")? 317 } 318 319 fn verify_block_vdfs_parallel<'a>(blocks: impl IntoIterator<Item = &'a Block>) -> Result<()> { 320 let blocks = blocks.into_iter().collect::<Vec<_>>(); 321 if blocks.is_empty() { 322 return Ok(()); 323 } 324 325 // Two chain candidates may be validated concurrently by the network coordinator. Giving 326 // each validation at most half the available CPUs prevents the pair from oversubscribing the 327 // machine, while the upper bound keeps untrusted batches from creating excessive threads. 328 let available = std::thread::available_parallelism() 329 .map(usize::from) 330 .unwrap_or(1); 331 let workers = available.div_ceil(2).clamp(1, 8).min(blocks.len()); 332 let chunk_size = blocks.len().div_ceil(workers); 333 let invalid_height = std::thread::scope(|scope| -> Result<Option<u64>> { 334 let handles = blocks 335 .chunks(chunk_size) 336 .map(|chunk| { 337 scope.spawn(move || { 338 chunk.iter().find_map(|block| { 339 (!verify_vdf(&block.vdf_seed(), block.vdf_rounds, &block.vdf_output)) 340 .then_some(block.height) 341 }) 342 }) 343 }) 344 .collect::<Vec<_>>(); 345 346 let mut invalid_height = None; 347 for handle in handles { 348 let height = handle 349 .join() 350 .map_err(|_| anyhow::anyhow!("VDF verification worker panicked"))?; 351 invalid_height = match (invalid_height, height) { 352 (Some(left), Some(right)) => Some(left.min(right)), 353 (height @ Some(_), None) | (None, height @ Some(_)) => height, 354 (None, None) => None, 355 }; 356 } 357 Ok(invalid_height) 358 })?; 359 if invalid_height.is_some() { 360 anyhow::bail!("block VDF output is invalid"); 361 } 362 Ok(()) 363 } 364 365 fn log_batch_validation_timing( 366 blocks: usize, 367 state_elapsed: std::time::Duration, 368 vdf_elapsed: std::time::Duration, 369 fork: bool, 370 ) { 371 if debug_logging_enabled() { 372 eprintln!( 373 "block batch validation: state={:.3}s vdf={:.3}s blocks={blocks} fork={fork}", 374 state_elapsed.as_secs_f64(), 375 vdf_elapsed.as_secs_f64(), 376 ); 377 } 378 } 379 380 pub(super) async fn network_adjusted_time_ms(network: &GossipNetwork) -> u64 { 381 let local_time_ms = now_ms(); 382 network 383 .inner 384 .peers 385 .lock() 386 .await 387 .adjusted_time_ms_at(local_time_ms) 388 } 389 390 pub(super) async fn verify_block_vdf(block: Block) -> Result<Block> { 391 let seed = block.vdf_seed(); 392 let rounds = block.vdf_rounds; 393 let solution = block.vdf_output.clone(); 394 let valid = tokio::task::spawn_blocking(move || verify_vdf(&seed, rounds, &solution)) 395 .await 396 .context("VDF verification worker failed")?; 397 if !valid { 398 anyhow::bail!("block VDF output is invalid"); 399 } 400 401 Ok(block) 402 } 403 404 #[cfg(test)] 405 mod tests { 406 use std::collections::BTreeMap; 407 408 use super::{join_client_hello, verify_block_vdfs_parallel}; 409 use crate::app::{GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION}; 410 use crate::domain::Ledger; 411 412 #[test] 413 fn snapshot_join_identifies_as_an_unannounced_setup_placeholder() { 414 let GossipEnvelope::Hello(hello) = join_client_hello() else { 415 panic!("join handshake must start with Hello"); 416 }; 417 418 assert_eq!(hello.protocol_version, PROTOCOL_VERSION); 419 assert_eq!(hello.network_id, NETWORK_ID); 420 assert_eq!(hello.height, 0); 421 assert_eq!(hello.genesis_hash, hello.tip_hash); 422 assert!(hello.listen_addr.is_none()); 423 assert!(hello.node_id.is_none()); 424 } 425 426 #[test] 427 fn parallel_vdf_verification_rejects_an_invalid_proof() { 428 let genesis = Ledger::new(BTreeMap::new(), 1).chain()[0].clone(); 429 let mut later = genesis.clone(); 430 later.height = 9; 431 later.vdf_output = "invalid-vdf".to_string(); 432 let mut earlier = genesis; 433 earlier.height = 3; 434 earlier.vdf_output = "also-invalid".to_string(); 435 436 let error = verify_block_vdfs_parallel([&later, &earlier]).unwrap_err(); 437 438 assert_eq!(error.to_string(), "block VDF output is invalid"); 439 } 440 }