Skip to main content

cloud_sdk_reqwest/blocking/
raw.rs

1use core::fmt;
2
3use cloud_sdk::transport::{
4    BlockingRawHttpExecutor, BoundTransport, EndpointIdentity, EndpointIdentityError,
5    RawResponsePolicy, ResponseStorageSanitizer, ResponseWriter, TransportFailure,
6    TransportRequest,
7};
8use cloud_sdk_sanitization::sanitize_bytes;
9use http::header::HeaderValue;
10
11use crate::shared::{HttpsEndpoint, RawHttpError, RawHyperClient, RawTransportFailure};
12
13/// Raw blocking HTTP executor with no implicit authentication or provider policy.
14#[derive(Clone)]
15pub struct RawBlockingClient {
16    inner: RawHyperClient,
17    endpoint: HttpsEndpoint,
18}
19
20impl RawBlockingClient {
21    pub(super) const fn new(inner: RawHyperClient, endpoint: HttpsEndpoint) -> Self {
22        Self { inner, endpoint }
23    }
24
25    fn execute_inner(
26        &self,
27        request: TransportRequest<'_>,
28        policy: RawResponsePolicy<'_>,
29        response_writer: &mut ResponseWriter<'_>,
30    ) -> Result<(), RawTransportFailure> {
31        if tokio::runtime::Handle::try_current().is_ok() {
32            return Err(TransportFailure::not_sent(
33                RawHttpError::BlockingRuntimeContext,
34            ));
35        }
36        let runtime = tokio::runtime::Builder::new_current_thread()
37            .enable_all()
38            .build()
39            .map_err(|_| TransportFailure::not_sent(RawHttpError::RuntimeInitializationFailed))?;
40        let mut attempt = response_writer
41            .begin_attempt()
42            .map_err(|_| TransportFailure::not_sent(RawHttpError::ResponseAlreadyCommitted))?;
43        let completion = runtime.block_on(self.inner.execute(request, policy, &mut attempt))?;
44        let status = completion.status();
45        attempt.commit_completion(completion).map_err(|_| {
46            TransportFailure::response_started_with_status(
47                status,
48                RawHttpError::ResponseCommitFailed,
49            )
50        })
51    }
52
53    pub(crate) fn execute_authenticated(
54        &self,
55        request: TransportRequest<'_>,
56        policy: RawResponsePolicy<'_>,
57        authorization: HeaderValue,
58        response_writer: &mut ResponseWriter<'_>,
59    ) -> Result<(), RawTransportFailure> {
60        if tokio::runtime::Handle::try_current().is_ok() {
61            return Err(TransportFailure::not_sent(
62                RawHttpError::BlockingRuntimeContext,
63            ));
64        }
65        let runtime = tokio::runtime::Builder::new_current_thread()
66            .enable_all()
67            .build()
68            .map_err(|_| TransportFailure::not_sent(RawHttpError::RuntimeInitializationFailed))?;
69        let mut attempt = response_writer
70            .begin_attempt()
71            .map_err(|_| TransportFailure::not_sent(RawHttpError::ResponseAlreadyCommitted))?;
72        let completion = runtime.block_on(self.inner.execute_authenticated(
73            request,
74            policy,
75            authorization,
76            &mut attempt,
77        ))?;
78        let status = completion.status();
79        attempt.commit_completion(completion).map_err(|_| {
80            TransportFailure::response_started_with_status(
81                status,
82                RawHttpError::ResponseCommitFailed,
83            )
84        })
85    }
86}
87
88impl BlockingRawHttpExecutor for RawBlockingClient {
89    type Error = RawTransportFailure;
90
91    fn execute(
92        &self,
93        request: TransportRequest<'_>,
94        policy: RawResponsePolicy<'_>,
95        response: &mut ResponseWriter<'_>,
96    ) -> Result<(), Self::Error> {
97        self.execute_inner(request, policy, response)
98    }
99}
100
101impl ResponseStorageSanitizer for RawBlockingClient {
102    fn sanitize_response_storage(&self, response_storage: &mut [u8]) {
103        sanitize_bytes(response_storage);
104    }
105}
106
107impl BoundTransport for RawBlockingClient {
108    fn endpoint_identity(&self) -> Result<EndpointIdentity<'_>, EndpointIdentityError> {
109        self.endpoint.identity()
110    }
111}
112
113impl fmt::Debug for RawBlockingClient {
114    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
115        formatter
116            .debug_struct("RawBlockingClient")
117            .field("endpoint", &"[redacted]")
118            .finish_non_exhaustive()
119    }
120}