diff --git a/benches/storage_bench.rs b/benches/storage_bench.rs index 384810f..a7e488f 100644 --- a/benches/storage_bench.rs +++ b/benches/storage_bench.rs @@ -286,8 +286,14 @@ fn ensure_server_started() -> &'static ServerHandle { let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); let network = twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity( + queueber::RPC_IO_BUFFER_BYTES, + reader, + ), + futures::io::BufWriter::with_capacity( + queueber::RPC_IO_BUFFER_BYTES, + writer, + ), rpc_twoparty_capnp::Side::Server, Default::default(), ); @@ -329,8 +335,8 @@ where stream.set_nodelay(true).unwrap(); let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); let rpc_network = Box::new(twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity(queueber::RPC_IO_BUFFER_BYTES, reader), + futures::io::BufWriter::with_capacity(queueber::RPC_IO_BUFFER_BYTES, writer), rpc_twoparty_capnp::Side::Client, Default::default(), )); diff --git a/docs/perf-notes.md b/docs/perf-notes.md index 2377c39..f5bd5a5 100644 --- a/docs/perf-notes.md +++ b/docs/perf-notes.md @@ -107,10 +107,9 @@ ### Measurement plan -- Use `benches/*` and `criterion` to measure CPU/latency; extend with scenarios covering add/poll/remove mixes. -- Run `stress.sh` under varying concurrency to validate coalescing and contention improvements. -- Validate worker and blocking pool utilization (`top -H`, tracing). -- Add CPU/heap profiling (`pprof`/jemalloc) to identify remaining hotspots. +### Tried optimizations and results + +- RPC IO buffers: Increased `futures::io::BufReader`/`futures::io::BufWriter` capacities to 64 KiB at all RPC endpoints (server, client, benches, tests). Result: no measurable improvement in throughput or latency in stress tests/benchmarks (within noise). PR: `https://github.com/ORG/REPO/pull/XXXX` ### References diff --git a/src/bin/client/main.rs b/src/bin/client/main.rs index edd047f..685f7a8 100644 --- a/src/bin/client/main.rs +++ b/src/bin/client/main.rs @@ -539,8 +539,8 @@ where stream.set_nodelay(true).unwrap(); let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); let rpc_network = Box::new(twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity(queueber::RPC_IO_BUFFER_BYTES, reader), + futures::io::BufWriter::with_capacity(queueber::RPC_IO_BUFFER_BYTES, writer), rpc_twoparty_capnp::Side::Client, Default::default(), )); diff --git a/src/bin/queueber/main.rs b/src/bin/queueber/main.rs index 5fe624b..fc5b947 100644 --- a/src/bin/queueber/main.rs +++ b/src/bin/queueber/main.rs @@ -138,8 +138,8 @@ async fn main() -> Result<()> { ) .split(); let network = twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity(queueber::RPC_IO_BUFFER_BYTES, reader), + futures::io::BufWriter::with_capacity(queueber::RPC_IO_BUFFER_BYTES, writer), rpc_twoparty_capnp::Side::Server, Default::default(), ); diff --git a/src/lib.rs b/src/lib.rs index e252a60..6e7071c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -10,3 +10,5 @@ pub mod storage; // Re-export commonly used types pub use crate::storage::RetriedStorage; pub use crate::storage::Storage; + +pub const RPC_IO_BUFFER_BYTES: usize = 64 * 1024; diff --git a/tests/server_poll.rs b/tests/server_poll.rs index b86bf39..ebfbf08 100644 --- a/tests/server_poll.rs +++ b/tests/server_poll.rs @@ -77,8 +77,8 @@ fn start_test_server() -> TestServerHandle { let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); let network = twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity(queueber::RPC_IO_BUFFER_BYTES, reader), + futures::io::BufWriter::with_capacity(queueber::RPC_IO_BUFFER_BYTES, writer), rpc_twoparty_capnp::Side::Server, Default::default(), ); @@ -116,8 +116,8 @@ where stream.set_nodelay(true).unwrap(); let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); let rpc_network = Box::new(twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity(queueber::RPC_IO_BUFFER_BYTES, reader), + futures::io::BufWriter::with_capacity(queueber::RPC_IO_BUFFER_BYTES, writer), rpc_twoparty_capnp::Side::Client, Default::default(), )); diff --git a/tests/startup_recovery.rs b/tests/startup_recovery.rs index 85c6a8a..9887c11 100644 --- a/tests/startup_recovery.rs +++ b/tests/startup_recovery.rs @@ -90,8 +90,8 @@ fn start_server_on(data_dir: std::path::PathBuf, addr: SocketAddr) -> TestServer stream.set_nodelay(true).unwrap(); let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); let network = twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity(queueber::RPC_IO_BUFFER_BYTES, reader), + futures::io::BufWriter::with_capacity(queueber::RPC_IO_BUFFER_BYTES, writer), rpc_twoparty_capnp::Side::Server, Default::default(), ); @@ -125,8 +125,8 @@ where stream.set_nodelay(true).unwrap(); let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split(); let rpc_network = Box::new(twoparty::VatNetwork::new( - futures::io::BufReader::new(reader), - futures::io::BufWriter::new(writer), + futures::io::BufReader::with_capacity(queueber::RPC_IO_BUFFER_BYTES, reader), + futures::io::BufWriter::with_capacity(queueber::RPC_IO_BUFFER_BYTES, writer), rpc_twoparty_capnp::Side::Client, Default::default(), ));