aws_smithy_http_client/client/pool/
client.rs1use 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#[derive(Clone)]
40pub struct Client {
41 pool: ConnectionPool,
43 partition: Arc<PartitionState>,
45}
46
47impl Client {
48 pub fn new(pool: &ConnectionPool) -> Result<Self, ClientBuildError> {
55 Self::from_partition(pool, PartitionId::ANONYMOUS)
56 }
57
58 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#[derive(Clone, Debug, Eq, PartialEq)]
130pub struct ClientBuildError {
131 kind: ClientBuildErrorKind,
133}
134
135impl ClientBuildError {
136 fn invalid_partition(partition: PartitionId) -> Self {
138 Self {
139 kind: ClientBuildErrorKind::InvalidPartition(partition),
140 }
141 }
142
143 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#[derive(Clone, Debug, Eq, PartialEq)]
168enum ClientBuildErrorKind {
169 InvalidPartition(PartitionId),
171}
172
173struct PoolConnector {
182 pool: ConnectionPool,
184 partition: Arc<PartitionState>,
186 connect_timeout: Option<Duration>,
190 read_timeout: Option<Duration>,
192 sleep: Option<SharedAsyncSleep>,
194 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
208async 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#[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}