1use core::fmt;
4use core::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
5
6use cloud_sdk::Method;
7use cloud_sdk::authentication::{
8 AsyncAuthenticatedTransport, AuthenticatedRequest, BlockingAuthenticatedTransport,
9};
10use cloud_sdk::transport::{
11 AsyncResponseStaging, AsyncTransport, BlockingTransport, BoundTransport, EndpointIdentity,
12 EndpointIdentityError, RequestHeaders, RequestTarget, ResponseCompletion, ResponseWriter,
13 TransportRequest,
14};
15
16use crate::mock::{MockResponseSink, stage_response};
17use crate::{MAX_DYNAMIC_RECORDS, MockError, RequestRecordSlot, ResponseFixture};
18
19#[derive(Clone, Copy)]
21pub struct DynamicRequest<'request> {
22 sequence: usize,
23 request: TransportRequest<'request>,
24}
25
26impl<'request> DynamicRequest<'request> {
27 const fn new(sequence: usize, request: TransportRequest<'request>) -> Self {
28 Self { sequence, request }
29 }
30
31 #[must_use]
33 pub const fn sequence(self) -> usize {
34 self.sequence
35 }
36
37 #[must_use]
39 pub const fn method(self) -> Method {
40 self.request.method()
41 }
42
43 #[must_use]
45 pub const fn target(self) -> RequestTarget<'request> {
46 self.request.target()
47 }
48
49 #[must_use]
51 pub const fn body(self) -> &'request [u8] {
52 self.request.body()
53 }
54
55 #[must_use]
57 pub const fn headers(self) -> RequestHeaders<'request> {
58 self.request.headers()
59 }
60}
61
62impl fmt::Debug for DynamicRequest<'_> {
63 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
64 formatter
65 .debug_struct("DynamicRequest")
66 .field("sequence", &self.sequence)
67 .field("method", &self.request.method())
68 .field("target", &"[redacted]")
69 .field("body", &"[redacted]")
70 .field("headers", &self.request.headers())
71 .finish()
72 }
73}
74
75pub trait ProviderFixtureBuilder<'fixture> {
77 type Error;
79
80 fn build<'request>(
82 &self,
83 request: DynamicRequest<'request>,
84 ) -> Result<&'fixture ResponseFixture<'fixture>, Self::Error>;
85}
86
87pub struct DynamicResponder<F> {
89 responder: F,
90}
91
92impl<F> DynamicResponder<F> {
93 #[must_use]
95 pub const fn new(responder: F) -> Self {
96 Self { responder }
97 }
98}
99
100impl<'fixture, E, F> ProviderFixtureBuilder<'fixture> for DynamicResponder<F>
101where
102 F: for<'request> Fn(DynamicRequest<'request>) -> Result<&'fixture ResponseFixture<'fixture>, E>,
103{
104 type Error = E;
105
106 fn build<'request>(
107 &self,
108 request: DynamicRequest<'request>,
109 ) -> Result<&'fixture ResponseFixture<'fixture>, Self::Error> {
110 (self.responder)(request)
111 }
112}
113
114#[derive(Clone, Copy, Debug, Eq, PartialEq)]
116pub enum DynamicMockConfigError {
117 NoRecordSlots,
119 TooManyRecordSlots,
121 DirtyRecordSlot,
123}
124
125impl_static_error!(DynamicMockConfigError,
126 Self::NoRecordSlots => "dynamic mock requires at least one record slot",
127 Self::TooManyRecordSlots => "dynamic mock recording capacity exceeds the limit",
128 Self::DirtyRecordSlot => "dynamic mock record slots are not empty",
129);
130
131pub enum DynamicMockError<E> {
133 Exhausted,
135 ConcurrentRequest,
137 Builder(E),
139 Fixture(MockError),
141 CursorOverflow,
143}
144
145impl<E> fmt::Debug for DynamicMockError<E> {
146 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
147 formatter.write_str(match self {
148 Self::Exhausted => "DynamicMockError::Exhausted",
149 Self::ConcurrentRequest => "DynamicMockError::ConcurrentRequest",
150 Self::Builder(_) => "DynamicMockError::Builder([redacted])",
151 Self::Fixture(_) => "DynamicMockError::Fixture([redacted])",
152 Self::CursorOverflow => "DynamicMockError::CursorOverflow",
153 })
154 }
155}
156
157impl<E> fmt::Display for DynamicMockError<E> {
158 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
159 formatter.write_str(match self {
160 Self::Exhausted => "dynamic mock recording capacity is exhausted",
161 Self::ConcurrentRequest => "dynamic mock request overlaps another request",
162 Self::Builder(_) => "dynamic fixture builder rejected the request",
163 Self::Fixture(_) => "dynamic response fixture could not be staged",
164 Self::CursorOverflow => "dynamic mock cursor overflowed",
165 })
166 }
167}
168
169impl<E> core::error::Error for DynamicMockError<E> {}
170
171pub struct DynamicMockTransport<'fixture, 'records, B> {
173 builder: B,
174 records: &'records [RequestRecordSlot],
175 cursor: AtomicUsize,
176 in_flight: AtomicBool,
177 endpoint: Option<EndpointIdentity<'fixture>>,
178}
179
180impl<'fixture, 'records, B> DynamicMockTransport<'fixture, 'records, B> {
181 pub fn new(
183 builder: B,
184 records: &'records [RequestRecordSlot],
185 ) -> Result<Self, DynamicMockConfigError> {
186 if records.is_empty() {
187 return Err(DynamicMockConfigError::NoRecordSlots);
188 }
189 if records.len() > MAX_DYNAMIC_RECORDS {
190 return Err(DynamicMockConfigError::TooManyRecordSlots);
191 }
192 if records.iter().any(|slot| !slot.is_empty()) {
193 return Err(DynamicMockConfigError::DirtyRecordSlot);
194 }
195 Ok(Self {
196 builder,
197 records,
198 cursor: AtomicUsize::new(0),
199 in_flight: AtomicBool::new(false),
200 endpoint: None,
201 })
202 }
203
204 #[must_use]
206 pub const fn with_endpoint(mut self, endpoint: EndpointIdentity<'fixture>) -> Self {
207 self.endpoint = Some(endpoint);
208 self
209 }
210
211 #[must_use]
213 pub fn recorded(&self) -> usize {
214 self.cursor.load(Ordering::Acquire)
215 }
216
217 #[must_use]
219 pub const fn capacity(&self) -> usize {
220 self.records.len()
221 }
222
223 #[must_use]
225 pub fn record(&self, index: usize) -> Option<crate::RecordedRequest> {
226 self.records.get(index).and_then(RequestRecordSlot::get)
227 }
228}
229
230struct InFlightGuard<'a>(&'a AtomicBool);
231
232impl Drop for InFlightGuard<'_> {
233 fn drop(&mut self) {
234 self.0.store(false, Ordering::Release);
235 }
236}
237
238impl<'fixture, B> DynamicMockTransport<'fixture, '_, B>
239where
240 B: ProviderFixtureBuilder<'fixture>,
241{
242 fn stage_inner<'buffer>(
243 &self,
244 request: TransportRequest<'_>,
245 response: &mut impl MockResponseSink<'buffer>,
246 ) -> Result<ResponseCompletion, DynamicMockError<B::Error>> {
247 self.in_flight
248 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
249 .map_err(|_| DynamicMockError::ConcurrentRequest)?;
250 let _guard = InFlightGuard(&self.in_flight);
251 let sequence = self.cursor.load(Ordering::Acquire);
252 let slot = self
253 .records
254 .get(sequence)
255 .ok_or(DynamicMockError::Exhausted)?;
256 let fixture = self
257 .builder
258 .build(DynamicRequest::new(sequence, request))
259 .map_err(DynamicMockError::Builder)?;
260 let completion = stage_response(fixture, response).map_err(DynamicMockError::Fixture)?;
261 let next = sequence
262 .checked_add(1)
263 .ok_or(DynamicMockError::CursorOverflow)?;
264 slot.commit(sequence, request, fixture.status());
265 self.cursor.store(next, Ordering::Release);
266 Ok(completion)
267 }
268
269 fn send_inner(
270 &self,
271 request: TransportRequest<'_>,
272 response: &mut ResponseWriter<'_>,
273 ) -> Result<(), DynamicMockError<B::Error>> {
274 if response.is_committed() {
275 return Err(DynamicMockError::Fixture(MockError::ResponseWriterRejected));
276 }
277 let mut attempt = response
278 .begin_attempt()
279 .map_err(|_| DynamicMockError::Fixture(MockError::ResponseWriterRejected))?;
280 let completion = self.stage_inner(request, &mut attempt)?;
281 attempt
282 .commit_completion(completion)
283 .map_err(|_| DynamicMockError::Fixture(MockError::ResponseWriterRejected))
284 }
285}
286
287impl<'fixture, B> BlockingTransport for DynamicMockTransport<'fixture, '_, B>
288where
289 B: ProviderFixtureBuilder<'fixture>,
290{
291 type Error = DynamicMockError<B::Error>;
292
293 fn send(
294 &self,
295 request: TransportRequest<'_>,
296 response: &mut ResponseWriter<'_>,
297 ) -> Result<(), Self::Error> {
298 self.send_inner(request, response)
299 }
300}
301
302impl<'fixture, B> BlockingAuthenticatedTransport for DynamicMockTransport<'fixture, '_, B>
303where
304 B: ProviderFixtureBuilder<'fixture>,
305{
306 type Error = DynamicMockError<B::Error>;
307
308 fn send_authenticated(
309 &self,
310 request: AuthenticatedRequest<'_, '_>,
311 response: &mut ResponseWriter<'_>,
312 ) -> Result<(), Self::Error> {
313 self.send_inner(request.transport_request(), response)
314 }
315}
316
317impl<'fixture, B> AsyncTransport for DynamicMockTransport<'fixture, '_, B>
318where
319 B: ProviderFixtureBuilder<'fixture> + Sync,
320{
321 type Error = DynamicMockError<B::Error>;
322
323 async fn send<'transport, 'request, 'writer, 'buffer>(
324 &'transport self,
325 request: TransportRequest<'request>,
326 mut response: AsyncResponseStaging<'writer, 'buffer>,
327 ) -> Result<ResponseCompletion, Self::Error>
328 where
329 'transport: 'writer,
330 'request: 'writer,
331 'buffer: 'writer,
332 {
333 self.stage_inner(request, &mut response)
334 }
335}
336
337impl<'fixture, B> AsyncAuthenticatedTransport for DynamicMockTransport<'fixture, '_, B>
338where
339 B: ProviderFixtureBuilder<'fixture> + Sync,
340{
341 type Error = DynamicMockError<B::Error>;
342
343 async fn send_authenticated<'transport, 'request, 'policy, 'writer, 'buffer>(
344 &'transport self,
345 request: AuthenticatedRequest<'request, 'policy>,
346 mut response: AsyncResponseStaging<'writer, 'buffer>,
347 ) -> Result<ResponseCompletion, Self::Error>
348 where
349 'transport: 'writer,
350 'request: 'writer,
351 'policy: 'writer,
352 'buffer: 'writer,
353 {
354 self.stage_inner(request.transport_request(), &mut response)
355 }
356}
357
358impl<B> BoundTransport for DynamicMockTransport<'_, '_, B> {
359 fn endpoint_identity(&self) -> Result<EndpointIdentity<'_>, EndpointIdentityError> {
360 self.endpoint.ok_or(EndpointIdentityError::UnboundTransport)
361 }
362}
363
364impl<B> fmt::Debug for DynamicMockTransport<'_, '_, B> {
365 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
366 formatter
367 .debug_struct("DynamicMockTransport")
368 .field("recorded", &self.recorded())
369 .field("capacity", &self.capacity())
370 .finish_non_exhaustive()
371 }
372}