mirror of
https://github.com/lexmount/moli.git
synced 2026-10-05 16:00:54 +00:00
perf(curl): reuse receive storage after would-block
This commit is contained in:
@@ -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);
|
||||
|
||||
@@ -38,6 +38,7 @@ struct Owner {
|
||||
sessions: HashMap<CurlTransferId, Session>,
|
||||
dns: CurlDnsOwnerResidence<CurlTransferId, Pending>,
|
||||
poll: SocketPoll,
|
||||
receive: Vec<u8>,
|
||||
}
|
||||
|
||||
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)),
|
||||
|
||||
@@ -113,7 +113,7 @@ impl Session {
|
||||
Ok(Step::Idle)
|
||||
}
|
||||
|
||||
pub(super) fn advance(&mut self) -> Result<Step, String> {
|
||||
pub(super) fn advance(&mut self, receive: &mut Vec<u8>) -> Result<Step, String> {
|
||||
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<Step, String> {
|
||||
fn read_chunk(&mut self, receive: &mut Vec<u8>) -> Result<Step, String> {
|
||||
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.
|
||||
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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::<u8>();
|
||||
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::<usize>()
|
||||
};
|
||||
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();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user