Our Personal Data Server from scratch!
0

Configure Feed

Select the types of activity you want to include in your feed.

ripple: transport backpressure & connect coalesing

Lewis: May this revision serve well! <lu5a@proton.me>

authored by

Lewis and committed by
Tangled
(Jun 14, 2026, 6:46 PM +0300) 562f970b 06fd6a1c

+500 -18
+16
crates/tranquil-ripple/src/metrics.rs
··· 50 50 "tranquil_ripple_gossip_delta_bytes", 51 51 "Size of CRDT delta chunks in bytes" 52 52 ); 53 + metrics::describe_counter!( 54 + "tranquil_ripple_transport_write_failures_total", 55 + "Total outbound frame writes that failed or timed out" 56 + ); 57 + metrics::describe_counter!( 58 + "tranquil_ripple_transport_inbound_dropped_total", 59 + "Total inbound frames dropped because the buffer budget was saturated" 60 + ); 53 61 } 54 62 55 63 pub fn record_cache_hit() { ··· 103 111 pub fn record_gossip_delta_bytes(bytes: usize) { 104 112 histogram!("tranquil_ripple_gossip_delta_bytes").record(bytes as f64); 105 113 } 114 + 115 + pub fn record_transport_write_failure() { 116 + counter!("tranquil_ripple_transport_write_failures_total").increment(1); 117 + } 118 + 119 + pub fn record_transport_inbound_dropped() { 120 + counter!("tranquil_ripple_transport_inbound_dropped_total").increment(1); 121 + }
+484 -18
crates/tranquil-ripple/src/transport.rs
··· 6 6 }; 7 7 use rustls::pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer}; 8 8 use std::collections::HashMap; 9 - use std::net::SocketAddr; 9 + use std::net::{IpAddr, SocketAddr}; 10 10 use std::sync::Arc; 11 11 use std::sync::atomic::{AtomicU64, Ordering}; 12 12 use std::time::Duration; 13 - use tokio::sync::{Semaphore, mpsc}; 13 + use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc, watch}; 14 14 use tokio_util::sync::CancellationToken; 15 15 16 16 pub(crate) const MAX_FRAME_SIZE: usize = 4 * 1024 * 1024; 17 17 const MAX_INBOUND_CONNECTIONS: usize = 512; 18 18 const MAX_OUTBOUND_CONNECTIONS: usize = 512; 19 + const MAX_QUEUED_WRITES: usize = 1024; 19 20 const MAX_CONCURRENT_UNI_STREAMS: u32 = 64; 21 + const INBOUND_BYTE_BUDGET: usize = 128 * 1024 * 1024; 22 + const MAX_READS_PER_PEER: usize = 32; 20 23 const READ_CHUNK_BYTES: usize = 256 * 1024; 21 24 const INCOMING_CHANNEL_DEPTH: usize = 1024; 22 25 const STREAM_RECEIVE_WINDOW: u32 = MAX_FRAME_SIZE as u32; ··· 24 27 const KEEPALIVE: Duration = Duration::from_secs(20); 25 28 const IDLE_TIMEOUT: Duration = Duration::from_secs(60); 26 29 const CONNECT_TIMEOUT: Duration = Duration::from_secs(5); 30 + const CONNECT_JOIN_TIMEOUT: Duration = Duration::from_secs(30); 27 31 const WRITE_TIMEOUT: Duration = Duration::from_secs(10); 28 32 const READ_TIMEOUT: Duration = Duration::from_secs(30); 29 33 const RIPPLE_ALPN: &[u8] = b"ripple/1"; ··· 56 60 pub from: SocketAddr, 57 61 pub tag: ChannelTag, 58 62 pub data: Vec<u8>, 63 + _budget: OwnedSemaphorePermit, 59 64 } 60 65 61 66 struct PeerConn { ··· 63 68 generation: u64, 64 69 } 65 70 71 + type ConnectingMap = Arc<parking_lot::Mutex<HashMap<SocketAddr, watch::Receiver<bool>>>>; 72 + 73 + struct ConnectingGuard { 74 + connecting: ConnectingMap, 75 + target: SocketAddr, 76 + } 77 + 78 + impl Drop for ConnectingGuard { 79 + fn drop(&mut self) { 80 + self.connecting.lock().remove(&self.target); 81 + } 82 + } 83 + 66 84 pub struct Transport { 67 85 endpoint: Endpoint, 68 86 local_addr: SocketAddr, 69 87 connections: Arc<parking_lot::Mutex<HashMap<SocketAddr, PeerConn>>>, 70 - connecting: Arc<parking_lot::Mutex<std::collections::HashSet<SocketAddr>>>, 88 + connecting: ConnectingMap, 71 89 conn_generation: Arc<AtomicU64>, 72 90 outbound_permits: Arc<Semaphore>, 91 + queue_permits: Arc<Semaphore>, 92 + inbound_byte_budget: Arc<Semaphore>, 93 + peer_read_limiter: PeerReadLimiter, 73 94 shutdown: CancellationToken, 74 95 incoming_tx: mpsc::Sender<IncomingFrame>, 75 96 } ··· 88 109 endpoint.set_default_client_config(client_config); 89 110 let local_addr = endpoint.local_addr()?; 90 111 let (incoming_tx, incoming_rx) = mpsc::channel(INCOMING_CHANNEL_DEPTH); 112 + let inbound_byte_budget = Arc::new(Semaphore::new(INBOUND_BYTE_BUDGET)); 113 + let peer_read_limiter = PeerReadLimiter::new(MAX_READS_PER_PEER); 91 114 92 115 let transport = Self { 93 116 endpoint: endpoint.clone(), 94 117 local_addr, 95 118 connections: Arc::new(parking_lot::Mutex::new(HashMap::new())), 96 - connecting: Arc::new(parking_lot::Mutex::new(std::collections::HashSet::new())), 119 + connecting: Arc::new(parking_lot::Mutex::new(HashMap::new())), 97 120 conn_generation: Arc::new(AtomicU64::new(0)), 98 121 outbound_permits: Arc::new(Semaphore::new(MAX_OUTBOUND_CONNECTIONS)), 122 + queue_permits: Arc::new(Semaphore::new(MAX_QUEUED_WRITES)), 123 + inbound_byte_budget: inbound_byte_budget.clone(), 124 + peer_read_limiter: peer_read_limiter.clone(), 99 125 shutdown: shutdown.clone(), 100 126 incoming_tx: incoming_tx.clone(), 101 127 }; ··· 122 148 continue; 123 149 }; 124 150 let tx = incoming_tx.clone(); 151 + let byte_budget = inbound_byte_budget.clone(); 152 + let limiter = peer_read_limiter.clone(); 125 153 let conn_cancel = cancel.child_token(); 126 154 tokio::spawn(async move { 127 155 let _permit = permit; ··· 129 157 Ok(conn) => { 130 158 let from = conn.remote_address(); 131 159 tracing::debug!(peer = %from, "accepted inbound connection"); 132 - run_conn_reader(conn, from, tx, conn_cancel).await; 160 + run_conn_reader(conn, from, tx, byte_budget, limiter, conn_cancel) 161 + .await; 133 162 } 134 163 Err(e) => tracing::warn!(error = %e, "inbound handshake failed"), 135 164 } ··· 151 180 pub fn try_queue(&self, target: SocketAddr, tag: ChannelTag, data: &[u8]) -> bool { 152 181 let conn = self.connections.lock().get(&target).map(|p| p.conn.clone()); 153 182 let Some(conn) = conn else { return false }; 183 + let Ok(permit) = self.queue_permits.clone().try_acquire_owned() else { 184 + return false; 185 + }; 154 186 let data = data.to_vec(); 155 187 tokio::spawn(async move { 188 + let _permit = permit; 156 189 if let Err(e) = write_frame(&conn, tag, &data).await { 157 190 tracing::debug!(error = %e, "queued write failed"); 158 191 } ··· 161 194 } 162 195 163 196 pub async fn send(&self, target: SocketAddr, tag: ChannelTag, data: &[u8]) { 197 + let Ok(permit) = self.queue_permits.clone().acquire_owned().await else { 198 + return; 199 + }; 164 200 let existing = self 165 201 .connections 166 202 .lock() ··· 184 220 } 185 221 } 186 222 } 223 + drop(permit); 187 224 self.connect_and_send(target, tag, data).await; 188 225 } 189 226 190 227 async fn connect_and_send(&self, target: SocketAddr, tag: ChannelTag, data: &[u8]) { 191 - { 228 + enum Role { 229 + Lead(watch::Sender<bool>), 230 + Join(watch::Receiver<bool>), 231 + } 232 + let role = { 192 233 let mut connecting = self.connecting.lock(); 193 - if connecting.contains(&target) { 194 - tracing::warn!(peer = %target, "connection already in-flight, dropping frame"); 195 - return; 234 + match connecting.get(&target) { 235 + Some(rx) => Role::Join(rx.clone()), 236 + None => { 237 + let (tx, rx) = watch::channel(false); 238 + connecting.insert(target, rx); 239 + Role::Lead(tx) 240 + } 196 241 } 197 - connecting.insert(target); 242 + }; 243 + match role { 244 + Role::Lead(done) => { 245 + let _guard = ConnectingGuard { 246 + connecting: self.connecting.clone(), 247 + target, 248 + }; 249 + let existing = self.connections.lock().get(&target).map(|p| p.conn.clone()); 250 + match existing { 251 + Some(conn) => { 252 + if let Err(e) = write_frame(&conn, tag, data).await { 253 + tracing::debug!(peer = %target, error = %e, "write on freshly established connection failed"); 254 + } 255 + } 256 + None => self.connect_and_send_inner(target, tag, data).await, 257 + } 258 + let _ = done.send(true); 259 + } 260 + Role::Join(mut done) => { 261 + let _ = 262 + tokio::time::timeout(CONNECT_JOIN_TIMEOUT, done.wait_for(|ready| *ready)).await; 263 + let conn = self.connections.lock().get(&target).map(|p| p.conn.clone()); 264 + match conn { 265 + Some(conn) => { 266 + if let Err(e) = write_frame(&conn, tag, data).await { 267 + tracing::debug!(peer = %target, error = %e, "write after joined connect failed"); 268 + } 269 + } 270 + None => { 271 + crate::metrics::record_transport_write_failure(); 272 + tracing::debug!(peer = %target, "connection attempt failed, dropping frame"); 273 + } 274 + } 275 + } 198 276 } 199 - 200 - self.connect_and_send_inner(target, tag, data).await; 201 - self.connecting.lock().remove(&target); 202 277 } 203 278 204 279 async fn connect_and_send_inner(&self, target: SocketAddr, tag: ChannelTag, data: &[u8]) { ··· 208 283 max = MAX_OUTBOUND_CONNECTIONS, 209 284 "outbound connection limit reached, dropping" 210 285 ); 286 + crate::metrics::record_transport_write_failure(); 211 287 return; 212 288 }; 213 289 let shutdown = self.shutdown.clone(); ··· 249 325 conn.clone(), 250 326 target, 251 327 self.incoming_tx.clone(), 328 + self.inbound_byte_budget.clone(), 329 + self.peer_read_limiter.clone(), 252 330 reader_cancel.clone(), 253 331 )); 254 332 ··· 270 348 tracing::debug!(peer = %target, "established outbound connection"); 271 349 } 272 350 Err(e) => { 351 + crate::metrics::record_transport_write_failure(); 273 352 tracing::warn!(peer = %target, error = %e, "failed to connect after retries"); 274 353 } 275 354 } ··· 280 359 conn: Connection, 281 360 from: SocketAddr, 282 361 incoming_tx: mpsc::Sender<IncomingFrame>, 362 + byte_budget: Arc<Semaphore>, 363 + peer_limiter: PeerReadLimiter, 283 364 cancel: CancellationToken, 284 365 ) { 285 366 loop { ··· 291 372 recv, 292 373 from, 293 374 incoming_tx.clone(), 375 + byte_budget.clone(), 376 + peer_limiter.clone(), 294 377 )); 295 378 } 296 379 Err(_) => break, ··· 301 384 302 385 enum FrameReadError { 303 386 Oversize, 387 + BudgetExhausted, 304 388 Stream(String), 305 389 } 306 390 307 - async fn read_frame(recv: RecvStream, from: SocketAddr, incoming_tx: mpsc::Sender<IncomingFrame>) { 391 + async fn read_frame( 392 + recv: RecvStream, 393 + from: SocketAddr, 394 + incoming_tx: mpsc::Sender<IncomingFrame>, 395 + byte_budget: Arc<Semaphore>, 396 + peer_limiter: PeerReadLimiter, 397 + ) { 398 + let Some(_peer_guard) = peer_limiter.try_acquire(from.ip()) else { 399 + crate::metrics::record_transport_inbound_dropped(); 400 + tracing::debug!(peer = %from, "per-peer inbound read limit reached, dropping frame"); 401 + return; 402 + }; 403 + 308 404 let read = async { 309 405 let mut recv = recv; 310 406 let mut tag_byte = [0u8; 1]; ··· 312 408 .await 313 409 .map_err(|e| FrameReadError::Stream(e.to_string()))?; 314 410 let mut data = Vec::with_capacity(READ_CHUNK_BYTES); 411 + let mut budget = byte_budget 412 + .clone() 413 + .try_acquire_many_owned(0) 414 + .expect("acquiring zero permits always succeeds"); 315 415 loop { 316 416 match recv.read_chunk(READ_CHUNK_BYTES, true).await { 317 417 Ok(Some(chunk)) => { ··· 319 419 if data.len() + len > MAX_FRAME_SIZE { 320 420 return Err(FrameReadError::Oversize); 321 421 } 422 + let permit = byte_budget 423 + .clone() 424 + .try_acquire_many_owned(len as u32) 425 + .map_err(|_| FrameReadError::BudgetExhausted)?; 426 + budget.merge(permit); 322 427 data.extend_from_slice(&chunk.bytes); 323 428 } 324 429 Ok(None) => break, 325 430 Err(e) => return Err(FrameReadError::Stream(e.to_string())), 326 431 } 327 432 } 328 - Ok::<(u8, Vec<u8>), FrameReadError>((tag_byte[0], data)) 433 + Ok::<(u8, Vec<u8>, OwnedSemaphorePermit), FrameReadError>((tag_byte[0], data, budget)) 329 434 }; 330 435 331 436 match tokio::time::timeout(READ_TIMEOUT, read).await { 332 - Ok(Ok((tag_byte, data))) => match ChannelTag::from_u8(tag_byte) { 437 + Ok(Ok((tag_byte, data, budget))) => match ChannelTag::from_u8(tag_byte) { 333 438 Some(tag) => { 334 - let frame = IncomingFrame { from, tag, data }; 439 + let frame = IncomingFrame { 440 + from, 441 + tag, 442 + data, 443 + _budget: budget, 444 + }; 335 445 if let Err(e) = incoming_tx.try_send(frame) { 336 446 tracing::warn!(peer = %from, error = %e, "incoming frame channel full, dropping frame"); 337 447 } ··· 339 449 None => tracing::debug!(tag = tag_byte, "unknown channel tag, dropping frame"), 340 450 }, 341 451 Ok(Err(FrameReadError::Oversize)) => { 452 + crate::metrics::record_transport_inbound_dropped(); 342 453 tracing::debug!(peer = %from, max = MAX_FRAME_SIZE, "inbound frame exceeds max size, dropping"); 343 454 } 455 + Ok(Err(FrameReadError::BudgetExhausted)) => { 456 + crate::metrics::record_transport_inbound_dropped(); 457 + tracing::debug!(peer = %from, "inbound byte budget saturated, dropping frame"); 458 + } 344 459 Ok(Err(FrameReadError::Stream(msg))) => { 345 460 tracing::debug!(peer = %from, error = %msg, "failed reading uni stream"); 346 461 } ··· 350 465 } 351 466 } 352 467 468 + #[derive(Clone)] 469 + struct PeerReadLimiter { 470 + counts: Arc<parking_lot::Mutex<HashMap<IpAddr, usize>>>, 471 + max_per_peer: usize, 472 + } 473 + 474 + impl PeerReadLimiter { 475 + fn new(max_per_peer: usize) -> Self { 476 + Self { 477 + counts: Arc::new(parking_lot::Mutex::new(HashMap::new())), 478 + max_per_peer, 479 + } 480 + } 481 + 482 + fn try_acquire(&self, peer: IpAddr) -> Option<PeerReadGuard> { 483 + let mut counts = self.counts.lock(); 484 + let count = counts.entry(peer).or_insert(0); 485 + match *count >= self.max_per_peer { 486 + true => None, 487 + false => { 488 + *count += 1; 489 + Some(PeerReadGuard { 490 + counts: self.counts.clone(), 491 + peer, 492 + }) 493 + } 494 + } 495 + } 496 + } 497 + 498 + struct PeerReadGuard { 499 + counts: Arc<parking_lot::Mutex<HashMap<IpAddr, usize>>>, 500 + peer: IpAddr, 501 + } 502 + 503 + impl Drop for PeerReadGuard { 504 + fn drop(&mut self) { 505 + let mut counts = self.counts.lock(); 506 + if let Some(count) = counts.get_mut(&self.peer) { 507 + *count -= 1; 508 + if *count == 0 { 509 + counts.remove(&self.peer); 510 + } 511 + } 512 + } 513 + } 514 + 353 515 async fn write_frame(conn: &Connection, tag: ChannelTag, data: &[u8]) -> std::io::Result<()> { 354 516 if data.len() > MAX_FRAME_SIZE { 355 517 tracing::warn!( ··· 357 519 max = MAX_FRAME_SIZE, 358 520 "refusing to send oversized frame" 359 521 ); 522 + crate::metrics::record_transport_write_failure(); 360 523 return Ok(()); 361 524 } 362 525 let timed_out = || std::io::Error::new(std::io::ErrorKind::TimedOut, "write timeout"); 363 526 let deadline = tokio::time::Instant::now() + WRITE_TIMEOUT; 364 - match tokio::time::timeout_at(deadline, conn.open_uni()).await { 527 + let result = match tokio::time::timeout_at(deadline, conn.open_uni()).await { 365 528 Ok(Ok(mut send)) => { 366 529 let write = async { 367 530 send.write_all(&[tag as u8]) ··· 385 548 } 386 549 Ok(Err(e)) => Err(std::io::Error::other(e)), 387 550 Err(_) => Err(timed_out()), 551 + }; 552 + if result.is_err() { 553 + crate::metrics::record_transport_write_failure(); 388 554 } 555 + result 389 556 } 390 557 391 558 fn transport_config() -> TransportConfig { ··· 556 723 } 557 724 558 725 #[tokio::test] 726 + async fn max_size_frame_roundtrips_and_oversize_is_refused() { 727 + let shutdown = CancellationToken::new(); 728 + let (sender, _rx_sender) = 729 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 730 + .await 731 + .expect("bind sender"); 732 + let (receiver, mut rx_receiver) = 733 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 734 + .await 735 + .expect("bind receiver"); 736 + let target = receiver.local_addr(); 737 + 738 + let big = vec![0xABu8; MAX_FRAME_SIZE]; 739 + sender.send(target, ChannelTag::CrdtSync, &big).await; 740 + let frame = tokio::time::timeout(Duration::from_secs(20), rx_receiver.recv()) 741 + .await 742 + .expect("max-size frame arrives before timeout") 743 + .expect("channel open"); 744 + assert_eq!(frame.data.len(), MAX_FRAME_SIZE); 745 + 746 + let oversize = vec![0u8; MAX_FRAME_SIZE + 1]; 747 + sender.send(target, ChannelTag::CrdtSync, &oversize).await; 748 + let res = tokio::time::timeout(Duration::from_secs(2), rx_receiver.recv()).await; 749 + assert!(res.is_err(), "oversize frame must be refused sender-side"); 750 + 751 + shutdown.cancel(); 752 + } 753 + 754 + #[tokio::test] 755 + async fn stalled_streams_capped_per_peer() { 756 + let shutdown = CancellationToken::new(); 757 + let (receiver, _rx_receiver) = 758 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 759 + .await 760 + .expect("bind receiver"); 761 + let target = receiver.local_addr(); 762 + 763 + let client_config = build_client_config().expect("client config"); 764 + let mut endpoint = 765 + Endpoint::client("127.0.0.1:0".parse().unwrap()).expect("client endpoint"); 766 + endpoint.set_default_client_config(client_config); 767 + let conn = endpoint 768 + .connect(target, RIPPLE_SERVER_NAME) 769 + .expect("connect") 770 + .await 771 + .expect("handshake"); 772 + 773 + let mut sends = Vec::new(); 774 + for _ in 0..(MAX_READS_PER_PEER + 16) { 775 + let mut s = conn.open_uni().await.expect("open uni"); 776 + s.write_all(&[ChannelTag::Gossip as u8]) 777 + .await 778 + .expect("write tag byte"); 779 + sends.push(s); 780 + } 781 + 782 + let peer_ip: IpAddr = "127.0.0.1".parse().unwrap(); 783 + let held = || { 784 + receiver 785 + .peer_read_limiter 786 + .counts 787 + .lock() 788 + .get(&peer_ip) 789 + .copied() 790 + .unwrap_or(0) 791 + }; 792 + let mut capped = false; 793 + for _ in 0..100 { 794 + if held() == MAX_READS_PER_PEER { 795 + capped = true; 796 + break; 797 + } 798 + tokio::time::sleep(Duration::from_millis(50)).await; 799 + } 800 + assert!( 801 + capped, 802 + "stalled streams from one peer must be capped at MAX_READS_PER_PEER; held={}", 803 + held() 804 + ); 805 + 806 + drop(sends); 807 + endpoint.close(0u32.into(), b"done"); 808 + shutdown.cancel(); 809 + } 810 + 811 + #[tokio::test] 812 + async fn inbound_byte_budget_held_until_consumed() { 813 + let shutdown = CancellationToken::new(); 814 + let (sender, _rx_sender) = 815 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 816 + .await 817 + .expect("bind sender"); 818 + let (receiver, _rx_receiver) = 819 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 820 + .await 821 + .expect("bind receiver"); 822 + let target = receiver.local_addr(); 823 + 824 + assert_eq!( 825 + receiver.inbound_byte_budget.available_permits(), 826 + INBOUND_BYTE_BUDGET 827 + ); 828 + 829 + let payload = vec![7u8; 1024 * 1024]; 830 + sender.send(target, ChannelTag::CrdtSync, &payload).await; 831 + 832 + let mut held = false; 833 + for _ in 0..100 { 834 + if receiver.inbound_byte_budget.available_permits() 835 + <= INBOUND_BYTE_BUDGET - payload.len() 836 + { 837 + held = true; 838 + break; 839 + } 840 + tokio::time::sleep(Duration::from_millis(50)).await; 841 + } 842 + assert!( 843 + held, 844 + "an undrained inbound frame must hold its byte budget; available={}", 845 + receiver.inbound_byte_budget.available_permits() 846 + ); 847 + 848 + shutdown.cancel(); 849 + } 850 + 851 + #[tokio::test] 559 852 async fn incoming_frame_from_matches_peer_listen_addr() { 560 853 let shutdown = CancellationToken::new(); 561 854 let (sender, _rx_sender) = ··· 581 874 "inbound frames must report the peer's canonical listen address" 582 875 ); 583 876 877 + shutdown.cancel(); 878 + } 879 + 880 + #[tokio::test] 881 + async fn inbound_budget_exhaustion_drops_excess_frames() { 882 + use futures::StreamExt; 883 + 884 + let shutdown = CancellationToken::new(); 885 + let (sender, _rx_sender) = 886 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 887 + .await 888 + .expect("bind sender"); 889 + let (receiver, rx_receiver) = 890 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 891 + .await 892 + .expect("bind receiver"); 893 + let target = receiver.local_addr(); 894 + 895 + let frame_count = INBOUND_BYTE_BUDGET / MAX_FRAME_SIZE + 1; 896 + let payload = vec![0x5Au8; MAX_FRAME_SIZE]; 897 + futures::stream::iter(0..frame_count) 898 + .for_each(|_| sender.send(target, ChannelTag::CrdtSync, &payload)) 899 + .await; 900 + 901 + let received: Vec<IncomingFrame> = futures::stream::unfold(rx_receiver, |mut rx| async { 902 + tokio::time::timeout(Duration::from_secs(5), rx.recv()) 903 + .await 904 + .ok() 905 + .flatten() 906 + .map(|frame| (frame, rx)) 907 + }) 908 + .collect() 909 + .await; 910 + 911 + assert!( 912 + !received.is_empty(), 913 + "frames within the budget must be delivered" 914 + ); 915 + assert!( 916 + received.len() < frame_count, 917 + "frames beyond the inbound byte budget must be dropped while earlier frames sit unconsumed; sent={frame_count} received={}", 918 + received.len() 919 + ); 920 + 921 + shutdown.cancel(); 922 + } 923 + 924 + #[tokio::test] 925 + async fn concurrent_sends_to_fresh_peer_all_delivered() { 926 + let shutdown = CancellationToken::new(); 927 + let (sender, _rx_sender) = 928 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 929 + .await 930 + .expect("bind sender"); 931 + let (receiver, rx_receiver) = 932 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 933 + .await 934 + .expect("bind receiver"); 935 + let target = receiver.local_addr(); 936 + 937 + let payloads: Vec<Vec<u8>> = (0..8u8).map(|i| vec![i; 16]).collect(); 938 + futures::future::join_all( 939 + payloads 940 + .iter() 941 + .map(|p| sender.send(target, ChannelTag::CrdtSync, p)), 942 + ) 943 + .await; 944 + 945 + use futures::StreamExt; 946 + let received: Vec<IncomingFrame> = futures::stream::unfold(rx_receiver, |mut rx| async { 947 + tokio::time::timeout(Duration::from_secs(2), rx.recv()) 948 + .await 949 + .ok() 950 + .flatten() 951 + .map(|frame| (frame, rx)) 952 + }) 953 + .collect() 954 + .await; 955 + 956 + let mut seen: Vec<Vec<u8>> = received.into_iter().map(|f| f.data).collect(); 957 + seen.sort(); 958 + assert_eq!( 959 + seen, payloads, 960 + "sends racing a fresh connection must join it and deliver every frame" 961 + ); 962 + 963 + shutdown.cancel(); 964 + } 965 + 966 + #[tokio::test] 967 + async fn timed_out_write_resets_stream_instead_of_truncating() { 968 + let shutdown = CancellationToken::new(); 969 + let (sender, _rx_sender) = 970 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 971 + .await 972 + .expect("bind sender"); 973 + 974 + let server_config = build_server_config().expect("server config"); 975 + let stalled_server = Endpoint::server(server_config, "127.0.0.1:0".parse().unwrap()) 976 + .expect("bind stalled server"); 977 + let target = stalled_server.local_addr().expect("local addr"); 978 + 979 + let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>(); 980 + let reader = tokio::spawn(async move { 981 + let incoming = stalled_server.accept().await.expect("incoming"); 982 + let conn = incoming.await.expect("handshake"); 983 + let mut recv = conn.accept_uni().await.expect("accept uni"); 984 + release_rx.await.expect("release signal"); 985 + recv.read_to_end(MAX_FRAME_SIZE * 2).await 986 + }); 987 + 988 + let payload = vec![0xEEu8; MAX_FRAME_SIZE]; 989 + let started = tokio::time::Instant::now(); 990 + sender.send(target, ChannelTag::CrdtSync, &payload).await; 991 + assert!( 992 + started.elapsed() >= WRITE_TIMEOUT, 993 + "a frame one byte over the stream receive window must stall until the write timeout" 994 + ); 995 + 996 + release_tx.send(()).expect("server task alive"); 997 + let read = reader.await.expect("server task"); 998 + assert!( 999 + read.is_err(), 1000 + "a timed-out partial write must reset the stream, not surface a truncated frame; read {} bytes", 1001 + read.map(|d| d.len()).unwrap_or(0) 1002 + ); 1003 + 1004 + shutdown.cancel(); 1005 + } 1006 + 1007 + #[tokio::test] 1008 + async fn write_timeout_keeps_connection() { 1009 + use futures::StreamExt; 1010 + 1011 + let shutdown = CancellationToken::new(); 1012 + let (sender, _rx_sender) = 1013 + Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) 1014 + .await 1015 + .expect("bind sender"); 1016 + 1017 + let server_config = build_server_config().expect("server config"); 1018 + let mute_server = Endpoint::server(server_config, "127.0.0.1:0".parse().unwrap()) 1019 + .expect("bind mute server"); 1020 + let target = mute_server.local_addr().expect("local addr"); 1021 + let mute = tokio::spawn(async move { 1022 + let incoming = mute_server.accept().await.expect("incoming"); 1023 + let conn = incoming.await.expect("handshake"); 1024 + std::future::pending::<()>().await; 1025 + drop(conn); 1026 + }); 1027 + 1028 + futures::stream::iter(0..u64::from(MAX_CONCURRENT_UNI_STREAMS)) 1029 + .for_each(|_| sender.send(target, ChannelTag::Gossip, b"fill stream credit")) 1030 + .await; 1031 + assert!( 1032 + sender.connections.lock().contains_key(&target), 1033 + "connection must be established before exhausting stream credit" 1034 + ); 1035 + 1036 + let started = tokio::time::Instant::now(); 1037 + sender 1038 + .send(target, ChannelTag::Gossip, b"blocked frame") 1039 + .await; 1040 + assert!( 1041 + started.elapsed() >= WRITE_TIMEOUT, 1042 + "send with exhausted stream credit must block until the write timeout" 1043 + ); 1044 + assert!( 1045 + sender.connections.lock().contains_key(&target), 1046 + "a timed-out write must keep the connection registered" 1047 + ); 1048 + 1049 + mute.abort(); 584 1050 shutdown.cancel(); 585 1051 } 586 1052 }