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    /// Executes once through a delivery-classified blocking transport.
256    pub fn execute_blocking<'buffer, T, C>(
257        mut self,
258        clock: &C,
259        transport: &T,
260        response_storage: &'buffer mut [u8],
261        response_header_storage: &'buffer mut [u8],
262    ) -> Result<CheckedResponseGuard<'buffer>, PermitExecutionError<T::Error>>
263    where
264        T: BlockingAuthenticatedTransport + BoundTransport,
265        T::Error: DeliveryClassified,
266        C: PermitClock + ?Sized,
267    {
268        sanitize_bytes(response_storage);
269        sanitize_bytes(response_header_storage);
270        self.ensure_fresh(clock.now(), response_storage, response_header_storage)?;
271        let result = self.subject.prepared().execute_blocking_authorized(
272            transport,
273            Some(self.subject.endpoint()),
274            response_storage,
275            response_header_storage,
276        );
277        self.finish_result(result)
278    }
279
280    /// Executes once through a delivery-classified Send-async transport.
281    #[allow(clippy::manual_async_fn)]
282    pub fn execute_async<'transport, 'buffer, T, C>(
283        mut self,
284        clock: &'transport C,
285        transport: &'transport T,
286        response_storage: &'buffer mut [u8],
287        response_header_storage: &'buffer mut [u8],
288    ) -> impl core::future::Future<
289        Output = Result<CheckedResponseGuard<'buffer>, PermitExecutionError<T::Error>>,
290    > + 'transport
291    where
292        T: AsyncAuthenticatedTransport + BoundTransport,
293        T::Error: DeliveryClassified,
294        C: PermitClock + Sync + ?Sized,
295        'request: 'transport,
296        'permit: 'transport,
297        'buffer: 'transport,
298    {
299        sanitize_bytes(response_storage);
300        sanitize_bytes(response_header_storage);
301        async move {
302            self.ensure_fresh(clock.now(), response_storage, response_header_storage)?;
303            let result = self
304                .subject
305                .prepared()
306                .execute_async_authorized(
307                    transport,
308                    Some(self.subject.endpoint()),
309                    response_storage,
310                    response_header_storage,
311                )
312                .await;
313            self.finish_result(result)
314        }
315    }
316
317    /// Executes once through a delivery-classified local-async transport.
318    #[allow(clippy::manual_async_fn)]
319    pub fn execute_local_async<'transport, 'buffer, T, C>(
320        mut self,
321        clock: &'transport C,
322        transport: &'transport T,
323        response_storage: &'buffer mut [u8],
324        response_header_storage: &'buffer mut [u8],
325    ) -> impl core::future::Future<
326        Output = Result<CheckedResponseGuard<'buffer>, PermitExecutionError<T::Error>>,
327    > + 'transport
328    where
329        T: LocalAsyncAuthenticatedTransport + BoundTransport,
330        T::Error: DeliveryClassified,
331        C: PermitClock + ?Sized,
332        'request: 'transport,
333        'permit: 'transport,
334        'buffer: 'transport,
335    {
336        sanitize_bytes(response_storage);
337        sanitize_bytes(response_header_storage);
338        async move {
339            self.ensure_fresh(clock.now(), response_storage, response_header_storage)?;
340            let result = self
341                .subject
342                .prepared()
343                .execute_local_async_authorized(
344                    transport,
345                    Some(self.subject.endpoint()),
346                    response_storage,
347                    response_header_storage,
348                )
349                .await;
350            self.finish_result(result)
351        }
352    }
353
354    fn ensure_fresh<E>(
355        &mut self,
356        now: PermitTimestamp,
357        response_storage: &mut [u8],
358        response_header_storage: &mut [u8],
359    ) -> Result<(), PermitExecutionError<E>> {
360        let observed = match &mut self.owner {
361            AttemptOwner::Direct(owner) => owner.observe(now),
362            AttemptOwner::Shared(owner) => {
363                owner.observe_attempt(self.subject, self.generation, now)
364            }
365        };
366        if let Err(error) = observed {
367            sanitize_bytes(response_storage);
368            sanitize_bytes(response_header_storage);
369            let disposition = self.finish(AttemptPhase::Rejected);
370            return Err(PermitExecutionError {
371                execution: PreparedExecutionError::AuthorizationInvalid(error),
372                disposition,
373            });
374        }
375        Ok(())
376    }
377
378    fn finish_result<'buffer, E: DeliveryClassified>(
379        &mut self,
380        result: Result<CheckedResponseGuard<'buffer>, PreparedExecutionError<E>>,
381    ) -> Result<CheckedResponseGuard<'buffer>, PermitExecutionError<E>> {
382        match result {
383            Ok(response) => {
384                let _ = self.finish(AttemptPhase::Applied);
385                Ok(response)
386            }
387            Err(execution) => {
388                let phase = match execution.delivery_phase() {
389                    DeliveryPhase::NotSent => AttemptPhase::NotSent,
390                    DeliveryPhase::PossiblySent | DeliveryPhase::ResponseStarted => {
391                        AttemptPhase::Uncertain
392                    }
393                };
394                let disposition = self.finish(phase);
395                Err(PermitExecutionError {
396                    execution,
397                    disposition,
398                })
399            }
400        }
401    }
402
403    fn finish(&mut self, phase: AttemptPhase) -> PermitDisposition {
404        self.finished = true;
405        match &mut self.owner {
406            AttemptOwner::Direct(owner) => owner.complete(self.generation, phase),
407            AttemptOwner::Shared(owner) => owner.complete(self.generation, phase),
408        }
409    }
410}
411
412impl Drop for PermitAttempt<'_, '_, '_> {
413    fn drop(&mut self) {
414        if !self.finished {
415            let _ = self.finish(AttemptPhase::Uncertain);
416        }
417    }
418}
419
420impl core::fmt::Debug for PermitAttempt<'_, '_, '_> {
421    fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
422        formatter
423            .debug_struct("PermitAttempt")
424            .field("generation", &self.generation)
425            .field("plan", &"[redacted]")
426            .finish()
427    }
428}