Skip to main content

cloud_sdk/operation/
prepared.rs

1//! Prepared operation storage, endpoint binding, and execution.
2
3use 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/// Complete request, endpoint, operation metadata, and response policy.
31#[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    /// Creates a complete prepared request after checking cross-policy invariants.
47    ///
48    /// # Errors
49    ///
50    /// Returns [`PreparedRequestPolicyError::MissingRequestIdHeader`] when
51    /// operation metadata protects or retains request IDs but the raw response
52    /// policy does not admit `x-request-id`. Returns
53    /// [`PreparedRequestPolicyError::ReadOnlyMethodMismatch`] when read-only
54    /// metadata is paired with any method other than `GET` or `HEAD`.
55    /// `body_sensitivity` is mandatory so provider implementations cannot omit
56    /// the confidential-body review and silently receive a public default.
57    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    /// Binds a validated provider operation identifier to this request.
95    #[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    /// Marks the immutable prepared body snapshot as byte-for-byte replayable.
102    ///
103    /// Providers must call this only after preparation has completed and when
104    /// the borrowed body bytes cannot change for the prepared request lifetime.
105    #[must_use]
106    pub const fn with_replayable_body(mut self) -> Self {
107        self.body_replayability = BodyReplayability::Replayable;
108        self
109    }
110
111    /// Upgrades the explicit body classification to sensitive.
112    ///
113    /// This cannot downgrade a sensitive body. Providers should classify the
114    /// body at construction; this helper supports reviewed wrapper policies.
115    #[must_use]
116    pub const fn with_sensitive_body(mut self) -> Self {
117        self.body_sensitivity = RequestBodySensitivity::Sensitive;
118        self
119    }
120
121    /// Requires provider-owned authorization evidence during plan construction.
122    ///
123    /// This marker can only tighten a prepared request. Generic plan builders
124    /// reject marked requests; provider wrappers must use the evidence-aware
125    /// digest builder and retain their typed dispatch validation.
126    #[must_use]
127    pub const fn with_required_authorization_evidence(mut self) -> Self {
128        self.authorization_evidence_required = true;
129        self
130    }
131
132    /// Returns the validated transport request.
133    #[must_use]
134    pub const fn transport_request(self) -> TransportRequest<'request> {
135        self.request
136    }
137
138    /// Returns the bound provider service.
139    #[must_use]
140    pub const fn service(self) -> ProviderService<'request> {
141        self.service
142    }
143
144    /// Returns complete safety and retry metadata.
145    #[must_use]
146    pub const fn metadata(self) -> OperationMetadata {
147        self.metadata
148    }
149
150    /// Returns complete checked-response policy.
151    #[must_use]
152    pub const fn response_policy(self) -> ResponsePolicy {
153        self.response_policy
154    }
155
156    /// Returns the complete provider-owned authentication-scope policy.
157    #[must_use]
158    pub const fn authentication_policy(self) -> AuthenticationScopePolicy<'request> {
159        self.authentication_policy
160    }
161
162    /// Returns the complete status-class raw response policy.
163    #[must_use]
164    pub const fn raw_response_policy(self) -> RawResponsePolicy<'request> {
165        self.raw_response_policy
166    }
167
168    /// Returns the request with its mandatory authentication and raw wire policy.
169    #[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    /// Returns the provider operation identifier when one was bound.
179    #[must_use]
180    pub const fn operation_id(self) -> Option<OperationId> {
181        self.operation_id
182    }
183
184    /// Returns the explicit request-body replay capability.
185    #[must_use]
186    pub const fn body_replayability(self) -> BodyReplayability {
187        self.body_replayability
188    }
189
190    /// Returns the provider-declared request-body sensitivity.
191    #[must_use]
192    pub const fn body_sensitivity(self) -> RequestBodySensitivity {
193        self.body_sensitivity
194    }
195
196    /// Reports whether plan construction requires provider-owned evidence.
197    #[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    /// Applies the complete prepared response policy without executing transport.
248    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    /// Applies operation-owned metadata policy before provider error decoding.
257    ///
258    /// This is the error-status counterpart to [`Self::validate_response`].
259    /// It extracts and protects, discards, or admits retention of the provider
260    /// request identifier without applying success-status or body policy.
261    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    /// Verifies endpoint identity, executes once, and validates the response.
269    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    /// Async equivalent of [`Self::execute_blocking`] without owning an executor.
346    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}