1use std::{
16 any::type_name,
17 borrow::Cow,
18 fmt::Debug,
19 num::NonZeroUsize,
20 sync::{
21 Arc,
22 atomic::{AtomicU64, Ordering},
23 },
24 time::Duration,
25};
26
27use bytes::{Bytes, BytesMut};
28use bytesize::ByteSize;
29use eyeball::SharedObservable;
30use http::Method;
31use matrix_sdk_base::SendOutsideWasm;
32use ruma::api::{
33 OutgoingRequest, OutgoingRequestExt, SupportedVersions,
34 auth_scheme::{self, AuthScheme, SendAccessToken},
35 error::{FromHttpResponseError, IntoHttpError},
36 path_builder,
37};
38use tokio::sync::{Semaphore, SemaphorePermit};
39use tracing::{Instrument, debug, error, field::debug, trace};
40
41use crate::{HttpResult, config::RequestConfig, error::HttpError};
42
43#[cfg(not(target_family = "wasm"))]
44mod native;
45#[cfg(target_family = "wasm")]
46mod wasm;
47
48#[cfg(not(target_family = "wasm"))]
49pub(crate) use native::HttpSettings;
50
51pub(crate) const DEFAULT_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
52
53#[derive(Clone, Debug)]
54struct MaybeSemaphore(Arc<Option<Semaphore>>);
55
56#[allow(dead_code)] struct MaybeSemaphorePermit<'a>(Option<SemaphorePermit<'a>>);
58
59impl MaybeSemaphore {
60 fn new(max: Option<NonZeroUsize>) -> Self {
61 let inner = max.map(|i| Semaphore::new(i.into()));
62 MaybeSemaphore(Arc::new(inner))
63 }
64
65 async fn acquire(&self) -> MaybeSemaphorePermit<'_> {
66 match self.0.as_ref() {
67 Some(inner) => {
68 MaybeSemaphorePermit(inner.acquire().await.ok())
71 }
72 None => MaybeSemaphorePermit(None),
73 }
74 }
75}
76
77#[derive(Clone, Debug)]
78pub(crate) struct HttpClient {
79 pub(crate) inner: reqwest::Client,
80 pub(crate) request_config: RequestConfig,
81 concurrent_request_semaphore: MaybeSemaphore,
82 next_request_id: Arc<AtomicU64>,
83}
84
85impl HttpClient {
86 pub(crate) fn new(inner: reqwest::Client, request_config: RequestConfig) -> Self {
87 HttpClient {
88 inner,
89 request_config,
90 concurrent_request_semaphore: MaybeSemaphore::new(
91 request_config.max_concurrent_requests,
92 ),
93 next_request_id: AtomicU64::new(0).into(),
94 }
95 }
96
97 fn get_request_id(&self) -> String {
98 let request_id = self.next_request_id.fetch_add(1, Ordering::SeqCst);
99 format!("REQ-{request_id}")
100 }
101
102 fn serialize_request<R>(
103 &self,
104 request: R,
105 config: RequestConfig,
106 homeserver: String,
107 access_token: Option<&str>,
108 path_builder_input: <R::PathBuilder as path_builder::PathBuilder>::Input<'_>,
109 ) -> Result<http::Request<Bytes>, IntoHttpError>
110 where
111 R: OutgoingRequest + Debug,
112 R::Authentication: SupportedAuthScheme,
113 {
114 trace!(request_type = type_name::<R>(), "Serializing request");
115
116 let send_access_token = match access_token {
117 Some(access_token) => match (config.force_auth, config.skip_auth) {
118 (true, true) | (true, false) => SendAccessToken::Always(access_token),
119 (false, true) => SendAccessToken::None,
120 (false, false) => SendAccessToken::IfRequired(access_token),
121 },
122 None => SendAccessToken::None,
123 };
124 let authentication_input = R::Authentication::authentication_input(send_access_token);
125
126 let request = request
127 .try_into_http_request::<BytesMut>(
128 &homeserver,
129 authentication_input,
130 path_builder_input,
131 )?
132 .map(|body| body.freeze());
133
134 Ok(request)
135 }
136
137 pub fn send<R>(
138 &self,
139 request: R,
140 config: Option<RequestConfig>,
141 homeserver: String,
142 access_token: Option<&str>,
143 path_builder_input: <R::PathBuilder as path_builder::PathBuilder>::Input<'_>,
144 send_progress: SharedObservable<TransmissionProgress>,
145 ) -> impl Future<Output = Result<R::IncomingResponse, HttpError>>
146 where
147 R: OutgoingRequest + Debug,
148 R::Authentication: SupportedAuthScheme,
149 HttpError: From<FromHttpResponseError<R::EndpointError>>,
150 {
151 fn make_span(client: &HttpClient, config: &RequestConfig) -> tracing::Span {
154 tracing::info_span!(
155 "send",
156 uri = tracing::field::Empty,
157 ?config,
158 method = tracing::field::Empty,
159 request_id = client.get_request_id(),
160 request_size = tracing::field::Empty,
161 request_duration = tracing::field::Empty,
162 status = tracing::field::Empty,
163 response_size = tracing::field::Empty,
164 sentry_event_id = tracing::field::Empty
165 )
166 }
167 fn record_request_uri_and_size(request: &http::Request<Bytes>) {
168 let method = request.method();
169
170 let mut uri_parts = request.uri().clone().into_parts();
171
172 if let Some(path_and_query) = &mut uri_parts.path_and_query {
175 *path_and_query =
176 path_and_query.path().try_into().expect("path is valid PathAndQuery");
177 }
178
179 let uri = http::Uri::from_parts(uri_parts).expect("created from valid URI");
180
181 let span = tracing::Span::current();
182 span.record("method", debug(method)).record("uri", uri.to_string());
183
184 if [Method::POST, Method::PUT, Method::PATCH].contains(method) {
187 let request_size = request.body().len().try_into().unwrap_or(u64::MAX);
188 span.record(
189 "request_size",
190 ByteSize(request_size).display().si_short().to_string(),
191 );
192 }
193 }
194 fn log_got_response() {
197 debug!("Got response");
198 }
199 fn log_error(e: &HttpError) {
200 error!("Error while sending request: {e:?}");
201 }
202
203 let config = match config {
204 Some(config) => config,
205 None => self.request_config,
206 };
207
208 async move {
209 let request = self
210 .serialize_request(request, config, homeserver, access_token, path_builder_input)
211 .map_err(HttpError::IntoHttp)?;
212 record_request_uri_and_size(&request);
213
214 let _handle = self.concurrent_request_semaphore.acquire().await;
216
217 match Box::pin(self.send_request::<R>(request, config, send_progress)).await {
220 Ok(response) => {
221 log_got_response();
222 Ok(response)
223 }
224 Err(e) => {
225 log_error(&e);
226 Err(e)
227 }
228 }
229 }
230 .instrument(make_span(self, &config))
231 }
232}
233
234#[derive(Clone, Copy, Debug, Default)]
236pub struct TransmissionProgress {
237 pub current: usize,
239 pub total: usize,
241}
242
243async fn response_to_http_response(
244 mut response: reqwest::Response,
245) -> Result<http::Response<Bytes>, reqwest::Error> {
246 let status = response.status();
247
248 let mut http_builder = http::Response::builder().status(status);
249 let headers = http_builder.headers_mut().expect("Can't get the response builder headers");
250
251 for (k, v) in response.headers_mut().drain() {
252 if let Some(key) = k {
253 headers.insert(key, v);
254 }
255 }
256
257 let body = response.bytes().await?;
258
259 Ok(http_builder.body(body).expect("Can't construct a response using the given body"))
260}
261
262pub trait SupportedAuthScheme: AuthScheme {
267 fn authentication_input(access_token: SendAccessToken<'_>) -> Self::Input<'_>;
269}
270
271impl SupportedAuthScheme for auth_scheme::NoAccessToken {
272 fn authentication_input(access_token: SendAccessToken<'_>) -> Self::Input<'_> {
273 access_token
274 }
275}
276
277impl SupportedAuthScheme for auth_scheme::AccessToken {
278 fn authentication_input(access_token: SendAccessToken<'_>) -> Self::Input<'_> {
279 access_token
280 }
281}
282
283impl SupportedAuthScheme for auth_scheme::AccessTokenOptional {
284 fn authentication_input(access_token: SendAccessToken<'_>) -> Self::Input<'_> {
285 access_token
286 }
287}
288
289impl SupportedAuthScheme for auth_scheme::AppserviceToken {
290 fn authentication_input(access_token: SendAccessToken<'_>) -> Self::Input<'_> {
291 access_token
292 }
293}
294
295impl SupportedAuthScheme for auth_scheme::AppserviceTokenOptional {
296 fn authentication_input(access_token: SendAccessToken<'_>) -> Self::Input<'_> {
297 access_token
298 }
299}
300
301impl SupportedAuthScheme for auth_scheme::NoAuthentication {
302 fn authentication_input(_access_token: SendAccessToken<'_>) -> Self::Input<'_> {}
303}
304
305pub trait SupportedPathBuilder: path_builder::PathBuilder {
311 fn get_path_builder_input(
314 client: &crate::Client,
315 skip_auth: bool,
316 ) -> impl Future<Output = HttpResult<Self::Input<'static>>> + SendOutsideWasm;
317}
318
319impl SupportedPathBuilder for path_builder::VersionHistory {
320 async fn get_path_builder_input(
321 client: &crate::Client,
322 skip_auth: bool,
323 ) -> HttpResult<Cow<'static, SupportedVersions>> {
324 if !client.auth_ctx().has_valid_access_token() {
329 if let Ok(Some(versions)) = client.supported_versions_cached().await {
331 return Ok(Cow::Owned(versions));
332 }
333
334 let response = client.fetch_server_versions_inner(true, None).await?;
337
338 Ok(Cow::Owned(response.as_supported_versions()))
339 } else if skip_auth {
340 let cached_versions = client.supported_versions_cached().await;
341
342 let versions = if let Ok(Some(versions)) = cached_versions {
343 versions
344 } else {
345 let request_config = RequestConfig::default().retry_limit(5).skip_auth();
348 let response =
349 client.fetch_server_versions_inner(true, Some(request_config)).await?;
350
351 response.as_supported_versions()
352 };
353
354 Ok(Cow::Owned(versions))
355 } else {
356 client.supported_versions_inner(true).await.map(Cow::Owned)
357 }
358 }
359}
360
361impl SupportedPathBuilder for path_builder::SinglePath {
362 async fn get_path_builder_input(_client: &crate::Client, _skip_auth: bool) -> HttpResult<()> {
363 Ok(())
364 }
365}
366
367#[cfg(all(test, not(target_family = "wasm")))]
368mod tests {
369 use std::{
370 num::NonZeroUsize,
371 sync::{
372 Arc,
373 atomic::{AtomicU8, Ordering},
374 },
375 time::Duration,
376 };
377
378 use matrix_sdk_common::executor::spawn;
379 use matrix_sdk_test::{async_test, test_json};
380 use wiremock::{
381 Mock, Request, ResponseTemplate,
382 matchers::{method, path},
383 };
384
385 use crate::{
386 http_client::RequestConfig,
387 test_utils::{set_client_session, test_client_builder_with_server},
388 };
389
390 #[async_test]
391 async fn test_ensure_concurrent_request_limit_is_observed() {
392 let (client_builder, server) = test_client_builder_with_server().await;
393 let client = client_builder
394 .request_config(RequestConfig::default().max_concurrent_requests(NonZeroUsize::new(5)))
395 .build()
396 .await
397 .unwrap();
398
399 set_client_session(&client).await;
400
401 let counter = Arc::new(AtomicU8::new(0));
402 let inner_counter = counter.clone();
403
404 Mock::given(method("GET"))
405 .and(path("/_matrix/client/versions"))
406 .respond_with(ResponseTemplate::new(200).set_body_json(&*test_json::VERSIONS))
407 .mount(&server)
408 .await;
409
410 Mock::given(method("GET"))
411 .and(path("_matrix/client/r0/account/whoami"))
412 .respond_with(move |_req: &Request| {
413 inner_counter.fetch_add(1, Ordering::SeqCst);
414 ResponseTemplate::new(200).set_delay(Duration::from_secs(60))
416 })
417 .mount(&server)
418 .await;
419
420 let bg_task = spawn(async move {
421 futures_util::future::join_all((0..10).map(|_| client.whoami())).await
422 });
423
424 tokio::time::sleep(Duration::from_millis(300)).await;
426
427 assert_eq!(
428 counter.load(Ordering::SeqCst),
429 5,
430 "More requests passed than the limit we configured"
431 );
432 bg_task.abort();
433 }
434
435 #[async_test]
436 async fn test_ensure_no_max_concurrent_request_does_not_limit() {
437 let (client_builder, server) = test_client_builder_with_server().await;
438 let client = client_builder
439 .request_config(RequestConfig::default().max_concurrent_requests(None))
440 .build()
441 .await
442 .unwrap();
443
444 set_client_session(&client).await;
445
446 let counter = Arc::new(AtomicU8::new(0));
447 let inner_counter = counter.clone();
448
449 Mock::given(method("GET"))
450 .and(path("/_matrix/client/versions"))
451 .respond_with(ResponseTemplate::new(200).set_body_json(&*test_json::VERSIONS))
452 .mount(&server)
453 .await;
454
455 Mock::given(method("GET"))
456 .and(path("_matrix/client/r0/account/whoami"))
457 .respond_with(move |_req: &Request| {
458 inner_counter.fetch_add(1, Ordering::SeqCst);
459 ResponseTemplate::new(200).set_delay(Duration::from_secs(60))
460 })
461 .mount(&server)
462 .await;
463
464 let bg_task = spawn(async move {
465 futures_util::future::join_all((0..254).map(|_| client.whoami())).await
466 });
467
468 tokio::time::sleep(Duration::from_secs(1)).await;
470
471 assert_eq!(counter.load(Ordering::SeqCst), 254, "Not all requests passed through");
472 bg_task.abort();
473 }
474}