toolkit_http/builder.rs
1use crate::config::{
2 ClientAuthConfig, HttpClientConfig, RateLimitConfig, RedirectConfig, RetryConfig, TlsConfig,
3 TlsRootConfig, TlsVersion, TransportSecurity,
4};
5use crate::error::HttpError;
6use crate::layers::{OtelLayer, RetryLayer, SecureRedirectPolicy, UserAgentLayer};
7use crate::response::ResponseBody;
8use crate::tls;
9use bytes::Bytes;
10use http::Response;
11use http_body_util::{BodyExt, Full};
12use hyper_rustls::HttpsConnector;
13use hyper_util::client::legacy::Client;
14use hyper_util::client::legacy::connect::HttpConnector;
15use hyper_util::rt::{TokioExecutor, TokioTimer};
16use std::time::Duration;
17use tower::buffer::Buffer;
18use tower::limit::ConcurrencyLimitLayer;
19use tower::load_shed::LoadShedLayer;
20use tower::timeout::TimeoutLayer;
21use tower::util::BoxCloneService;
22use tower::{ServiceBuilder, ServiceExt};
23use tower_http::decompression::DecompressionLayer;
24use tower_http::follow_redirect::FollowRedirectLayer;
25
26/// Type-erased inner service between layer composition steps in [`HttpClientBuilder::build`].
27type InnerService =
28 BoxCloneService<http::Request<Full<Bytes>>, http::Response<ResponseBody>, HttpError>;
29
30/// Builder for constructing an [`HttpClient`] with a layered tower middleware stack.
31pub struct HttpClientBuilder {
32 config: HttpClientConfig,
33 auth_layer: Option<Box<dyn FnOnce(InnerService) -> InnerService + Send>>,
34 metrics_layer: Option<Box<dyn FnOnce(InnerService) -> InnerService + Send>>,
35}
36
37impl HttpClientBuilder {
38 /// Create a new builder with default configuration
39 #[must_use]
40 pub fn new() -> Self {
41 Self {
42 config: HttpClientConfig::default(),
43 auth_layer: None,
44 metrics_layer: None,
45 }
46 }
47
48 /// Create a builder with a specific configuration
49 #[must_use]
50 pub fn with_config(config: HttpClientConfig) -> Self {
51 Self {
52 config,
53 auth_layer: None,
54 metrics_layer: None,
55 }
56 }
57
58 /// Set the per-request timeout
59 ///
60 /// This timeout applies to each individual HTTP request/attempt.
61 /// If retries are enabled, each retry attempt gets its own timeout.
62 #[must_use]
63 pub fn timeout(mut self, timeout: Duration) -> Self {
64 self.config.request_timeout = timeout;
65 self
66 }
67
68 /// Set the total timeout spanning all retry attempts
69 ///
70 /// When set, the entire operation (including all retries and backoff delays)
71 /// must complete within this duration. If the deadline is exceeded,
72 /// the request fails with `HttpError::DeadlineExceeded(total_timeout)`.
73 #[must_use]
74 pub fn total_timeout(mut self, timeout: Duration) -> Self {
75 self.config.total_timeout = Some(timeout);
76 self
77 }
78
79 /// Set the user agent string
80 #[must_use]
81 pub fn user_agent(mut self, user_agent: impl Into<String>) -> Self {
82 self.config.user_agent = user_agent.into();
83 self
84 }
85
86 /// Set the retry configuration
87 #[must_use]
88 pub fn retry(mut self, retry: Option<RetryConfig>) -> Self {
89 self.config.retry = retry;
90 self
91 }
92
93 /// Set the maximum response body size
94 #[must_use]
95 pub fn max_body_size(mut self, size: usize) -> Self {
96 self.config.max_body_size = size;
97 self
98 }
99
100 /// Set the concurrency-limit configuration.
101 ///
102 /// This installs a *concurrency* cap (max in-flight requests), not a
103 /// requests-per-second rate limit. `None` disables the limiter entirely.
104 /// `Some(cfg)` installs a load-shedding
105 /// [`ConcurrencyLimitLayer`](tower::limit::ConcurrencyLimitLayer) capped at
106 /// `cfg.max_concurrent_requests`; requests beyond the cap fail fast with
107 /// [`HttpError::Overloaded`](crate::error::HttpError::Overloaded) rather
108 /// than queueing. A `max_concurrent_requests` of `usize::MAX` is treated as
109 /// unlimited (the layer is skipped); `0` is clamped to `1` at
110 /// [`build`](Self::build) time so the client can never wedge shedding every
111 /// request.
112 #[must_use]
113 pub fn concurrency_limit(mut self, rate_limit: Option<RateLimitConfig>) -> Self {
114 self.config.rate_limit = rate_limit;
115 self
116 }
117
118 /// Set transport security mode
119 ///
120 /// Use `TransportSecurity::TlsOnly` to enforce HTTPS for all connections.
121 #[must_use]
122 pub fn transport(mut self, transport: TransportSecurity) -> Self {
123 self.config.transport = transport;
124 self
125 }
126
127 /// Deny insecure HTTP connections, enforcing TLS for all traffic
128 ///
129 /// Equivalent to `.transport(TransportSecurity::TlsOnly)`.
130 ///
131 /// Use this when TLS enforcement is required (e.g., production environments).
132 #[must_use]
133 pub fn deny_insecure_http(mut self) -> Self {
134 tracing::debug!(
135 target: "toolkit_http::security",
136 "deny_insecure_http() called - enforcing TLS for all connections"
137 );
138 self.config.transport = TransportSecurity::TlsOnly;
139 self
140 }
141
142 /// Set the TLS handshake configuration (minimum protocol version and
143 /// optional mutual-TLS client identity).
144 ///
145 /// Replaces the entire [`TlsConfig`]; use [`crate::config::TlsConfig`]'s
146 /// fields to set `min_version` and `client_auth`. The root-trust strategy
147 /// is configured separately via the `tls_roots` config field.
148 ///
149 /// Mutual-TLS PEM material referenced by `client_auth` is read and parsed
150 /// in [`HttpClientBuilder::build`]; IO/parse failures surface there as
151 /// [`HttpError::Tls`].
152 ///
153 /// [`HttpClientBuilder::build`]: crate::builder::HttpClientBuilder::build
154 #[must_use]
155 pub fn tls(mut self, tls: TlsConfig) -> Self {
156 self.config.tls = tls;
157 self
158 }
159
160 /// Set the minimum TLS protocol version, leaving the rest of the
161 /// [`TlsConfig`] (e.g. `client_auth`) untouched.
162 ///
163 /// Convenience shortcut for mutating `config.tls.min_version` without
164 /// rebuilding the whole [`TlsConfig`] via [`HttpClientBuilder::tls`].
165 ///
166 /// [`HttpClientBuilder::tls`]: crate::builder::HttpClientBuilder::tls
167 #[must_use]
168 pub fn tls_min_version(mut self, min_version: TlsVersion) -> Self {
169 self.config.tls.min_version = min_version;
170 self
171 }
172
173 /// Set the mutual-TLS client identity, leaving the rest of the
174 /// [`TlsConfig`] (e.g. `min_version`) untouched.
175 ///
176 /// Convenience shortcut for mutating `config.tls.client_auth` without
177 /// rebuilding the whole [`TlsConfig`] via [`HttpClientBuilder::tls`]. The
178 /// referenced PEM material is read and parsed in
179 /// [`HttpClientBuilder::build`]; IO/parse failures surface there as
180 /// [`HttpError::Tls`].
181 ///
182 /// [`HttpClientBuilder::tls`]: crate::builder::HttpClientBuilder::tls
183 /// [`HttpClientBuilder::build`]: crate::builder::HttpClientBuilder::build
184 #[must_use]
185 pub fn client_auth(mut self, client_auth: ClientAuthConfig) -> Self {
186 self.config.tls.client_auth = Some(client_auth);
187 self
188 }
189
190 /// Enable OpenTelemetry tracing layer
191 ///
192 /// When enabled, creates spans for outbound requests with HTTP metadata
193 /// and injects W3C trace context headers (when `otel` feature is enabled).
194 #[must_use]
195 pub fn with_otel(mut self) -> Self {
196 self.config.otel = true;
197 self
198 }
199
200 /// Insert an optional auth layer between retry and timeout in the stack.
201 ///
202 /// Stack position: `… → Retry → **this layer** → Timeout → …`
203 ///
204 /// The layer sits inside the retry loop so each attempt re-executes it
205 /// (e.g. re-reads a refreshed bearer token). Only one auth layer can be
206 /// set; a second call replaces the first.
207 #[must_use]
208 pub fn with_auth_layer(
209 mut self,
210 wrap: impl FnOnce(InnerService) -> InnerService + Send + 'static,
211 ) -> Self {
212 self.auth_layer = Some(Box::new(wrap));
213 self
214 }
215
216 /// Insert a metrics layer between the rate-limit and retry layers.
217 ///
218 /// Stack position: `… → RateLimit → **this layer** → Retry → Auth → Timeout → …`
219 ///
220 /// The layer sits outside the retry loop so it observes a single logical
221 /// request once, regardless of how many transport-level retries the retry
222 /// layer issues. If the use case is "count every attempt", the equivalent
223 /// observation can be made with [`with_otel`](Self::with_otel) (one span
224 /// per attempt) and a `tracing` → metrics bridge.
225 ///
226 /// Only one metrics layer can be set; a second call replaces the first.
227 #[must_use]
228 pub fn with_metrics_layer(
229 mut self,
230 wrap: impl FnOnce(InnerService) -> InnerService + Send + 'static,
231 ) -> Self {
232 self.metrics_layer = Some(Box::new(wrap));
233 self
234 }
235
236 /// Record OpenTelemetry client metrics using the default route classifier.
237 ///
238 /// Emits `http.client.request.duration` (a histogram, in seconds) for each
239 /// logical request. `client_type` names the instrumentation scope (the
240 /// meter), like the server-side gear name. The default classifier labels
241 /// each request `"METHOD host"` via the `http.route` attribute — a bounded
242 /// value that never contains a raw path, so metric cardinality stays
243 /// controlled. Use [`with_metrics_by`](Self::with_metrics_by)
244 /// to supply route templates.
245 ///
246 /// Like [`with_metrics_layer`](Self::with_metrics_layer), the layer sits
247 /// outside the retry loop, so it observes one logical request, not each
248 /// attempt.
249 #[cfg(feature = "otel")]
250 #[must_use]
251 pub fn with_metrics(self, client_type: impl Into<String>) -> Self {
252 self.with_metrics_by(client_type, crate::layers::default_classify)
253 }
254
255 /// Record OpenTelemetry client metrics with a custom route classifier.
256 ///
257 /// Like [`with_metrics`](Self::with_metrics) but accepts a classifier — `classify` produces the
258 /// `http.route` attribute for each request (the Rust analogue of go-appkit's
259 /// `ClassifyRequest`). It MUST return a bounded set of values (e.g.
260 /// `"GET /users/{id}"`), never raw paths, to avoid unbounded metric
261 /// cardinality.
262 #[cfg(feature = "otel")]
263 #[must_use]
264 pub fn with_metrics_by(
265 self,
266 client_type: impl Into<String>,
267 classify: impl Fn(&http::Request<Full<Bytes>>) -> std::borrow::Cow<'static, str>
268 + Send
269 + Sync
270 + 'static,
271 ) -> Self {
272 let client_type = client_type.into();
273 let classify: crate::layers::ClassifyFn = std::sync::Arc::new(classify);
274 let layer = crate::layers::MetricsLayer::new(&client_type, classify);
275 self.with_metrics_layer(move |svc| {
276 ServiceBuilder::new()
277 .layer(layer)
278 .service(svc)
279 .boxed_clone()
280 })
281 }
282
283 /// Set the buffer capacity for concurrent request handling
284 ///
285 /// The HTTP client uses an internal buffer to allow concurrent requests
286 /// without external locking. This sets the maximum number of requests
287 /// that can be queued.
288 ///
289 /// **Note**: A capacity of 0 is invalid and will be clamped to 1.
290 /// Tower's Buffer panics with capacity=0, so we enforce minimum of 1.
291 #[must_use]
292 pub fn buffer_capacity(mut self, capacity: usize) -> Self {
293 // Clamp to at least 1 - tower::Buffer panics with capacity=0
294 self.config.buffer_capacity = capacity.max(1);
295 self
296 }
297
298 /// Set the maximum number of redirects to follow
299 ///
300 /// Set to `0` to disable redirect following (3xx responses pass through as-is).
301 /// Default: 10
302 #[must_use]
303 pub fn max_redirects(mut self, max_redirects: usize) -> Self {
304 self.config.redirect.max_redirects = max_redirects;
305 self
306 }
307
308 /// Disable redirect following
309 ///
310 /// Equivalent to `.max_redirects(0)`. When disabled, 3xx responses are
311 /// returned to the caller without following the `Location` header.
312 #[must_use]
313 pub fn no_redirects(mut self) -> Self {
314 self.config.redirect = RedirectConfig::disabled();
315 self
316 }
317
318 /// Set the redirect policy configuration
319 ///
320 /// Use this to configure redirect security settings:
321 /// - `same_origin_only`: Only follow redirects to the same host
322 /// - `strip_sensitive_headers`: Remove `Authorization`/`Cookie` on cross-origin
323 /// - `allow_https_downgrade`: Allow HTTPS → HTTP redirects (not recommended)
324 ///
325 /// # Example
326 ///
327 /// ```rust,ignore
328 /// let client = HttpClient::builder()
329 /// .redirect(RedirectConfig::permissive()) // Allow all redirects with header stripping
330 /// .build()?;
331 /// ```
332 #[must_use]
333 pub fn redirect(mut self, config: RedirectConfig) -> Self {
334 self.config.redirect = config;
335 self
336 }
337
338 /// Set the idle connection timeout for the connection pool
339 ///
340 /// Connections that remain idle for longer than this duration will be
341 /// closed and removed from the pool. Default: 90 seconds.
342 ///
343 /// Set to `None` to disable idle timeout (connections kept indefinitely).
344 #[must_use]
345 pub fn pool_idle_timeout(mut self, timeout: Option<Duration>) -> Self {
346 self.config.pool_idle_timeout = timeout;
347 self
348 }
349
350 /// Set the maximum number of idle connections per host
351 ///
352 /// Limits how many idle connections are kept in the pool for each host.
353 /// Default: 32.
354 ///
355 /// - Setting to `0` disables connection reuse entirely
356 /// - Setting too high may waste resources on rarely-used connections
357 #[must_use]
358 pub fn pool_max_idle_per_host(mut self, max: usize) -> Self {
359 self.config.pool_max_idle_per_host = max;
360 self
361 }
362
363 /// Build the HTTP client with all configured layers
364 ///
365 /// # Errors
366 /// Returns an error if TLS initialization fails or configuration is invalid.
367 ///
368 /// Under `--features fips`, returns [`HttpError::InsecureTransport`] when
369 /// `config.transport == TransportSecurity::AllowInsecureHttp`. FIPS builds
370 /// must use [`TransportSecurity::TlsOnly`] — there is no opt-out.
371 pub fn build(self) -> Result<crate::HttpClient, HttpError> {
372 // Reject AllowInsecureHttp under --features fips before any TLS work.
373 // The check lives here (rather than in HttpClientConfig) so that
374 // constructing a config with AllowInsecureHttp is still cheap and
375 // infallible; only actually building a client fails closed.
376 #[cfg(feature = "fips")]
377 if self.config.transport == TransportSecurity::AllowInsecureHttp {
378 tracing::warn!(
379 target: "toolkit_http::security",
380 "rejecting AllowInsecureHttp under --features fips: returning HttpError::InsecureTransport"
381 );
382 return Err(HttpError::InsecureTransport);
383 }
384
385 let timeout = self.config.request_timeout;
386 let total_timeout = self.config.total_timeout;
387
388 // Build the HTTPS connector (may fail for Native roots if no valid certs)
389 let https = build_https_connector(
390 self.config.tls_roots,
391 self.config.transport,
392 &self.config.tls,
393 )?;
394
395 // Create the base hyper client with HTTP/2 support and connection pool settings
396 let mut client_builder = Client::builder(TokioExecutor::new());
397
398 // Configure connection pool
399 // CRITICAL: pool_timer is required for pool_idle_timeout to work!
400 client_builder
401 .pool_timer(TokioTimer::new())
402 .pool_max_idle_per_host(self.config.pool_max_idle_per_host)
403 .http2_only(false); // Allow both HTTP/1 and HTTP/2 via ALPN
404
405 // Set idle timeout (None = no timeout, connections kept indefinitely)
406 if let Some(idle_timeout) = self.config.pool_idle_timeout {
407 client_builder.pool_idle_timeout(idle_timeout);
408 }
409
410 let hyper_client = client_builder.build::<_, Full<Bytes>>(https);
411
412 // Parse user agent header (may fail)
413 let ua_layer = UserAgentLayer::try_new(&self.config.user_agent)?;
414
415 // =======================================================================
416 // Tower Layer Stack (outer to inner)
417 // =======================================================================
418 //
419 // Request flow (outer → inner):
420 // Buffer → OtelLayer → [MetricsLayer?] → LoadShed/Concurrency →
421 // RetryLayer → [AuthLayer?] → ErrorMapping → Timeout → UserAgent →
422 // Decompression → FollowRedirect → hyper_client
423 //
424 // AuthLayer (if set via with_auth_layer) sits inside the retry
425 // loop so each retry re-acquires credentials (e.g. refreshed
426 // bearer token).
427 //
428 // MetricsLayer (if set via with_metrics_layer) sits outside the
429 // retry loop so it observes one logical request, not per-attempt, and
430 // outside the concurrency limiter so a load-shed rejection is still
431 // counted (recorded with error.type = "overloaded").
432 //
433 // Response flow (inner → outer):
434 // hyper_client → FollowRedirect → Decompression → UserAgent →
435 // Timeout → ErrorMapping → [AuthLayer?] → RetryLayer →
436 // LoadShed/Concurrency → [MetricsLayer?] → OtelLayer → Buffer
437 //
438 // Key semantics (reqwest-like):
439 // - send() returns Ok(Response) for ALL HTTP statuses (including 4xx/5xx)
440 // - send() returns Err only for transport/timeout/TLS errors
441 // - Non-2xx converted to error ONLY via error_for_status()
442 // - RetryLayer handles both Err (transport) and Ok(Response) (status)
443 // retries internally, draining body before retry for connection reuse
444 // - FollowRedirect handles 3xx responses internally with security protections:
445 // * Same-origin enforcement (default) - blocks SSRF attacks
446 // * Sensitive header stripping on cross-origin redirects
447 // * HTTPS downgrade protection
448 //
449 // =======================================================================
450 //
451 let redirect_policy = SecureRedirectPolicy::new(self.config.redirect.clone());
452
453 // Build the service stack with secure redirect following
454 let service = ServiceBuilder::new()
455 .layer(TimeoutLayer::new(timeout))
456 .layer(ua_layer)
457 .layer(DecompressionLayer::new())
458 .layer(FollowRedirectLayer::with_policy(redirect_policy))
459 .service(hyper_client);
460
461 // Map the decompression body to our boxed ResponseBody type.
462 // This converts Response<DecompressionBody<Incoming>> to Response<ResponseBody>.
463 //
464 // The decompression body's error type is tower_http::BoxError, which we convert
465 // to our boxed error type for consistency with the ResponseBody definition.
466 let service = service.map_response(map_decompression_response);
467
468 // Map errors to HttpError with proper timeout duration
469 let service = service.map_err(move |e: tower::BoxError| map_tower_error(e, timeout));
470
471 // Box the service for type erasure
472 let mut boxed_service = service.boxed_clone();
473
474 // Apply auth layer (between timeout and retry).
475 // Inside retry so each retry attempt re-acquires the token.
476 if let Some(wrap) = self.auth_layer {
477 boxed_service = wrap(boxed_service);
478 }
479
480 // Conditionally wrap with RetryLayer
481 //
482 // RetryLayer handles retries for both:
483 // - Err(HttpError::Transport/Timeout) - transport-level failures
484 // - Ok(Response) with retryable status codes (429, 5xx for GET, etc.)
485 //
486 // When retrying on status codes, RetryLayer drains the response body
487 // (up to configured limit) to allow connection reuse.
488 //
489 // If total_timeout is set, the entire operation (including all retries)
490 // must complete within that duration.
491 if let Some(ref retry_config) = self.config.retry {
492 let retry_layer = RetryLayer::with_total_timeout(retry_config.clone(), total_timeout);
493 let retry_service = ServiceBuilder::new()
494 .layer(retry_layer)
495 .service(boxed_service);
496 boxed_service = retry_service.boxed_clone();
497 }
498
499 // Conditionally wrap with concurrency limit + load shedding.
500 // LoadShedLayer returns error immediately when ConcurrencyLimitLayer is
501 // saturated instead of waiting indefinitely (Poll::Pending).
502 //
503 // Applied BEFORE the metrics layer (i.e. inner to it) on purpose: a shed
504 // request then propagates back out through MetricsLayer and is recorded
505 // in `http.client.request.duration` with `error.type = "overloaded"`,
506 // rather than being rejected outside metrics and left invisible.
507 if let Some(rate_limit) = self.config.rate_limit
508 && rate_limit.max_concurrent_requests < usize::MAX
509 {
510 let limited_service = ServiceBuilder::new()
511 .layer(LoadShedLayer::new())
512 // `.max(1)`: a cap of 0 never grants a permit (sheds everything);
513 // mirrors the `buffer_capacity` clamp below.
514 .layer(ConcurrencyLimitLayer::new(
515 rate_limit.max_concurrent_requests.max(1),
516 ))
517 .service(boxed_service);
518 // Map load shed errors to HttpError::Overloaded
519 let limited_service = limited_service.map_err(map_load_shed_error);
520 boxed_service = limited_service.boxed_clone();
521 }
522
523 // Apply metrics layer (outside both retry and the concurrency limiter).
524 // Outside retry: observes one logical request, not per-attempt. Outside
525 // the limiter: so load-shed rejections are counted (see above).
526 if let Some(wrap) = self.metrics_layer {
527 boxed_service = wrap(boxed_service);
528 }
529
530 // Conditionally wrap with OTEL tracing layer (outermost layer before buffer)
531 // Applied last so it sees the final request after UserAgent and other modifications.
532 // Creates spans, records status, and injects trace context headers.
533 if self.config.otel {
534 let otel_service = ServiceBuilder::new()
535 .layer(OtelLayer::new())
536 .service(boxed_service);
537 boxed_service = otel_service.boxed_clone();
538 }
539
540 // Wrap in Buffer as the final step for true concurrent access
541 // Buffer spawns a background task that processes requests from a channel,
542 // providing Clone + Send + Sync without any mutex serialization.
543 let buffer_capacity = self.config.buffer_capacity.max(1);
544 let buffered_service: crate::client::BufferedService =
545 Buffer::new(boxed_service, buffer_capacity);
546
547 Ok(crate::HttpClient {
548 service: buffered_service,
549 max_body_size: self.config.max_body_size,
550 transport_security: self.config.transport,
551 })
552 }
553}
554
555#[cfg(test)]
556impl HttpClientBuilder {
557 /// Build an `HttpClient` with a custom inner service replacing the
558 /// hyper connector. The full middleware stack (Retry, Concurrency,
559 /// Buffer, etc.) is applied on top.
560 ///
561 /// The inner service must handle `Request<Full<Bytes>>` and return
562 /// `Response<ResponseBody>`. Use this to inject a fake slow service
563 /// for cancellation testing without needing a real HTTP server.
564 fn build_with_inner_service(self, inner: InnerService) -> crate::HttpClient {
565 let mut boxed_service = inner;
566
567 if let Some(ref retry_config) = self.config.retry {
568 let retry_layer =
569 RetryLayer::with_total_timeout(retry_config.clone(), self.config.total_timeout);
570 let retry_service = ServiceBuilder::new()
571 .layer(retry_layer)
572 .service(boxed_service);
573 boxed_service = retry_service.boxed_clone();
574 }
575
576 if let Some(rate_limit) = self.config.rate_limit
577 && rate_limit.max_concurrent_requests < usize::MAX
578 {
579 let limited_service = ServiceBuilder::new()
580 .layer(LoadShedLayer::new())
581 .layer(ConcurrencyLimitLayer::new(
582 rate_limit.max_concurrent_requests.max(1),
583 ))
584 .service(boxed_service);
585 let limited_service = limited_service.map_err(map_load_shed_error);
586 boxed_service = limited_service.boxed_clone();
587 }
588
589 let buffer_capacity = self.config.buffer_capacity.max(1);
590 let buffered_service: crate::client::BufferedService =
591 Buffer::new(boxed_service, buffer_capacity);
592
593 crate::HttpClient {
594 service: buffered_service,
595 max_body_size: self.config.max_body_size,
596 transport_security: self.config.transport,
597 }
598 }
599}
600
601impl Default for HttpClientBuilder {
602 fn default() -> Self {
603 Self::new()
604 }
605}
606
607/// Map tower errors to `HttpError` with actual timeout duration
608///
609/// Attempts to extract existing `HttpError` from the boxed error before
610/// wrapping as `Transport`. This preserves typed errors like `Overloaded`
611/// and `ServiceClosed` that may have been boxed by tower middleware.
612fn map_tower_error(err: tower::BoxError, timeout: Duration) -> HttpError {
613 if err.is::<tower::timeout::error::Elapsed>() {
614 return HttpError::Timeout(timeout);
615 }
616
617 // Try to extract existing HttpError before wrapping as Transport
618 match err.downcast::<HttpError>() {
619 Ok(http_err) => *http_err,
620 Err(other) => HttpError::Transport(other),
621 }
622}
623
624/// Map load shed errors to `HttpError::Overloaded`.
625///
626/// A shed request is observable via metrics: the concurrency limiter is inner to
627/// [`MetricsLayer`](crate::layers::metrics), so the rejection is recorded in
628/// `http.client.request.duration` with `error.type = "overloaded"`.
629fn map_load_shed_error(err: tower::BoxError) -> HttpError {
630 if err.is::<tower::load_shed::error::Overloaded>() {
631 HttpError::Overloaded
632 } else {
633 // Pass through other HttpError types (from inner service)
634 match err.downcast::<HttpError>() {
635 Ok(http_err) => *http_err,
636 Err(err) => HttpError::Transport(err),
637 }
638 }
639}
640
641/// Map the decompression response to our boxed response body type.
642///
643/// This converts `Response<DecompressionBody<Incoming>>` to `Response<ResponseBody>`
644/// by boxing the body with appropriate error type mapping.
645fn map_decompression_response<B>(response: Response<B>) -> Response<ResponseBody>
646where
647 B: hyper::body::Body<Data = Bytes> + Send + Sync + 'static,
648 B::Error: Into<Box<dyn std::error::Error + Send + Sync>>,
649{
650 let (parts, body) = response.into_parts();
651 // Convert the decompression body errors to our boxed error type.
652 // tower-http's DecompressionBody uses tower_http::BoxError which is
653 // compatible with our Box<dyn Error + Send + Sync> via Into.
654 let boxed_body: ResponseBody = body.map_err(Into::into).boxed();
655 Response::from_parts(parts, boxed_body)
656}
657
658/// Build the HTTPS connector with the specified TLS root configuration.
659///
660/// For `TlsRootConfig::Native`, uses cached native root certificates to avoid
661/// repeated OS certificate store lookups on each `build()` call.
662///
663/// HTTP/2 is enabled via `enable_all_versions()` which configures ALPN to
664/// advertise both h2 and http/1.1. Protocol selection happens during TLS
665/// handshake based on server support.
666///
667/// # Errors
668///
669/// Returns `HttpError::Tls` if `TlsRootConfig::Native` is requested but no
670/// valid root certificates are available from the OS certificate store.
671fn build_https_connector(
672 tls_roots: TlsRootConfig,
673 transport: TransportSecurity,
674 tls: &TlsConfig,
675) -> Result<HttpsConnector<HttpConnector>, HttpError> {
676 let allow_http = transport == TransportSecurity::AllowInsecureHttp;
677
678 // Both branches build a `ClientConfig` ourselves (rather than using
679 // `with_provider_and_webpki_roots`) so we can apply `require_ems = true`
680 // under the `fips` feature — see `tls::build_client_config`. The
681 // functions return `tls::TlsConfigError` (an enum implementing
682 // `std::error::Error`); boxing it into `HttpError::Tls`'s
683 // `Box<dyn Error + Send + Sync>` preserves the source chain via
684 // `TlsConfigError`'s own `Error::source()` impl.
685 let client_config = match tls_roots {
686 TlsRootConfig::WebPki => tls::webpki_roots_client_config(tls),
687 TlsRootConfig::Native => tls::native_roots_client_config(tls),
688 }
689 .map_err(|e| HttpError::Tls(Box::new(e)))?;
690
691 let builder = hyper_rustls::HttpsConnectorBuilder::new().with_tls_config(client_config);
692 let connector = if allow_http {
693 builder.https_or_http().enable_all_versions().build()
694 } else {
695 builder.https_only().enable_all_versions().build()
696 };
697 Ok(connector)
698}
699
700#[cfg(test)]
701#[cfg_attr(coverage_nightly, coverage(off))]
702mod tests {
703 use super::*;
704 use crate::config::DEFAULT_USER_AGENT;
705
706 #[test]
707 fn test_builder_default() {
708 let builder = HttpClientBuilder::new();
709 assert_eq!(builder.config.request_timeout, Duration::from_secs(30));
710 assert_eq!(builder.config.user_agent, DEFAULT_USER_AGENT);
711 assert!(builder.config.retry.is_some());
712 assert_eq!(builder.config.buffer_capacity, 1024);
713 }
714
715 #[test]
716 fn test_builder_with_config() {
717 let config = HttpClientConfig::minimal();
718 let builder = HttpClientBuilder::with_config(config);
719 assert_eq!(builder.config.request_timeout, Duration::from_secs(10));
720 }
721
722 #[test]
723 fn test_builder_timeout() {
724 let builder = HttpClientBuilder::new().timeout(Duration::from_mins(1));
725 assert_eq!(builder.config.request_timeout, Duration::from_mins(1));
726 }
727
728 #[test]
729 fn test_builder_user_agent() {
730 let builder = HttpClientBuilder::new().user_agent("custom/1.0");
731 assert_eq!(builder.config.user_agent, "custom/1.0");
732 }
733
734 #[test]
735 fn test_builder_retry() {
736 let builder = HttpClientBuilder::new().retry(None);
737 assert!(builder.config.retry.is_none());
738 }
739
740 #[test]
741 fn test_builder_max_body_size() {
742 let builder = HttpClientBuilder::new().max_body_size(1024);
743 assert_eq!(builder.config.max_body_size, 1024);
744 }
745
746 #[test]
747 fn test_builder_transport_security() {
748 let builder = HttpClientBuilder::new().transport(TransportSecurity::TlsOnly);
749 assert_eq!(builder.config.transport, TransportSecurity::TlsOnly);
750
751 let builder = HttpClientBuilder::new().deny_insecure_http();
752 assert_eq!(builder.config.transport, TransportSecurity::TlsOnly);
753
754 let builder = HttpClientBuilder::new();
755 #[cfg(not(feature = "fips"))]
756 assert_eq!(
757 builder.config.transport,
758 TransportSecurity::AllowInsecureHttp
759 );
760 #[cfg(feature = "fips")]
761 assert_eq!(builder.config.transport, TransportSecurity::TlsOnly);
762 }
763
764 #[test]
765 fn test_builder_otel() {
766 let builder = HttpClientBuilder::new().with_otel();
767 assert!(builder.config.otel);
768 }
769
770 #[test]
771 fn test_builder_buffer_capacity() {
772 let builder = HttpClientBuilder::new().buffer_capacity(512);
773 assert_eq!(builder.config.buffer_capacity, 512);
774 }
775
776 /// Test that `buffer_capacity=0` is clamped to 1 to prevent panic.
777 ///
778 /// Tower's Buffer panics with capacity=0, so we enforce minimum of 1.
779 #[test]
780 fn test_builder_buffer_capacity_zero_clamped() {
781 let builder = HttpClientBuilder::new().buffer_capacity(0);
782 assert_eq!(
783 builder.config.buffer_capacity, 1,
784 "buffer_capacity=0 should be clamped to 1"
785 );
786 }
787
788 /// Test that `buffer_capacity=0` via config is clamped during `build()`.
789 #[tokio::test]
790 async fn test_builder_buffer_capacity_zero_in_config_clamped() {
791 let config = HttpClientConfig {
792 buffer_capacity: 0, // Invalid - should be clamped in build()
793 ..Default::default()
794 };
795 let result = HttpClientBuilder::with_config(config).build();
796 // Should succeed (clamped to 1), not panic
797 assert!(
798 result.is_ok(),
799 "build() should succeed with capacity clamped to 1"
800 );
801 }
802
803 #[test]
804 fn test_builder_concurrency_limit() {
805 let builder = HttpClientBuilder::new().concurrency_limit(Some(RateLimitConfig {
806 max_concurrent_requests: 7,
807 }));
808 let rate_limit = builder.config.rate_limit.expect("rate_limit should be set");
809 assert_eq!(rate_limit.max_concurrent_requests, 7);
810 }
811
812 #[test]
813 fn test_builder_concurrency_limit_none_disables() {
814 let builder = HttpClientBuilder::new().concurrency_limit(None);
815 assert!(
816 builder.config.rate_limit.is_none(),
817 "None should disable the limiter"
818 );
819 }
820
821 /// `max_concurrent_requests: usize::MAX` is treated as unlimited — the layer
822 /// is skipped in `build()`. Building must still succeed.
823 #[tokio::test]
824 async fn test_builder_concurrency_limit_unlimited_builds() {
825 let client = HttpClientBuilder::new()
826 .concurrency_limit(Some(RateLimitConfig::unlimited()))
827 .build();
828 assert!(client.is_ok());
829 }
830
831 /// Smoke test: building with `max_concurrent_requests: 0` must not fail or
832 /// panic. This does *not* prove the `0 → 1` clamp works (a dead 0-permit
833 /// limiter also builds fine) — that is covered behaviourally in
834 /// toolkit-contract's `concurrency_limit::cap_zero_is_clamped_and_still_serves`,
835 /// which builds a cap-0 client and asserts a request is still served.
836 #[tokio::test]
837 async fn test_builder_concurrency_limit_zero_builds() {
838 let client = HttpClientBuilder::new()
839 .concurrency_limit(Some(RateLimitConfig {
840 max_concurrent_requests: 0,
841 }))
842 .build();
843 assert!(client.is_ok(), "build() should succeed with a 0 cap");
844 }
845
846 #[tokio::test]
847 async fn test_builder_build_with_otel() {
848 let client = HttpClientBuilder::new().with_otel().build();
849 assert!(client.is_ok());
850 }
851
852 #[tokio::test]
853 async fn test_builder_with_auth_layer() {
854 let client = HttpClientBuilder::new()
855 .with_auth_layer(|svc| svc) // identity transform
856 .build();
857 assert!(client.is_ok());
858 }
859
860 #[tokio::test]
861 async fn test_builder_with_metrics_layer() {
862 let client = HttpClientBuilder::new()
863 .with_metrics_layer(|svc| svc) // identity transform
864 .build();
865 assert!(client.is_ok());
866 }
867
868 #[tokio::test]
869 async fn test_builder_with_metrics_layer_second_call_replaces_first() {
870 use std::sync::Arc;
871 use std::sync::atomic::{AtomicUsize, Ordering};
872
873 let call_count = Arc::new(AtomicUsize::new(0));
874 let call_count2 = call_count.clone();
875
876 // Second call should replace the first; only one layer is applied.
877 let client = HttpClientBuilder::new()
878 .with_metrics_layer(|_svc| {
879 // This closure should NOT be called (replaced by the second).
880 panic!("first metrics layer should have been replaced");
881 })
882 .with_metrics_layer(move |svc| {
883 call_count2.fetch_add(1, Ordering::SeqCst);
884 svc
885 })
886 .build();
887
888 assert!(client.is_ok());
889 assert_eq!(
890 call_count.load(Ordering::SeqCst),
891 1,
892 "second metrics layer must be applied exactly once"
893 );
894 }
895
896 #[tokio::test]
897 async fn test_builder_build() {
898 let client = HttpClientBuilder::new().build();
899 assert!(client.is_ok());
900 }
901
902 #[tokio::test]
903 async fn test_builder_build_with_deny_insecure_http() {
904 let client = HttpClientBuilder::new().deny_insecure_http().build();
905 assert!(client.is_ok());
906 }
907
908 #[tokio::test]
909 async fn test_builder_build_with_sse_config() {
910 use crate::config::HttpClientConfig;
911 let config = HttpClientConfig::sse();
912 let client = HttpClientBuilder::with_config(config).build();
913 assert!(client.is_ok(), "SSE config should build successfully");
914 }
915
916 #[tokio::test]
917 async fn test_builder_build_invalid_user_agent() {
918 let client = HttpClientBuilder::new()
919 .user_agent("invalid\x00agent")
920 .build();
921 assert!(client.is_err());
922 }
923
924 #[tokio::test]
925 async fn test_builder_default_uses_webpki_roots() {
926 let builder = HttpClientBuilder::new();
927 assert_eq!(builder.config.tls_roots, TlsRootConfig::WebPki);
928 // Build should succeed without OS native roots
929 let client = builder.build();
930 assert!(client.is_ok());
931 }
932
933 #[tokio::test]
934 async fn test_builder_native_roots() {
935 let config = HttpClientConfig {
936 tls_roots: TlsRootConfig::Native,
937 ..Default::default()
938 };
939 let result = HttpClientBuilder::with_config(config).build();
940
941 // Native roots may succeed or fail depending on OS certificate availability.
942 // On systems with certs: Ok(_)
943 // On minimal containers without certs: Err(HttpError::Tls(_))
944 match &result {
945 Ok(_) => {
946 // Success on systems with native certs
947 }
948 Err(HttpError::Tls(err)) => {
949 // Expected failure on systems without native certs
950 let msg = err.to_string();
951 assert!(
952 msg.contains("native root") || msg.contains("certificate"),
953 "TLS error should mention certificates: {msg}"
954 );
955 }
956 Err(other) => {
957 panic!("Unexpected error type: {other:?}");
958 }
959 }
960 }
961
962 #[tokio::test]
963 async fn test_builder_webpki_roots_https_only() {
964 let config = HttpClientConfig {
965 tls_roots: TlsRootConfig::WebPki,
966 transport: TransportSecurity::TlsOnly,
967 ..Default::default()
968 };
969 let client = HttpClientBuilder::with_config(config).build();
970 assert!(client.is_ok());
971 }
972
973 /// Verify HTTP/2 is enabled for all TLS root configurations.
974 ///
975 /// HTTP/2 support is configured via `enable_all_versions()` on the connector,
976 /// which sets up ALPN to negotiate h2 or http/1.1 during TLS handshake.
977 /// The hyper client uses `http2_only(false)` to allow both protocols.
978 ///
979 /// The `AllowInsecureHttp` sub-cases are skipped under `--features fips`
980 /// because the FIPS guard in `build()` rejects insecure transport.
981 #[tokio::test]
982 async fn test_http2_enabled_for_all_configurations() {
983 // Test WebPki with default transport (AllowInsecureHttp without fips, TlsOnly under fips)
984 let client = HttpClientBuilder::new().build();
985 assert!(
986 client.is_ok(),
987 "WebPki + default transport should build with HTTP/2 enabled"
988 );
989
990 // Test WebPki with TlsOnly (HTTPS only)
991 let client = HttpClientBuilder::new()
992 .transport(TransportSecurity::TlsOnly)
993 .build();
994 assert!(
995 client.is_ok(),
996 "WebPki + TlsOnly should build with HTTP/2 enabled"
997 );
998
999 // Test Native roots with AllowInsecureHttp (non-fips only)
1000 #[cfg(not(feature = "fips"))]
1001 {
1002 let config = HttpClientConfig {
1003 tls_roots: TlsRootConfig::Native,
1004 transport: TransportSecurity::AllowInsecureHttp,
1005 ..Default::default()
1006 };
1007 let client = HttpClientBuilder::with_config(config).build();
1008 assert!(
1009 client.is_ok(),
1010 "Native + AllowInsecureHttp should build with HTTP/2 enabled"
1011 );
1012 }
1013
1014 // Test Native roots with TlsOnly (HTTPS only)
1015 let config = HttpClientConfig {
1016 tls_roots: TlsRootConfig::Native,
1017 transport: TransportSecurity::TlsOnly,
1018 ..Default::default()
1019 };
1020 let client = HttpClientBuilder::with_config(config).build();
1021 assert!(
1022 client.is_ok(),
1023 "Native + TlsOnly should build with HTTP/2 enabled"
1024 );
1025 }
1026
1027 /// Test that concurrency limit uses fail-fast behavior (C2).
1028 ///
1029 /// `LoadShedLayer` + `ConcurrencyLimitLayer` combination returns Overloaded error
1030 /// immediately when capacity is exhausted, instead of blocking indefinitely.
1031 #[tokio::test]
1032 async fn test_load_shedding_returns_overloaded_error() {
1033 use bytes::Bytes;
1034 use http::{Request, Response};
1035 use http_body_util::Full;
1036 use std::future::Future;
1037 use std::pin::Pin;
1038 use std::sync::Arc;
1039 use std::sync::atomic::{AtomicUsize, Ordering};
1040 use std::task::{Context, Poll};
1041 use tower::Service;
1042 use tower::ServiceExt;
1043
1044 // A service that holds a slot forever once called
1045 #[derive(Clone)]
1046 struct SlotHoldingService {
1047 active: Arc<AtomicUsize>,
1048 }
1049
1050 impl Service<Request<Full<Bytes>>> for SlotHoldingService {
1051 type Response = Response<Full<Bytes>>;
1052 type Error = HttpError;
1053 type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
1054
1055 fn poll_ready(&mut self, _: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1056 Poll::Ready(Ok(()))
1057 }
1058
1059 fn call(&mut self, _: Request<Full<Bytes>>) -> Self::Future {
1060 self.active.fetch_add(1, Ordering::SeqCst);
1061 // Never complete - holds the slot
1062 Box::pin(std::future::pending())
1063 }
1064 }
1065
1066 let active = Arc::new(AtomicUsize::new(0));
1067
1068 // Build a service with load shedding and concurrency limit of 1
1069 let service = tower::ServiceBuilder::new()
1070 .layer(LoadShedLayer::new())
1071 .layer(ConcurrencyLimitLayer::new(1))
1072 .service(SlotHoldingService {
1073 active: active.clone(),
1074 });
1075
1076 let service = service.map_err(map_load_shed_error);
1077
1078 // First request: will occupy the single slot
1079 let req1 = Request::builder()
1080 .uri("http://test")
1081 .body(Full::new(Bytes::new()))
1082 .unwrap();
1083 let mut svc1 = service.clone();
1084
1085 let svc1_ready = svc1.ready().await.unwrap();
1086 let _pending_fut = svc1_ready.call(req1);
1087
1088 // Wait for the slot to be occupied
1089 tokio::time::sleep(Duration::from_millis(10)).await;
1090 assert_eq!(
1091 active.load(Ordering::SeqCst),
1092 1,
1093 "First request should be active"
1094 );
1095
1096 // Second request: LoadShedLayer should reject because ConcurrencyLimit is at capacity
1097 let req2 = Request::builder()
1098 .uri("http://test")
1099 .body(Full::new(Bytes::new()))
1100 .unwrap();
1101
1102 let mut svc2 = service.clone();
1103
1104 // LoadShedLayer checks poll_ready and returns Overloaded if inner service is not ready
1105 let result = tokio::time::timeout(Duration::from_millis(100), async {
1106 // poll_ready should return quickly with error (not block)
1107 match svc2.ready().await {
1108 Ok(ready_svc) => ready_svc.call(req2).await,
1109 Err(e) => Err(e),
1110 }
1111 })
1112 .await;
1113
1114 // Should complete within timeout (not hang) and return Overloaded
1115 assert!(result.is_ok(), "Request should not hang");
1116 let err = result.unwrap().unwrap_err();
1117 assert!(
1118 matches!(err, HttpError::Overloaded),
1119 "Expected Overloaded error, got: {err:?}"
1120 );
1121 }
1122
1123 // ==========================================================================
1124 // map_tower_error Tests
1125 // ==========================================================================
1126
1127 /// Test that `map_tower_error` preserves `HttpError::Overloaded` when wrapped in `BoxError`
1128 #[test]
1129 fn test_map_tower_error_preserves_overloaded() {
1130 let http_err = HttpError::Overloaded;
1131 let boxed: tower::BoxError = Box::new(http_err);
1132 let result = map_tower_error(boxed, Duration::from_secs(30));
1133
1134 assert!(
1135 matches!(result, HttpError::Overloaded),
1136 "Should preserve HttpError::Overloaded, got: {result:?}"
1137 );
1138 }
1139
1140 /// Test that `map_tower_error` preserves `HttpError::ServiceClosed` when wrapped in `BoxError`
1141 #[test]
1142 fn test_map_tower_error_preserves_service_closed() {
1143 let http_err = HttpError::ServiceClosed;
1144 let boxed: tower::BoxError = Box::new(http_err);
1145 let result = map_tower_error(boxed, Duration::from_secs(30));
1146
1147 assert!(
1148 matches!(result, HttpError::ServiceClosed),
1149 "Should preserve HttpError::ServiceClosed, got: {result:?}"
1150 );
1151 }
1152
1153 /// Test that `map_tower_error` preserves `HttpError::Timeout` with original duration
1154 #[test]
1155 fn test_map_tower_error_preserves_timeout_attempt() {
1156 let original_duration = Duration::from_secs(5);
1157 let http_err = HttpError::Timeout(original_duration);
1158 let boxed: tower::BoxError = Box::new(http_err);
1159 // Pass a different timeout to verify original is preserved
1160 let result = map_tower_error(boxed, Duration::from_secs(30));
1161
1162 match result {
1163 HttpError::Timeout(d) => {
1164 assert_eq!(
1165 d, original_duration,
1166 "Should preserve original timeout duration"
1167 );
1168 }
1169 other => panic!("Should preserve HttpError::Timeout, got: {other:?}"),
1170 }
1171 }
1172
1173 /// Test that `map_tower_error` wraps unknown errors as Transport
1174 #[test]
1175 fn test_map_tower_error_wraps_unknown_as_transport() {
1176 let other_err: tower::BoxError = Box::new(std::io::Error::new(
1177 std::io::ErrorKind::ConnectionRefused,
1178 "connection refused",
1179 ));
1180 let result = map_tower_error(other_err, Duration::from_secs(30));
1181
1182 assert!(
1183 matches!(result, HttpError::Transport(_)),
1184 "Should wrap unknown errors as Transport, got: {result:?}"
1185 );
1186 }
1187
1188 // ==========================================================================
1189 // Cancellation chain test
1190 //
1191 // Proves that dropping the response future from HttpClient cancels the
1192 // inner service future through the toolkit-http middleware stack
1193 // (Buffer → Concurrency → inner service). Retry is disabled to
1194 // isolate the cancellation path.
1195 //
1196 // Uses build_with_inner_service() to inject a fake slow service at the
1197 // bottom of the real tower stack - no HTTP server needed.
1198 // ==========================================================================
1199
1200 /// Dropping the `HttpClient::send()` future must cancel the inner
1201 /// service future through the full middleware stack.
1202 ///
1203 /// Injects a fake service via `build_with_inner_service()` that
1204 /// blocks on a `Notify` (never completes) and signals a second
1205 /// `Notify` from its `Drop` impl. No sleeps - purely notification-based.
1206 #[tokio::test]
1207 async fn test_cancellation_propagates_through_full_stack() {
1208 use crate::response::ResponseBody;
1209 use std::future::Future;
1210 use std::pin::Pin;
1211 use std::sync::Arc;
1212 use std::sync::atomic::{AtomicBool, Ordering};
1213 use std::task::{Context, Poll};
1214 use tower::Service;
1215
1216 #[derive(Clone)]
1217 struct PendingService {
1218 completed: Arc<AtomicBool>,
1219 drop_notifier: Arc<tokio::sync::Notify>,
1220 started_notifier: Arc<tokio::sync::Notify>,
1221 }
1222
1223 struct FutureGuard {
1224 completed: Arc<AtomicBool>,
1225 drop_notifier: Arc<tokio::sync::Notify>,
1226 }
1227
1228 impl Drop for FutureGuard {
1229 fn drop(&mut self) {
1230 if !self.completed.load(Ordering::SeqCst) {
1231 self.drop_notifier.notify_one();
1232 }
1233 }
1234 }
1235
1236 impl Service<http::Request<Full<Bytes>>> for PendingService {
1237 type Response = http::Response<ResponseBody>;
1238 type Error = HttpError;
1239 type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
1240
1241 fn poll_ready(&mut self, _: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1242 Poll::Ready(Ok(()))
1243 }
1244
1245 fn call(&mut self, _: http::Request<Full<Bytes>>) -> Self::Future {
1246 let completed = self.completed.clone();
1247 let drop_notifier = self.drop_notifier.clone();
1248 let started_notifier = self.started_notifier.clone();
1249 Box::pin(async move {
1250 let _guard = FutureGuard {
1251 completed: completed.clone(),
1252 drop_notifier,
1253 };
1254 // Signal that the request reached the inner service
1255 started_notifier.notify_one();
1256 // Block forever - only completes via drop
1257 std::future::pending::<()>().await;
1258 completed.store(true, Ordering::SeqCst);
1259 unreachable!()
1260 })
1261 }
1262 }
1263
1264 let inner_completed = Arc::new(AtomicBool::new(false));
1265 let drop_notifier = Arc::new(tokio::sync::Notify::new());
1266 let started_notifier = Arc::new(tokio::sync::Notify::new());
1267
1268 let inner = PendingService {
1269 completed: inner_completed.clone(),
1270 drop_notifier: drop_notifier.clone(),
1271 started_notifier: started_notifier.clone(),
1272 };
1273
1274 // Build the real HttpClient stack with our fake service at the bottom.
1275 // Retry disabled to isolate cancellation. Tests: Buffer → Concurrency → PendingService
1276 let client = HttpClientBuilder::new()
1277 .timeout(Duration::from_secs(30))
1278 .retry(None)
1279 .build_with_inner_service(inner.boxed_clone());
1280
1281 // Spawn the request so we can drop it explicitly.
1282 // URL uses https:// so the scheme validation in RequestBuilder passes
1283 // under both the non-fips default (AllowInsecureHttp) and the fips
1284 // default (TlsOnly); the connector here is the injected PendingService,
1285 // so no real TLS handshake happens.
1286 let send_handle = tokio::spawn({
1287 let client = client.clone();
1288 async move { client.get("https://fake/slow").send().await }
1289 });
1290
1291 // Wait for the request to reach the inner service
1292 started_notifier.notified().await;
1293
1294 // Drop the in-flight request by aborting the task
1295 send_handle.abort();
1296
1297 // Wait for the drop notification - no sleep, pure notification
1298 tokio::time::timeout(Duration::from_secs(5), drop_notifier.notified())
1299 .await
1300 .expect(
1301 "Inner service future should have been dropped within 5s - \
1302 the full toolkit-http stack must propagate cancellation",
1303 );
1304
1305 assert!(
1306 !inner_completed.load(Ordering::SeqCst),
1307 "Inner service future should NOT have completed"
1308 );
1309 }
1310}