cloud_sdk_reqwest/blocking/
raw.rs1use 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#[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}