1use super::{
4 ExecutionPermitError, PermitClock, PermitDisposition, PermitExecutionError,
5 PermitIdempotencyKey, PermitScope, PermitState, PermitTimestamp, PlanSubject,
6 ReconciliationToken, RecoveryToken, ReplayPolicy,
7};
8use crate::authentication::{
9 AsyncAuthenticatedTransport, BlockingAuthenticatedTransport, LocalAsyncAuthenticatedTransport,
10};
11use crate::operation::{CheckedResponseGuard, PreparedExecutionError};
12use crate::transport::{BoundTransport, DeliveryClassified, DeliveryPhase};
13use cloud_sdk_sanitization::sanitize_bytes;
14
15pub(super) struct DirectState<'request, 'fingerprint> {
16 subject: PlanSubject<'request, 'fingerprint>,
17 state: PermitState,
18 generation: u16,
19 remaining: u16,
20 last_offset: u32,
21}
22
23impl<'request, 'fingerprint> DirectState<'request, 'fingerprint> {
24 pub(super) fn new(
25 subject: PlanSubject<'request, 'fingerprint>,
26 expected: PermitScope,
27 now: PermitTimestamp,
28 ) -> Result<Self, ExecutionPermitError> {
29 if subject.scope() != expected {
30 return Err(ExecutionPermitError::ScopeMismatch);
31 }
32 let last_offset = subject.validity().offset(now)?;
33 Ok(Self {
34 subject,
35 state: PermitState::Ready,
36 generation: 0,
37 remaining: subject.attempt_budget().get(),
38 last_offset,
39 })
40 }
41
42 pub(super) const fn state(&self) -> PermitState {
43 self.state
44 }
45
46 pub(super) fn begin(
47 &mut self,
48 now: PermitTimestamp,
49 ) -> Result<PermitAttempt<'_, 'request, 'fingerprint>, ExecutionPermitError> {
50 self.begin_for(self.subject, now)
51 }
52
53 pub(super) fn begin_for(
54 &mut self,
55 candidate: PlanSubject<'_, '_>,
56 now: PermitTimestamp,
57 ) -> Result<PermitAttempt<'_, 'request, 'fingerprint>, ExecutionPermitError> {
58 if !self.subject.fingerprint().matches(candidate.fingerprint()) {
59 return Err(ExecutionPermitError::FingerprintMismatch);
60 }
61 self.observe(now)?;
62 match self.state {
63 PermitState::Ready if self.remaining != 0 => {}
64 PermitState::Ready | PermitState::Spent => return Err(ExecutionPermitError::Spent),
65 PermitState::InFlight => return Err(ExecutionPermitError::AttemptInFlight),
66 PermitState::Recoverable => return Err(ExecutionPermitError::RecoveryRequired),
67 PermitState::PendingReconciliation => {
68 return Err(ExecutionPermitError::ReconciliationRequired);
69 }
70 }
71 self.remaining = self
72 .remaining
73 .checked_sub(1)
74 .ok_or(ExecutionPermitError::Spent)?;
75 self.state = PermitState::InFlight;
76 Ok(PermitAttempt::direct(self, self.generation))
77 }
78
79 pub(super) fn recover_not_sent(
80 &mut self,
81 token: RecoveryToken,
82 now: PermitTimestamp,
83 ) -> Result<(), ExecutionPermitError> {
84 self.observe(now)?;
85 if self.state != PermitState::Recoverable || token.0 != self.generation {
86 return Err(ExecutionPermitError::StaleGeneration);
87 }
88 if self.subject.replay_policy() == ReplayPolicy::SingleAttempt {
89 return Err(ExecutionPermitError::ReplayForbidden);
90 }
91 self.rearm()
92 }
93
94 pub(super) fn reconcile_not_applied(
95 &mut self,
96 token: ReconciliationToken,
97 candidate: PlanSubject<'_, '_>,
98 idempotency: PermitIdempotencyKey<'_>,
99 now: PermitTimestamp,
100 ) -> Result<(), ExecutionPermitError> {
101 self.observe(now)?;
102 if self.state != PermitState::PendingReconciliation || token.0 != self.generation {
103 return Err(ExecutionPermitError::StaleGeneration);
104 }
105 if self.subject.replay_policy() != ReplayPolicy::ReconcileThenRetry {
106 return Err(ExecutionPermitError::ReplayForbidden);
107 }
108 if !self.subject.fingerprint().matches(candidate.fingerprint()) {
109 return Err(ExecutionPermitError::FingerprintMismatch);
110 }
111 if !self
112 .subject
113 .idempotency()
114 .is_some_and(|expected| expected.matches(idempotency))
115 {
116 return Err(ExecutionPermitError::IdempotencyMismatch);
117 }
118 self.rearm()
119 }
120
121 fn observe(&mut self, now: PermitTimestamp) -> Result<(), ExecutionPermitError> {
122 let offset = match self.subject.validity().offset(now) {
123 Ok(offset) => offset,
124 Err(error) => {
125 self.state = PermitState::Spent;
126 self.remaining = 0;
127 return Err(error);
128 }
129 };
130 if offset < self.last_offset {
131 self.state = PermitState::Spent;
132 self.remaining = 0;
133 return Err(ExecutionPermitError::ClockRollback);
134 }
135 self.last_offset = offset;
136 Ok(())
137 }
138
139 fn rearm(&mut self) -> Result<(), ExecutionPermitError> {
140 if self.remaining == 0 {
141 self.state = PermitState::Spent;
142 return Err(ExecutionPermitError::Spent);
143 }
144 let Some(generation) = self.generation.checked_add(1) else {
145 self.state = PermitState::Spent;
146 return Err(ExecutionPermitError::GenerationExhausted);
147 };
148 self.generation = generation;
149 self.state = PermitState::Ready;
150 Ok(())
151 }
152
153 fn complete(&mut self, generation: u16, phase: AttemptPhase) -> PermitDisposition {
154 if self.state != PermitState::InFlight || self.generation != generation {
155 self.state = PermitState::Spent;
156 return PermitDisposition::Spent;
157 }
158 match phase {
159 AttemptPhase::Applied | AttemptPhase::Rejected => {
160 self.state = PermitState::Spent;
161 PermitDisposition::Spent
162 }
163 AttemptPhase::NotSent if self.remaining == 0 => {
164 self.state = PermitState::Spent;
165 PermitDisposition::Spent
166 }
167 AttemptPhase::NotSent => {
168 self.state = PermitState::Recoverable;
169 PermitDisposition::Recoverable(RecoveryToken(generation))
170 }
171 AttemptPhase::Uncertain => {
172 self.state = PermitState::PendingReconciliation;
173 PermitDisposition::PendingReconciliation(ReconciliationToken(generation))
174 }
175 }
176 }
177}
178
179#[derive(Clone, Copy)]
180pub(super) enum AttemptPhase {
181 Applied,
182 Rejected,
183 NotSent,
184 Uncertain,
185}
186
187enum AttemptOwner<'permit, 'request, 'fingerprint> {
188 Direct(&'permit mut DirectState<'request, 'fingerprint>),
189 Shared(&'permit super::shared::SharedPermitState),
190}
191
192#[must_use]
204pub struct PermitAttempt<'permit, 'request, 'fingerprint> {
205 owner: AttemptOwner<'permit, 'request, 'fingerprint>,
206 subject: PlanSubject<'request, 'fingerprint>,
207 generation: u16,
208 finished: bool,
209}
210
211impl<'permit, 'request, 'fingerprint> PermitAttempt<'permit, 'request, 'fingerprint> {
212 pub(super) fn direct(
213 owner: &'permit mut DirectState<'request, 'fingerprint>,
214 generation: u16,
215 ) -> Self {
216 Self {
217 subject: owner.subject,
218 owner: AttemptOwner::Direct(owner),
219 generation,
220 finished: false,
221 }
222 }
223
224 pub(super) fn shared(
225 owner: &'permit super::shared::SharedPermitState,
226 subject: PlanSubject<'request, 'fingerprint>,
227 generation: u16,
228 ) -> Self {
229 Self {
230 owner: AttemptOwner::Shared(owner),
231 subject,
232 generation,
233 finished: false,
234 }
235 }
236
237 pub fn complete(mut self, phase: DeliveryPhase) -> PermitDisposition {
243 let phase = match phase {
244 DeliveryPhase::NotSent => AttemptPhase::NotSent,
245 DeliveryPhase::PossiblySent | DeliveryPhase::ResponseStarted => AttemptPhase::Uncertain,
246 };
247 self.finish(phase)
248 }
249
250 pub fn complete_applied(mut self) -> PermitDisposition {
252 self.finish(AttemptPhase::Applied)
253 }
254
255 pub fn reject_authorization<E>(
260 mut self,
261 error: ExecutionPermitError,
262 response_storage: &mut [u8],
263 response_header_storage: &mut [u8],
264 ) -> PermitExecutionError<E> {
265 sanitize_bytes(response_storage);
266 sanitize_bytes(response_header_storage);
267 let disposition = self.finish(AttemptPhase::Rejected);
268 PermitExecutionError {
269 execution: PreparedExecutionError::AuthorizationInvalid(error),
270 disposition,
271 }
272 }
273
274 pub fn execute_blocking<'buffer, T, C>(
276 mut self,
277 clock: &C,
278 transport: &T,
279 response_storage: &'buffer mut [u8],
280 response_header_storage: &'buffer mut [u8],
281 ) -> Result<CheckedResponseGuard<'buffer>, PermitExecutionError<T::Error>>
282 where
283 T: BlockingAuthenticatedTransport + BoundTransport,
284 T::Error: DeliveryClassified,
285 C: PermitClock + ?Sized,
286 {
287 sanitize_bytes(response_storage);
288 sanitize_bytes(response_header_storage);
289 self.ensure_fresh(clock.now(), response_storage, response_header_storage)?;
290 let result = self.subject.prepared().execute_blocking_authorized(
291 transport,
292 Some(self.subject.endpoint()),
293 response_storage,
294 response_header_storage,
295 );
296 self.finish_result(result)
297 }
298
299 #[allow(clippy::manual_async_fn)]
301 pub fn execute_async<'transport, 'buffer, T, C>(
302 mut self,
303 clock: &'transport C,
304 transport: &'transport T,
305 response_storage: &'buffer mut [u8],
306 response_header_storage: &'buffer mut [u8],
307 ) -> impl core::future::Future<
308 Output = Result<CheckedResponseGuard<'buffer>, PermitExecutionError<T::Error>>,
309 > + 'transport
310 where
311 T: AsyncAuthenticatedTransport + BoundTransport,
312 T::Error: DeliveryClassified,
313 C: PermitClock + Sync + ?Sized,
314 'request: 'transport,
315 'permit: 'transport,
316 'buffer: 'transport,
317 {
318 sanitize_bytes(response_storage);
319 sanitize_bytes(response_header_storage);
320 async move {
321 self.ensure_fresh(clock.now(), response_storage, response_header_storage)?;
322 let result = self
323 .subject
324 .prepared()
325 .execute_async_authorized(
326 transport,
327 Some(self.subject.endpoint()),
328 response_storage,
329 response_header_storage,
330 )
331 .await;
332 self.finish_result(result)
333 }
334 }
335
336 #[allow(clippy::manual_async_fn)]
338 pub fn execute_local_async<'transport, 'buffer, T, C>(
339 mut self,
340 clock: &'transport C,
341 transport: &'transport T,
342 response_storage: &'buffer mut [u8],
343 response_header_storage: &'buffer mut [u8],
344 ) -> impl core::future::Future<
345 Output = Result<CheckedResponseGuard<'buffer>, PermitExecutionError<T::Error>>,
346 > + 'transport
347 where
348 T: LocalAsyncAuthenticatedTransport + BoundTransport,
349 T::Error: DeliveryClassified,
350 C: PermitClock + ?Sized,
351 'request: 'transport,
352 'permit: 'transport,
353 'buffer: 'transport,
354 {
355 sanitize_bytes(response_storage);
356 sanitize_bytes(response_header_storage);
357 async move {
358 self.ensure_fresh(clock.now(), response_storage, response_header_storage)?;
359 let result = self
360 .subject
361 .prepared()
362 .execute_local_async_authorized(
363 transport,
364 Some(self.subject.endpoint()),
365 response_storage,
366 response_header_storage,
367 )
368 .await;
369 self.finish_result(result)
370 }
371 }
372
373 fn ensure_fresh<E>(
374 &mut self,
375 now: PermitTimestamp,
376 response_storage: &mut [u8],
377 response_header_storage: &mut [u8],
378 ) -> Result<(), PermitExecutionError<E>> {
379 let observed = match &mut self.owner {
380 AttemptOwner::Direct(owner) => owner.observe(now),
381 AttemptOwner::Shared(owner) => {
382 owner.observe_attempt(self.subject, self.generation, now)
383 }
384 };
385 if let Err(error) = observed {
386 sanitize_bytes(response_storage);
387 sanitize_bytes(response_header_storage);
388 let disposition = self.finish(AttemptPhase::Rejected);
389 return Err(PermitExecutionError {
390 execution: PreparedExecutionError::AuthorizationInvalid(error),
391 disposition,
392 });
393 }
394 Ok(())
395 }
396
397 fn finish_result<'buffer, E: DeliveryClassified>(
398 &mut self,
399 result: Result<CheckedResponseGuard<'buffer>, PreparedExecutionError<E>>,
400 ) -> Result<CheckedResponseGuard<'buffer>, PermitExecutionError<E>> {
401 match result {
402 Ok(response) => {
403 let _ = self.finish(AttemptPhase::Applied);
404 Ok(response)
405 }
406 Err(execution) => {
407 let phase = match execution.delivery_phase() {
408 DeliveryPhase::NotSent => AttemptPhase::NotSent,
409 DeliveryPhase::PossiblySent | DeliveryPhase::ResponseStarted => {
410 AttemptPhase::Uncertain
411 }
412 };
413 let disposition = self.finish(phase);
414 Err(PermitExecutionError {
415 execution,
416 disposition,
417 })
418 }
419 }
420 }
421
422 fn finish(&mut self, phase: AttemptPhase) -> PermitDisposition {
423 self.finished = true;
424 match &mut self.owner {
425 AttemptOwner::Direct(owner) => owner.complete(self.generation, phase),
426 AttemptOwner::Shared(owner) => owner.complete(self.generation, phase),
427 }
428 }
429}
430
431impl Drop for PermitAttempt<'_, '_, '_> {
432 fn drop(&mut self) {
433 if !self.finished {
434 let _ = self.finish(AttemptPhase::Uncertain);
435 }
436 }
437}
438
439impl core::fmt::Debug for PermitAttempt<'_, '_, '_> {
440 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
441 formatter
442 .debug_struct("PermitAttempt")
443 .field("generation", &self.generation)
444 .field("plan", &"[redacted]")
445 .finish()
446 }
447}