From 852a1b1e5c78a494630ff87cbcb1902372430542 Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 08:25:12 +0200 Subject: [PATCH 01/13] test: define opaque pool API contracts Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/opaque_memory.rs | 8 ++++++++ crates/fetch/src/custom.rs | 8 ++++++++ crates/http_extensions/src/body/builder.rs | 8 ++++++++ 3 files changed, 24 insertions(+) diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index 8c374c466..f6cd12072 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -82,6 +82,14 @@ mod tests { assert_impl_all!(OpaqueMemory: MemoryShared); + #[test] + #[ignore = "stub"] + fn opaque_pool_wraps_custom_provider() { + // Arrange a custom MemoryShared provider. + // Wrap it in OpaquePool and reserve memory. + // Assert the custom provider handled the reservation. + } + #[test] fn wraps_inner() { let provider = GlobalPool::new(); diff --git a/crates/fetch/src/custom.rs b/crates/fetch/src/custom.rs index e1bf7fff7..f1d43fb77 100644 --- a/crates/fetch/src/custom.rs +++ b/crates/fetch/src/custom.rs @@ -311,6 +311,14 @@ mod tests { FakeHandler::from_fn(|_req| HttpResponseBuilder::new_fake().status(StatusCode::OK).build()) } + #[test] + #[ignore = "stub"] + fn custom_deps_accept_custom_opaque_pool() { + // Arrange CustomDeps with a non-GlobalPool provider wrapped in OpaquePool. + // Build a custom client and execute a request. + // Assert the request succeeds through the custom pipeline. + } + #[cfg_attr(miri, ignore)] #[tokio::test] async fn create_builder_serves_requests_through_custom_pipeline() { diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 9f8f38ffe..b03255f5d 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -399,6 +399,14 @@ mod tests { assert_impl_all!(HttpBodyBuilder: Send, Sync, AsRef, std::fmt::Debug); } + #[test] + #[ignore = "stub"] + fn new_accepts_opaque_pool_with_custom_memory() { + // Arrange a custom MemoryShared provider wrapped in OpaquePool. + // Create an HttpBodyBuilder with HttpBodyBuilder::new. + // Assert body creation reserves through the custom provider. + } + #[test] fn new_with_global_memory() { let clock = Clock::new_frozen(); From 6024805e59e6d617a5bf0f97243c460b1fd96f63 Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 08:27:03 +0200 Subject: [PATCH 02/13] feat(bytesbuf): add opaque memory pool Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/mod.rs | 2 +- crates/bytesbuf/src/mem/opaque_memory.rs | 27 +++++++------ crates/fetch/Cargo.toml | 2 +- crates/fetch/examples/http_client_app.rs | 6 +-- crates/fetch/examples/http_client_custom.rs | 4 +- crates/fetch/src/custom.rs | 18 ++++----- crates/fetch/src/fake.rs | 2 +- crates/fetch/src/tokio.rs | 8 ++-- crates/fetch/tests/telemetry_scope.rs | 4 +- crates/fetch_hyper/src/builder.rs | 4 +- crates/http_extensions/Cargo.toml | 2 +- .../http_extensions/examples/custom_client.rs | 4 +- .../http_extensions/examples/custom_server.rs | 4 +- crates/http_extensions/src/body/builder.rs | 39 +++++-------------- .../http_extensions/src/body/timeout_body.rs | 21 +++++----- 15 files changed, 66 insertions(+), 81 deletions(-) diff --git a/crates/bytesbuf/src/mem/mod.rs b/crates/bytesbuf/src/mem/mod.rs index 5af516f84..39a04dd5f 100644 --- a/crates/bytesbuf/src/mem/mod.rs +++ b/crates/bytesbuf/src/mem/mod.rs @@ -56,7 +56,7 @@ pub use global::GlobalPool; pub use has_memory::HasMemory; pub use memory::Memory; pub use memory_shared::MemoryShared; -pub use opaque_memory::OpaqueMemory; +pub use opaque_memory::{OpaqueMemory, OpaquePool}; #[cfg(any(test, feature = "test-util"))] pub mod testing; diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index f6cd12072..52e8e08aa 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -8,7 +8,7 @@ use thread_aware::ThreadAware; use crate::mem::{Memory, MemoryShared}; -/// Adapter to erase the type of a [`MemoryShared`] implementation. +/// Type-erased memory pool backed by any [`MemoryShared`] implementation. /// /// This adapter adds some inefficiency due to additional indirection overhead for /// every memory reservation, so avoid this adapter if you can tolerate alternatives (generics). @@ -18,12 +18,12 @@ use crate::mem::{Memory, MemoryShared}; /// provider. Cloning the adapter clones the wrapped provider; whether the clones then share any /// state is up to that provider. #[derive(Debug, ThreadAware)] -pub struct OpaqueMemory { +pub struct OpaquePool { inner: Box, } -impl OpaqueMemory { - /// Creates a new instance of the adapter. +impl OpaquePool { + /// Creates a type-erased pool backed by `inner`. #[must_use] pub fn new(inner: impl MemoryShared) -> Self { Self { inner: Box::new(inner) } @@ -53,7 +53,7 @@ impl OpaqueMemory { } } -impl Clone for OpaqueMemory { +impl Clone for OpaquePool { fn clone(&self) -> Self { Self { inner: self.inner.clone_boxed(), @@ -61,13 +61,16 @@ impl Clone for OpaqueMemory { } } -impl Memory for OpaqueMemory { +impl Memory for OpaquePool { #[cfg_attr(test, mutants::skip)] // Trivial forwarder. fn reserve(&self, min_bytes: usize) -> crate::BytesBuf { self.reserve(min_bytes) } } +/// Compatibility alias for the former name of [`OpaquePool`]. +pub type OpaqueMemory = OpaquePool; + #[cfg_attr(coverage_nightly, coverage(off))] #[cfg(all(test, feature = "std"))] mod tests { @@ -80,7 +83,7 @@ mod tests { use super::*; use crate::mem::GlobalPool; - assert_impl_all!(OpaqueMemory: MemoryShared); + assert_impl_all!(OpaquePool: MemoryShared); #[test] #[ignore = "stub"] @@ -93,7 +96,7 @@ mod tests { #[test] fn wraps_inner() { let provider = GlobalPool::new(); - let memory = OpaqueMemory::new(provider); + let memory = OpaquePool::new(provider); let builder = memory.reserve(1024); assert!(builder.capacity() >= 1024); @@ -102,7 +105,7 @@ mod tests { #[test] fn memory_trait() { let provider = GlobalPool::new(); - let memory = OpaqueMemory::new(provider); + let memory = OpaquePool::new(provider); // Call reserve via the Memory trait to verify the impl block let builder = Memory::reserve(&memory, 1024); @@ -111,7 +114,7 @@ mod tests { #[test] fn relocate_does_not_break_reservation() { - let mut memory = OpaqueMemory::new(GlobalPool::new()); + let mut memory = OpaquePool::new(GlobalPool::new()); let affinities = pinned_affinities(&[2]); memory.relocate(Some(affinities[0]), affinities[1]); @@ -144,7 +147,7 @@ mod tests { } let relocated = Arc::new(AtomicUsize::new(0)); - let mut memory = OpaqueMemory::new(TrackingMemory { + let mut memory = OpaquePool::new(TrackingMemory { relocated: Arc::clone(&relocated), inner: GlobalPool::new(), }); @@ -157,7 +160,7 @@ mod tests { #[test] fn clone_is_usable_independently() { - let memory = OpaqueMemory::new(GlobalPool::new()); + let memory = OpaquePool::new(GlobalPool::new()); let mut clone = memory.clone(); // Relocating the clone must leave both the clone and the original usable. diff --git a/crates/fetch/Cargo.toml b/crates/fetch/Cargo.toml index 72ea9c287..07b078507 100644 --- a/crates/fetch/Cargo.toml +++ b/crates/fetch/Cargo.toml @@ -19,9 +19,9 @@ repository = "https://github.com/microsoft/oxidizer/tree/main/crates/fetch" [package.metadata.cargo_check_external_types] allowed_external_types = [ - "bytesbuf::mem::global::GlobalPool", "bytesbuf::mem::has_memory::HasMemory", "bytesbuf::mem::memory::Memory", + "bytesbuf::mem::opaque_memory::OpaquePool", "data_privacy::redaction_engine::RedactionEngine", "fetch_options::*", "fetch_tls::*", diff --git a/crates/fetch/examples/http_client_app.rs b/crates/fetch/examples/http_client_app.rs index 9d7846a6c..5244d5707 100644 --- a/crates/fetch/examples/http_client_app.rs +++ b/crates/fetch/examples/http_client_app.rs @@ -6,7 +6,7 @@ use std::sync::Arc; -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaquePool}; use fetch::HttpClient; use fetch::tls::TlsOptions; use ohno::ErrorExt; @@ -18,7 +18,7 @@ use tick::Clock; #[fundle::bundle] struct App { clock: Clock, - global_pool: GlobalPool, + memory_pool: OpaquePool, client: HttpClient, } @@ -42,7 +42,7 @@ async fn main() -> Result<(), ohno::AppError> { // Initialize and set up the App instance; fundle ensures all fields are properly constructed. let app = App::builder() .clock(|_| Clock::new_tokio()) - .global_pool(|_| GlobalPool::new()) + .memory_pool(|_| OpaquePool::new(GlobalPool::new())) .client({ let meter_provider = meter_provider.clone(); move |x| { diff --git a/crates/fetch/examples/http_client_custom.rs b/crates/fetch/examples/http_client_custom.rs index dea4f0dc8..04de20479 100644 --- a/crates/fetch/examples/http_client_custom.rs +++ b/crates/fetch/examples/http_client_custom.rs @@ -4,7 +4,7 @@ //! Plugs a custom `EchoHandler` into [`fetch::custom::create_builder`] as the transport //! handler. Every request's body is returned verbatim in the response. -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaquePool}; use fetch::custom::{CustomDeps, Isolation, create_builder}; use fetch::{HttpRequest, HttpResponse, HttpResponseBuilder}; use http::StatusCode; @@ -16,7 +16,7 @@ use tick::Clock; async fn main() -> Result<(), ohno::AppError> { let deps = CustomDeps { clock: Clock::new_tokio(), - global_pool: GlobalPool::new(), + memory_pool: OpaquePool::new(GlobalPool::new()), extras: (), }; diff --git a/crates/fetch/src/custom.rs b/crates/fetch/src/custom.rs index f1d43fb77..b2879b68f 100644 --- a/crates/fetch/src/custom.rs +++ b/crates/fetch/src/custom.rs @@ -17,7 +17,7 @@ use std::borrow::Cow; use std::fmt::Debug; use std::sync::Arc; -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::OpaquePool; use http_extensions::{HttpBodyBuilder, RequestHandler}; use opentelemetry::metrics::Meter; use thread_aware::{PerCore, ThreadAware, unaware}; @@ -53,7 +53,7 @@ where /// Clock for timing operations and timeouts. pub clock: Clock, /// Memory pool for usage-neutral memory allocations. - pub global_pool: GlobalPool, + pub memory_pool: OpaquePool, /// Extra dependencies forwarded verbatim to [`CustomContext::extras`]. pub extras: Extras, } @@ -210,12 +210,12 @@ impl HttpClient { runtime_name: runtime.into(), name: transport.into(), clock: deps.clock.clone(), - global_pool: deps.global_pool.clone(), + memory_pool: deps.memory_pool.clone(), isolation, inner: thread_aware::Arc::new_with((deps, unaware(factory)), |(deps, factory)| { Arc::new(move |options, meter, pool_index| { let context = CustomContext { - body_builder: create_body_builder(&deps.global_pool, &deps.clock, &options), + body_builder: create_body_builder(&deps.memory_pool, &deps.clock, &options), clock: deps.clock.clone(), pool_index, extras: deps.extras.clone(), @@ -242,7 +242,7 @@ pub(crate) struct Transport { name: Cow<'static, str>, inner: thread_aware::Arc, clock: Clock, - global_pool: GlobalPool, + memory_pool: OpaquePool, isolation: Isolation, } @@ -268,7 +268,7 @@ impl Transport { } pub(crate) fn create_body_builder(&self, options: &ClientOptions) -> HttpBodyBuilder { - create_body_builder(&self.global_pool, &self.clock, options) + create_body_builder(&self.memory_pool, &self.clock, options) } } @@ -278,7 +278,7 @@ impl Debug for Transport { } } -pub(crate) fn create_body_builder(pool: &GlobalPool, clock: &Clock, options: &ClientOptions) -> HttpBodyBuilder { +pub(crate) fn create_body_builder(pool: &OpaquePool, clock: &Clock, options: &ClientOptions) -> HttpBodyBuilder { HttpBodyBuilder::new(pool.clone(), clock).with_options(options.response_body_options) } @@ -301,7 +301,7 @@ mod tests { fn custom_deps() -> CustomDeps { CustomDeps { clock: FakeDeps::default().clock, - global_pool: bytesbuf::mem::GlobalPool::new(), + memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: (), } } @@ -356,7 +356,7 @@ mod tests { let counter = Arc::new(AtomicUsize::new(0)); let deps = CustomDeps { clock: FakeDeps::default().clock, - global_pool: bytesbuf::mem::GlobalPool::new(), + memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: unaware(Arc::clone(&counter)), }; diff --git a/crates/fetch/src/fake.rs b/crates/fetch/src/fake.rs index 2669a4c13..91acef55f 100644 --- a/crates/fetch/src/fake.rs +++ b/crates/fetch/src/fake.rs @@ -84,7 +84,7 @@ impl HttpClient { Isolation::Shared, CustomDeps { clock: deps.clock, - global_pool: bytesbuf::mem::GlobalPool::new(), + memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: handler, }, ) diff --git a/crates/fetch/src/tokio.rs b/crates/fetch/src/tokio.rs index 71a00d6f4..6558b3a4f 100644 --- a/crates/fetch/src/tokio.rs +++ b/crates/fetch/src/tokio.rs @@ -33,7 +33,7 @@ pub struct TokioDeps { /// Clock for timing operations and timeouts. pub clock: Clock, /// Memory pool for usage-neutral memory allocations. - pub global_pool: bytesbuf::mem::GlobalPool, + pub memory_pool: bytesbuf::mem::OpaquePool, } impl Default for TokioDeps { @@ -47,7 +47,7 @@ impl TokioDeps { #[must_use] pub fn with_clock(clock: &Clock) -> Self { Self { - global_pool: bytesbuf::mem::GlobalPool::new(), + memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), clock: clock.clone(), } } @@ -64,7 +64,7 @@ impl HttpClient { pub fn builder_tokio(deps: impl Into) -> HttpClientBuilder { let deps = deps.into(); let clock = deps.clock.clone(); - let global_pool = deps.global_pool.clone(); + let memory_pool = deps.memory_pool.clone(); // Re-layer on top of the in-crate `builder_custom_internal` path: the // full `TokioDeps` rides through `CustomDeps::extras` so that the @@ -77,7 +77,7 @@ impl HttpClient { Isolation::Shared, CustomDeps { clock, - global_pool, + memory_pool, extras: deps, }, ) diff --git a/crates/fetch/tests/telemetry_scope.rs b/crates/fetch/tests/telemetry_scope.rs index b3e09999e..042498301 100644 --- a/crates/fetch/tests/telemetry_scope.rs +++ b/crates/fetch/tests/telemetry_scope.rs @@ -94,7 +94,7 @@ async fn custom_transport_scope_attribute() { let deps = CustomDeps { clock: Clock::new_frozen(), - global_pool: bytesbuf::mem::GlobalPool::new(), + memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; @@ -128,7 +128,7 @@ async fn custom_transport_instrument_inherits_scope() { let deps = CustomDeps { clock: Clock::new_frozen(), - global_pool: bytesbuf::mem::GlobalPool::new(), + memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; diff --git a/crates/fetch_hyper/src/builder.rs b/crates/fetch_hyper/src/builder.rs index db7803ac5..e22e4d9c1 100644 --- a/crates/fetch_hyper/src/builder.rs +++ b/crates/fetch_hyper/src/builder.rs @@ -8,7 +8,7 @@ use std::fmt; use std::marker::PhantomData; use anyspawn::Spawner; -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaquePool}; use fetch_options::{ConnectionIdleTimeout, ConnectionKeepAlive, ConnectionPoolOptions, Http2Options, PoolIndex, TransportOptions}; use fetch_tls::TlsBackend; use http::Version; @@ -205,7 +205,7 @@ where let body_builder = self .body_builder .clone() - .unwrap_or_else(|| HttpBodyBuilder::new(GlobalPool::new(), &self.clock)); + .unwrap_or_else(|| HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &self.clock)); HyperTransport::new(build_hyper_handler(self, tls, body_builder, &meter).into_dynamic()) } diff --git a/crates/http_extensions/Cargo.toml b/crates/http_extensions/Cargo.toml index 5434d2dc2..152374fc5 100644 --- a/crates/http_extensions/Cargo.toml +++ b/crates/http_extensions/Cargo.toml @@ -25,10 +25,10 @@ ignored = ["uuid"] allowed_external_types = [ "futures_core::stream::Stream", "bytes::bytes::Bytes", - "bytesbuf::mem::global::GlobalPool", "bytesbuf::mem::has_memory::HasMemory", "bytesbuf::mem::memory::Memory", "bytesbuf::mem::memory_shared::MemoryShared", + "bytesbuf::mem::opaque_memory::OpaquePool", "bytesbuf::view::BytesView", "http::*", "http_body::Body", diff --git a/crates/http_extensions/examples/custom_client.rs b/crates/http_extensions/examples/custom_client.rs index 0663b4124..7ac778d9d 100644 --- a/crates/http_extensions/examples/custom_client.rs +++ b/crates/http_extensions/examples/custom_client.rs @@ -6,7 +6,7 @@ //! This example demonstrates how to create a simple HTTP client that just echoes back the //! data it receives. -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaquePool}; use http_extensions::{HttpBodyBuilder, HttpRequest, HttpRequestBuilderExt, HttpResponse, HttpResponseBuilder, StatusExt}; use layered::Service; use tick::Clock; @@ -46,7 +46,7 @@ impl AsRef for CustomClient { impl Default for CustomClient { fn default() -> Self { Self { - builder: HttpBodyBuilder::new(GlobalPool::new(), &Clock::new_tokio()), + builder: HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &Clock::new_tokio()), } } } diff --git a/crates/http_extensions/examples/custom_server.rs b/crates/http_extensions/examples/custom_server.rs index 1b75e64cf..8f9e723be 100644 --- a/crates/http_extensions/examples/custom_server.rs +++ b/crates/http_extensions/examples/custom_server.rs @@ -11,7 +11,7 @@ use std::sync::Arc; use std::time::Duration; use bytesbuf::BytesView; -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaquePool}; use futures::TryStreamExt; use http::Request; use http_body_util::BodyExt; @@ -31,7 +31,7 @@ async fn main() -> Result<(), ohno::AppError> { let clock = Clock::new_tokio(); // In a real application, the application framework would provide the global memory pool. - let body_builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let body_builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); let body_builder_clone = body_builder.clone(); // Define an execution stack of middleware diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index b03255f5d..6a2b36184 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -1,7 +1,7 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT License. -use bytesbuf::mem::{GlobalPool, HasMemory, Memory, MemoryShared, OpaqueMemory}; +use bytesbuf::mem::{GlobalPool, HasMemory, Memory, MemoryShared, OpaquePool}; use bytesbuf::{BytesBuf, BytesView}; use futures::{Stream, TryStreamExt}; use http_body::{Body, Frame}; @@ -44,7 +44,7 @@ use crate::{HttpError, Result}; /// With the `test-util` feature enabled, you can create a test instance using `HttpBodyBuilder::new_fake()`. #[derive(Debug, Clone, ThreadAware)] pub struct HttpBodyBuilder { - memory: MemoryWrapper, + memory: OpaquePool, clock: Clock, pub(super) options: HttpBodyOptions, } @@ -71,16 +71,16 @@ impl HttpBodyBuilder { #[cfg(any(feature = "test-util", test))] #[must_use] pub fn new_fake() -> Self { - Self::new(GlobalPool::new(), &Clock::new_frozen()) + Self::new(OpaquePool::new(GlobalPool::new()), &Clock::new_frozen()) } /// Creates a new instance of [`HttpBodyBuilder`]. /// - /// This method uses a per-thread memory pool from [`GlobalPool`]. + /// The provided pool can wrap any [`MemoryShared`] implementation. #[must_use] - pub fn new(memory: GlobalPool, clock: &Clock) -> Self { + pub fn new(memory: OpaquePool, clock: &Clock) -> Self { Self { - memory: MemoryWrapper::Global(memory), + memory, clock: clock.clone(), options: HttpBodyOptions::default(), } @@ -94,11 +94,7 @@ impl HttpBodyBuilder { /// relocated along with it. #[must_use] pub fn with_custom_memory(memory: impl MemoryShared, clock: &Clock) -> Self { - Self { - memory: MemoryWrapper::Opaque(OpaqueMemory::new(memory)), - clock: clock.clone(), - options: HttpBodyOptions::default(), - } + Self::new(OpaquePool::new(memory), clock) } /// Sets default [`HttpBodyOptions`] for all bodies created by this builder. @@ -357,21 +353,6 @@ impl HasMemory for HttpBodyBuilder { } } -#[derive(Debug, Clone, ThreadAware)] -enum MemoryWrapper { - Global(GlobalPool), - Opaque(OpaqueMemory), -} - -impl Memory for MemoryWrapper { - fn reserve(&self, min_bytes: usize) -> BytesBuf { - match self { - Self::Global(pool) => pool.reserve(min_bytes), - Self::Opaque(memory) => memory.reserve(min_bytes), - } - } -} - impl AsRef for HttpBodyBuilder { fn as_ref(&self) -> &Clock { &self.clock @@ -410,7 +391,7 @@ mod tests { #[test] fn new_with_global_memory() { let clock = Clock::new_frozen(); - let memory = GlobalPool::new(); + let memory = OpaquePool::new(GlobalPool::new()); let builder = HttpBodyBuilder::new(memory, &clock); let body = builder.text("test"); assert_eq!(body.content_length(), Some(4)); @@ -551,7 +532,7 @@ mod tests { #[test] fn stream_with_timeout_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); let chunks: Vec> = [b"hello " as &[u8], b"world"] .iter() .map(|c| Ok(BytesView::copied_from_slice(c, &builder))) @@ -627,7 +608,7 @@ mod tests { fn builder_merges_per_call_options_with_defaults() { let clock = Clock::new_frozen(); let builder_options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(builder_options); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock).with_options(builder_options); // Per-call options override the builder-level default. let per_call = HttpBodyOptions::default().timeout(Duration::from_secs(5)); diff --git a/crates/http_extensions/src/body/timeout_body.rs b/crates/http_extensions/src/body/timeout_body.rs index c3f59fadc..b8e450736 100644 --- a/crates/http_extensions/src/body/timeout_body.rs +++ b/crates/http_extensions/src/body/timeout_body.rs @@ -96,7 +96,7 @@ mod tests { use std::time::Duration; use bytesbuf::BytesView; - use bytesbuf::mem::GlobalPool; + use bytesbuf::mem::{GlobalPool, OpaquePool}; use futures::executor::block_on; use http_body::{Body, Frame}; use tick::ClockControl; @@ -107,7 +107,7 @@ mod tests { #[test] fn stream_body_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); // Stream yields data immediately — well within the generous timeout, // exercising the TimeoutBody happy path via stream with timeout options. @@ -121,7 +121,7 @@ mod tests { #[test] fn stream_body_times_out_when_pending() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); // A body that never yields data. let options = HttpBodyOptions::default().timeout(Duration::from_millis(100)); @@ -136,7 +136,8 @@ mod tests { #[test] fn body_timeout_chains_with_buffer_limit() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); + let builder = + HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); assert_eq!(builder.options, HttpBodyOptions::default().buffer_limit(1024)); @@ -160,7 +161,7 @@ mod tests { #[test] fn size_hint_delegates_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); // Full body has an exact size hint; verify it passes through TimeoutBody. let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); @@ -176,7 +177,7 @@ mod tests { #[test] fn is_end_stream_true_when_inner_is_empty() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Empty::new(), &options); @@ -186,7 +187,7 @@ mod tests { #[test] fn is_end_stream_false_when_inner_has_data() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Full::new(BytesView::copied_from_slice(b"data", &builder)), &options); @@ -196,7 +197,7 @@ mod tests { #[test] fn poll_frame_returns_data_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); let body = builder.body( @@ -210,7 +211,7 @@ mod tests { #[test] fn poll_frame_times_out_when_pending_with_short_timeout() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); // A body that never yields data with a very short timeout. let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); @@ -227,7 +228,7 @@ mod tests { fn poll_frame_returns_data_even_when_clock_advanced_past_timeout() { let control = ClockControl::new(); let clock = control.to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); // Use a body that has data immediately available (Full is always ready). let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); From f679e1d8a9c9fca02d373ddc77b021a0130dcb3f Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 08:29:09 +0200 Subject: [PATCH 03/13] test: cover custom opaque pools Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/opaque_memory.rs | 10 +++++---- crates/fetch/src/custom.rs | 25 ++++++++++++++++------ crates/http_extensions/src/body/builder.rs | 11 ++++++---- 3 files changed, 32 insertions(+), 14 deletions(-) diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index 52e8e08aa..1015b931c 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -82,15 +82,17 @@ mod tests { use super::*; use crate::mem::GlobalPool; + use crate::mem::testing::TransparentMemory; assert_impl_all!(OpaquePool: MemoryShared); #[test] - #[ignore = "stub"] fn opaque_pool_wraps_custom_provider() { - // Arrange a custom MemoryShared provider. - // Wrap it in OpaquePool and reserve memory. - // Assert the custom provider handled the reservation. + let pool = OpaquePool::new(TransparentMemory::new()); + + let buffer = pool.reserve(1024); + + assert_eq!(buffer.capacity(), 1024); } #[test] diff --git a/crates/fetch/src/custom.rs b/crates/fetch/src/custom.rs index b2879b68f..1d9184458 100644 --- a/crates/fetch/src/custom.rs +++ b/crates/fetch/src/custom.rs @@ -293,6 +293,8 @@ mod tests { use thread_aware::unaware; use super::{CustomContext, CustomDeps, Isolation, create_builder}; + use bytesbuf::mem::OpaquePool; + use bytesbuf::mem::testing::TransparentMemory; use crate::HttpResponseBuilder; use crate::fake::FakeDeps; use crate::pipeline::Pipeline; @@ -311,12 +313,23 @@ mod tests { FakeHandler::from_fn(|_req| HttpResponseBuilder::new_fake().status(StatusCode::OK).build()) } - #[test] - #[ignore = "stub"] - fn custom_deps_accept_custom_opaque_pool() { - // Arrange CustomDeps with a non-GlobalPool provider wrapped in OpaquePool. - // Build a custom client and execute a request. - // Assert the request succeeds through the custom pipeline. + #[cfg_attr(miri, ignore)] + #[tokio::test] + async fn custom_deps_accept_custom_opaque_pool() { + let deps = CustomDeps { + clock: FakeDeps::default().clock, + memory_pool: OpaquePool::new(TransparentMemory::new()), + extras: (), + }; + + let client = create_builder("test-runtime", "test", ok_factory, Isolation::Shared, deps) + .insecure_allow_http() + .minimal_pipeline() + .build(); + + let response = client.post("http://example.com").text("custom pool").fetch().await.unwrap(); + + assert_eq!(response.status(), StatusCode::OK); } #[cfg_attr(miri, ignore)] diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 6a2b36184..1e8833d17 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -381,11 +381,14 @@ mod tests { } #[test] - #[ignore = "stub"] fn new_accepts_opaque_pool_with_custom_memory() { - // Arrange a custom MemoryShared provider wrapped in OpaquePool. - // Create an HttpBodyBuilder with HttpBodyBuilder::new. - // Assert body creation reserves through the custom provider. + let clock = Clock::new_frozen(); + let pool = OpaquePool::new(TransparentMemory::new()); + + let builder = HttpBodyBuilder::new(pool, &clock); + let body = builder.text("custom pool"); + + assert_eq!(body.content_length(), Some(11)); } #[test] From 517079842b3413d46e0b2057ab7983766430d83d Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 08:38:40 +0200 Subject: [PATCH 04/13] test(fetch): avoid test-util pool dependency Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/fetch/src/custom.rs | 19 +++++++++++++++---- 1 file changed, 15 insertions(+), 4 deletions(-) diff --git a/crates/fetch/src/custom.rs b/crates/fetch/src/custom.rs index 1d9184458..0704f427d 100644 --- a/crates/fetch/src/custom.rs +++ b/crates/fetch/src/custom.rs @@ -288,13 +288,13 @@ mod tests { use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; + use bytesbuf::BytesBuf; + use bytesbuf::mem::{GlobalPool, Memory, OpaquePool}; use http::StatusCode; use http_extensions::FakeHandler; - use thread_aware::unaware; + use thread_aware::{ThreadAware, unaware}; use super::{CustomContext, CustomDeps, Isolation, create_builder}; - use bytesbuf::mem::OpaquePool; - use bytesbuf::mem::testing::TransparentMemory; use crate::HttpResponseBuilder; use crate::fake::FakeDeps; use crate::pipeline::Pipeline; @@ -313,12 +313,23 @@ mod tests { FakeHandler::from_fn(|_req| HttpResponseBuilder::new_fake().status(StatusCode::OK).build()) } + #[derive(Clone, Debug, ThreadAware)] + struct CustomMemory { + inner: GlobalPool, + } + + impl Memory for CustomMemory { + fn reserve(&self, min_bytes: usize) -> BytesBuf { + self.inner.reserve(min_bytes) + } + } + #[cfg_attr(miri, ignore)] #[tokio::test] async fn custom_deps_accept_custom_opaque_pool() { let deps = CustomDeps { clock: FakeDeps::default().clock, - memory_pool: OpaquePool::new(TransparentMemory::new()), + memory_pool: OpaquePool::new(CustomMemory { inner: GlobalPool::new() }), extras: (), }; From d465abad80820d6011894704c96178ab61bcf1e8 Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 08:57:48 +0200 Subject: [PATCH 05/13] fix: preserve memory pool API compatibility Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/opaque_memory.rs | 7 +++++++ crates/fetch/examples/http_client_app.rs | 4 ++-- crates/fetch/examples/http_client_custom.rs | 2 +- crates/fetch/src/custom.rs | 16 +++++++------- crates/fetch/src/fake.rs | 2 +- crates/fetch/src/tokio.rs | 8 +++---- crates/fetch/tests/telemetry_scope.rs | 4 ++-- crates/fetch_hyper/src/builder.rs | 4 ++-- .../http_extensions/examples/custom_client.rs | 4 ++-- .../http_extensions/examples/custom_server.rs | 4 ++-- crates/http_extensions/src/body/builder.rs | 16 +++++++------- .../http_extensions/src/body/timeout_body.rs | 21 +++++++++---------- 12 files changed, 50 insertions(+), 42 deletions(-) diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index 1015b931c..a9b694977 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -53,6 +53,13 @@ impl OpaquePool { } } +#[cfg(feature = "std")] +impl From for OpaquePool { + fn from(pool: crate::mem::GlobalPool) -> Self { + Self::new(pool) + } +} + impl Clone for OpaquePool { fn clone(&self) -> Self { Self { diff --git a/crates/fetch/examples/http_client_app.rs b/crates/fetch/examples/http_client_app.rs index 5244d5707..fd66cb417 100644 --- a/crates/fetch/examples/http_client_app.rs +++ b/crates/fetch/examples/http_client_app.rs @@ -18,7 +18,7 @@ use tick::Clock; #[fundle::bundle] struct App { clock: Clock, - memory_pool: OpaquePool, + global_pool: OpaquePool, client: HttpClient, } @@ -42,7 +42,7 @@ async fn main() -> Result<(), ohno::AppError> { // Initialize and set up the App instance; fundle ensures all fields are properly constructed. let app = App::builder() .clock(|_| Clock::new_tokio()) - .memory_pool(|_| OpaquePool::new(GlobalPool::new())) + .global_pool(|_| OpaquePool::new(GlobalPool::new())) .client({ let meter_provider = meter_provider.clone(); move |x| { diff --git a/crates/fetch/examples/http_client_custom.rs b/crates/fetch/examples/http_client_custom.rs index 04de20479..456dec40b 100644 --- a/crates/fetch/examples/http_client_custom.rs +++ b/crates/fetch/examples/http_client_custom.rs @@ -16,7 +16,7 @@ use tick::Clock; async fn main() -> Result<(), ohno::AppError> { let deps = CustomDeps { clock: Clock::new_tokio(), - memory_pool: OpaquePool::new(GlobalPool::new()), + global_pool: OpaquePool::new(GlobalPool::new()), extras: (), }; diff --git a/crates/fetch/src/custom.rs b/crates/fetch/src/custom.rs index 0704f427d..2a5309f05 100644 --- a/crates/fetch/src/custom.rs +++ b/crates/fetch/src/custom.rs @@ -53,7 +53,7 @@ where /// Clock for timing operations and timeouts. pub clock: Clock, /// Memory pool for usage-neutral memory allocations. - pub memory_pool: OpaquePool, + pub global_pool: OpaquePool, /// Extra dependencies forwarded verbatim to [`CustomContext::extras`]. pub extras: Extras, } @@ -210,12 +210,12 @@ impl HttpClient { runtime_name: runtime.into(), name: transport.into(), clock: deps.clock.clone(), - memory_pool: deps.memory_pool.clone(), + global_pool: deps.global_pool.clone(), isolation, inner: thread_aware::Arc::new_with((deps, unaware(factory)), |(deps, factory)| { Arc::new(move |options, meter, pool_index| { let context = CustomContext { - body_builder: create_body_builder(&deps.memory_pool, &deps.clock, &options), + body_builder: create_body_builder(&deps.global_pool, &deps.clock, &options), clock: deps.clock.clone(), pool_index, extras: deps.extras.clone(), @@ -242,7 +242,7 @@ pub(crate) struct Transport { name: Cow<'static, str>, inner: thread_aware::Arc, clock: Clock, - memory_pool: OpaquePool, + global_pool: OpaquePool, isolation: Isolation, } @@ -268,7 +268,7 @@ impl Transport { } pub(crate) fn create_body_builder(&self, options: &ClientOptions) -> HttpBodyBuilder { - create_body_builder(&self.memory_pool, &self.clock, options) + create_body_builder(&self.global_pool, &self.clock, options) } } @@ -303,7 +303,7 @@ mod tests { fn custom_deps() -> CustomDeps { CustomDeps { clock: FakeDeps::default().clock, - memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: (), } } @@ -329,7 +329,7 @@ mod tests { async fn custom_deps_accept_custom_opaque_pool() { let deps = CustomDeps { clock: FakeDeps::default().clock, - memory_pool: OpaquePool::new(CustomMemory { inner: GlobalPool::new() }), + global_pool: OpaquePool::new(CustomMemory { inner: GlobalPool::new() }), extras: (), }; @@ -380,7 +380,7 @@ mod tests { let counter = Arc::new(AtomicUsize::new(0)); let deps = CustomDeps { clock: FakeDeps::default().clock, - memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: unaware(Arc::clone(&counter)), }; diff --git a/crates/fetch/src/fake.rs b/crates/fetch/src/fake.rs index 91acef55f..7d52a7c47 100644 --- a/crates/fetch/src/fake.rs +++ b/crates/fetch/src/fake.rs @@ -84,7 +84,7 @@ impl HttpClient { Isolation::Shared, CustomDeps { clock: deps.clock, - memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: handler, }, ) diff --git a/crates/fetch/src/tokio.rs b/crates/fetch/src/tokio.rs index 6558b3a4f..493053bde 100644 --- a/crates/fetch/src/tokio.rs +++ b/crates/fetch/src/tokio.rs @@ -33,7 +33,7 @@ pub struct TokioDeps { /// Clock for timing operations and timeouts. pub clock: Clock, /// Memory pool for usage-neutral memory allocations. - pub memory_pool: bytesbuf::mem::OpaquePool, + pub global_pool: bytesbuf::mem::OpaquePool, } impl Default for TokioDeps { @@ -47,7 +47,7 @@ impl TokioDeps { #[must_use] pub fn with_clock(clock: &Clock) -> Self { Self { - memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), clock: clock.clone(), } } @@ -64,7 +64,7 @@ impl HttpClient { pub fn builder_tokio(deps: impl Into) -> HttpClientBuilder { let deps = deps.into(); let clock = deps.clock.clone(); - let memory_pool = deps.memory_pool.clone(); + let global_pool = deps.global_pool.clone(); // Re-layer on top of the in-crate `builder_custom_internal` path: the // full `TokioDeps` rides through `CustomDeps::extras` so that the @@ -77,7 +77,7 @@ impl HttpClient { Isolation::Shared, CustomDeps { clock, - memory_pool, + global_pool, extras: deps, }, ) diff --git a/crates/fetch/tests/telemetry_scope.rs b/crates/fetch/tests/telemetry_scope.rs index 042498301..cbecafa24 100644 --- a/crates/fetch/tests/telemetry_scope.rs +++ b/crates/fetch/tests/telemetry_scope.rs @@ -94,7 +94,7 @@ async fn custom_transport_scope_attribute() { let deps = CustomDeps { clock: Clock::new_frozen(), - memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; @@ -128,7 +128,7 @@ async fn custom_transport_instrument_inherits_scope() { let deps = CustomDeps { clock: Clock::new_frozen(), - memory_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; diff --git a/crates/fetch_hyper/src/builder.rs b/crates/fetch_hyper/src/builder.rs index e22e4d9c1..db7803ac5 100644 --- a/crates/fetch_hyper/src/builder.rs +++ b/crates/fetch_hyper/src/builder.rs @@ -8,7 +8,7 @@ use std::fmt; use std::marker::PhantomData; use anyspawn::Spawner; -use bytesbuf::mem::{GlobalPool, OpaquePool}; +use bytesbuf::mem::GlobalPool; use fetch_options::{ConnectionIdleTimeout, ConnectionKeepAlive, ConnectionPoolOptions, Http2Options, PoolIndex, TransportOptions}; use fetch_tls::TlsBackend; use http::Version; @@ -205,7 +205,7 @@ where let body_builder = self .body_builder .clone() - .unwrap_or_else(|| HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &self.clock)); + .unwrap_or_else(|| HttpBodyBuilder::new(GlobalPool::new(), &self.clock)); HyperTransport::new(build_hyper_handler(self, tls, body_builder, &meter).into_dynamic()) } diff --git a/crates/http_extensions/examples/custom_client.rs b/crates/http_extensions/examples/custom_client.rs index 7ac778d9d..0663b4124 100644 --- a/crates/http_extensions/examples/custom_client.rs +++ b/crates/http_extensions/examples/custom_client.rs @@ -6,7 +6,7 @@ //! This example demonstrates how to create a simple HTTP client that just echoes back the //! data it receives. -use bytesbuf::mem::{GlobalPool, OpaquePool}; +use bytesbuf::mem::GlobalPool; use http_extensions::{HttpBodyBuilder, HttpRequest, HttpRequestBuilderExt, HttpResponse, HttpResponseBuilder, StatusExt}; use layered::Service; use tick::Clock; @@ -46,7 +46,7 @@ impl AsRef for CustomClient { impl Default for CustomClient { fn default() -> Self { Self { - builder: HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &Clock::new_tokio()), + builder: HttpBodyBuilder::new(GlobalPool::new(), &Clock::new_tokio()), } } } diff --git a/crates/http_extensions/examples/custom_server.rs b/crates/http_extensions/examples/custom_server.rs index 8f9e723be..1b75e64cf 100644 --- a/crates/http_extensions/examples/custom_server.rs +++ b/crates/http_extensions/examples/custom_server.rs @@ -11,7 +11,7 @@ use std::sync::Arc; use std::time::Duration; use bytesbuf::BytesView; -use bytesbuf::mem::{GlobalPool, OpaquePool}; +use bytesbuf::mem::GlobalPool; use futures::TryStreamExt; use http::Request; use http_body_util::BodyExt; @@ -31,7 +31,7 @@ async fn main() -> Result<(), ohno::AppError> { let clock = Clock::new_tokio(); // In a real application, the application framework would provide the global memory pool. - let body_builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let body_builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let body_builder_clone = body_builder.clone(); // Define an execution stack of middleware diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 1e8833d17..470c5a60e 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -71,16 +71,18 @@ impl HttpBodyBuilder { #[cfg(any(feature = "test-util", test))] #[must_use] pub fn new_fake() -> Self { - Self::new(OpaquePool::new(GlobalPool::new()), &Clock::new_frozen()) + Self::new(GlobalPool::new(), &Clock::new_frozen()) } /// Creates a new instance of [`HttpBodyBuilder`]. /// - /// The provided pool can wrap any [`MemoryShared`] implementation. + /// Accepts either an [`OpaquePool`] or the default [`GlobalPool`]. + /// Use [`with_custom_memory`](Self::with_custom_memory) to wrap another [`MemoryShared`] + /// implementation. #[must_use] - pub fn new(memory: OpaquePool, clock: &Clock) -> Self { + pub fn new(memory: impl Into, clock: &Clock) -> Self { Self { - memory, + memory: memory.into(), clock: clock.clone(), options: HttpBodyOptions::default(), } @@ -394,7 +396,7 @@ mod tests { #[test] fn new_with_global_memory() { let clock = Clock::new_frozen(); - let memory = OpaquePool::new(GlobalPool::new()); + let memory = GlobalPool::new(); let builder = HttpBodyBuilder::new(memory, &clock); let body = builder.text("test"); assert_eq!(body.content_length(), Some(4)); @@ -535,7 +537,7 @@ mod tests { #[test] fn stream_with_timeout_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let chunks: Vec> = [b"hello " as &[u8], b"world"] .iter() .map(|c| Ok(BytesView::copied_from_slice(c, &builder))) @@ -611,7 +613,7 @@ mod tests { fn builder_merges_per_call_options_with_defaults() { let clock = Clock::new_frozen(); let builder_options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock).with_options(builder_options); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(builder_options); // Per-call options override the builder-level default. let per_call = HttpBodyOptions::default().timeout(Duration::from_secs(5)); diff --git a/crates/http_extensions/src/body/timeout_body.rs b/crates/http_extensions/src/body/timeout_body.rs index b8e450736..c3f59fadc 100644 --- a/crates/http_extensions/src/body/timeout_body.rs +++ b/crates/http_extensions/src/body/timeout_body.rs @@ -96,7 +96,7 @@ mod tests { use std::time::Duration; use bytesbuf::BytesView; - use bytesbuf::mem::{GlobalPool, OpaquePool}; + use bytesbuf::mem::GlobalPool; use futures::executor::block_on; use http_body::{Body, Frame}; use tick::ClockControl; @@ -107,7 +107,7 @@ mod tests { #[test] fn stream_body_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // Stream yields data immediately — well within the generous timeout, // exercising the TimeoutBody happy path via stream with timeout options. @@ -121,7 +121,7 @@ mod tests { #[test] fn stream_body_times_out_when_pending() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // A body that never yields data. let options = HttpBodyOptions::default().timeout(Duration::from_millis(100)); @@ -136,8 +136,7 @@ mod tests { #[test] fn body_timeout_chains_with_buffer_limit() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = - HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); assert_eq!(builder.options, HttpBodyOptions::default().buffer_limit(1024)); @@ -161,7 +160,7 @@ mod tests { #[test] fn size_hint_delegates_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // Full body has an exact size hint; verify it passes through TimeoutBody. let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); @@ -177,7 +176,7 @@ mod tests { #[test] fn is_end_stream_true_when_inner_is_empty() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Empty::new(), &options); @@ -187,7 +186,7 @@ mod tests { #[test] fn is_end_stream_false_when_inner_has_data() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Full::new(BytesView::copied_from_slice(b"data", &builder)), &options); @@ -197,7 +196,7 @@ mod tests { #[test] fn poll_frame_returns_data_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); let body = builder.body( @@ -211,7 +210,7 @@ mod tests { #[test] fn poll_frame_times_out_when_pending_with_short_timeout() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // A body that never yields data with a very short timeout. let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); @@ -228,7 +227,7 @@ mod tests { fn poll_frame_returns_data_even_when_clock_advanced_past_timeout() { let control = ClockControl::new(); let clock = control.to_clock(); - let builder = HttpBodyBuilder::new(OpaquePool::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // Use a body that has data immediately available (Full is always ready). let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); From d4b2b88ef3fd8abada9af4141b98840081965926 Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 09:08:10 +0200 Subject: [PATCH 06/13] fix(http_extensions): gate test-only pool import Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/http_extensions/src/body/builder.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 470c5a60e..84647fac2 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -1,7 +1,9 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT License. -use bytesbuf::mem::{GlobalPool, HasMemory, Memory, MemoryShared, OpaquePool}; +#[cfg(any(feature = "test-util", test))] +use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{HasMemory, Memory, MemoryShared, OpaquePool}; use bytesbuf::{BytesBuf, BytesView}; use futures::{Stream, TryStreamExt}; use http_body::{Body, Frame}; From 5da02f308aa02008f7ebe905659968273ff513c8 Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 09:53:12 +0200 Subject: [PATCH 07/13] review: use existing opaque memory adapter Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/mod.rs | 2 +- crates/bytesbuf/src/mem/opaque_memory.rs | 44 +++++-------------- crates/fetch/Cargo.toml | 2 +- crates/fetch/examples/http_client_app.rs | 6 +-- crates/fetch/examples/http_client_custom.rs | 4 +- crates/fetch/src/custom.rs | 18 ++++---- crates/fetch/src/fake.rs | 2 +- crates/fetch/src/tokio.rs | 4 +- crates/fetch/tests/telemetry_scope.rs | 4 +- crates/fetch_hyper/src/builder.rs | 4 +- crates/http_extensions/Cargo.toml | 2 +- .../http_extensions/examples/custom_client.rs | 4 +- .../http_extensions/examples/custom_server.rs | 4 +- crates/http_extensions/src/body/builder.rs | 28 ++++++------ .../http_extensions/src/body/timeout_body.rs | 21 ++++----- 15 files changed, 64 insertions(+), 85 deletions(-) diff --git a/crates/bytesbuf/src/mem/mod.rs b/crates/bytesbuf/src/mem/mod.rs index 39a04dd5f..5af516f84 100644 --- a/crates/bytesbuf/src/mem/mod.rs +++ b/crates/bytesbuf/src/mem/mod.rs @@ -56,7 +56,7 @@ pub use global::GlobalPool; pub use has_memory::HasMemory; pub use memory::Memory; pub use memory_shared::MemoryShared; -pub use opaque_memory::{OpaqueMemory, OpaquePool}; +pub use opaque_memory::OpaqueMemory; #[cfg(any(test, feature = "test-util"))] pub mod testing; diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index a9b694977..8c374c466 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -8,7 +8,7 @@ use thread_aware::ThreadAware; use crate::mem::{Memory, MemoryShared}; -/// Type-erased memory pool backed by any [`MemoryShared`] implementation. +/// Adapter to erase the type of a [`MemoryShared`] implementation. /// /// This adapter adds some inefficiency due to additional indirection overhead for /// every memory reservation, so avoid this adapter if you can tolerate alternatives (generics). @@ -18,12 +18,12 @@ use crate::mem::{Memory, MemoryShared}; /// provider. Cloning the adapter clones the wrapped provider; whether the clones then share any /// state is up to that provider. #[derive(Debug, ThreadAware)] -pub struct OpaquePool { +pub struct OpaqueMemory { inner: Box, } -impl OpaquePool { - /// Creates a type-erased pool backed by `inner`. +impl OpaqueMemory { + /// Creates a new instance of the adapter. #[must_use] pub fn new(inner: impl MemoryShared) -> Self { Self { inner: Box::new(inner) } @@ -53,14 +53,7 @@ impl OpaquePool { } } -#[cfg(feature = "std")] -impl From for OpaquePool { - fn from(pool: crate::mem::GlobalPool) -> Self { - Self::new(pool) - } -} - -impl Clone for OpaquePool { +impl Clone for OpaqueMemory { fn clone(&self) -> Self { Self { inner: self.inner.clone_boxed(), @@ -68,16 +61,13 @@ impl Clone for OpaquePool { } } -impl Memory for OpaquePool { +impl Memory for OpaqueMemory { #[cfg_attr(test, mutants::skip)] // Trivial forwarder. fn reserve(&self, min_bytes: usize) -> crate::BytesBuf { self.reserve(min_bytes) } } -/// Compatibility alias for the former name of [`OpaquePool`]. -pub type OpaqueMemory = OpaquePool; - #[cfg_attr(coverage_nightly, coverage(off))] #[cfg(all(test, feature = "std"))] mod tests { @@ -89,23 +79,13 @@ mod tests { use super::*; use crate::mem::GlobalPool; - use crate::mem::testing::TransparentMemory; - - assert_impl_all!(OpaquePool: MemoryShared); - #[test] - fn opaque_pool_wraps_custom_provider() { - let pool = OpaquePool::new(TransparentMemory::new()); - - let buffer = pool.reserve(1024); - - assert_eq!(buffer.capacity(), 1024); - } + assert_impl_all!(OpaqueMemory: MemoryShared); #[test] fn wraps_inner() { let provider = GlobalPool::new(); - let memory = OpaquePool::new(provider); + let memory = OpaqueMemory::new(provider); let builder = memory.reserve(1024); assert!(builder.capacity() >= 1024); @@ -114,7 +94,7 @@ mod tests { #[test] fn memory_trait() { let provider = GlobalPool::new(); - let memory = OpaquePool::new(provider); + let memory = OpaqueMemory::new(provider); // Call reserve via the Memory trait to verify the impl block let builder = Memory::reserve(&memory, 1024); @@ -123,7 +103,7 @@ mod tests { #[test] fn relocate_does_not_break_reservation() { - let mut memory = OpaquePool::new(GlobalPool::new()); + let mut memory = OpaqueMemory::new(GlobalPool::new()); let affinities = pinned_affinities(&[2]); memory.relocate(Some(affinities[0]), affinities[1]); @@ -156,7 +136,7 @@ mod tests { } let relocated = Arc::new(AtomicUsize::new(0)); - let mut memory = OpaquePool::new(TrackingMemory { + let mut memory = OpaqueMemory::new(TrackingMemory { relocated: Arc::clone(&relocated), inner: GlobalPool::new(), }); @@ -169,7 +149,7 @@ mod tests { #[test] fn clone_is_usable_independently() { - let memory = OpaquePool::new(GlobalPool::new()); + let memory = OpaqueMemory::new(GlobalPool::new()); let mut clone = memory.clone(); // Relocating the clone must leave both the clone and the original usable. diff --git a/crates/fetch/Cargo.toml b/crates/fetch/Cargo.toml index 07b078507..2743edf96 100644 --- a/crates/fetch/Cargo.toml +++ b/crates/fetch/Cargo.toml @@ -21,7 +21,7 @@ repository = "https://github.com/microsoft/oxidizer/tree/main/crates/fetch" allowed_external_types = [ "bytesbuf::mem::has_memory::HasMemory", "bytesbuf::mem::memory::Memory", - "bytesbuf::mem::opaque_memory::OpaquePool", + "bytesbuf::mem::opaque_memory::OpaqueMemory", "data_privacy::redaction_engine::RedactionEngine", "fetch_options::*", "fetch_tls::*", diff --git a/crates/fetch/examples/http_client_app.rs b/crates/fetch/examples/http_client_app.rs index fd66cb417..90c301d92 100644 --- a/crates/fetch/examples/http_client_app.rs +++ b/crates/fetch/examples/http_client_app.rs @@ -6,7 +6,7 @@ use std::sync::Arc; -use bytesbuf::mem::{GlobalPool, OpaquePool}; +use bytesbuf::mem::{GlobalPool, OpaqueMemory}; use fetch::HttpClient; use fetch::tls::TlsOptions; use ohno::ErrorExt; @@ -18,7 +18,7 @@ use tick::Clock; #[fundle::bundle] struct App { clock: Clock, - global_pool: OpaquePool, + global_pool: OpaqueMemory, client: HttpClient, } @@ -42,7 +42,7 @@ async fn main() -> Result<(), ohno::AppError> { // Initialize and set up the App instance; fundle ensures all fields are properly constructed. let app = App::builder() .clock(|_| Clock::new_tokio()) - .global_pool(|_| OpaquePool::new(GlobalPool::new())) + .global_pool(|_| OpaqueMemory::new(GlobalPool::new())) .client({ let meter_provider = meter_provider.clone(); move |x| { diff --git a/crates/fetch/examples/http_client_custom.rs b/crates/fetch/examples/http_client_custom.rs index 456dec40b..725aacf24 100644 --- a/crates/fetch/examples/http_client_custom.rs +++ b/crates/fetch/examples/http_client_custom.rs @@ -4,7 +4,7 @@ //! Plugs a custom `EchoHandler` into [`fetch::custom::create_builder`] as the transport //! handler. Every request's body is returned verbatim in the response. -use bytesbuf::mem::{GlobalPool, OpaquePool}; +use bytesbuf::mem::{GlobalPool, OpaqueMemory}; use fetch::custom::{CustomDeps, Isolation, create_builder}; use fetch::{HttpRequest, HttpResponse, HttpResponseBuilder}; use http::StatusCode; @@ -16,7 +16,7 @@ use tick::Clock; async fn main() -> Result<(), ohno::AppError> { let deps = CustomDeps { clock: Clock::new_tokio(), - global_pool: OpaquePool::new(GlobalPool::new()), + global_pool: OpaqueMemory::new(GlobalPool::new()), extras: (), }; diff --git a/crates/fetch/src/custom.rs b/crates/fetch/src/custom.rs index 2a5309f05..ffb2b9849 100644 --- a/crates/fetch/src/custom.rs +++ b/crates/fetch/src/custom.rs @@ -17,7 +17,7 @@ use std::borrow::Cow; use std::fmt::Debug; use std::sync::Arc; -use bytesbuf::mem::OpaquePool; +use bytesbuf::mem::OpaqueMemory; use http_extensions::{HttpBodyBuilder, RequestHandler}; use opentelemetry::metrics::Meter; use thread_aware::{PerCore, ThreadAware, unaware}; @@ -53,7 +53,7 @@ where /// Clock for timing operations and timeouts. pub clock: Clock, /// Memory pool for usage-neutral memory allocations. - pub global_pool: OpaquePool, + pub global_pool: OpaqueMemory, /// Extra dependencies forwarded verbatim to [`CustomContext::extras`]. pub extras: Extras, } @@ -242,7 +242,7 @@ pub(crate) struct Transport { name: Cow<'static, str>, inner: thread_aware::Arc, clock: Clock, - global_pool: OpaquePool, + global_pool: OpaqueMemory, isolation: Isolation, } @@ -278,7 +278,7 @@ impl Debug for Transport { } } -pub(crate) fn create_body_builder(pool: &OpaquePool, clock: &Clock, options: &ClientOptions) -> HttpBodyBuilder { +pub(crate) fn create_body_builder(pool: &OpaqueMemory, clock: &Clock, options: &ClientOptions) -> HttpBodyBuilder { HttpBodyBuilder::new(pool.clone(), clock).with_options(options.response_body_options) } @@ -289,7 +289,7 @@ mod tests { use std::sync::atomic::{AtomicUsize, Ordering}; use bytesbuf::BytesBuf; - use bytesbuf::mem::{GlobalPool, Memory, OpaquePool}; + use bytesbuf::mem::{GlobalPool, Memory, OpaqueMemory}; use http::StatusCode; use http_extensions::FakeHandler; use thread_aware::{ThreadAware, unaware}; @@ -303,7 +303,7 @@ mod tests { fn custom_deps() -> CustomDeps { CustomDeps { clock: FakeDeps::default().clock, - global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: (), } } @@ -326,10 +326,10 @@ mod tests { #[cfg_attr(miri, ignore)] #[tokio::test] - async fn custom_deps_accept_custom_opaque_pool() { + async fn custom_deps_accept_custom_opaque_memory() { let deps = CustomDeps { clock: FakeDeps::default().clock, - global_pool: OpaquePool::new(CustomMemory { inner: GlobalPool::new() }), + global_pool: OpaqueMemory::new(CustomMemory { inner: GlobalPool::new() }), extras: (), }; @@ -380,7 +380,7 @@ mod tests { let counter = Arc::new(AtomicUsize::new(0)); let deps = CustomDeps { clock: FakeDeps::default().clock, - global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: unaware(Arc::clone(&counter)), }; diff --git a/crates/fetch/src/fake.rs b/crates/fetch/src/fake.rs index 7d52a7c47..7923fcf8d 100644 --- a/crates/fetch/src/fake.rs +++ b/crates/fetch/src/fake.rs @@ -84,7 +84,7 @@ impl HttpClient { Isolation::Shared, CustomDeps { clock: deps.clock, - global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: handler, }, ) diff --git a/crates/fetch/src/tokio.rs b/crates/fetch/src/tokio.rs index 493053bde..754b3fa7f 100644 --- a/crates/fetch/src/tokio.rs +++ b/crates/fetch/src/tokio.rs @@ -33,7 +33,7 @@ pub struct TokioDeps { /// Clock for timing operations and timeouts. pub clock: Clock, /// Memory pool for usage-neutral memory allocations. - pub global_pool: bytesbuf::mem::OpaquePool, + pub global_pool: bytesbuf::mem::OpaqueMemory, } impl Default for TokioDeps { @@ -47,7 +47,7 @@ impl TokioDeps { #[must_use] pub fn with_clock(clock: &Clock) -> Self { Self { - global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), clock: clock.clone(), } } diff --git a/crates/fetch/tests/telemetry_scope.rs b/crates/fetch/tests/telemetry_scope.rs index cbecafa24..4dfca2287 100644 --- a/crates/fetch/tests/telemetry_scope.rs +++ b/crates/fetch/tests/telemetry_scope.rs @@ -94,7 +94,7 @@ async fn custom_transport_scope_attribute() { let deps = CustomDeps { clock: Clock::new_frozen(), - global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; @@ -128,7 +128,7 @@ async fn custom_transport_instrument_inherits_scope() { let deps = CustomDeps { clock: Clock::new_frozen(), - global_pool: bytesbuf::mem::OpaquePool::new(bytesbuf::mem::GlobalPool::new()), + global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; diff --git a/crates/fetch_hyper/src/builder.rs b/crates/fetch_hyper/src/builder.rs index db7803ac5..c7547b3a9 100644 --- a/crates/fetch_hyper/src/builder.rs +++ b/crates/fetch_hyper/src/builder.rs @@ -8,7 +8,7 @@ use std::fmt; use std::marker::PhantomData; use anyspawn::Spawner; -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaqueMemory}; use fetch_options::{ConnectionIdleTimeout, ConnectionKeepAlive, ConnectionPoolOptions, Http2Options, PoolIndex, TransportOptions}; use fetch_tls::TlsBackend; use http::Version; @@ -205,7 +205,7 @@ where let body_builder = self .body_builder .clone() - .unwrap_or_else(|| HttpBodyBuilder::new(GlobalPool::new(), &self.clock)); + .unwrap_or_else(|| HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &self.clock)); HyperTransport::new(build_hyper_handler(self, tls, body_builder, &meter).into_dynamic()) } diff --git a/crates/http_extensions/Cargo.toml b/crates/http_extensions/Cargo.toml index 152374fc5..19b9536e1 100644 --- a/crates/http_extensions/Cargo.toml +++ b/crates/http_extensions/Cargo.toml @@ -28,7 +28,7 @@ allowed_external_types = [ "bytesbuf::mem::has_memory::HasMemory", "bytesbuf::mem::memory::Memory", "bytesbuf::mem::memory_shared::MemoryShared", - "bytesbuf::mem::opaque_memory::OpaquePool", + "bytesbuf::mem::opaque_memory::OpaqueMemory", "bytesbuf::view::BytesView", "http::*", "http_body::Body", diff --git a/crates/http_extensions/examples/custom_client.rs b/crates/http_extensions/examples/custom_client.rs index 0663b4124..dce546605 100644 --- a/crates/http_extensions/examples/custom_client.rs +++ b/crates/http_extensions/examples/custom_client.rs @@ -6,7 +6,7 @@ //! This example demonstrates how to create a simple HTTP client that just echoes back the //! data it receives. -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaqueMemory}; use http_extensions::{HttpBodyBuilder, HttpRequest, HttpRequestBuilderExt, HttpResponse, HttpResponseBuilder, StatusExt}; use layered::Service; use tick::Clock; @@ -46,7 +46,7 @@ impl AsRef for CustomClient { impl Default for CustomClient { fn default() -> Self { Self { - builder: HttpBodyBuilder::new(GlobalPool::new(), &Clock::new_tokio()), + builder: HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &Clock::new_tokio()), } } } diff --git a/crates/http_extensions/examples/custom_server.rs b/crates/http_extensions/examples/custom_server.rs index 1b75e64cf..e698dc5ae 100644 --- a/crates/http_extensions/examples/custom_server.rs +++ b/crates/http_extensions/examples/custom_server.rs @@ -11,7 +11,7 @@ use std::sync::Arc; use std::time::Duration; use bytesbuf::BytesView; -use bytesbuf::mem::GlobalPool; +use bytesbuf::mem::{GlobalPool, OpaqueMemory}; use futures::TryStreamExt; use http::Request; use http_body_util::BodyExt; @@ -31,7 +31,7 @@ async fn main() -> Result<(), ohno::AppError> { let clock = Clock::new_tokio(); // In a real application, the application framework would provide the global memory pool. - let body_builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let body_builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); let body_builder_clone = body_builder.clone(); // Define an execution stack of middleware diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 84647fac2..04ad5cfe0 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -3,7 +3,7 @@ #[cfg(any(feature = "test-util", test))] use bytesbuf::mem::GlobalPool; -use bytesbuf::mem::{HasMemory, Memory, MemoryShared, OpaquePool}; +use bytesbuf::mem::{HasMemory, Memory, MemoryShared, OpaqueMemory}; use bytesbuf::{BytesBuf, BytesView}; use futures::{Stream, TryStreamExt}; use http_body::{Body, Frame}; @@ -46,7 +46,7 @@ use crate::{HttpError, Result}; /// With the `test-util` feature enabled, you can create a test instance using `HttpBodyBuilder::new_fake()`. #[derive(Debug, Clone, ThreadAware)] pub struct HttpBodyBuilder { - memory: OpaquePool, + memory: OpaqueMemory, clock: Clock, pub(super) options: HttpBodyOptions, } @@ -73,18 +73,16 @@ impl HttpBodyBuilder { #[cfg(any(feature = "test-util", test))] #[must_use] pub fn new_fake() -> Self { - Self::new(GlobalPool::new(), &Clock::new_frozen()) + Self::new(OpaqueMemory::new(GlobalPool::new()), &Clock::new_frozen()) } /// Creates a new instance of [`HttpBodyBuilder`]. /// - /// Accepts either an [`OpaquePool`] or the default [`GlobalPool`]. - /// Use [`with_custom_memory`](Self::with_custom_memory) to wrap another [`MemoryShared`] - /// implementation. + /// The provided type-erased memory can wrap any [`MemoryShared`] implementation. #[must_use] - pub fn new(memory: impl Into, clock: &Clock) -> Self { + pub fn new(memory: OpaqueMemory, clock: &Clock) -> Self { Self { - memory: memory.into(), + memory, clock: clock.clone(), options: HttpBodyOptions::default(), } @@ -98,7 +96,7 @@ impl HttpBodyBuilder { /// relocated along with it. #[must_use] pub fn with_custom_memory(memory: impl MemoryShared, clock: &Clock) -> Self { - Self::new(OpaquePool::new(memory), clock) + Self::new(OpaqueMemory::new(memory), clock) } /// Sets default [`HttpBodyOptions`] for all bodies created by this builder. @@ -385,11 +383,11 @@ mod tests { } #[test] - fn new_accepts_opaque_pool_with_custom_memory() { + fn new_accepts_opaque_memory_with_custom_provider() { let clock = Clock::new_frozen(); - let pool = OpaquePool::new(TransparentMemory::new()); + let memory = OpaqueMemory::new(TransparentMemory::new()); - let builder = HttpBodyBuilder::new(pool, &clock); + let builder = HttpBodyBuilder::new(memory, &clock); let body = builder.text("custom pool"); assert_eq!(body.content_length(), Some(11)); @@ -398,7 +396,7 @@ mod tests { #[test] fn new_with_global_memory() { let clock = Clock::new_frozen(); - let memory = GlobalPool::new(); + let memory = OpaqueMemory::new(GlobalPool::new()); let builder = HttpBodyBuilder::new(memory, &clock); let body = builder.text("test"); assert_eq!(body.content_length(), Some(4)); @@ -539,7 +537,7 @@ mod tests { #[test] fn stream_with_timeout_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); let chunks: Vec> = [b"hello " as &[u8], b"world"] .iter() .map(|c| Ok(BytesView::copied_from_slice(c, &builder))) @@ -615,7 +613,7 @@ mod tests { fn builder_merges_per_call_options_with_defaults() { let clock = Clock::new_frozen(); let builder_options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(builder_options); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock).with_options(builder_options); // Per-call options override the builder-level default. let per_call = HttpBodyOptions::default().timeout(Duration::from_secs(5)); diff --git a/crates/http_extensions/src/body/timeout_body.rs b/crates/http_extensions/src/body/timeout_body.rs index c3f59fadc..cea4358a8 100644 --- a/crates/http_extensions/src/body/timeout_body.rs +++ b/crates/http_extensions/src/body/timeout_body.rs @@ -96,7 +96,7 @@ mod tests { use std::time::Duration; use bytesbuf::BytesView; - use bytesbuf::mem::GlobalPool; + use bytesbuf::mem::{GlobalPool, OpaqueMemory}; use futures::executor::block_on; use http_body::{Body, Frame}; use tick::ClockControl; @@ -107,7 +107,7 @@ mod tests { #[test] fn stream_body_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); // Stream yields data immediately — well within the generous timeout, // exercising the TimeoutBody happy path via stream with timeout options. @@ -121,7 +121,7 @@ mod tests { #[test] fn stream_body_times_out_when_pending() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); // A body that never yields data. let options = HttpBodyOptions::default().timeout(Duration::from_millis(100)); @@ -136,7 +136,8 @@ mod tests { #[test] fn body_timeout_chains_with_buffer_limit() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); + let builder = + HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); assert_eq!(builder.options, HttpBodyOptions::default().buffer_limit(1024)); @@ -160,7 +161,7 @@ mod tests { #[test] fn size_hint_delegates_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); // Full body has an exact size hint; verify it passes through TimeoutBody. let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); @@ -176,7 +177,7 @@ mod tests { #[test] fn is_end_stream_true_when_inner_is_empty() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Empty::new(), &options); @@ -186,7 +187,7 @@ mod tests { #[test] fn is_end_stream_false_when_inner_has_data() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Full::new(BytesView::copied_from_slice(b"data", &builder)), &options); @@ -196,7 +197,7 @@ mod tests { #[test] fn poll_frame_returns_data_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); let body = builder.body( @@ -210,7 +211,7 @@ mod tests { #[test] fn poll_frame_times_out_when_pending_with_short_timeout() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); // A body that never yields data with a very short timeout. let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); @@ -227,7 +228,7 @@ mod tests { fn poll_frame_returns_data_even_when_clock_advanced_past_timeout() { let control = ClockControl::new(); let clock = control.to_clock(); - let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); // Use a body that has data immediately available (Full is always ready). let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); From 81c97f5188d9c4cb86258ca93ff21fd201d77e36 Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 09:59:20 +0200 Subject: [PATCH 08/13] review: simplify opaque memory APIs Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/fetch/examples/http_client_custom.rs | 2 +- crates/fetch/src/custom.rs | 16 +++++++------- crates/fetch/src/fake.rs | 2 +- crates/fetch/src/tokio.rs | 4 ++-- crates/fetch/tests/telemetry_scope.rs | 4 ++-- .../benches/http_request_builder.rs | 3 ++- crates/http_extensions/src/body/builder.rs | 22 +------------------ 7 files changed, 17 insertions(+), 36 deletions(-) diff --git a/crates/fetch/examples/http_client_custom.rs b/crates/fetch/examples/http_client_custom.rs index 725aacf24..1889cd426 100644 --- a/crates/fetch/examples/http_client_custom.rs +++ b/crates/fetch/examples/http_client_custom.rs @@ -16,7 +16,7 @@ use tick::Clock; async fn main() -> Result<(), ohno::AppError> { let deps = CustomDeps { clock: Clock::new_tokio(), - global_pool: OpaqueMemory::new(GlobalPool::new()), + memory: OpaqueMemory::new(GlobalPool::new()), extras: (), }; diff --git a/crates/fetch/src/custom.rs b/crates/fetch/src/custom.rs index ffb2b9849..c850a4019 100644 --- a/crates/fetch/src/custom.rs +++ b/crates/fetch/src/custom.rs @@ -53,7 +53,7 @@ where /// Clock for timing operations and timeouts. pub clock: Clock, /// Memory pool for usage-neutral memory allocations. - pub global_pool: OpaqueMemory, + pub memory: OpaqueMemory, /// Extra dependencies forwarded verbatim to [`CustomContext::extras`]. pub extras: Extras, } @@ -210,12 +210,12 @@ impl HttpClient { runtime_name: runtime.into(), name: transport.into(), clock: deps.clock.clone(), - global_pool: deps.global_pool.clone(), + memory: deps.memory.clone(), isolation, inner: thread_aware::Arc::new_with((deps, unaware(factory)), |(deps, factory)| { Arc::new(move |options, meter, pool_index| { let context = CustomContext { - body_builder: create_body_builder(&deps.global_pool, &deps.clock, &options), + body_builder: create_body_builder(&deps.memory, &deps.clock, &options), clock: deps.clock.clone(), pool_index, extras: deps.extras.clone(), @@ -242,7 +242,7 @@ pub(crate) struct Transport { name: Cow<'static, str>, inner: thread_aware::Arc, clock: Clock, - global_pool: OpaqueMemory, + memory: OpaqueMemory, isolation: Isolation, } @@ -268,7 +268,7 @@ impl Transport { } pub(crate) fn create_body_builder(&self, options: &ClientOptions) -> HttpBodyBuilder { - create_body_builder(&self.global_pool, &self.clock, options) + create_body_builder(&self.memory, &self.clock, options) } } @@ -303,7 +303,7 @@ mod tests { fn custom_deps() -> CustomDeps { CustomDeps { clock: FakeDeps::default().clock, - global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), + memory: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: (), } } @@ -329,7 +329,7 @@ mod tests { async fn custom_deps_accept_custom_opaque_memory() { let deps = CustomDeps { clock: FakeDeps::default().clock, - global_pool: OpaqueMemory::new(CustomMemory { inner: GlobalPool::new() }), + memory: OpaqueMemory::new(CustomMemory { inner: GlobalPool::new() }), extras: (), }; @@ -380,7 +380,7 @@ mod tests { let counter = Arc::new(AtomicUsize::new(0)); let deps = CustomDeps { clock: FakeDeps::default().clock, - global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), + memory: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: unaware(Arc::clone(&counter)), }; diff --git a/crates/fetch/src/fake.rs b/crates/fetch/src/fake.rs index 7923fcf8d..6aecf5023 100644 --- a/crates/fetch/src/fake.rs +++ b/crates/fetch/src/fake.rs @@ -84,7 +84,7 @@ impl HttpClient { Isolation::Shared, CustomDeps { clock: deps.clock, - global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), + memory: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: handler, }, ) diff --git a/crates/fetch/src/tokio.rs b/crates/fetch/src/tokio.rs index 754b3fa7f..74bb0f6a4 100644 --- a/crates/fetch/src/tokio.rs +++ b/crates/fetch/src/tokio.rs @@ -64,7 +64,7 @@ impl HttpClient { pub fn builder_tokio(deps: impl Into) -> HttpClientBuilder { let deps = deps.into(); let clock = deps.clock.clone(); - let global_pool = deps.global_pool.clone(); + let memory = deps.global_pool.clone(); // Re-layer on top of the in-crate `builder_custom_internal` path: the // full `TokioDeps` rides through `CustomDeps::extras` so that the @@ -77,7 +77,7 @@ impl HttpClient { Isolation::Shared, CustomDeps { clock, - global_pool, + memory, extras: deps, }, ) diff --git a/crates/fetch/tests/telemetry_scope.rs b/crates/fetch/tests/telemetry_scope.rs index 4dfca2287..685f988cc 100644 --- a/crates/fetch/tests/telemetry_scope.rs +++ b/crates/fetch/tests/telemetry_scope.rs @@ -94,7 +94,7 @@ async fn custom_transport_scope_attribute() { let deps = CustomDeps { clock: Clock::new_frozen(), - global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), + memory: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; @@ -128,7 +128,7 @@ async fn custom_transport_instrument_inherits_scope() { let deps = CustomDeps { clock: Clock::new_frozen(), - global_pool: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), + memory: bytesbuf::mem::OpaqueMemory::new(bytesbuf::mem::GlobalPool::new()), extras: (), }; diff --git a/crates/http_extensions/benches/http_request_builder.rs b/crates/http_extensions/benches/http_request_builder.rs index cbf6d28b5..97d9acea8 100644 --- a/crates/http_extensions/benches/http_request_builder.rs +++ b/crates/http_extensions/benches/http_request_builder.rs @@ -11,6 +11,7 @@ use alloc_tracker::{Allocator, Session}; use benchmarking::time_sample; +use bytesbuf::mem::OpaqueMemory; use bytesbuf::mem::testing::TransparentMemory; use criterion::{Criterion, criterion_group, criterion_main}; use http::header::CONTENT_TYPE; @@ -153,7 +154,7 @@ fn entry(c: &mut Criterion) { // Use TransparentMemory instead of GlobalPool so that every reserve() call from the // serde_json writer results in a real heap allocation. This makes alloc_tracker report // the true number of memory reservations, which GlobalPool would otherwise absorb. - let transparent_body_builder = HttpBodyBuilder::with_custom_memory(TransparentMemory::new(), &tick::Clock::new_frozen()); + let transparent_body_builder = HttpBodyBuilder::new(OpaqueMemory::new(TransparentMemory::new()), &tick::Clock::new_frozen()); let operation = session.operation("json_body_large_transparent"); group.bench_function("json_body_large_transparent", |b| { b.iter_custom(|iters| { diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 04ad5cfe0..3d323e209 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -88,17 +88,6 @@ impl HttpBodyBuilder { } } - /// Creates a new instance of [`HttpBodyBuilder`] with custom memory. - /// - /// The provided memory provider is type-erased and used in place of the global per-thread - /// memory used by [`HttpBodyBuilder::new`]. It remains thread-aware: when the builder is moved - /// between threads via a thread-aware runtime mechanism, the provider's thread-affine state is - /// relocated along with it. - #[must_use] - pub fn with_custom_memory(memory: impl MemoryShared, clock: &Clock) -> Self { - Self::new(OpaqueMemory::new(memory), clock) - } - /// Sets default [`HttpBodyOptions`] for all bodies created by this builder. /// /// Per-call options passed to [`body`](Self::body) or [`stream`](Self::stream) are @@ -405,15 +394,6 @@ mod tests { let _clock: &Clock = builder.as_ref(); } - #[test] - fn with_custom_memory() { - let clock = Clock::new_frozen(); - let builder = HttpBodyBuilder::with_custom_memory(TransparentMemory::new(), &clock); - let body = builder.text("hello"); - let data = BytesView::try_from(body).unwrap(); - assert_eq!(data.len(), 5); - } - #[test] fn with_options_sets_buffer_limit() { let options = HttpBodyOptions::default().buffer_limit(1024); @@ -596,7 +576,7 @@ mod tests { ); let clock = Clock::new_frozen(); - let builder = HttpBodyBuilder::with_custom_memory(TransparentMemory::new(), &clock); + let builder = HttpBodyBuilder::new(OpaqueMemory::new(TransparentMemory::new()), &clock); let body = builder.json(&payload).unwrap(); let bytes_view = body.into_bytes_no_buffering().unwrap(); From 6a0d2069ceb2cb293b16f08ea3151e04521ab5dc Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 11:16:24 +0200 Subject: [PATCH 09/13] review: avoid double-wrapping opaque memory Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/fetch_hyper/src/builder.rs | 4 +-- .../benches/http_request_builder.rs | 3 +- .../http_extensions/examples/custom_client.rs | 4 +-- .../http_extensions/examples/custom_server.rs | 4 +-- crates/http_extensions/src/body/builder.rs | 30 +++++++++++++------ .../http_extensions/src/body/timeout_body.rs | 21 +++++++------ 6 files changed, 38 insertions(+), 28 deletions(-) diff --git a/crates/fetch_hyper/src/builder.rs b/crates/fetch_hyper/src/builder.rs index c7547b3a9..db7803ac5 100644 --- a/crates/fetch_hyper/src/builder.rs +++ b/crates/fetch_hyper/src/builder.rs @@ -8,7 +8,7 @@ use std::fmt; use std::marker::PhantomData; use anyspawn::Spawner; -use bytesbuf::mem::{GlobalPool, OpaqueMemory}; +use bytesbuf::mem::GlobalPool; use fetch_options::{ConnectionIdleTimeout, ConnectionKeepAlive, ConnectionPoolOptions, Http2Options, PoolIndex, TransportOptions}; use fetch_tls::TlsBackend; use http::Version; @@ -205,7 +205,7 @@ where let body_builder = self .body_builder .clone() - .unwrap_or_else(|| HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &self.clock)); + .unwrap_or_else(|| HttpBodyBuilder::new(GlobalPool::new(), &self.clock)); HyperTransport::new(build_hyper_handler(self, tls, body_builder, &meter).into_dynamic()) } diff --git a/crates/http_extensions/benches/http_request_builder.rs b/crates/http_extensions/benches/http_request_builder.rs index 97d9acea8..6be54e5a5 100644 --- a/crates/http_extensions/benches/http_request_builder.rs +++ b/crates/http_extensions/benches/http_request_builder.rs @@ -11,7 +11,6 @@ use alloc_tracker::{Allocator, Session}; use benchmarking::time_sample; -use bytesbuf::mem::OpaqueMemory; use bytesbuf::mem::testing::TransparentMemory; use criterion::{Criterion, criterion_group, criterion_main}; use http::header::CONTENT_TYPE; @@ -154,7 +153,7 @@ fn entry(c: &mut Criterion) { // Use TransparentMemory instead of GlobalPool so that every reserve() call from the // serde_json writer results in a real heap allocation. This makes alloc_tracker report // the true number of memory reservations, which GlobalPool would otherwise absorb. - let transparent_body_builder = HttpBodyBuilder::new(OpaqueMemory::new(TransparentMemory::new()), &tick::Clock::new_frozen()); + let transparent_body_builder = HttpBodyBuilder::new(TransparentMemory::new(), &tick::Clock::new_frozen()); let operation = session.operation("json_body_large_transparent"); group.bench_function("json_body_large_transparent", |b| { b.iter_custom(|iters| { diff --git a/crates/http_extensions/examples/custom_client.rs b/crates/http_extensions/examples/custom_client.rs index dce546605..0663b4124 100644 --- a/crates/http_extensions/examples/custom_client.rs +++ b/crates/http_extensions/examples/custom_client.rs @@ -6,7 +6,7 @@ //! This example demonstrates how to create a simple HTTP client that just echoes back the //! data it receives. -use bytesbuf::mem::{GlobalPool, OpaqueMemory}; +use bytesbuf::mem::GlobalPool; use http_extensions::{HttpBodyBuilder, HttpRequest, HttpRequestBuilderExt, HttpResponse, HttpResponseBuilder, StatusExt}; use layered::Service; use tick::Clock; @@ -46,7 +46,7 @@ impl AsRef for CustomClient { impl Default for CustomClient { fn default() -> Self { Self { - builder: HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &Clock::new_tokio()), + builder: HttpBodyBuilder::new(GlobalPool::new(), &Clock::new_tokio()), } } } diff --git a/crates/http_extensions/examples/custom_server.rs b/crates/http_extensions/examples/custom_server.rs index e698dc5ae..1b75e64cf 100644 --- a/crates/http_extensions/examples/custom_server.rs +++ b/crates/http_extensions/examples/custom_server.rs @@ -11,7 +11,7 @@ use std::sync::Arc; use std::time::Duration; use bytesbuf::BytesView; -use bytesbuf::mem::{GlobalPool, OpaqueMemory}; +use bytesbuf::mem::GlobalPool; use futures::TryStreamExt; use http::Request; use http_body_util::BodyExt; @@ -31,7 +31,7 @@ async fn main() -> Result<(), ohno::AppError> { let clock = Clock::new_tokio(); // In a real application, the application framework would provide the global memory pool. - let body_builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let body_builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let body_builder_clone = body_builder.clone(); // Define an execution stack of middleware diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 3d323e209..83ee56b11 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -1,6 +1,8 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT License. +use std::any::{Any, TypeId}; + #[cfg(any(feature = "test-util", test))] use bytesbuf::mem::GlobalPool; use bytesbuf::mem::{HasMemory, Memory, MemoryShared, OpaqueMemory}; @@ -51,6 +53,17 @@ pub struct HttpBodyBuilder { pub(super) options: HttpBodyOptions, } +fn into_opaque_memory(memory: M) -> OpaqueMemory { + if TypeId::of::() == TypeId::of::() { + let memory: Box = Box::new(memory); + *memory + .downcast::() + .expect("the concrete type was verified as OpaqueMemory above") + } else { + OpaqueMemory::new(memory) + } +} + impl HttpBodyBuilder { /// Creates a test-friendly [`HttpBodyBuilder`] instance. /// @@ -73,16 +86,16 @@ impl HttpBodyBuilder { #[cfg(any(feature = "test-util", test))] #[must_use] pub fn new_fake() -> Self { - Self::new(OpaqueMemory::new(GlobalPool::new()), &Clock::new_frozen()) + Self::new(GlobalPool::new(), &Clock::new_frozen()) } /// Creates a new instance of [`HttpBodyBuilder`]. /// - /// The provided type-erased memory can wrap any [`MemoryShared`] implementation. + /// The provided memory is type-erased unless it is already an [`OpaqueMemory`]. #[must_use] - pub fn new(memory: OpaqueMemory, clock: &Clock) -> Self { + pub fn new(memory: impl MemoryShared, clock: &Clock) -> Self { Self { - memory, + memory: into_opaque_memory(memory), clock: clock.clone(), options: HttpBodyOptions::default(), } @@ -385,8 +398,7 @@ mod tests { #[test] fn new_with_global_memory() { let clock = Clock::new_frozen(); - let memory = OpaqueMemory::new(GlobalPool::new()); - let builder = HttpBodyBuilder::new(memory, &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let body = builder.text("test"); assert_eq!(body.content_length(), Some(4)); @@ -517,7 +529,7 @@ mod tests { #[test] fn stream_with_timeout_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let chunks: Vec> = [b"hello " as &[u8], b"world"] .iter() .map(|c| Ok(BytesView::copied_from_slice(c, &builder))) @@ -576,7 +588,7 @@ mod tests { ); let clock = Clock::new_frozen(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(TransparentMemory::new()), &clock); + let builder = HttpBodyBuilder::new(TransparentMemory::new(), &clock); let body = builder.json(&payload).unwrap(); let bytes_view = body.into_bytes_no_buffering().unwrap(); @@ -593,7 +605,7 @@ mod tests { fn builder_merges_per_call_options_with_defaults() { let clock = Clock::new_frozen(); let builder_options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock).with_options(builder_options); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(builder_options); // Per-call options override the builder-level default. let per_call = HttpBodyOptions::default().timeout(Duration::from_secs(5)); diff --git a/crates/http_extensions/src/body/timeout_body.rs b/crates/http_extensions/src/body/timeout_body.rs index cea4358a8..c3f59fadc 100644 --- a/crates/http_extensions/src/body/timeout_body.rs +++ b/crates/http_extensions/src/body/timeout_body.rs @@ -96,7 +96,7 @@ mod tests { use std::time::Duration; use bytesbuf::BytesView; - use bytesbuf::mem::{GlobalPool, OpaqueMemory}; + use bytesbuf::mem::GlobalPool; use futures::executor::block_on; use http_body::{Body, Frame}; use tick::ClockControl; @@ -107,7 +107,7 @@ mod tests { #[test] fn stream_body_returns_data_before_timeout() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // Stream yields data immediately — well within the generous timeout, // exercising the TimeoutBody happy path via stream with timeout options. @@ -121,7 +121,7 @@ mod tests { #[test] fn stream_body_times_out_when_pending() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // A body that never yields data. let options = HttpBodyOptions::default().timeout(Duration::from_millis(100)); @@ -136,8 +136,7 @@ mod tests { #[test] fn body_timeout_chains_with_buffer_limit() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = - HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock).with_options(HttpBodyOptions::default().buffer_limit(1024)); assert_eq!(builder.options, HttpBodyOptions::default().buffer_limit(1024)); @@ -161,7 +160,7 @@ mod tests { #[test] fn size_hint_delegates_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // Full body has an exact size hint; verify it passes through TimeoutBody. let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); @@ -177,7 +176,7 @@ mod tests { #[test] fn is_end_stream_true_when_inner_is_empty() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Empty::new(), &options); @@ -187,7 +186,7 @@ mod tests { #[test] fn is_end_stream_false_when_inner_has_data() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(1)); let body = builder.body(http_body_util::Full::new(BytesView::copied_from_slice(b"data", &builder)), &options); @@ -197,7 +196,7 @@ mod tests { #[test] fn poll_frame_returns_data_through_timeout_body() { let clock = ClockControl::new().to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); let options = HttpBodyOptions::default().timeout(Duration::from_secs(30)); let body = builder.body( @@ -211,7 +210,7 @@ mod tests { #[test] fn poll_frame_times_out_when_pending_with_short_timeout() { let clock = ClockControl::new().auto_advance_timers(true).to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // A body that never yields data with a very short timeout. let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); @@ -228,7 +227,7 @@ mod tests { fn poll_frame_returns_data_even_when_clock_advanced_past_timeout() { let control = ClockControl::new(); let clock = control.to_clock(); - let builder = HttpBodyBuilder::new(OpaqueMemory::new(GlobalPool::new()), &clock); + let builder = HttpBodyBuilder::new(GlobalPool::new(), &clock); // Use a body that has data immediately available (Full is always ready). let options = HttpBodyOptions::default().timeout(Duration::from_millis(1)); From ca604d3a1c7936772cf28e59d78e8ef41af58dae Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 11:22:07 +0200 Subject: [PATCH 10/13] review(bytesbuf): avoid nested opaque memory Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/opaque_memory.rs | 23 ++++++++++++++++++++-- crates/http_extensions/src/body/builder.rs | 15 +------------- 2 files changed, 22 insertions(+), 16 deletions(-) diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index 8c374c466..2d2e472cd 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -1,6 +1,8 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT License. +use core::any::{Any, TypeId}; + #[cfg(not(test))] use alloc::boxed::Box; @@ -25,8 +27,15 @@ pub struct OpaqueMemory { impl OpaqueMemory { /// Creates a new instance of the adapter. #[must_use] - pub fn new(inner: impl MemoryShared) -> Self { - Self { inner: Box::new(inner) } + pub fn new(inner: M) -> Self { + if TypeId::of::() == TypeId::of::() { + let inner: Box = Box::new(inner); + *inner + .downcast::() + .expect("the concrete type was verified as OpaqueMemory above") + } else { + Self { inner: Box::new(inner) } + } } /// Reserves at least `min_bytes` bytes of memory capacity. @@ -91,6 +100,16 @@ mod tests { assert!(builder.capacity() >= 1024); } + #[test] + fn accepts_existing_opaque_memory() { + let memory = OpaqueMemory::new(GlobalPool::new()); + let memory = OpaqueMemory::new(memory); + + let builder = memory.reserve(1024); + + assert!(builder.capacity() >= 1024); + } + #[test] fn memory_trait() { let provider = GlobalPool::new(); diff --git a/crates/http_extensions/src/body/builder.rs b/crates/http_extensions/src/body/builder.rs index 83ee56b11..bcb6f8b03 100644 --- a/crates/http_extensions/src/body/builder.rs +++ b/crates/http_extensions/src/body/builder.rs @@ -1,8 +1,6 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT License. -use std::any::{Any, TypeId}; - #[cfg(any(feature = "test-util", test))] use bytesbuf::mem::GlobalPool; use bytesbuf::mem::{HasMemory, Memory, MemoryShared, OpaqueMemory}; @@ -53,17 +51,6 @@ pub struct HttpBodyBuilder { pub(super) options: HttpBodyOptions, } -fn into_opaque_memory(memory: M) -> OpaqueMemory { - if TypeId::of::() == TypeId::of::() { - let memory: Box = Box::new(memory); - *memory - .downcast::() - .expect("the concrete type was verified as OpaqueMemory above") - } else { - OpaqueMemory::new(memory) - } -} - impl HttpBodyBuilder { /// Creates a test-friendly [`HttpBodyBuilder`] instance. /// @@ -95,7 +82,7 @@ impl HttpBodyBuilder { #[must_use] pub fn new(memory: impl MemoryShared, clock: &Clock) -> Self { Self { - memory: into_opaque_memory(memory), + memory: OpaqueMemory::new(memory), clock: clock.clone(), options: HttpBodyOptions::default(), } From 2876de09e0050073c9664d5dcab678672ee1fd60 Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 11:30:13 +0200 Subject: [PATCH 11/13] fix(bytesbuf): document opaque constructor panic Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/opaque_memory.rs | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index 2d2e472cd..9af85f62c 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -26,6 +26,11 @@ pub struct OpaqueMemory { impl OpaqueMemory { /// Creates a new instance of the adapter. + /// + /// # Panics + /// + /// Panics only if runtime type identification reports [`OpaqueMemory`] but downcasting the + /// same value to [`OpaqueMemory`] fails, which would indicate a standard library defect. #[must_use] pub fn new(inner: M) -> Self { if TypeId::of::() == TypeId::of::() { From 0090821168c3e7334fa13069bb1b74a76d600b3d Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 11:38:02 +0200 Subject: [PATCH 12/13] fix(bytesbuf): apply nightly import formatting Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/opaque_memory.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index 9af85f62c..c818bd375 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -1,10 +1,9 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT License. -use core::any::{Any, TypeId}; - #[cfg(not(test))] use alloc::boxed::Box; +use core::any::{Any, TypeId}; use thread_aware::ThreadAware; From ceac638a9e79cbe0eb40bee0b87c6c5af289a2ab Mon Sep 17 00:00:00 2001 From: Martin Tomka Date: Fri, 21 Aug 2026 11:48:43 +0200 Subject: [PATCH 13/13] fix(bytesbuf): use accepted downcast spelling Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: c98d886c-c9b6-42aa-ac1e-990eb589d5c0 --- crates/bytesbuf/src/mem/opaque_memory.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/bytesbuf/src/mem/opaque_memory.rs b/crates/bytesbuf/src/mem/opaque_memory.rs index c818bd375..4f699d6ba 100644 --- a/crates/bytesbuf/src/mem/opaque_memory.rs +++ b/crates/bytesbuf/src/mem/opaque_memory.rs @@ -28,8 +28,8 @@ impl OpaqueMemory { /// /// # Panics /// - /// Panics only if runtime type identification reports [`OpaqueMemory`] but downcasting the - /// same value to [`OpaqueMemory`] fails, which would indicate a standard library defect. + /// Panics only if runtime type identification reports [`OpaqueMemory`] but the downcast of + /// the same value to [`OpaqueMemory`] fails, which would indicate a standard library defect. #[must_use] pub fn new(inner: M) -> Self { if TypeId::of::() == TypeId::of::() {