Skip to main content

cloud_sdk/pagination/header_cursor/execution/
asynchronous.rs

1use cloud_sdk_sanitization::SecretBuffer;
2
3use crate::authentication::AsyncAuthenticatedTransport;
4use crate::buffer::write_u64;
5use crate::operation::PreparedExecutionError;
6use crate::pagination::{PaginationCursor, PaginationError, PaginationLimits};
7use crate::transport::{
8    BoundTransport, EndpointIdentity, MAX_REQUEST_HEADERS, RequestHeader, RequestHeaders,
9};
10
11use super::{
12    HeaderCursorContinuation, HeaderCursorExecutionError, HeaderCursorPage, HeaderCursorSession,
13    clear_execution_buffers, finish_page,
14};
15
16impl HeaderCursorSession<'_, '_> {
17    /// Executes the initial executor-neutral async request and decodes its response.
18    #[allow(clippy::too_many_arguments)]
19    pub async fn execute_async<'response, 'cursor, 'endpoint, T>(
20        &self,
21        transport: &'endpoint T,
22        response_storage: &'response mut [u8],
23        response_header_storage: &'response mut [u8],
24        decimal_scratch: &mut [u8],
25        transfer_scratch: &mut [u8],
26        cursor_destination: &'cursor mut [u8],
27        limits: PaginationLimits,
28    ) -> Result<
29        HeaderCursorPage<'response, 'cursor, 'endpoint, '_, '_, '_>,
30        HeaderCursorExecutionError<T::Error>,
31    >
32    where
33        T: AsyncAuthenticatedTransport + BoundTransport,
34    {
35        clear_execution_buffers(
36            response_storage,
37            response_header_storage,
38            decimal_scratch,
39            transfer_scratch,
40            cursor_destination,
41        );
42        let endpoint = transport.endpoint_identity().map_err(|error| {
43            HeaderCursorExecutionError::Prepared(PreparedExecutionError::EndpointIdentity(error))
44        })?;
45        execute_async(
46            self,
47            None,
48            endpoint,
49            transport,
50            response_storage,
51            response_header_storage,
52            decimal_scratch,
53            transfer_scratch,
54            cursor_destination,
55            limits,
56        )
57        .await
58    }
59}
60
61impl<'cursor, 'endpoint, 'session, 'request, 'policy>
62    HeaderCursorContinuation<'cursor, 'endpoint, 'session, 'request, 'policy>
63{
64    /// Executes the next executor-neutral async request with retained context.
65    #[allow(clippy::too_many_arguments)]
66    pub async fn execute_async<'response, 'next, T>(
67        &self,
68        transport: &T,
69        response_storage: &'response mut [u8],
70        response_header_storage: &'response mut [u8],
71        decimal_scratch: &mut [u8],
72        transfer_scratch: &mut [u8],
73        cursor_destination: &'next mut [u8],
74        limits: PaginationLimits,
75    ) -> Result<
76        HeaderCursorPage<'response, 'next, 'endpoint, 'session, 'request, 'policy>,
77        HeaderCursorExecutionError<T::Error>,
78    >
79    where
80        T: AsyncAuthenticatedTransport + BoundTransport,
81    {
82        clear_execution_buffers(
83            response_storage,
84            response_header_storage,
85            decimal_scratch,
86            transfer_scratch,
87            cursor_destination,
88        );
89        let endpoint = transport.endpoint_identity().map_err(|error| {
90            HeaderCursorExecutionError::Prepared(PreparedExecutionError::EndpointIdentity(error))
91        })?;
92        if endpoint != self.endpoint {
93            return Err(HeaderCursorExecutionError::Pagination(
94                PaginationError::EndpointMismatch,
95            ));
96        }
97        execute_async(
98            self.session,
99            Some(&self.cursor),
100            self.endpoint,
101            transport,
102            response_storage,
103            response_header_storage,
104            decimal_scratch,
105            transfer_scratch,
106            cursor_destination,
107            limits,
108        )
109        .await
110    }
111}
112
113#[allow(clippy::too_many_arguments)]
114async fn execute_async<'response, 'cursor, 'endpoint, 'session, 'request, 'policy, T>(
115    session: &'session HeaderCursorSession<'request, 'policy>,
116    cursor: Option<&PaginationCursor<'_>>,
117    endpoint: EndpointIdentity<'endpoint>,
118    transport: &T,
119    response_storage: &'response mut [u8],
120    response_header_storage: &'response mut [u8],
121    decimal_scratch: &mut [u8],
122    transfer_scratch: &mut [u8],
123    cursor_destination: &'cursor mut [u8],
124    limits: PaginationLimits,
125) -> Result<
126    HeaderCursorPage<'response, 'cursor, 'endpoint, 'session, 'request, 'policy>,
127    HeaderCursorExecutionError<T::Error>,
128>
129where
130    T: AsyncAuthenticatedTransport + BoundTransport,
131{
132    let response = match cursor {
133        Some(cursor) => {
134            execute_with_value(
135                session,
136                Some(cursor.as_bytes()),
137                transport,
138                response_storage,
139                response_header_storage,
140                decimal_scratch,
141            )
142            .await
143        }
144        None => {
145            execute_with_value(
146                session,
147                None,
148                transport,
149                response_storage,
150                response_header_storage,
151                decimal_scratch,
152            )
153            .await
154        }
155    }?;
156    finish_page(
157        session,
158        endpoint,
159        response,
160        transfer_scratch,
161        cursor_destination,
162        limits,
163    )
164}
165
166async fn execute_with_value<'response, T>(
167    session: &HeaderCursorSession<'_, '_>,
168    cursor: Option<&[u8]>,
169    transport: &T,
170    response_storage: &'response mut [u8],
171    response_header_storage: &'response mut [u8],
172    decimal_scratch: &mut [u8],
173) -> Result<crate::operation::CheckedResponseGuard<'response>, HeaderCursorExecutionError<T::Error>>
174where
175    T: AsyncAuthenticatedTransport + BoundTransport,
176{
177    let mut decimal = SecretBuffer::new(decimal_scratch);
178    let mut len = 0_usize;
179    write_u64(
180        decimal.as_mut_slice(),
181        &mut len,
182        session.policy.page_size(),
183        PaginationError::OutputTooSmall,
184    )
185    .map_err(HeaderCursorExecutionError::Pagination)?;
186    let size = core::str::from_utf8(
187        decimal
188            .as_slice()
189            .get(..len)
190            .ok_or(PaginationError::OutputTooSmall)
191            .map_err(HeaderCursorExecutionError::Pagination)?,
192    )
193    .map_err(|_| HeaderCursorExecutionError::Pagination(PaginationError::InvalidHeaderState))?;
194    let size = RequestHeader::new(session.policy.size_request().as_str(), size)
195        .map_err(|_| HeaderCursorExecutionError::Pagination(PaginationError::InvalidHeaderState))?;
196    let cursor_header = match cursor {
197        Some(value) => {
198            let value = core::str::from_utf8(value).map_err(|_| {
199                HeaderCursorExecutionError::Pagination(PaginationError::InvalidHeaderState)
200            })?;
201            Some(
202                RequestHeader::sensitive(session.policy.cursor_request().as_str(), value).map_err(
203                    |_| HeaderCursorExecutionError::Pagination(PaginationError::InvalidHeaderState),
204                )?,
205            )
206        }
207        None => None,
208    };
209    let pagination_entries = [size, cursor_header.unwrap_or(size)];
210    let pagination_len = if cursor_header.is_some() { 2 } else { 1 };
211    let pagination = pagination_entries
212        .get(..pagination_len)
213        .ok_or(PaginationError::InvalidHeaderState)
214        .map_err(HeaderCursorExecutionError::Pagination)?;
215    let base_headers = session.prepared.transport_request().headers();
216    let base = base_headers.as_slice();
217    let count = base
218        .len()
219        .checked_add(pagination.len())
220        .ok_or(PaginationError::RequestHeaderConflict)
221        .map_err(HeaderCursorExecutionError::Pagination)?;
222    if count > MAX_REQUEST_HEADERS {
223        return Err(HeaderCursorExecutionError::Pagination(
224            PaginationError::RequestHeaderConflict,
225        ));
226    }
227    let mut entries = [size; MAX_REQUEST_HEADERS];
228    entries
229        .get_mut(..base.len())
230        .ok_or(PaginationError::RequestHeaderConflict)
231        .map_err(HeaderCursorExecutionError::Pagination)?
232        .copy_from_slice(base);
233    entries
234        .get_mut(base.len()..count)
235        .ok_or(PaginationError::RequestHeaderConflict)
236        .map_err(HeaderCursorExecutionError::Pagination)?
237        .copy_from_slice(pagination);
238    let selected = entries
239        .get(..count)
240        .ok_or(PaginationError::RequestHeaderConflict)
241        .map_err(HeaderCursorExecutionError::Pagination)?;
242    let headers = RequestHeaders::new(selected).map_err(|_| {
243        HeaderCursorExecutionError::Pagination(PaginationError::RequestHeaderConflict)
244    })?;
245    let prepared = session.prepared.with_request_headers(headers);
246    prepared
247        .execute_async(transport, response_storage, response_header_storage)
248        .await
249        .map_err(HeaderCursorExecutionError::Prepared)
250}