Skip to main content

aws_smithy_http_client/client/pool/
client.rs

1/*
2 * Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
3 * SPDX-License-Identifier: Apache-2.0
4 */
5
6//! Partition-bound handles for Smithy HTTP operations.
7
8use super::partition::PartitionId;
9use super::registry::PartitionState;
10use super::ConnectionPool;
11use crate::client::timeout::{self, TimeoutKind};
12use crate::client::{downcast_error, ConnectorCallTimer};
13use crate::sync::Arc;
14use aws_smithy_async::rt::sleep::{default_async_sleep, SharedAsyncSleep};
15use aws_smithy_async::time::SharedTimeSource;
16use aws_smithy_runtime_api::box_error::BoxError;
17use aws_smithy_runtime_api::client::connector_metadata::ConnectorMetadata;
18use aws_smithy_runtime_api::client::http::telemetry::CaptureHttpAttemptTelemetry;
19use aws_smithy_runtime_api::client::http::{
20    HttpClient, HttpConnector, HttpConnectorFuture, HttpConnectorSettings, SharedHttpConnector,
21};
22use aws_smithy_runtime_api::client::orchestrator::{HttpRequest, HttpResponse};
23use aws_smithy_runtime_api::client::result::ConnectorError;
24use aws_smithy_runtime_api::client::runtime_components::{
25    RuntimeComponents, RuntimeComponentsBuilder,
26};
27use aws_smithy_types::config_bag::ConfigBag;
28use std::borrow::Cow;
29use std::error::Error;
30use std::fmt;
31use std::time::Duration;
32
33/// Smithy HTTP client bound to one connection-pool partition.
34///
35/// Construction resolves the partition once. Every request uses that partition
36/// for establishment placement and local reuse, while the configured reuse
37/// scope may permit dispatch through a peer partition's connection. Cloning a
38/// client shares the pool, resolved partition state, and reusable connections.
39#[derive(Clone)]
40pub struct Client {
41    /// Shared connection pool and immutable policy.
42    pool: ConnectionPool,
43    /// Resolved runtime and network placement for this client.
44    partition: Arc<PartitionState>,
45}
46
47impl Client {
48    /// Creates a client for a pool built without explicit partitions.
49    ///
50    /// # Errors
51    ///
52    /// Returns [`ClientBuildError`] when the pool contains only explicit
53    /// partitions.
54    pub fn new(pool: &ConnectionPool) -> Result<Self, ClientBuildError> {
55        Self::from_partition(pool, PartitionId::ANONYMOUS)
56    }
57
58    /// Creates a client bound to `id`.
59    ///
60    /// The identity may name a declared partition or the anonymous partition
61    /// retained by a pool built without declarations.
62    ///
63    /// # Errors
64    ///
65    /// Returns [`ClientBuildError`] when the pool does not contain `id`.
66    pub fn from_partition(
67        pool: &ConnectionPool,
68        id: PartitionId,
69    ) -> Result<Self, ClientBuildError> {
70        let partition = pool
71            .inner
72            .registry
73            .partition(id)
74            .ok_or_else(|| ClientBuildError::invalid_partition(id))?;
75        Ok(Self {
76            pool: pool.clone(),
77            partition,
78        })
79    }
80}
81
82impl fmt::Debug for Client {
83    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
84        f.debug_struct("Client")
85            .field("partition", &self.partition.id())
86            .field("pool", &self.pool)
87            .finish_non_exhaustive()
88    }
89}
90
91impl HttpClient for Client {
92    fn validate_base_client_config(
93        &self,
94        _: &RuntimeComponentsBuilder,
95        _: &ConfigBag,
96    ) -> Result<(), BoxError> {
97        self.pool
98            .inner
99            .transport
100            .initialize_for_partition(&self.partition)
101    }
102
103    fn http_connector(
104        &self,
105        settings: &HttpConnectorSettings,
106        components: &RuntimeComponents,
107    ) -> SharedHttpConnector {
108        let connect_timeout = settings.connect_timeout();
109        let read_timeout = settings.read_timeout();
110        let sleep = components.sleep_impl().or_else(default_async_sleep);
111        let time_source = components.time_source().unwrap_or_default();
112
113        SharedHttpConnector::new(PoolConnector {
114            pool: self.pool.clone(),
115            partition: self.partition.clone(),
116            connect_timeout,
117            read_timeout,
118            sleep,
119            time_source,
120        })
121    }
122
123    fn connector_metadata(&self) -> Option<ConnectorMetadata> {
124        Some(ConnectorMetadata::new("hyper", Some(Cow::Borrowed("1.x"))))
125    }
126}
127
128/// Error returned when a [`Client`] cannot be resolved from a pool.
129#[derive(Clone, Debug, Eq, PartialEq)]
130pub struct ClientBuildError {
131    /// Construction failure represented without exposing the internal enum.
132    kind: ClientBuildErrorKind,
133}
134
135impl ClientBuildError {
136    /// Creates an unresolved-partition error.
137    fn invalid_partition(partition: PartitionId) -> Self {
138        Self {
139            kind: ClientBuildErrorKind::InvalidPartition(partition),
140        }
141    }
142
143    /// Returns the unresolved partition identity, when applicable.
144    pub fn partition(&self) -> Option<PartitionId> {
145        match &self.kind {
146            ClientBuildErrorKind::InvalidPartition(partition) => Some(*partition),
147        }
148    }
149}
150
151impl fmt::Display for ClientBuildError {
152    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
153        match &self.kind {
154            ClientBuildErrorKind::InvalidPartition(partition) => {
155                write!(
156                    f,
157                    "connection pool does not contain partition {partition:?}"
158                )
159            }
160        }
161    }
162}
163
164impl Error for ClientBuildError {}
165
166/// Internal categories represented by [`ClientBuildError`].
167#[derive(Clone, Debug, Eq, PartialEq)]
168enum ClientBuildErrorKind {
169    /// The pool retained no partition with this identity.
170    InvalidPartition(PartitionId),
171}
172
173/// Per-operation Smithy adapter over one resolved pool partition.
174///
175/// [`HttpClient::http_connector`] combines the retained pool and partition
176/// with that operation's timeout settings. Each call validates timeout timer
177/// availability, converts the Smithy request, and routes it through the shared
178/// pool. The connect timeout covers a newly started transport operation; the
179/// read timeout covers pool acquisition and dispatch through response headers.
180/// The adapter neither creates another pool nor resolves the partition again.
181struct PoolConnector {
182    /// Shared pool used for acquisition and dispatch.
183    pool: ConnectionPool,
184    /// Partition selected when the client was constructed.
185    partition: Arc<PartitionState>,
186    /// Maximum duration of the transport connection operation.
187    ///
188    /// Connector readiness completes before this timeout starts.
189    connect_timeout: Option<Duration>,
190    /// Maximum duration through response headers.
191    read_timeout: Option<Duration>,
192    /// Runtime timer used by operation timeouts.
193    sleep: Option<SharedAsyncSleep>,
194    /// Operation runtime clock used for request-attempt telemetry.
195    time_source: SharedTimeSource,
196}
197
198impl fmt::Debug for PoolConnector {
199    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
200        f.debug_struct("PoolConnector")
201            .field("partition", &self.partition.id())
202            .field("connect_timeout", &self.connect_timeout)
203            .field("read_timeout", &self.read_timeout)
204            .finish_non_exhaustive()
205    }
206}
207
208/// Converts and dispatches one request through the partitioned pool.
209async fn send_pool_request(
210    pool: &ConnectionPool,
211    partition: Arc<PartitionState>,
212    request: HttpRequest,
213    connect_timeout: Option<Duration>,
214    read_timeout: Option<Duration>,
215    sleep: Option<SharedAsyncSleep>,
216    attempt_telemetry: Option<super::dispatch::AttemptTelemetryInput>,
217) -> Result<HttpResponse, ConnectorError> {
218    if (connect_timeout.is_some() || read_timeout.is_some()) && sleep.is_none() {
219        return Err(ConnectorError::user(MissingAsyncSleep.into()));
220    }
221
222    let request = request
223        .try_into_http1x()
224        .map_err(|error| ConnectorError::user(error.into()))?;
225    let options = super::dispatch::RequestOptions::new(
226        connect_timeout
227            .zip(sleep.clone())
228            .map(|(duration, sleep)| super::establish::TransportTimeout::new(duration, sleep)),
229        attempt_telemetry,
230    );
231    let send = pool.send_request(partition, request, options);
232    let response =
233        timeout::maybe_timeout_future(send, read_timeout, sleep.as_ref(), TimeoutKind::Read)
234            .await
235            .map_err(downcast_error)?;
236    HttpResponse::try_from(response).map_err(|error| ConnectorError::other(error.into(), None))
237}
238
239impl HttpConnector for PoolConnector {
240    fn call(&self, request: HttpRequest) -> HttpConnectorFuture {
241        let attempt_capture = request.extension::<CaptureHttpAttemptTelemetry>().cloned();
242        let connector_call_timer = ConnectorCallTimer::start(&request, &self.time_source);
243        let attempt_telemetry = attempt_capture.as_ref().map(|capture| {
244            super::dispatch::AttemptTelemetryInput::new(capture.clone(), self.time_source.clone())
245        });
246        let pool = self.pool.clone();
247        let partition = self.partition.clone();
248        let connect_timeout = self.connect_timeout;
249        let read_timeout = self.read_timeout;
250        let sleep = self.sleep.clone();
251        HttpConnectorFuture::new(async move {
252            let result = send_pool_request(
253                &pool,
254                partition,
255                request,
256                connect_timeout,
257                read_timeout,
258                sleep,
259                attempt_telemetry,
260            )
261            .await;
262            if let Some(timer) = connector_call_timer {
263                timer.finish();
264            }
265            result
266        })
267    }
268}
269
270/// Operation timeouts were configured without a runtime timer.
271#[derive(Debug)]
272struct MissingAsyncSleep;
273
274impl fmt::Display for MissingAsyncSleep {
275    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
276        f.write_str("an async sleep implementation is required for HTTP connect or read timeouts")
277    }
278}
279
280impl Error for MissingAsyncSleep {}
281
282#[cfg(all(test, feature = "rt-tokio", not(smithy_http_client_loom)))]
283mod tests {
284    use super::*;
285
286    #[tokio::test]
287    async fn configured_timeout_without_sleep_returns_a_user_error() {
288        let pool = ConnectionPool::builder()
289            .idle_timeout(None)
290            .build_http()
291            .unwrap();
292        let client = Client::new(&pool).unwrap();
293
294        for (connect_timeout, read_timeout) in [
295            (Some(Duration::from_secs(1)), None),
296            (None, Some(Duration::from_secs(1))),
297        ] {
298            let connector = PoolConnector {
299                pool: client.pool.clone(),
300                partition: client.partition.clone(),
301                connect_timeout,
302                read_timeout,
303                sleep: None,
304                time_source: SharedTimeSource::default(),
305            };
306
307            let error = connector
308                .call(HttpRequest::get("http://example.com/").unwrap())
309                .await
310                .expect_err("a timeout without sleep unexpectedly started a request");
311            assert!(error.is_user(), "unexpected connector error: {error:?}");
312            assert_eq!(
313                "an async sleep implementation is required for HTTP connect or read timeouts",
314                std::error::Error::source(&error).unwrap().to_string()
315            );
316        }
317    }
318}