From feabdf3db116bf004291459aebf4ecdcaebfc9e2 Mon Sep 17 00:00:00 2001 From: ldm0 Date: Thu, 10 Sep 2026 16:11:47 +0800 Subject: [PATCH] perf(curl): reuse receive storage after would-block --- moli-curl/src/websocket.rs | 4 ++ moli-curl/src/websocket/owner.rs | 4 +- moli-curl/src/websocket/session.rs | 20 +++++-- moli-curl/src/websocket/tests.rs | 1 + moli-curl/src/websocket/tests/storage.rs | 72 ++++++++++++++++++++++++ 5 files changed, 95 insertions(+), 6 deletions(-) create mode 100644 moli-curl/src/websocket/tests/storage.rs diff --git a/moli-curl/src/websocket.rs b/moli-curl/src/websocket.rs index ba13fe4cba..f4ad2bc4da 100644 --- a/moli-curl/src/websocket.rs +++ b/moli-curl/src/websocket.rs @@ -134,6 +134,8 @@ struct Control { #[cfg(test)] read_attempts: std::sync::atomic::AtomicUsize, #[cfg(test)] + receive_allocations: std::sync::atomic::AtomicUsize, + #[cfg(test)] read_waiting: tokio::sync::Notify, } @@ -343,6 +345,8 @@ impl CurlWebSocketRuntime { #[cfg(test)] read_attempts: std::sync::atomic::AtomicUsize::new(0), #[cfg(test)] + receive_allocations: std::sync::atomic::AtomicUsize::new(0), + #[cfg(test)] read_waiting: tokio::sync::Notify::new(), }); let (event_tx, events) = mpsc::channel(MAX_PENDING_EVENTS); diff --git a/moli-curl/src/websocket/owner.rs b/moli-curl/src/websocket/owner.rs index 249fb54a96..ee18175bf4 100644 --- a/moli-curl/src/websocket/owner.rs +++ b/moli-curl/src/websocket/owner.rs @@ -38,6 +38,7 @@ struct Owner { sessions: HashMap, dns: CurlDnsOwnerResidence, poll: SocketPoll, + receive: Vec, } pub(super) fn run( @@ -50,6 +51,7 @@ pub(super) fn run( sessions: HashMap::new(), dns: CurlDnsOwnerResidence::default(), poll: SocketPoll::default(), + receive: Vec::new(), }; let _ = waker_tx.send(owner.multi.waker()); while !shutdown.load(Ordering::Acquire) { @@ -169,7 +171,7 @@ impl Owner { let mut retired = Vec::new(); let mut progressed = false; for (id, session) in &mut self.sessions { - match session.advance() { + match session.advance(&mut self.receive) { Ok(Step::Progress) => progressed = true, Ok(Step::Idle) => {} terminal => retired.push((*id, terminal)), diff --git a/moli-curl/src/websocket/session.rs b/moli-curl/src/websocket/session.rs index 6a8d4add6b..28c38f43d7 100644 --- a/moli-curl/src/websocket/session.rs +++ b/moli-curl/src/websocket/session.rs @@ -113,7 +113,7 @@ impl Session { Ok(Step::Idle) } - pub(super) fn advance(&mut self) -> Result { + pub(super) fn advance(&mut self, receive: &mut Vec) -> Result { if self.io.cancelled() { return Ok(Step::Closed); } @@ -131,7 +131,7 @@ impl Session { } // Receive delivery must never hold up a pending write. progressed |= self.write_pending()?; - match self.read_chunk()? { + match self.read_chunk(receive)? { Step::Progress => progressed = true, Step::Idle => break, Step::Closed => return Ok(Step::Closed), @@ -180,7 +180,7 @@ impl Session { matches!(self.phase, Phase::Open) && self.io.control.reading.load(Ordering::Acquire) } - fn read_chunk(&mut self) -> Result { + fn read_chunk(&mut self, receive: &mut Vec) -> Result { if !self.reading.can_run(self.reading_enabled()) { return Ok(Step::Idle); } @@ -194,14 +194,24 @@ impl Session { } Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => return Ok(Step::Closed), }; - let mut data = vec![0; CHUNK_BYTES]; + // The owner lends one spare buffer across connections. AGAIN keeps it; + // success transfers ownership to the event without copying the payload. + #[cfg(test)] + if receive.capacity() < CHUNK_BYTES { + self.io + .control + .receive_allocations + .fetch_add(1, Ordering::Relaxed); + } + receive.resize(CHUNK_BYTES, 0); #[cfg(test)] self.io .control .read_attempts .fetch_add(1, Ordering::Relaxed); - match self.handle.ws_recv(&mut data) { + match self.handle.ws_recv(receive) { Ok((count, frame)) => { + let mut data = std::mem::take(receive); data.truncate(count); // Let the caller answer Close before probing EOF. A peer can // half-close TCP in the same packet as its Close frame. diff --git a/moli-curl/src/websocket/tests.rs b/moli-curl/src/websocket/tests.rs index 8bebdd9ded..5b6d9e8d54 100644 --- a/moli-curl/src/websocket/tests.rs +++ b/moli-curl/src/websocket/tests.rs @@ -10,6 +10,7 @@ use tokio_tungstenite::tungstenite::{self, Message, handshake::derive_accept_key use super::*; mod readiness; +mod storage; const DEADLINE: Duration = Duration::from_secs(10); diff --git a/moli-curl/src/websocket/tests/storage.rs b/moli-curl/src/websocket/tests/storage.rs new file mode 100644 index 0000000000..3e8458c811 --- /dev/null +++ b/moli-curl/src/websocket/tests/storage.rs @@ -0,0 +1,72 @@ +use super::*; + +#[tokio::test] +async fn native_receive_storage_is_shared_across_idle_probes_and_transferred_to_chunks() { + let runtime = CurlWebSocketRuntime::new().unwrap(); + let mut connections = Vec::new(); + let mut peers = Vec::new(); + let mut writers = Vec::new(); + for _ in 0..3 { + let (write, writes) = std::sync::mpsc::channel::(); + let (url, peer) = server(move |mut stream| { + upgrade(&mut stream, &[]); + while let Ok(byte) = writes.recv_timeout(DEADLINE) { + stream.write_all(&[0x82, 1, byte]).unwrap(); + } + assert_eq!(stream.read(&mut [0]).unwrap(), 0); + }); + let mut connection = runtime.connect(CurlWebSocketRequest::new(url)).unwrap(); + opened(&mut connection).await; + connection.sender().set_reading(true); + timeout(DEADLINE, connection.sender.control.read_waiting.notified()) + .await + .unwrap(); + connections.push(connection); + peers.push(peer); + writers.push(write); + } + let allocations = |connections: &[CurlWebSocketConnection]| { + connections + .iter() + .map(|connection| { + connection + .sender + .control + .receive_allocations + .load(Ordering::Acquire) + }) + .sum::() + }; + assert_eq!( + allocations(&connections), + 1, + "idle connections must share the owner's spare receive storage" + ); + + let mut retained = Vec::new(); + for byte in 0..12 { + let index = usize::from(byte) % connections.len(); + writers[index].send(byte).unwrap(); + let CurlWebSocketEvent::Chunk { data, .. } = event(&mut connections[index]).await else { + panic!("expected payload"); + }; + retained.push(data); + timeout( + DEADLINE, + connections[index].sender.control.read_waiting.notified(), + ) + .await + .unwrap(); + // Success transfers the previous allocation; the following AGAIN only + // needs one replacement. Keep all prior chunks alive to detect aliasing. + assert_eq!(allocations(&connections), usize::from(byte) + 2); + for (expected, data) in retained.iter().enumerate() { + assert_eq!(data, &[expected as u8]); + } + } + drop(writers); + drop(connections); + for peer in peers { + peer.join().unwrap(); + } +}