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