Skip to main content

cloud_sdk/operation/permit/
state.rs

1//! Direct permit transitions and execution-attempt ownership.
2
3use 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/// One in-flight attempt. Dropping it records uncertain delivery.
193///
194/// The prepared request cannot be extracted from this capability.
195///
196/// ```compile_fail
197/// use cloud_sdk::operation::PermitAttempt;
198///
199/// fn extract(attempt: &PermitAttempt<'_, '_, '_>) {
200///     let _ = attempt.prepared();
201/// }
202/// ```
203#[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    /// Completes a manually driven attempt with conservative delivery state.
238    ///
239    /// `NotSent` is sound only when a delivery-aware transport boundary proves
240    /// that no request bytes reached the peer. Unknown state must be reported
241    /// as `PossiblySent`; ordinary callers should prefer the execute methods.
242    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    /// Marks a successful checked provider response and spends authority.
251    pub fn complete_applied(mut self) -> PermitDisposition {
252        self.finish(AttemptPhase::Applied)
253    }
254
255    /// Rejects this in-flight attempt before transport dispatch.
256    ///
257    /// Provider wrappers use this after validating request-bound evidence with
258    /// the same clock sample supplied to the generic permit check.
259    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    /// Executes once through a delivery-classified blocking transport.
275    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    /// Executes once through a delivery-classified Send-async transport.
300    #[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    /// Executes once through a delivery-classified local-async transport.
337    #[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}