iuna

iuna

iuna - experimental mainnet-candidate protocol
git clone https://getiuna.org/git/iuna.git
Log | Files | Refs | README | LICENSE

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 }