line_codec.rs (19811B)
1 use anyhow::{Context, Result}; 2 use tokio::{ 3 io::{AsyncBufReadExt, AsyncRead, BufReader}, 4 net::tcp::OwnedReadHalf, 5 }; 6 7 use crate::{ 8 app::{GossipEnvelope, TRANSACTION_BATCH_LIMIT}, 9 domain::BURN_COMMITTEE_SIZE, 10 }; 11 12 use super::{ 13 GossipNetwork, MAX_BLOCK_BATCH, MAX_BLOCK_LOCATOR_HASHES, MAX_GOSSIP_LINE_BYTES, 14 MAX_INVENTORY_ITEMS, MAX_OBJECT_REQUESTS, MAX_PEER_LIST, metrics::P2pMetricsCounters, 15 }; 16 17 pub(super) struct LimitedLineReader<R> { 18 reader: BufReader<R>, 19 pending: Vec<u8>, 20 } 21 22 impl<R: AsyncRead + Unpin> LimitedLineReader<R> { 23 pub(super) fn new(reader: R) -> Self { 24 Self { 25 reader: BufReader::new(reader), 26 pending: Vec::new(), 27 } 28 } 29 30 pub(super) async fn read_line(&mut self) -> Result<Option<String>> { 31 loop { 32 let available = self.reader.fill_buf().await?; 33 if available.is_empty() { 34 if self.pending.is_empty() { 35 return Ok(None); 36 } 37 anyhow::bail!("peer closed before completing a gossip message"); 38 } 39 40 if let Some(newline) = available.iter().position(|byte| *byte == b'\n') { 41 if self.pending.len() + newline > MAX_GOSSIP_LINE_BYTES { 42 anyhow::bail!("p2p message exceeds {MAX_GOSSIP_LINE_BYTES} byte limit"); 43 } 44 self.pending.extend_from_slice(&available[..newline]); 45 self.reader.consume(newline + 1); 46 if self.pending.ends_with(b"\r") { 47 self.pending.pop(); 48 } 49 let bytes = std::mem::take(&mut self.pending); 50 return String::from_utf8(bytes) 51 .context("p2p message is not valid UTF-8") 52 .map(Some); 53 } 54 55 if self.pending.len() + available.len() > MAX_GOSSIP_LINE_BYTES { 56 anyhow::bail!("p2p message exceeds {MAX_GOSSIP_LINE_BYTES} byte limit"); 57 } 58 let consumed = available.len(); 59 self.pending.extend_from_slice(available); 60 self.reader.consume(consumed); 61 } 62 } 63 } 64 65 pub(super) async fn read_session_envelope( 66 network: &GossipNetwork, 67 connection_label: &str, 68 reader: &mut LimitedLineReader<OwnedReadHalf>, 69 ) -> Result<Option<GossipEnvelope>> { 70 let Some(line) = reader.read_line().await? else { 71 return Ok(None); 72 }; 73 P2pMetricsCounters::add(&network.inner.metrics.bytes_received, line.len() as u64 + 1); 74 if line.trim().is_empty() { 75 P2pMetricsCounters::inc(&network.inner.metrics.empty_frames); 76 P2pMetricsCounters::set_last( 77 &network.inner.metrics.last_empty_frame_remote, 78 connection_label.to_string(), 79 ); 80 anyhow::bail!("empty p2p envelope"); 81 } 82 83 match parse_envelope(&line) { 84 Ok(envelope) => { 85 P2pMetricsCounters::inc(&network.inner.metrics.envelopes_received); 86 record_received_envelope_kind(&network.inner.metrics, &envelope); 87 Ok(Some(envelope)) 88 } 89 Err(error) => { 90 P2pMetricsCounters::inc(&network.inner.metrics.parse_errors); 91 P2pMetricsCounters::set_last( 92 &network.inner.metrics.last_parse_error, 93 format!("{connection_label}: {error:#}"), 94 ); 95 Err(error) 96 } 97 } 98 } 99 100 pub(super) fn record_received_envelope_kind( 101 metrics: &P2pMetricsCounters, 102 envelope: &GossipEnvelope, 103 ) { 104 match envelope { 105 GossipEnvelope::Hello(_) => { 106 P2pMetricsCounters::inc(&metrics.hello_envelopes_received); 107 } 108 GossipEnvelope::PeerStatus { .. } => { 109 P2pMetricsCounters::inc(&metrics.peer_status_envelopes_received); 110 } 111 GossipEnvelope::Inventory { .. } => { 112 P2pMetricsCounters::inc(&metrics.inventory_envelopes_received); 113 } 114 GossipEnvelope::Transaction(_) => { 115 P2pMetricsCounters::inc(&metrics.data_envelopes_received); 116 P2pMetricsCounters::inc(&metrics.transaction_envelopes_received); 117 P2pMetricsCounters::inc(&metrics.transactions_received); 118 } 119 GossipEnvelope::Transactions { transactions } => { 120 P2pMetricsCounters::inc(&metrics.data_envelopes_received); 121 P2pMetricsCounters::inc(&metrics.transaction_envelopes_received); 122 P2pMetricsCounters::add(&metrics.transactions_received, transactions.len() as u64); 123 } 124 GossipEnvelope::TransactionV2 { .. } => { 125 P2pMetricsCounters::inc(&metrics.data_envelopes_received); 126 P2pMetricsCounters::inc(&metrics.transaction_envelopes_received); 127 P2pMetricsCounters::inc(&metrics.transactions_received); 128 } 129 GossipEnvelope::TransactionsV2 { envelopes } => { 130 P2pMetricsCounters::inc(&metrics.data_envelopes_received); 131 P2pMetricsCounters::inc(&metrics.transaction_envelopes_received); 132 P2pMetricsCounters::add(&metrics.transactions_received, envelopes.len() as u64); 133 } 134 GossipEnvelope::BurnBundle(_) => { 135 P2pMetricsCounters::inc(&metrics.data_envelopes_received); 136 P2pMetricsCounters::inc(&metrics.burn_bundle_envelopes_received); 137 P2pMetricsCounters::inc(&metrics.burn_bundles_received); 138 } 139 GossipEnvelope::BurnBundles { bundles } => { 140 P2pMetricsCounters::inc(&metrics.data_envelopes_received); 141 P2pMetricsCounters::inc(&metrics.burn_bundle_envelopes_received); 142 P2pMetricsCounters::add(&metrics.burn_bundles_received, bundles.len() as u64); 143 } 144 GossipEnvelope::Block(_) 145 | GossipEnvelope::Blocks { .. } 146 | GossipEnvelope::ChainBootstrap(_) => { 147 P2pMetricsCounters::inc(&metrics.data_envelopes_received); 148 } 149 GossipEnvelope::ChainBootstrapRequest 150 | GossipEnvelope::BlockLocatorRequest { .. } 151 | GossipEnvelope::BlockRangeRequest { .. } 152 | GossipEnvelope::BlockRequest { .. } 153 | GossipEnvelope::BurnBundleRequest { .. } 154 | GossipEnvelope::PeerAnnouncement { .. } 155 | GossipEnvelope::PeerVerificationChallenge { .. } 156 | GossipEnvelope::PeerVerificationResponse { .. } 157 | GossipEnvelope::PeerList { .. } => { 158 P2pMetricsCounters::inc(&metrics.control_envelopes_received); 159 } 160 } 161 } 162 163 pub(super) fn parse_envelope(line: &str) -> Result<GossipEnvelope> { 164 if line.trim().is_empty() { 165 anyhow::bail!("empty p2p envelope"); 166 } 167 let envelope = serde_json::from_str(line).context("invalid p2p envelope JSON")?; 168 validate_envelope_limits(&envelope)?; 169 Ok(envelope) 170 } 171 172 pub(super) fn validate_envelope_limits(envelope: &GossipEnvelope) -> Result<()> { 173 match envelope { 174 GossipEnvelope::BlockRangeRequest { limit, .. } => { 175 ensure_len("block range request", *limit, MAX_BLOCK_BATCH)?; 176 } 177 GossipEnvelope::BlockLocatorRequest { locator, limit } => { 178 ensure_len("block locator", locator.len(), MAX_BLOCK_LOCATOR_HASHES)?; 179 ensure_len("block locator request", *limit, MAX_BLOCK_BATCH)?; 180 } 181 GossipEnvelope::BlockRequest { hashes } => { 182 ensure_len("block request", hashes.len(), MAX_OBJECT_REQUESTS)?; 183 } 184 GossipEnvelope::Inventory { blocks } => { 185 ensure_len("block inventory", blocks.len(), MAX_INVENTORY_ITEMS)?; 186 } 187 GossipEnvelope::Transactions { transactions } => { 188 ensure_len( 189 "transaction batch", 190 transactions.len(), 191 TRANSACTION_BATCH_LIMIT, 192 )?; 193 } 194 GossipEnvelope::TransactionsV2 { envelopes } => { 195 ensure_len( 196 "transaction v2 batch", 197 envelopes.len(), 198 TRANSACTION_BATCH_LIMIT, 199 )?; 200 } 201 GossipEnvelope::BurnBundles { bundles } => { 202 ensure_len("burn bundle batch", bundles.len(), TRANSACTION_BATCH_LIMIT)?; 203 } 204 GossipEnvelope::BurnBundleRequest { slots, .. } => { 205 ensure_len("burn bundle request", slots.len(), BURN_COMMITTEE_SIZE)?; 206 } 207 GossipEnvelope::Blocks { blocks } => { 208 ensure_len("block batch", blocks.len(), MAX_BLOCK_BATCH)?; 209 } 210 GossipEnvelope::PeerList { peers } => { 211 ensure_len("peer list", peers.len(), MAX_PEER_LIST)?; 212 } 213 GossipEnvelope::Hello(_) 214 | GossipEnvelope::ChainBootstrapRequest 215 | GossipEnvelope::ChainBootstrap(_) 216 | GossipEnvelope::PeerStatus { .. } 217 | GossipEnvelope::Transaction(_) 218 | GossipEnvelope::TransactionV2 { .. } 219 | GossipEnvelope::BurnBundle(_) 220 | GossipEnvelope::Block(_) 221 | GossipEnvelope::PeerAnnouncement { .. } 222 | GossipEnvelope::PeerVerificationChallenge { .. } 223 | GossipEnvelope::PeerVerificationResponse { .. } => {} 224 } 225 Ok(()) 226 } 227 228 fn ensure_len(label: &str, len: usize, max: usize) -> Result<()> { 229 if len > max { 230 anyhow::bail!("{label} has {len} items, exceeding limit {max}"); 231 } 232 Ok(()) 233 } 234 235 #[cfg(test)] 236 mod tests { 237 use tokio::io::AsyncWriteExt; 238 239 use crate::{ 240 adapters::p2p::metrics::P2pMetricsCounters, 241 app::{BlockInventory, GossipEnvelope, TRANSACTION_BATCH_LIMIT}, 242 domain::{ 243 BURN_COMMITTEE_SIZE, Block, BurnBundle, BurnBundleSection, FinalizerMode, OutPoint, 244 Transaction, TxInput, TxOutput, 245 }, 246 }; 247 248 use super::{ 249 LimitedLineReader, MAX_BLOCK_BATCH, MAX_BLOCK_LOCATOR_HASHES, MAX_GOSSIP_LINE_BYTES, 250 MAX_INVENTORY_ITEMS, MAX_OBJECT_REQUESTS, MAX_PEER_LIST, parse_envelope, 251 record_received_envelope_kind, validate_envelope_limits, 252 }; 253 254 fn burn(signature: &str) -> Transaction { 255 Transaction::Burn { 256 inputs: vec![TxInput { 257 outpoint: OutPoint { 258 txid: format!("{signature:0<64}"), 259 index: 0, 260 }, 261 owner: "owner".to_string(), 262 signature: signature.to_string(), 263 }], 264 change: vec![TxOutput { 265 address: "owner".to_string(), 266 amount: 1, 267 }], 268 amount: 1, 269 fee: 1, 270 anchor: None, 271 signature: signature.to_string(), 272 } 273 } 274 275 fn burn_bundle(slot: u8, signature: &str) -> BurnBundle { 276 BurnBundle { 277 height: 1, 278 prev_hash: "parent".to_string(), 279 slot, 280 member: format!("member-{slot}"), 281 reward_address: None, 282 burns: vec![burn(signature)], 283 burns_v2: Vec::new(), 284 signature: format!("bundle-{signature}"), 285 } 286 } 287 288 fn dummy_block(height: u64) -> Block { 289 Block { 290 height, 291 prev_hash: "0".repeat(64), 292 timestamp_ms: height, 293 miner: "0".repeat(64), 294 reward_address: None, 295 reward_address_signature: None, 296 finalizer_mode: FinalizerMode::Ticket, 297 finalizer_rank: 0, 298 reward: 0, 299 vdf_rounds: 0, 300 vdf_output: "0:0".to_string(), 301 leader_proof: None, 302 burn_bundle_section: BurnBundleSection::default(), 303 transactions: Vec::new(), 304 transactions_v2: Vec::new(), 305 hash: format!("{height:064x}"), 306 } 307 } 308 309 #[test] 310 fn metrics_count_transaction_and_burn_bundle_batches() { 311 let metrics = P2pMetricsCounters::default(); 312 313 record_received_envelope_kind( 314 &metrics, 315 &GossipEnvelope::Transactions { 316 transactions: vec![burn("a"), burn("b")], 317 }, 318 ); 319 record_received_envelope_kind( 320 &metrics, 321 &GossipEnvelope::BurnBundles { 322 bundles: vec![burn_bundle(1, "c"), burn_bundle(2, "d")], 323 }, 324 ); 325 record_received_envelope_kind(&metrics, &GossipEnvelope::Transaction(burn("e"))); 326 record_received_envelope_kind( 327 &metrics, 328 &GossipEnvelope::TransactionsV2 { 329 envelopes: vec!["00".to_string(), "01".to_string()], 330 }, 331 ); 332 record_received_envelope_kind( 333 &metrics, 334 &GossipEnvelope::TransactionV2 { 335 envelope: "02".to_string(), 336 }, 337 ); 338 record_received_envelope_kind(&metrics, &GossipEnvelope::BurnBundle(burn_bundle(1, "f"))); 339 340 let snapshot = metrics.snapshot(); 341 assert_eq!(snapshot.data_envelopes_received, 6); 342 assert_eq!(snapshot.transaction_envelopes_received, 4); 343 assert_eq!(snapshot.transactions_received, 6); 344 assert_eq!(snapshot.burn_bundle_envelopes_received, 2); 345 assert_eq!(snapshot.burn_bundles_received, 3); 346 } 347 348 #[test] 349 fn parser_accepts_burn_bundle_envelopes() { 350 let envelope = GossipEnvelope::BurnBundles { 351 bundles: vec![burn_bundle(1, "a")], 352 }; 353 let line = serde_json::to_string(&envelope).unwrap(); 354 355 assert_eq!(parse_envelope(&line).unwrap(), envelope); 356 357 let envelope = GossipEnvelope::TransactionV2 { 358 envelope: "000102ff".to_string(), 359 }; 360 let line = serde_json::to_string(&envelope).unwrap(); 361 assert_eq!(parse_envelope(&line).unwrap(), envelope); 362 } 363 364 #[test] 365 fn envelope_item_limits_reject_only_above_the_boundary() { 366 assert!( 367 validate_envelope_limits(&GossipEnvelope::BlockRangeRequest { 368 from_height: 1, 369 limit: MAX_BLOCK_BATCH 370 }) 371 .is_ok() 372 ); 373 assert!( 374 validate_envelope_limits(&GossipEnvelope::BlockRangeRequest { 375 from_height: 1, 376 limit: MAX_BLOCK_BATCH + 1 377 }) 378 .is_err() 379 ); 380 assert!( 381 validate_envelope_limits(&GossipEnvelope::BlockRequest { 382 hashes: vec!["0".repeat(64); MAX_OBJECT_REQUESTS] 383 }) 384 .is_ok() 385 ); 386 assert!( 387 validate_envelope_limits(&GossipEnvelope::BlockRequest { 388 hashes: vec!["0".repeat(64); MAX_OBJECT_REQUESTS + 1] 389 }) 390 .is_err() 391 ); 392 assert!( 393 validate_envelope_limits(&GossipEnvelope::Inventory { 394 blocks: vec![ 395 BlockInventory { 396 height: 1, 397 hash: "0".repeat(64) 398 }; 399 MAX_INVENTORY_ITEMS 400 ] 401 }) 402 .is_ok() 403 ); 404 assert!( 405 validate_envelope_limits(&GossipEnvelope::Inventory { 406 blocks: vec![ 407 BlockInventory { 408 height: 1, 409 hash: "0".repeat(64) 410 }; 411 MAX_INVENTORY_ITEMS + 1 412 ] 413 }) 414 .is_err() 415 ); 416 assert!( 417 validate_envelope_limits(&GossipEnvelope::Transactions { 418 transactions: vec![burn("a"); TRANSACTION_BATCH_LIMIT] 419 }) 420 .is_ok() 421 ); 422 assert!( 423 validate_envelope_limits(&GossipEnvelope::Transactions { 424 transactions: vec![burn("a"); TRANSACTION_BATCH_LIMIT + 1] 425 }) 426 .is_err() 427 ); 428 assert!( 429 validate_envelope_limits(&GossipEnvelope::TransactionsV2 { 430 envelopes: vec!["00".to_string(); TRANSACTION_BATCH_LIMIT] 431 }) 432 .is_ok() 433 ); 434 assert!( 435 validate_envelope_limits(&GossipEnvelope::TransactionsV2 { 436 envelopes: vec!["00".to_string(); TRANSACTION_BATCH_LIMIT + 1] 437 }) 438 .is_err() 439 ); 440 assert!( 441 validate_envelope_limits(&GossipEnvelope::BurnBundles { 442 bundles: vec![burn_bundle(1, "a"); TRANSACTION_BATCH_LIMIT] 443 }) 444 .is_ok() 445 ); 446 assert!( 447 validate_envelope_limits(&GossipEnvelope::BurnBundles { 448 bundles: vec![burn_bundle(1, "a"); TRANSACTION_BATCH_LIMIT + 1] 449 }) 450 .is_err() 451 ); 452 assert!( 453 validate_envelope_limits(&GossipEnvelope::BurnBundleRequest { 454 height: 1, 455 prev_hash: "0".repeat(64), 456 slots: vec![1; BURN_COMMITTEE_SIZE] 457 }) 458 .is_ok() 459 ); 460 assert!( 461 validate_envelope_limits(&GossipEnvelope::BurnBundleRequest { 462 height: 1, 463 prev_hash: "0".repeat(64), 464 slots: vec![1; BURN_COMMITTEE_SIZE + 1] 465 }) 466 .is_err() 467 ); 468 assert!( 469 validate_envelope_limits(&GossipEnvelope::Blocks { 470 blocks: vec![dummy_block(1); MAX_BLOCK_BATCH] 471 }) 472 .is_ok() 473 ); 474 assert!( 475 validate_envelope_limits(&GossipEnvelope::Blocks { 476 blocks: vec![dummy_block(1); MAX_BLOCK_BATCH + 1] 477 }) 478 .is_err() 479 ); 480 assert!( 481 validate_envelope_limits(&GossipEnvelope::BlockLocatorRequest { 482 locator: vec!["0".repeat(64); MAX_BLOCK_LOCATOR_HASHES], 483 limit: MAX_BLOCK_BATCH, 484 }) 485 .is_ok() 486 ); 487 assert!( 488 validate_envelope_limits(&GossipEnvelope::BlockLocatorRequest { 489 locator: vec!["0".repeat(64); MAX_BLOCK_LOCATOR_HASHES + 1], 490 limit: MAX_BLOCK_BATCH, 491 }) 492 .is_err() 493 ); 494 assert!( 495 validate_envelope_limits(&GossipEnvelope::PeerList { 496 peers: vec!["127.0.0.1:9444".to_string(); MAX_PEER_LIST] 497 }) 498 .is_ok() 499 ); 500 assert!( 501 validate_envelope_limits(&GossipEnvelope::PeerList { 502 peers: vec!["127.0.0.1:9444".to_string(); MAX_PEER_LIST + 1] 503 }) 504 .is_err() 505 ); 506 } 507 508 #[tokio::test] 509 async fn performance_budget_p2p_line_reader_enforces_message_size() { 510 let (mut client, server) = tokio::io::duplex(MAX_GOSSIP_LINE_BYTES + 1); 511 let mut reader = LimitedLineReader::new(server); 512 let line = vec![b'a'; MAX_GOSSIP_LINE_BYTES]; 513 client.write_all(&line).await.unwrap(); 514 client.write_all(b"\n").await.unwrap(); 515 516 let read = reader.read_line().await.unwrap().unwrap(); 517 518 assert_eq!(read.len(), MAX_GOSSIP_LINE_BYTES); 519 520 let (mut client, server) = tokio::io::duplex(MAX_GOSSIP_LINE_BYTES + 2); 521 let mut reader = LimitedLineReader::new(server); 522 let line = vec![b'a'; MAX_GOSSIP_LINE_BYTES + 1]; 523 client.write_all(&line).await.unwrap(); 524 client.write_all(b"\n").await.unwrap(); 525 526 let error = reader.read_line().await.unwrap_err(); 527 528 assert!(error.to_string().contains("p2p message exceeds")); 529 } 530 531 #[test] 532 fn performance_budget_p2p_batch_parser_enforces_item_limits() { 533 let at_budget = GossipEnvelope::PeerList { 534 peers: vec!["127.0.0.1:9444".to_string(); MAX_PEER_LIST], 535 }; 536 let line = serde_json::to_string(&at_budget).unwrap(); 537 538 assert_eq!(parse_envelope(&line).unwrap(), at_budget); 539 540 let over_budget = GossipEnvelope::PeerList { 541 peers: vec!["127.0.0.1:9444".to_string(); MAX_PEER_LIST + 1], 542 }; 543 let line = serde_json::to_string(&over_budget).unwrap(); 544 let error = parse_envelope(&line).unwrap_err(); 545 546 assert!(error.to_string().contains("peer list has")); 547 } 548 }