cloud_sdk/pagination/header_cursor/execution/
local_async.rs1use cloud_sdk_sanitization::SecretBuffer;
2
3use crate::authentication::LocalAsyncAuthenticatedTransport;
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 #[allow(clippy::too_many_arguments)]
19 pub async fn execute_local_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: LocalAsyncAuthenticatedTransport + 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_local_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 #[allow(clippy::too_many_arguments)]
66 pub async fn execute_local_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: LocalAsyncAuthenticatedTransport + 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_local_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_local_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: LocalAsyncAuthenticatedTransport + 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: LocalAsyncAuthenticatedTransport + 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_local_async(transport, response_storage, response_header_storage)
248 .await
249 .map_err(HeaderCursorExecutionError::Prepared)
250}