1use core::fmt;
4
5use cloud_sdk_sanitization::sanitize_bytes;
6
7use crate::authentication::{
8 AsyncAuthenticatedTransport, AuthenticatedRequest, AuthenticationScopePolicy,
9 BlockingAuthenticatedTransport, drive_async_authenticated,
10};
11use crate::operation::{
12 CheckedResponseGuard, OperationId, OperationImpact, OperationMetadata, RequestIdPolicy,
13 ResponsePolicy, ResponsePolicyError,
14};
15use crate::transport::{
16 BoundTransport, EndpointIdentity, RawResponsePolicy, RequestHeaders, ResponseBuffer,
17 TransportRequest,
18};
19
20mod body;
21mod error;
22mod service;
23mod storage;
24pub use body::{BodyReplayability, RequestBodySensitivity};
25use error::{EndpointCheckError, map_endpoint_error};
26pub use error::{PreparedExecutionError, PreparedRequestPolicyError};
27pub use service::ProviderService;
28pub use storage::{PreparationStorage, PrepareOperation};
29
30#[derive(Clone, Copy)]
32pub struct PreparedRequest<'request> {
33 request: TransportRequest<'request>,
34 service: ProviderService<'request>,
35 metadata: OperationMetadata,
36 response_policy: ResponsePolicy,
37 authentication_policy: AuthenticationScopePolicy<'request>,
38 raw_response_policy: RawResponsePolicy<'request>,
39 operation_id: Option<OperationId>,
40 body_replayability: BodyReplayability,
41 body_sensitivity: RequestBodySensitivity,
42 authorization_evidence_required: bool,
43}
44
45impl<'request> PreparedRequest<'request> {
46 pub fn new(
58 request: TransportRequest<'request>,
59 service: ProviderService<'request>,
60 metadata: OperationMetadata,
61 response_policy: ResponsePolicy,
62 authentication_policy: AuthenticationScopePolicy<'request>,
63 raw_response_policy: RawResponsePolicy<'request>,
64 body_sensitivity: RequestBodySensitivity,
65 ) -> Result<Self, PreparedRequestPolicyError> {
66 if matches!(metadata.impact(), OperationImpact::ReadOnly)
67 && !request.method().permits_direct_read_only()
68 {
69 return Err(PreparedRequestPolicyError::ReadOnlyMethodMismatch);
70 }
71 if metadata.request_id_policy() != RequestIdPolicy::Discard
72 && !raw_response_policy.admits_header("x-request-id")
73 {
74 return Err(PreparedRequestPolicyError::MissingRequestIdHeader);
75 }
76 Ok(Self {
77 request,
78 service,
79 metadata,
80 response_policy,
81 authentication_policy,
82 raw_response_policy,
83 operation_id: None,
84 body_replayability: if request.body().is_empty() {
85 BodyReplayability::Replayable
86 } else {
87 BodyReplayability::NotReplayable
88 },
89 body_sensitivity,
90 authorization_evidence_required: false,
91 })
92 }
93
94 #[must_use]
96 pub const fn with_operation_id(mut self, operation_id: OperationId) -> Self {
97 self.operation_id = Some(operation_id);
98 self
99 }
100
101 #[must_use]
106 pub const fn with_replayable_body(mut self) -> Self {
107 self.body_replayability = BodyReplayability::Replayable;
108 self
109 }
110
111 #[must_use]
116 pub const fn with_sensitive_body(mut self) -> Self {
117 self.body_sensitivity = RequestBodySensitivity::Sensitive;
118 self
119 }
120
121 #[must_use]
127 pub const fn with_required_authorization_evidence(mut self) -> Self {
128 self.authorization_evidence_required = true;
129 self
130 }
131
132 #[must_use]
134 pub const fn transport_request(self) -> TransportRequest<'request> {
135 self.request
136 }
137
138 #[must_use]
140 pub const fn service(self) -> ProviderService<'request> {
141 self.service
142 }
143
144 #[must_use]
146 pub const fn metadata(self) -> OperationMetadata {
147 self.metadata
148 }
149
150 #[must_use]
152 pub const fn response_policy(self) -> ResponsePolicy {
153 self.response_policy
154 }
155
156 #[must_use]
158 pub const fn authentication_policy(self) -> AuthenticationScopePolicy<'request> {
159 self.authentication_policy
160 }
161
162 #[must_use]
164 pub const fn raw_response_policy(self) -> RawResponsePolicy<'request> {
165 self.raw_response_policy
166 }
167
168 #[must_use]
170 pub(crate) const fn authenticated_request(self) -> AuthenticatedRequest<'request, 'request> {
171 AuthenticatedRequest::new(
172 self.request,
173 self.authentication_policy,
174 &self.raw_response_policy,
175 )
176 }
177
178 #[must_use]
180 pub const fn operation_id(self) -> Option<OperationId> {
181 self.operation_id
182 }
183
184 #[must_use]
186 pub const fn body_replayability(self) -> BodyReplayability {
187 self.body_replayability
188 }
189
190 #[must_use]
192 pub const fn body_sensitivity(self) -> RequestBodySensitivity {
193 self.body_sensitivity
194 }
195
196 #[must_use]
198 pub const fn authorization_evidence_required(self) -> bool {
199 self.authorization_evidence_required
200 }
201
202 pub(crate) fn with_request_headers<'headers>(
203 self,
204 headers: RequestHeaders<'headers>,
205 ) -> PreparedRequest<'headers>
206 where
207 'request: 'headers,
208 {
209 let request: TransportRequest<'headers> = self.request;
210 PreparedRequest {
211 request: request.with_headers(headers),
212 service: self.service,
213 metadata: self.metadata,
214 response_policy: self.response_policy,
215 authentication_policy: self.authentication_policy,
216 raw_response_policy: self.raw_response_policy,
217 operation_id: self.operation_id,
218 body_replayability: self.body_replayability,
219 body_sensitivity: self.body_sensitivity,
220 authorization_evidence_required: self.authorization_evidence_required,
221 }
222 }
223
224 pub(crate) fn has_same_retry_policy(&self, other: &Self) -> bool {
225 self.service == other.service
226 && self.metadata == other.metadata
227 && self.response_policy == other.response_policy
228 && self.authentication_policy == other.authentication_policy
229 && self.raw_response_policy == other.raw_response_policy
230 && self.operation_id == other.operation_id
231 && self.body_replayability == other.body_replayability
232 && self.body_sensitivity == other.body_sensitivity
233 && self.authorization_evidence_required == other.authorization_evidence_required
234 && self.has_same_header_policy(other)
235 }
236
237 fn has_same_header_policy(&self, other: &Self) -> bool {
238 let left = self.request.headers().as_slice();
239 let right = other.request.headers().as_slice();
240 left.len() == right.len()
241 && left
242 .iter()
243 .zip(right)
244 .all(|(left, right)| left.sensitivity() == right.sensitivity())
245 }
246
247 pub fn validate_response<'buffer>(
249 self,
250 response: ResponseBuffer<'buffer>,
251 ) -> Result<CheckedResponseGuard<'buffer>, ResponsePolicyError> {
252 self.response_policy
253 .validate(response, self.metadata.request_id_policy())
254 }
255
256 pub fn apply_response_metadata_policy(
262 self,
263 response: &mut ResponseBuffer<'_>,
264 ) -> Result<(), ResponsePolicyError> {
265 super::policy::apply_request_id_policy(response, self.metadata.request_id_policy())
266 }
267
268 pub fn execute_blocking<'buffer, T>(
270 self,
271 transport: &T,
272 response_storage: &'buffer mut [u8],
273 response_header_storage: &'buffer mut [u8],
274 ) -> Result<CheckedResponseGuard<'buffer>, PreparedExecutionError<T::Error>>
275 where
276 T: BlockingAuthenticatedTransport + BoundTransport,
277 {
278 let response = self.send_blocking(transport, response_storage, response_header_storage)?;
279 self.response_policy
280 .validate(response, self.metadata.request_id_policy())
281 .map_err(PreparedExecutionError::ResponsePolicy)
282 }
283
284 pub(crate) fn send_blocking<'buffer, T>(
285 self,
286 transport: &T,
287 response_storage: &'buffer mut [u8],
288 response_header_storage: &'buffer mut [u8],
289 ) -> Result<ResponseBuffer<'buffer>, PreparedExecutionError<T::Error>>
290 where
291 T: BlockingAuthenticatedTransport + BoundTransport,
292 {
293 if self.requires_execution_permit() {
294 sanitize_bytes(response_storage);
295 sanitize_bytes(response_header_storage);
296 return Err(PreparedExecutionError::AuthorizationRequired);
297 }
298 self.send_blocking_authorized(transport, None, response_storage, response_header_storage)
299 }
300
301 pub(crate) fn execute_blocking_authorized<'buffer, T>(
302 self,
303 transport: &T,
304 confirmed_endpoint: Option<EndpointIdentity<'_>>,
305 response_storage: &'buffer mut [u8],
306 response_header_storage: &'buffer mut [u8],
307 ) -> Result<CheckedResponseGuard<'buffer>, PreparedExecutionError<T::Error>>
308 where
309 T: BlockingAuthenticatedTransport + BoundTransport,
310 {
311 let response = self.send_blocking_authorized(
312 transport,
313 confirmed_endpoint,
314 response_storage,
315 response_header_storage,
316 )?;
317 self.response_policy
318 .validate(response, self.metadata.request_id_policy())
319 .map_err(PreparedExecutionError::ResponsePolicy)
320 }
321
322 pub(crate) fn send_blocking_authorized<'buffer, T>(
323 self,
324 transport: &T,
325 confirmed_endpoint: Option<EndpointIdentity<'_>>,
326 response_storage: &'buffer mut [u8],
327 response_header_storage: &'buffer mut [u8],
328 ) -> Result<ResponseBuffer<'buffer>, PreparedExecutionError<T::Error>>
329 where
330 T: BlockingAuthenticatedTransport + BoundTransport,
331 {
332 let mut response = ResponseBuffer::new(
333 response_storage,
334 self.raw_response_policy.max_body_bytes(),
335 response_header_storage,
336 );
337 self.verify_endpoint(transport, confirmed_endpoint)
338 .map_err(map_endpoint_error)?;
339 transport
340 .send_authenticated(self.authenticated_request(), response.writer())
341 .map_err(PreparedExecutionError::Transport)?;
342 Ok(response)
343 }
344
345 pub async fn execute_async<'transport, 'buffer, T>(
347 &'transport self,
348 transport: &'transport T,
349 response_storage: &'buffer mut [u8],
350 response_header_storage: &'buffer mut [u8],
351 ) -> Result<CheckedResponseGuard<'buffer>, PreparedExecutionError<T::Error>>
352 where
353 T: AsyncAuthenticatedTransport + BoundTransport,
354 'request: 'transport,
355 {
356 let response = self
357 .send_async(transport, response_storage, response_header_storage)
358 .await?;
359 self.response_policy
360 .validate(response, self.metadata.request_id_policy())
361 .map_err(PreparedExecutionError::ResponsePolicy)
362 }
363
364 pub(crate) async fn send_async<'transport, 'buffer, T>(
365 &'transport self,
366 transport: &'transport T,
367 response_storage: &'buffer mut [u8],
368 response_header_storage: &'buffer mut [u8],
369 ) -> Result<ResponseBuffer<'buffer>, PreparedExecutionError<T::Error>>
370 where
371 T: AsyncAuthenticatedTransport + BoundTransport,
372 'request: 'transport,
373 {
374 if self.requires_execution_permit() {
375 sanitize_bytes(response_storage);
376 sanitize_bytes(response_header_storage);
377 return Err(PreparedExecutionError::AuthorizationRequired);
378 }
379 self.send_async_authorized(transport, None, response_storage, response_header_storage)
380 .await
381 }
382
383 pub(crate) async fn execute_async_authorized<'transport, 'buffer, T>(
384 &'transport self,
385 transport: &'transport T,
386 confirmed_endpoint: Option<EndpointIdentity<'_>>,
387 response_storage: &'buffer mut [u8],
388 response_header_storage: &'buffer mut [u8],
389 ) -> Result<CheckedResponseGuard<'buffer>, PreparedExecutionError<T::Error>>
390 where
391 T: AsyncAuthenticatedTransport + BoundTransport,
392 'request: 'transport,
393 {
394 let response = self
395 .send_async_authorized(
396 transport,
397 confirmed_endpoint,
398 response_storage,
399 response_header_storage,
400 )
401 .await?;
402 self.response_policy
403 .validate(response, self.metadata.request_id_policy())
404 .map_err(PreparedExecutionError::ResponsePolicy)
405 }
406
407 pub(crate) async fn send_async_authorized<'transport, 'buffer, T>(
408 &'transport self,
409 transport: &'transport T,
410 confirmed_endpoint: Option<EndpointIdentity<'_>>,
411 response_storage: &'buffer mut [u8],
412 response_header_storage: &'buffer mut [u8],
413 ) -> Result<ResponseBuffer<'buffer>, PreparedExecutionError<T::Error>>
414 where
415 T: AsyncAuthenticatedTransport + BoundTransport,
416 'request: 'transport,
417 {
418 let mut response = ResponseBuffer::new(
419 response_storage,
420 self.raw_response_policy.max_body_bytes(),
421 response_header_storage,
422 );
423 self.verify_endpoint(transport, confirmed_endpoint)
424 .map_err(map_endpoint_error)?;
425 drive_async_authenticated(transport, self.authenticated_request(), response.writer())
426 .await
427 .map_err(|error| match error {
428 crate::transport::AsyncExecutionError::Transport(error) => {
429 PreparedExecutionError::Transport(error)
430 }
431 crate::transport::AsyncExecutionError::Response(error) => {
432 PreparedExecutionError::ResponseWriter(error)
433 }
434 })?;
435 Ok(response)
436 }
437
438 pub(crate) const fn requires_execution_permit(self) -> bool {
439 !self.request.method().permits_direct_read_only()
440 || !matches!(self.metadata.impact(), OperationImpact::ReadOnly)
441 || matches!(self.metadata.cost_intent(), super::CostIntent::MayIncurCost)
442 }
443
444 fn verify_endpoint<T>(
445 self,
446 transport: &T,
447 confirmed_endpoint: Option<EndpointIdentity<'_>>,
448 ) -> Result<(), EndpointCheckError>
449 where
450 T: BoundTransport,
451 {
452 let actual = transport
453 .endpoint_identity()
454 .map_err(EndpointCheckError::Invalid)?;
455 match confirmed_endpoint {
456 Some(expected) if actual == expected => Ok(()),
457 Some(_) => Err(EndpointCheckError::Mismatch),
458 None => self
459 .service
460 .endpoint_policy()
461 .verify(actual)
462 .map_err(|_| EndpointCheckError::Mismatch),
463 }
464 }
465}
466
467impl fmt::Debug for PreparedRequest<'_> {
468 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
469 formatter
470 .debug_struct("PreparedRequest")
471 .field("request", &self.request)
472 .field("service", &self.service)
473 .field("metadata", &self.metadata)
474 .field("response_policy", &self.response_policy)
475 .field("authentication_policy", &self.authentication_policy)
476 .field("raw_response_policy", &self.raw_response_policy)
477 .field("operation_id", &self.operation_id)
478 .field("body_replayability", &self.body_replayability)
479 .field("body_sensitivity", &self.body_sensitivity)
480 .field(
481 "authorization_evidence_required",
482 &self.authorization_evidence_required,
483 )
484 .finish()
485 }
486}