Skip to main content

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}