writer.rs (3909B)
1 use anyhow::Result; 2 use tokio::{io::AsyncWriteExt, net::tcp::OwnedWriteHalf, time::timeout}; 3 4 use crate::app::GossipEnvelope; 5 6 use super::{MAX_GOSSIP_LINE_BYTES, WRITE_TIMEOUT}; 7 8 pub(super) async fn write_payload( 9 writer: &mut OwnedWriteHalf, 10 payload: &[GossipEnvelope], 11 ) -> Result<()> { 12 timeout(WRITE_TIMEOUT, async { 13 for envelope in payload { 14 write_envelope_inner(writer, envelope).await?; 15 } 16 Ok::<(), anyhow::Error>(()) 17 }) 18 .await 19 .map_err(|_| anyhow::anyhow!("p2p batch write timed out after {WRITE_TIMEOUT:?}"))??; 20 Ok(()) 21 } 22 23 pub(super) async fn write_envelope( 24 writer: &mut OwnedWriteHalf, 25 envelope: &GossipEnvelope, 26 ) -> Result<()> { 27 timeout(WRITE_TIMEOUT, write_envelope_inner(writer, envelope)) 28 .await 29 .map_err(|_| anyhow::anyhow!("p2p write timed out after {WRITE_TIMEOUT:?}"))??; 30 Ok(()) 31 } 32 33 async fn write_envelope_inner( 34 writer: &mut OwnedWriteHalf, 35 envelope: &GossipEnvelope, 36 ) -> Result<()> { 37 let line = serde_json::to_string(envelope)?; 38 if line.len() > MAX_GOSSIP_LINE_BYTES { 39 anyhow::bail!( 40 "p2p message is {} bytes, exceeding {} byte limit", 41 line.len(), 42 MAX_GOSSIP_LINE_BYTES 43 ); 44 } 45 writer.write_all(line.as_bytes()).await?; 46 writer.write_all(b"\n").await?; 47 Ok(()) 48 } 49 50 pub(super) fn byte_bounded_block_page( 51 blocks: Vec<crate::domain::Block>, 52 ) -> Vec<crate::domain::Block> { 53 let mut page = Vec::new(); 54 let mut encoded_len = serde_json::to_vec(&GossipEnvelope::Blocks { blocks: Vec::new() }) 55 .expect("block envelope serialization cannot fail") 56 .len(); 57 for block in blocks { 58 let block_len = serde_json::to_vec(&block) 59 .expect("block serialization cannot fail") 60 .len(); 61 let separator_len = usize::from(!page.is_empty()); 62 let Some(candidate_len) = encoded_len 63 .checked_add(separator_len) 64 .and_then(|len| len.checked_add(block_len)) 65 else { 66 break; 67 }; 68 if candidate_len > MAX_GOSSIP_LINE_BYTES { 69 break; 70 } 71 encoded_len = candidate_len; 72 page.push(block); 73 } 74 page 75 } 76 77 #[cfg(test)] 78 mod tests { 79 use crate::{ 80 app::GossipEnvelope, 81 domain::{Block, BurnBundleSection, FinalizerMode}, 82 }; 83 84 use super::{MAX_GOSSIP_LINE_BYTES, byte_bounded_block_page}; 85 86 fn large_block(height: u64) -> Block { 87 Block { 88 height, 89 prev_hash: format!("{:064x}", height.saturating_sub(1)), 90 timestamp_ms: height, 91 miner: "0".repeat(64), 92 reward_address: None, 93 reward_address_signature: None, 94 finalizer_mode: FinalizerMode::Ticket, 95 finalizer_rank: 0, 96 reward: 0, 97 vdf_rounds: 0, 98 vdf_output: "x".repeat(512 * 1024), 99 leader_proof: None, 100 burn_bundle_section: BurnBundleSection::default(), 101 transactions: Vec::new(), 102 transactions_v2: Vec::new(), 103 hash: format!("{height:064x}"), 104 } 105 } 106 107 #[test] 108 fn block_page_is_bounded_by_wire_bytes_not_only_item_count() { 109 let blocks = (1..=32).map(large_block).collect::<Vec<_>>(); 110 let page = byte_bounded_block_page(blocks.clone()); 111 let encoded = serde_json::to_vec(&GossipEnvelope::Blocks { 112 blocks: page.clone(), 113 }) 114 .unwrap(); 115 116 assert!(!page.is_empty()); 117 assert!(page.len() < blocks.len()); 118 assert!(encoded.len() <= MAX_GOSSIP_LINE_BYTES); 119 120 let mut one_more = page; 121 one_more.push(blocks[one_more.len()].clone()); 122 assert!( 123 serde_json::to_vec(&GossipEnvelope::Blocks { blocks: one_more }) 124 .unwrap() 125 .len() 126 > MAX_GOSSIP_LINE_BYTES 127 ); 128 } 129 }