1use super::{
2 core_resource_failure, fmt, invalid_resource, new_deferred_device_cleanup_domain, watch, Arc,
3 AtomicBool, AtomicU8, BTreeMap, BTreeSet, BufferDescriptor, CoreOwnedAllocation,
4 DeviceCapacityClaim, DriverCommitAcknowledgement, DynamicPoolMaintenanceController,
5 DynamicPoolSet, ErasedPlanStaticDriver, Mutex, Ordering, PhantomData, PlanRuntimeResources,
6 PlanRuntimeStatic, PlanStaticResources, RefCell, ResourceAbandonSignal, ResourceActionCursor,
7 ResourceCommitView, ResourceCompensationRecord, ResourceDriverFailure, ResourceFailureId,
8 ResourceFailurePoint, ResourceFailureReceipt, ResourceId, ResourceLeaseAction,
9 ResourceLeaseState, ResourceLeaseTransitionReceipt, ResourceLeaseValidationContext,
10 ResourceLedgerEntrySnapshot, ResourceLedgerSnapshot, ResourceOwnedBuffer,
11 ResourceOwnershipReason, ResourcePoolOwnership, ResourceRecoveryStrategy,
12 ResourceReservationBatch, ResourceTransactionAction, ResourceTransactionContext,
13 ResourceTransactionDriver, ResourceTransactionIdentity, ResourceTransactionState,
14 ResourceTransitionReceipt, ResourceTransitionRecord, ResourceTransitionValidationContext,
15 RwLock, StaticProvisioningBinding, StaticProvisioningLease, StaticProvisioningPermit,
16 VNextError, PLAN_RUNTIME_OPEN,
17};
18
19mod sealed {
20 pub trait Sealed {}
21}
22
23pub trait TransactionStage: sealed::Sealed {
24 const STATE: ResourceTransactionState;
25 const TERMINAL: bool;
26}
27
28pub struct TransactionNew;
29pub struct TransactionReserved;
30pub struct TransactionCommitted;
31pub struct TransactionRolledBack;
32pub struct TransactionReleased;
33pub struct TransactionQuarantined;
34
35impl sealed::Sealed for TransactionNew {}
36impl sealed::Sealed for TransactionReserved {}
37impl sealed::Sealed for TransactionCommitted {}
38impl sealed::Sealed for TransactionRolledBack {}
39impl sealed::Sealed for TransactionReleased {}
40impl sealed::Sealed for TransactionQuarantined {}
41
42impl TransactionStage for TransactionNew {
43 const STATE: ResourceTransactionState = ResourceTransactionState::New;
44 const TERMINAL: bool = false;
45}
46impl TransactionStage for TransactionReserved {
47 const STATE: ResourceTransactionState = ResourceTransactionState::Reserved;
48 const TERMINAL: bool = false;
49}
50impl TransactionStage for TransactionCommitted {
51 const STATE: ResourceTransactionState = ResourceTransactionState::Committed;
52 const TERMINAL: bool = false;
53}
54impl TransactionStage for TransactionRolledBack {
55 const STATE: ResourceTransactionState = ResourceTransactionState::RolledBack;
56 const TERMINAL: bool = true;
57}
58impl TransactionStage for TransactionReleased {
59 const STATE: ResourceTransactionState = ResourceTransactionState::Released;
60 const TERMINAL: bool = true;
61}
62impl TransactionStage for TransactionQuarantined {
63 const STATE: ResourceTransactionState = ResourceTransactionState::Quarantined;
64 const TERMINAL: bool = true;
65}
66
67struct PendingSubsetRelease {
68 target_orders: Vec<usize>,
69 failure: ResourceFailureReceipt,
70}
71
72#[must_use = "dropping a nonterminal resource transaction emits an abandon signal"]
73pub struct ResourceTransaction<D: ResourceTransactionDriver, S: TransactionStage> {
74 driver: Option<D>,
75 maintenance_controller: Option<DynamicPoolMaintenanceController<D::Runtime>>,
76 dynamic_pools: Option<Arc<DynamicPoolSet<D::Runtime>>>,
77 identity: ResourceTransactionIdentity,
78 admission: StaticProvisioningBinding,
79 reservations: ResourceReservationBatch,
80 capacity_claim: Option<DeviceCapacityClaim>,
81 allocation_issued: Vec<AtomicBool>,
82 pending_allocations: Vec<RefCell<Option<CoreOwnedAllocation<D::Buffer>>>>,
83 states: Vec<ResourceTransactionState>,
84 lease: Option<StaticProvisioningLease<D::Runtime>>,
85 receipts: Vec<ResourceTransitionReceipt>,
86 lease_receipts: Vec<ResourceLeaseTransitionReceipt>,
87 recovery_history: Vec<ResourceFailureReceipt>,
88 latest_transition_context: Option<ResourceTransitionValidationContext>,
89 latest_lease_context: Option<ResourceLeaseValidationContext>,
90 pending_failure: Option<ResourceFailureReceipt>,
91 pending_subset_release: Option<PendingSubsetRelease>,
92 next_failure_id: u64,
93 finalized: bool,
94 stage: PhantomData<S>,
95}
96
97impl<D: ResourceTransactionDriver, S: TransactionStage> ResourceTransaction<D, S> {
98 pub fn maintenance_controller(&self) -> &DynamicPoolMaintenanceController<D::Runtime> {
99 self.maintenance_controller
100 .as_ref()
101 .expect("live resource transaction owns dynamic-pool maintenance authority")
102 }
103
104 pub fn identity(&self) -> &ResourceTransactionIdentity {
105 &self.identity
106 }
107
108 pub fn admission(&self) -> &StaticProvisioningBinding {
109 &self.admission
110 }
111
112 pub fn reservations(&self) -> &ResourceReservationBatch {
113 &self.reservations
114 }
115
116 pub fn receipts(&self) -> &[ResourceTransitionReceipt] {
117 &self.receipts
118 }
119
120 pub fn lease_receipts(&self) -> &[ResourceLeaseTransitionReceipt] {
121 &self.lease_receipts
122 }
123
124 pub fn recovery_history(&self) -> &[ResourceFailureReceipt] {
125 &self.recovery_history
126 }
127
128 pub const fn state(&self) -> ResourceTransactionState {
129 S::STATE
130 }
131
132 pub fn actual_states(&self) -> &[ResourceTransactionState] {
133 &self.states
134 }
135
136 pub fn ledger_snapshot(&self) -> ResourceLedgerSnapshot {
137 ResourceLedgerSnapshot {
138 identity: self.identity.clone(),
139 admission: self.admission.clone(),
140 entries: self.ledger_snapshot_entries(),
141 }
142 }
143
144 pub fn latest_transition_validation_context(
145 &self,
146 ) -> Option<&ResourceTransitionValidationContext> {
147 self.latest_transition_context.as_ref()
148 }
149
150 pub fn latest_lease_validation_context(&self) -> Option<&ResourceLeaseValidationContext> {
151 self.latest_lease_context.as_ref()
152 }
153
154 fn ledger_snapshot_entries(&self) -> Vec<ResourceLedgerEntrySnapshot> {
155 let lease = self
156 .lease
157 .as_ref()
158 .expect("every resource transaction owns its batch lease ledger");
159 lease
160 .slots
161 .iter()
162 .zip(&self.pending_allocations)
163 .zip(&self.states)
164 .map(|((slot, pending), &transaction_state)| {
165 let pending = pending.borrow();
166 ResourceLedgerEntrySnapshot {
167 entry: slot.entry.clone(),
168 transaction_state,
169 buffer_present: slot.buffer.is_some() || pending.is_some(),
170 actual_resource_id: slot
171 .actual_resource_id
172 .clone()
173 .or_else(|| pending.as_ref().map(|value| value.resource_id.clone())),
174 actual_generation: slot
175 .actual_generation
176 .or_else(|| pending.as_ref().map(|value| value.generation)),
177 actual_descriptor: slot
178 .descriptor
179 .clone()
180 .or_else(|| pending.as_ref().map(|value| value.descriptor.clone())),
181 }
182 })
183 .collect()
184 }
185
186 fn driver_and_context(
187 &mut self,
188 order: usize,
189 action: ResourceTransactionAction,
190 ) -> (&mut D, ResourceTransactionContext<'_, D::Runtime>) {
191 let before = self.states[order];
192 let driver = self
193 .driver
194 .as_mut()
195 .expect("live transaction owns its resource driver");
196 let context = ResourceTransactionContext {
197 runtime: &self
198 .lease
199 .as_ref()
200 .expect("transaction owns lease ledger")
201 .runtime,
202 identity: &self.identity,
203 binding: &self.admission,
204 reservations: &self.reservations,
205 cursor: Some(ResourceActionCursor {
206 order,
207 action,
208 before,
209 allocation_authorized: action == ResourceTransactionAction::Commit,
210 }),
211 allocation_authority: (action == ResourceTransactionAction::Commit)
212 .then_some(&self.allocation_issued[order]),
213 pending_allocation: (action == ResourceTransactionAction::Commit)
214 .then_some(&self.pending_allocations[order]),
215 };
216 (driver, context)
217 }
218
219 fn make_record(
220 &self,
221 order: usize,
222 action: ResourceTransactionAction,
223 before: ResourceTransactionState,
224 after: ResourceTransactionState,
225 ) -> ResourceTransitionRecord {
226 ResourceTransitionRecord::from_reservation(
227 &self.identity,
228 &self.admission,
229 &self.reservations.reservations[order],
230 action,
231 before,
232 after,
233 order,
234 )
235 }
236
237 fn transition_receipt(
238 &mut self,
239 action: ResourceTransactionAction,
240 before: Vec<ResourceLedgerEntrySnapshot>,
241 records: Vec<ResourceTransitionRecord>,
242 ) -> Result<ResourceTransitionReceipt, VNextError> {
243 let context = ResourceTransitionValidationContext {
244 identity: self.identity.clone(),
245 admission: self.admission.clone(),
246 action,
247 before,
248 after: self.ledger_snapshot_entries(),
249 };
250 let receipt = ResourceTransitionReceipt::from_context(&context, records)?;
251 self.latest_transition_context = Some(context);
252 Ok(receipt)
253 }
254
255 fn lease_transition(
256 &mut self,
257 resource_ids: &[ResourceId],
258 action: ResourceLeaseAction,
259 ) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
260 if self.pending_subset_release.is_some() {
261 return Err(invalid_resource(
262 "lease policy cannot change during pending release recovery",
263 ));
264 }
265 let orders = self.orders_for_ids(resource_ids)?;
266 if orders
267 .iter()
268 .any(|&order| self.states[order] != ResourceTransactionState::Committed)
269 {
270 return Err(invalid_resource(
271 "lease policy targets a resource that is not actually committed",
272 ));
273 }
274 let before_ledger = self.ledger_snapshot_entries();
275 let (before, after, entries) = self
276 .lease
277 .as_mut()
278 .expect("transaction owns lease ledger")
279 .transition_subset(&orders, action)?;
280 let context = ResourceLeaseValidationContext {
281 identity: self.identity.clone(),
282 admission: self.admission.clone(),
283 action,
284 before: before_ledger,
285 after: self.ledger_snapshot_entries(),
286 };
287 let receipt =
288 ResourceLeaseTransitionReceipt::from_context(&context, before, after, entries)?;
289 self.latest_lease_context = Some(context);
290 self.lease_receipts.push(receipt.clone());
291 Ok(receipt)
292 }
293
294 fn orders_for_ids(&self, resource_ids: &[ResourceId]) -> Result<Vec<usize>, VNextError> {
295 if resource_ids.is_empty() {
296 return Err(invalid_resource("resource subset must not be empty"));
297 }
298 let mut unique = BTreeSet::new();
299 let by_id = self
300 .reservations
301 .reservations()
302 .iter()
303 .enumerate()
304 .map(|(order, reservation)| (reservation.resource_id(), order))
305 .collect::<BTreeMap<_, _>>();
306 let mut orders = Vec::with_capacity(resource_ids.len());
307 for resource_id in resource_ids {
308 if !unique.insert(resource_id) {
309 return Err(invalid_resource("resource subset contains duplicates"));
310 }
311 orders.push(
312 *by_id
313 .get(resource_id)
314 .ok_or_else(|| invalid_resource("resource subset is not admitted"))?,
315 );
316 }
317 orders.sort_unstable();
318 Ok(orders)
319 }
320
321 fn set_pending_failure(&mut self, failure: &ResourceFailureReceipt) {
322 self.pending_failure = Some(failure.clone());
323 }
324
325 fn clear_pending_failure(&mut self) {
326 self.pending_failure = None;
327 }
328
329 fn issue_failure_id(&mut self) -> ResourceFailureId {
330 let current = self.next_failure_id;
331 self.next_failure_id = current
332 .checked_add(1)
333 .expect("resource transaction failure id space is exhausted");
334 ResourceFailureId::try_from(current).expect("transaction failure ids start at one")
335 }
336
337 fn advance<T: TransactionStage>(
338 mut self,
339 receipt: Option<ResourceTransitionReceipt>,
340 ) -> ResourceTransaction<D, T> {
341 self.finalized = true;
342 if let Some(receipt) = receipt {
343 self.receipts.push(receipt);
344 }
345 ResourceTransaction {
346 driver: self.driver.take(),
347 maintenance_controller: self.maintenance_controller.take(),
348 dynamic_pools: self.dynamic_pools.take(),
349 identity: self.identity.clone(),
350 admission: self.admission.clone(),
351 reservations: self.reservations.clone(),
352 capacity_claim: if T::TERMINAL {
353 if let Some(mut claim) = self.capacity_claim.take() {
354 claim.release();
355 }
356 None
357 } else {
358 self.capacity_claim.take()
359 },
360 allocation_issued: std::mem::take(&mut self.allocation_issued),
361 pending_allocations: std::mem::take(&mut self.pending_allocations),
362 states: std::mem::take(&mut self.states),
363 lease: self.lease.take(),
364 receipts: std::mem::take(&mut self.receipts),
365 lease_receipts: std::mem::take(&mut self.lease_receipts),
366 recovery_history: std::mem::take(&mut self.recovery_history),
367 latest_transition_context: self.latest_transition_context.take(),
368 latest_lease_context: self.latest_lease_context.take(),
369 pending_failure: None,
370 pending_subset_release: None,
371 next_failure_id: self.next_failure_id,
372 finalized: T::TERMINAL,
373 stage: PhantomData,
374 }
375 }
376
377 fn abandon_signal(&self) -> ResourceAbandonSignal {
378 ResourceAbandonSignal {
379 identity: self.identity.clone(),
380 admission: self.admission.clone(),
381 state: S::STATE,
382 pending_action: self
383 .pending_failure
384 .as_ref()
385 .map(ResourceFailureReceipt::action),
386 ledger: self.ledger_snapshot_entries(),
387 active_sequence_slots: Vec::new(),
388 poisoned_sequence_slots: Vec::new(),
389 undrained_sequence_slots: Vec::new(),
390 failure: self
391 .pending_failure
392 .as_ref()
393 .map(|receipt| receipt.failure.clone()),
394 }
395 }
396
397 fn take_all_owned_buffers(&mut self) -> Vec<ResourceOwnedBuffer<D::Buffer>> {
398 let mut buffers = self
399 .lease
400 .as_mut()
401 .expect("transaction owns lease ledger")
402 .take_owned_buffers(&self.reservations);
403 for (order, pending) in self.pending_allocations.iter_mut().enumerate() {
404 let Some(allocation) = pending.get_mut().take() else {
405 continue;
406 };
407 let reservation = &self.reservations.reservations[order];
408 buffers.push(ResourceOwnedBuffer {
409 order,
410 expected_resource_id: reservation.resource_id.clone(),
411 actual_resource_id: allocation.resource_id,
412 expected_generation: reservation.generation,
413 actual_generation: allocation.generation,
414 expected_descriptor: BufferDescriptor {
415 resource_id: reservation.resource_id.clone(),
416 size_bytes: reservation.size_bytes,
417 alignment_bytes: reservation.alignment_bytes,
418 usage: reservation.usage,
419 element_type: reservation.element_type,
420 },
421 actual_descriptor: allocation.descriptor,
422 buffer: allocation.buffer,
423 });
424 }
425 buffers.sort_by_key(|buffer| buffer.order);
426 buffers
427 }
428
429 fn quarantine_live(&mut self) -> Result<Vec<ResourceTransitionRecord>, ResourceDriverFailure> {
430 let buffers = self.take_all_owned_buffers();
431 let ownership = ResourcePoolOwnership {
432 runtime: Arc::clone(
433 &self
434 .lease
435 .as_ref()
436 .expect("transaction owns lease ledger")
437 .runtime,
438 ),
439 pool_identity: self.admission.pool_identity.clone(),
440 reason: ResourceOwnershipReason::Quarantine,
441 signal: None,
442 buffers,
443 capacity_claim: self.capacity_claim.take(),
444 };
445 let result = {
446 let driver = self
447 .driver
448 .as_mut()
449 .expect("live transaction owns its resource driver");
450 let context = ResourceTransactionContext {
451 runtime: &self
452 .lease
453 .as_ref()
454 .expect("transaction owns lease ledger")
455 .runtime,
456 identity: &self.identity,
457 binding: &self.admission,
458 reservations: &self.reservations,
459 cursor: None,
460 allocation_authority: None,
461 pending_allocation: None,
462 };
463 driver.quarantine_transaction(&context, ownership)
464 };
465 if let Err(failure) = result {
466 let (failure, mut ownership) = failure.into_parts();
467 let expected_claimed_bytes = self
468 .states
469 .iter()
470 .zip(self.reservations.reservations())
471 .filter(|(state, _)| state.is_live())
472 .map(|(_, reservation)| reservation.size_bytes())
473 .sum::<u64>();
474 if ownership.pool_identity != self.admission.pool_identity
475 || ownership.claimed_bytes() != expected_claimed_bytes
476 {
477 std::mem::forget(ownership);
478 return Err(ResourceDriverFailure::new(core_resource_failure(
479 "ownership_transfer_identity_mismatch",
480 "quarantine failure returned ownership for a different resource pool",
481 false,
482 ))
483 .expect("core failure has resource domain"));
484 }
485 self.lease
486 .as_mut()
487 .expect("transaction owns lease ledger")
488 .restore_owned_buffers(std::mem::take(&mut ownership.buffers));
489 self.capacity_claim = ownership.capacity_claim.take();
490 return Err(failure);
491 }
492
493 let mut records = Vec::new();
494 for order in 0..self.states.len() {
495 let before = self.states[order];
496 if before.is_live() {
497 self.states[order] = ResourceTransactionState::Quarantined;
498 records.push(self.make_record(
499 order,
500 ResourceTransactionAction::Quarantine,
501 before,
502 ResourceTransactionState::Quarantined,
503 ));
504 }
505 }
506 Ok(records)
507 }
508}
509
510impl<D: ResourceTransactionDriver, S: TransactionStage> Drop for ResourceTransaction<D, S> {
511 fn drop(&mut self) {
512 if self.finalized || S::TERMINAL {
513 return;
514 }
515 let signal = self.abandon_signal();
516 drop(self.maintenance_controller.take());
517 drop(self.dynamic_pools.take());
518 let buffers = self.take_all_owned_buffers();
519 let ownership = ResourcePoolOwnership {
520 runtime: Arc::clone(
521 &self
522 .lease
523 .as_ref()
524 .expect("transaction owns lease ledger")
525 .runtime,
526 ),
527 pool_identity: self.admission.pool_identity.clone(),
528 reason: ResourceOwnershipReason::Abandon,
529 signal: Some(signal),
530 buffers,
531 capacity_claim: self.capacity_claim.take(),
532 };
533 if let Some(driver) = self.driver.as_mut() {
534 let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
538 driver.abandon_transaction(ownership);
539 }));
540 } else {
541 std::mem::forget(ownership);
544 }
545 self.finalized = true;
546 }
547}
548
549impl<D: ResourceTransactionDriver> ResourceTransaction<D, TransactionNew> {
550 pub fn begin(
551 driver: D,
552 identity: ResourceTransactionIdentity,
553 permit: StaticProvisioningPermit<D::Runtime>,
554 ) -> Result<Self, VNextError> {
555 if identity.pool_id != permit.binding.pool_id()
556 || identity.request_id != permit.binding.request_id
557 || permit.binding.pool_identity.plan_id != permit.binding.plan_id
558 || permit.binding.pool_identity.plan_hash != permit.binding.plan_hash
559 || permit.binding.pool_identity.device_id != permit.binding.device_id
560 || permit
561 .binding
562 .pool_identity
563 .device_runtime_implementation_fingerprint
564 != permit.binding.device_runtime_implementation_fingerprint
565 || permit.binding.pool_identity.admission_generation
566 != permit.binding.admission_generation
567 || permit.binding.request_id != *permit.reservations.request_id()
568 || permit.binding.admitted_bytes != permit.reservations.total_size_bytes()
569 || permit.binding.plan_static_bytes != permit.reservations.plan_static_size_bytes()
570 || permit.binding.admitted_bytes != permit.binding.plan_static_bytes
571 || permit.binding.plan_static_bytes > permit.binding.usable_capacity_bytes
572 || permit.binding.usable_capacity_bytes > permit.binding.device_capacity_bytes
573 || driver.device_id() != &permit.binding.device_id
574 || driver.device_runtime_implementation_fingerprint()
575 != permit.binding.device_runtime_implementation_fingerprint
576 || driver.device_capacity_bytes() != permit.binding.device_capacity_bytes
577 || !Arc::ptr_eq(driver.runtime(), &permit.runtime)
578 {
579 return Err(invalid_resource(
580 "transaction identity, driver device, or admitted capacity does not match its permit",
581 ));
582 }
583 let _ = permit.seal;
584 if permit.capacity_claim.bytes() != permit.binding.admitted_bytes {
585 return Err(invalid_resource(
586 "device capacity claim differs from admitted bytes",
587 ));
588 }
589 let states = vec![ResourceTransactionState::New; permit.reservations.reservations.len()];
590 let allocation_issued = (0..states.len()).map(|_| AtomicBool::new(false)).collect();
591 let pending_allocations = (0..states.len()).map(|_| RefCell::new(None)).collect();
592 let lease = StaticProvisioningLease::new(
593 Arc::clone(&permit.runtime),
594 &identity,
595 &permit.binding,
596 &permit.reservations,
597 );
598 Ok(Self {
599 driver: Some(driver),
600 maintenance_controller: Some(permit.maintenance_controller),
601 dynamic_pools: Some(permit.dynamic_pools),
602 identity,
603 admission: permit.binding,
604 reservations: permit.reservations,
605 capacity_claim: Some(permit.capacity_claim),
606 allocation_issued,
607 pending_allocations,
608 states,
609 lease: Some(lease),
610 receipts: Vec::new(),
611 lease_receipts: Vec::new(),
612 recovery_history: Vec::new(),
613 latest_transition_context: None,
614 latest_lease_context: None,
615 pending_failure: None,
616 pending_subset_release: None,
617 next_failure_id: 1,
618 finalized: false,
619 stage: PhantomData,
620 })
621 }
622
623 pub fn reserve(
624 mut self,
625 ) -> Result<
626 ResourceTransaction<D, TransactionReserved>,
627 ResourcePrepareTransitionError<D, TransactionNew>,
628 > {
629 debug_assert!(
630 self.states
631 .iter()
632 .all(|state| *state == ResourceTransactionState::New),
633 "new transaction ledger must be uniformly new"
634 );
635 let ledger_before = self.ledger_snapshot_entries();
636 let mut completed = Vec::new();
637 for order in 0..self.states.len() {
638 let result = {
639 let reservation = self.reservations.reservations[order].clone();
640 let (driver, context) =
641 self.driver_and_context(order, ResourceTransactionAction::Reserve);
642 driver.reserve_resource(&context, &reservation)
643 };
644 match result {
645 Ok(()) => {
646 let before = self.states[order];
647 let after = before
648 .transition(
649 &self.reservations.reservations[order].resource_id,
650 ResourceTransactionAction::Reserve,
651 )
652 .expect("reserve was preflight validated");
653 self.states[order] = after;
654 completed.push(self.make_record(
655 order,
656 ResourceTransactionAction::Reserve,
657 before,
658 after,
659 ));
660 }
661 Err(failure) => {
662 let failure_id = self.issue_failure_id();
663 let receipt = ResourceFailureReceipt::new(
664 failure_id,
665 &self.identity,
666 &self.admission,
667 ResourceTransactionAction::Reserve,
668 failure.into_failure(),
669 Some(ResourceFailurePoint::new(
670 &self.reservations.reservations[order],
671 order,
672 self.states[order],
673 )),
674 completed,
675 ResourceRecoveryStrategy::ReverseCompensation,
676 ledger_before,
677 self.ledger_snapshot_entries(),
678 );
679 self.set_pending_failure(&receipt);
680 return Err(ResourcePrepareTransitionError {
681 transaction: Some(self),
682 failure: receipt,
683 });
684 }
685 }
686 }
687 let receipt = self
688 .transition_receipt(ResourceTransactionAction::Reserve, ledger_before, completed)
689 .expect("core reserve journal must validate");
690 Ok(self.advance(Some(receipt)))
691 }
692}
693
694#[derive(Debug)]
695pub enum ResourceCommitTransitionError<D: ResourceTransactionDriver> {
696 Recoverable(ResourcePrepareTransitionError<D, TransactionReserved>),
697 Poisoned(ResourcePoisonedTransaction<D>),
698}
699
700impl<D: ResourceTransactionDriver> ResourceTransaction<D, TransactionReserved> {
701 pub fn commit(
702 mut self,
703 ) -> Result<ResourceTransaction<D, TransactionCommitted>, ResourceCommitTransitionError<D>>
704 {
705 debug_assert!(
706 self.states
707 .iter()
708 .all(|state| *state == ResourceTransactionState::Reserved),
709 "reserved transaction must be uniformly reserved before commit"
710 );
711 let ledger_before = self.ledger_snapshot_entries();
712 let mut completed = Vec::new();
713 for order in 0..self.states.len() {
714 let driver_result = {
715 let reservation = self.reservations.reservations[order].clone();
716 let (driver, context) =
717 self.driver_and_context(order, ResourceTransactionAction::Commit);
718 driver
719 .commit_resource(&context, &reservation)
720 .map(|receipt| DriverCommitAcknowledgement::from_receipt(&receipt))
721 };
722 let allocation = self.pending_allocations[order].get_mut().take();
723 match (driver_result, allocation) {
724 (Ok(acknowledgement), Some(allocation))
725 if acknowledgement.matches(&allocation)
726 && allocation.matches(&self.reservations.reservations[order]) =>
727 {
728 self.lease
729 .as_mut()
730 .expect("transaction owns lease ledger")
731 .install(order, allocation);
732 let before = self.states[order];
733 let after = before
734 .transition(
735 &self.reservations.reservations[order].resource_id,
736 ResourceTransactionAction::Commit,
737 )
738 .expect("commit was preflight validated");
739 self.states[order] = after;
740 completed.push(self.make_record(
741 order,
742 ResourceTransactionAction::Commit,
743 before,
744 after,
745 ));
746 }
747 (driver_result, Some(poisoned)) => {
748 let actual_resource_id = poisoned.resource_id.clone();
749 let actual_generation = poisoned.generation;
750 let failure_envelope = match driver_result {
751 Ok(_) => core_resource_failure(
752 "invalid_commit_outcome",
753 format!(
754 "core-owned allocation `{}` generation {} does not match allocation `{}` generation {}",
755 actual_resource_id,
756 actual_generation,
757 self.reservations.reservations[order].resource_id(),
758 self.reservations.reservations[order].generation(),
759 ),
760 false,
761 ),
762 Err(failure) => failure.into_failure(),
763 };
764 self.lease
765 .as_mut()
766 .expect("transaction owns lease ledger")
767 .install(order, poisoned);
768 let failure_id = self.issue_failure_id();
769 let failure = ResourceFailureReceipt::new(
770 failure_id,
771 &self.identity,
772 &self.admission,
773 ResourceTransactionAction::Commit,
774 failure_envelope,
775 Some(ResourceFailurePoint::new(
776 &self.reservations.reservations[order],
777 order,
778 self.states[order],
779 )),
780 completed,
781 ResourceRecoveryStrategy::ReconcileOrQuarantine,
782 ledger_before,
783 self.ledger_snapshot_entries(),
784 );
785 self.set_pending_failure(&failure);
786 return Err(ResourceCommitTransitionError::Poisoned(
787 ResourcePoisonedTransaction {
788 transaction: Some(self),
789 failure,
790 poisoned_order: order,
791 },
792 ));
793 }
794 (driver_result, None) => {
795 let failure = match driver_result {
796 Ok(_) => core_resource_failure(
797 "commit_acknowledged_without_allocation",
798 "driver acknowledged commit without a core-owned allocation",
799 false,
800 ),
801 Err(failure) => failure.into_failure(),
802 };
803 let failure_id = self.issue_failure_id();
804 let receipt = ResourceFailureReceipt::new(
805 failure_id,
806 &self.identity,
807 &self.admission,
808 ResourceTransactionAction::Commit,
809 failure,
810 Some(ResourceFailurePoint::new(
811 &self.reservations.reservations[order],
812 order,
813 self.states[order],
814 )),
815 completed,
816 ResourceRecoveryStrategy::ReverseCompensation,
817 ledger_before,
818 self.ledger_snapshot_entries(),
819 );
820 self.set_pending_failure(&receipt);
821 return Err(ResourceCommitTransitionError::Recoverable(
822 ResourcePrepareTransitionError {
823 transaction: Some(self),
824 failure: receipt,
825 },
826 ));
827 }
828 }
829 }
830 let receipt = self
831 .transition_receipt(ResourceTransactionAction::Commit, ledger_before, completed)
832 .expect("core commit journal must validate");
833 Ok(self.advance(Some(receipt)))
834 }
835
836 pub fn rollback(
837 mut self,
838 ) -> Result<ResourceTransaction<D, TransactionRolledBack>, ResourceRollbackTransitionError<D>>
839 {
840 debug_assert!(
841 self.states
842 .iter()
843 .all(|state| *state == ResourceTransactionState::Reserved),
844 "rollback starts from a uniformly reserved transaction"
845 );
846 let ledger_before = self.ledger_snapshot_entries();
847 let mut completed = Vec::new();
848 for order in 0..self.states.len() {
849 let result = {
850 let reservation = self.reservations.reservations[order].clone();
851 let (driver, context) =
852 self.driver_and_context(order, ResourceTransactionAction::Rollback);
853 driver.rollback_resource(&context, &reservation)
854 };
855 match result {
856 Ok(()) => {
857 let before = self.states[order];
858 self.states[order] = ResourceTransactionState::RolledBack;
859 completed.push(self.make_record(
860 order,
861 ResourceTransactionAction::Rollback,
862 before,
863 ResourceTransactionState::RolledBack,
864 ));
865 }
866 Err(failure) => {
867 let failure_id = self.issue_failure_id();
868 let receipt = ResourceFailureReceipt::new(
869 failure_id,
870 &self.identity,
871 &self.admission,
872 ResourceTransactionAction::Rollback,
873 failure.into_failure(),
874 Some(ResourceFailurePoint::new(
875 &self.reservations.reservations[order],
876 order,
877 self.states[order],
878 )),
879 completed,
880 ResourceRecoveryStrategy::ForwardCompletion,
881 ledger_before,
882 self.ledger_snapshot_entries(),
883 );
884 self.set_pending_failure(&receipt);
885 return Err(ResourceRollbackTransitionError {
886 transaction: Some(self),
887 failure: receipt,
888 });
889 }
890 }
891 }
892 let receipt = self
893 .transition_receipt(
894 ResourceTransactionAction::Rollback,
895 ledger_before,
896 completed,
897 )
898 .expect("core rollback journal must validate");
899 Ok(self.advance(Some(receipt)))
900 }
901}
902
903impl<D: ResourceTransactionDriver> ResourceTransaction<D, TransactionCommitted> {
904 pub fn into_plan_runtime(
905 mut self,
906 ) -> Result<Arc<PlanRuntimeResources<D::Runtime>>, PlanRuntimeHandoffError<D>>
907 where
908 D: 'static,
909 {
910 let valid = self
911 .states
912 .iter()
913 .all(|state| *state == ResourceTransactionState::Committed)
914 && self.pending_failure.is_none()
915 && self.pending_subset_release.is_none()
916 && self
917 .pending_allocations
918 .iter()
919 .all(|allocation| allocation.borrow().is_none())
920 && self.lease.as_ref().is_some_and(|lease| {
921 !lease.slots.is_empty()
922 && lease.slots.iter().all(|slot| {
923 slot.buffer.is_some()
924 && slot.descriptor.is_some()
925 && slot.entry.state == ResourceLeaseState::Active
926 })
927 })
928 && self.driver.is_some()
929 && self.maintenance_controller.is_some()
930 && self.dynamic_pools.is_some()
931 && self.capacity_claim.is_some();
932 if !valid {
933 return Err(PlanRuntimeHandoffError {
934 error: invalid_resource(
935 "plan runtime handoff requires one complete active committed transaction",
936 ),
937 transaction: Some(self),
938 });
939 }
940 let lease = self
941 .lease
942 .take()
943 .expect("validated plan runtime handoff owns its static lease");
944 let runtime = Arc::clone(&lease.runtime);
945 let dynamic_pools = self
946 .dynamic_pools
947 .take()
948 .expect("validated plan runtime handoff owns its dynamic pools");
949 let driver: Box<dyn ErasedPlanStaticDriver<D::Runtime>> = Box::new(
950 self.driver
951 .take()
952 .expect("validated plan runtime handoff owns its driver"),
953 );
954 let static_resources = PlanStaticResources {
955 driver: Mutex::new(Some(driver)),
956 identity: self.identity.clone(),
957 admission: self.admission.clone(),
958 reservations: self.reservations.clone(),
959 states: std::mem::take(&mut self.states),
960 capacity_claim: self.capacity_claim.take(),
961 lease: Some(lease),
962 finalized: false,
963 };
964 let maintenance_controller = self
965 .maintenance_controller
966 .take()
967 .expect("validated plan runtime handoff owns maintenance authority");
968 let (lifecycle_tx, _) = watch::channel(PLAN_RUNTIME_OPEN);
969 self.finalized = true;
970 Ok(Arc::new(PlanRuntimeResources {
971 lifecycle: RwLock::new(()),
972 phase: AtomicU8::new(PLAN_RUNTIME_OPEN),
973 lifecycle_tx,
974 maintenance_controller,
975 dynamic_pools,
976 static_resources: PlanRuntimeStatic::Static(static_resources),
977 runtime,
978 deferred_cleanup_domain: new_deferred_device_cleanup_domain(),
979 }))
980 }
981
982 pub fn lease(&self) -> &StaticProvisioningLease<D::Runtime> {
983 self.lease
984 .as_ref()
985 .expect("committed transaction owns its batch lease")
986 }
987
988 fn defer_lease(
989 &mut self,
990 resource_ids: &[ResourceId],
991 ) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
992 self.lease_transition(resource_ids, ResourceLeaseAction::Defer)
993 }
994
995 fn resume_lease(
996 &mut self,
997 resource_ids: &[ResourceId],
998 ) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
999 self.lease_transition(resource_ids, ResourceLeaseAction::Resume)
1000 }
1001
1002 fn cancel_lease(
1003 &mut self,
1004 resource_ids: &[ResourceId],
1005 ) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
1006 self.lease_transition(resource_ids, ResourceLeaseAction::Cancel)
1007 }
1008
1009 pub fn defer_all(&mut self) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
1010 let ids = self
1011 .reservations
1012 .resource_ids()
1013 .cloned()
1014 .collect::<Vec<_>>();
1015 self.defer_lease(&ids)
1016 }
1017
1018 pub fn resume_all(&mut self) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
1019 let ids = self
1020 .reservations
1021 .resource_ids()
1022 .cloned()
1023 .collect::<Vec<_>>();
1024 self.resume_lease(&ids)
1025 }
1026
1027 pub fn cancel_all(&mut self) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
1028 let ids = self
1029 .reservations
1030 .resource_ids()
1031 .cloned()
1032 .collect::<Vec<_>>();
1033 self.cancel_lease(&ids)
1034 }
1035
1036 fn release_all_resources(
1040 &mut self,
1041 ) -> Result<ResourceTransitionReceipt, ResourceFailureReceipt> {
1042 if let Some(pending) = &self.pending_subset_release {
1043 return Err(pending.failure.clone());
1044 }
1045 let orders = (0..self.states.len()).collect::<Vec<_>>();
1046 if let Some(&order) = orders
1047 .iter()
1048 .find(|&&order| self.states[order] != ResourceTransactionState::Committed)
1049 {
1050 return Err(self.local_release_failure(format!(
1051 "resource `{}` is not actually committed",
1052 self.reservations.reservations[order].resource_id()
1053 )));
1054 }
1055 let ledger_before = self.ledger_snapshot_entries();
1056 let mut completed = Vec::new();
1057 for &order in &orders {
1058 let result = {
1059 let reservation = self.reservations.reservations[order].clone();
1060 let buffer = self
1061 .lease
1062 .as_ref()
1063 .expect("committed transaction owns lease")
1064 .buffer(order)
1065 .expect("committed ledger entry owns its buffer");
1066 let driver = self.driver.as_mut().expect("live transaction owns driver");
1067 let context = ResourceTransactionContext {
1068 runtime: &self.lease.as_ref().expect("transaction owns lease").runtime,
1069 identity: &self.identity,
1070 binding: &self.admission,
1071 reservations: &self.reservations,
1072 cursor: Some(ResourceActionCursor {
1073 order,
1074 action: ResourceTransactionAction::Release,
1075 before: self.states[order],
1076 allocation_authorized: false,
1077 }),
1078 allocation_authority: None,
1079 pending_allocation: None,
1080 };
1081 driver.release_resource(&context, &reservation, buffer)
1082 };
1083 match result {
1084 Ok(()) => {
1085 let before = self.states[order];
1086 self.lease
1087 .as_mut()
1088 .expect("transaction owns lease ledger")
1089 .clear(order);
1090 self.states[order] = ResourceTransactionState::Released;
1091 self.capacity_claim
1092 .as_mut()
1093 .expect("live transaction owns its static capacity claim")
1094 .release_bytes(self.reservations.reservations[order].size_bytes());
1095 completed.push(self.make_record(
1096 order,
1097 ResourceTransactionAction::Release,
1098 before,
1099 ResourceTransactionState::Released,
1100 ));
1101 }
1102 Err(failure) => {
1103 let failure_id = self.issue_failure_id();
1104 let receipt = ResourceFailureReceipt::new(
1105 failure_id,
1106 &self.identity,
1107 &self.admission,
1108 ResourceTransactionAction::Release,
1109 failure.into_failure(),
1110 Some(ResourceFailurePoint::new(
1111 &self.reservations.reservations[order],
1112 order,
1113 self.states[order],
1114 )),
1115 completed,
1116 ResourceRecoveryStrategy::ForwardCompletion,
1117 ledger_before,
1118 self.ledger_snapshot_entries(),
1119 );
1120 self.set_pending_failure(&receipt);
1121 self.pending_subset_release = Some(PendingSubsetRelease {
1122 target_orders: orders,
1123 failure: receipt.clone(),
1124 });
1125 return Err(receipt);
1126 }
1127 }
1128 }
1129 let receipt = self
1130 .transition_receipt(ResourceTransactionAction::Release, ledger_before, completed)
1131 .expect("core subset release journal must validate");
1132 self.receipts.push(receipt.clone());
1133 Ok(receipt)
1134 }
1135
1136 pub fn complete_pending_release(
1137 &mut self,
1138 ) -> Result<ResourceTransitionReceipt, ResourceFailureReceipt> {
1139 let Some(mut pending) = self.pending_subset_release.take() else {
1140 return Err(self.local_release_failure(
1141 "there is no pending subset release to complete".to_owned(),
1142 ));
1143 };
1144 for &order in &pending.target_orders {
1145 match self.states[order] {
1146 ResourceTransactionState::Released => continue,
1147 ResourceTransactionState::Committed => {}
1148 actual => {
1149 let failure = ResourceDriverFailure::new(core_resource_failure(
1150 "release_ledger_diverged",
1151 format!(
1152 "resource `{}` has unexpected actual state {}",
1153 self.reservations.reservations[order].resource_id(),
1154 actual.as_str()
1155 ),
1156 false,
1157 ))
1158 .expect("core failure has resource domain");
1159 pending.failure.record_recovery_failure(
1160 failure,
1161 Some(ResourceFailurePoint::new(
1162 &self.reservations.reservations[order],
1163 order,
1164 actual,
1165 )),
1166 );
1167 pending.failure.ledger_after = self.ledger_snapshot_entries();
1168 self.set_pending_failure(&pending.failure);
1169 self.pending_subset_release = Some(pending);
1170 return Err(self
1171 .pending_subset_release
1172 .as_ref()
1173 .expect("pending release was restored")
1174 .failure
1175 .clone());
1176 }
1177 }
1178 let result = {
1179 let reservation = self.reservations.reservations[order].clone();
1180 let buffer = self
1181 .lease
1182 .as_ref()
1183 .expect("committed transaction owns lease")
1184 .buffer(order)
1185 .expect("outstanding release owns its buffer");
1186 let driver = self.driver.as_mut().expect("live transaction owns driver");
1187 let context = ResourceTransactionContext {
1188 runtime: &self.lease.as_ref().expect("transaction owns lease").runtime,
1189 identity: &self.identity,
1190 binding: &self.admission,
1191 reservations: &self.reservations,
1192 cursor: Some(ResourceActionCursor {
1193 order,
1194 action: ResourceTransactionAction::Release,
1195 before: self.states[order],
1196 allocation_authorized: false,
1197 }),
1198 allocation_authority: None,
1199 pending_allocation: None,
1200 };
1201 driver.release_resource(&context, &reservation, buffer)
1202 };
1203 match result {
1204 Ok(()) => {
1205 let before = self.states[order];
1206 self.lease
1207 .as_mut()
1208 .expect("transaction owns lease ledger")
1209 .clear(order);
1210 self.states[order] = ResourceTransactionState::Released;
1211 self.capacity_claim
1212 .as_mut()
1213 .expect("live transaction owns its static capacity claim")
1214 .release_bytes(self.reservations.reservations[order].size_bytes());
1215 pending.failure.completed.push(self.make_record(
1216 order,
1217 ResourceTransactionAction::Release,
1218 before,
1219 ResourceTransactionState::Released,
1220 ));
1221 }
1222 Err(failure) => {
1223 pending.failure.record_recovery_failure(
1224 failure,
1225 Some(ResourceFailurePoint::new(
1226 &self.reservations.reservations[order],
1227 order,
1228 self.states[order],
1229 )),
1230 );
1231 pending.failure.ledger_after = self.ledger_snapshot_entries();
1232 self.set_pending_failure(&pending.failure);
1233 let snapshot = pending.failure.clone();
1234 self.pending_subset_release = Some(pending);
1235 return Err(snapshot);
1236 }
1237 }
1238 }
1239 pending.failure.recovery_complete = true;
1240 pending.failure.ledger_after = self.ledger_snapshot_entries();
1241 let receipt = self
1242 .transition_receipt(
1243 ResourceTransactionAction::Release,
1244 pending.failure.ledger_before.clone(),
1245 pending.failure.completed.clone(),
1246 )
1247 .expect("forward-completed release journal must validate");
1248 self.clear_pending_failure();
1249 self.recovery_history.push(pending.failure);
1250 self.receipts.push(receipt.clone());
1251 Ok(receipt)
1252 }
1253
1254 fn local_release_failure(&mut self, message: String) -> ResourceFailureReceipt {
1255 let ledger = self.ledger_snapshot_entries();
1256 let failure_id = self.issue_failure_id();
1257 ResourceFailureReceipt::new(
1258 failure_id,
1259 &self.identity,
1260 &self.admission,
1261 ResourceTransactionAction::Release,
1262 core_resource_failure("invalid_release_request", message, false),
1263 None,
1264 Vec::new(),
1265 ResourceRecoveryStrategy::ForwardCompletion,
1266 ledger.clone(),
1267 ledger,
1268 )
1269 }
1270
1271 pub fn release(
1272 mut self,
1273 ) -> Result<ResourceTransaction<D, TransactionReleased>, ResourceReleaseTransitionError<D>>
1274 {
1275 if self.pending_subset_release.is_some() {
1276 let failure = self
1277 .pending_subset_release
1278 .as_ref()
1279 .expect("pending release exists")
1280 .failure
1281 .clone();
1282 return Err(ResourceReleaseTransitionError {
1283 transaction: Some(self),
1284 failure,
1285 });
1286 }
1287 if self
1288 .states
1289 .iter()
1290 .all(|state| *state == ResourceTransactionState::Released)
1291 {
1292 return Ok(self.advance(None));
1293 }
1294 if !self
1295 .states
1296 .iter()
1297 .all(|state| *state == ResourceTransactionState::Committed)
1298 {
1299 let failure = self.local_release_failure(
1300 "transaction contains non-releasable actual states".to_owned(),
1301 );
1302 return Err(ResourceReleaseTransitionError {
1303 transaction: Some(self),
1304 failure,
1305 });
1306 }
1307 match self.release_all_resources() {
1308 Ok(_) => Ok(self.advance(None)),
1309 Err(failure) => Err(ResourceReleaseTransitionError {
1310 transaction: Some(self),
1311 failure,
1312 }),
1313 }
1314 }
1315}
1316
1317#[must_use = "failed plan runtime handoff retains the committed transaction"]
1318pub struct PlanRuntimeHandoffError<D>
1319where
1320 D: ResourceTransactionDriver,
1321{
1322 error: VNextError,
1323 transaction: Option<ResourceTransaction<D, TransactionCommitted>>,
1324}
1325
1326impl<D> PlanRuntimeHandoffError<D>
1327where
1328 D: ResourceTransactionDriver,
1329{
1330 pub fn error(&self) -> &VNextError {
1331 &self.error
1332 }
1333
1334 pub fn into_transaction(mut self) -> ResourceTransaction<D, TransactionCommitted> {
1335 self.transaction
1336 .take()
1337 .expect("plan runtime handoff error owns its transaction")
1338 }
1339}
1340
1341#[must_use = "reverse recovery must complete or ownership must be quarantined"]
1342pub struct ResourcePrepareTransitionError<D: ResourceTransactionDriver, S: TransactionStage> {
1343 transaction: Option<ResourceTransaction<D, S>>,
1344 failure: ResourceFailureReceipt,
1345}
1346
1347impl<D: ResourceTransactionDriver, S: TransactionStage> fmt::Debug
1348 for ResourcePrepareTransitionError<D, S>
1349{
1350 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1351 formatter
1352 .debug_struct("ResourcePrepareTransitionError")
1353 .field("stage", &S::STATE)
1354 .field("failure", &self.failure)
1355 .finish_non_exhaustive()
1356 }
1357}
1358
1359impl<D: ResourceTransactionDriver, S: TransactionStage> ResourcePrepareTransitionError<D, S> {
1360 pub fn failure(&self) -> &ResourceFailureReceipt {
1361 &self.failure
1362 }
1363
1364 pub fn recover(mut self) -> Result<ResourceTransaction<D, S>, Self> {
1365 while self.failure.compensation.len() < self.failure.completed.len() {
1366 let record_index = self.failure.completed.len() - 1 - self.failure.compensation.len();
1367 let attempted = self.failure.completed[record_index].clone();
1368 let order = attempted.order as usize;
1369 let transaction = self
1370 .transaction
1371 .as_mut()
1372 .expect("prepare recovery owns transaction");
1373 if transaction.states[order] != attempted.after {
1374 let failure = ResourceDriverFailure::new(core_resource_failure(
1375 "compensation_ledger_diverged",
1376 format!(
1377 "resource `{}` expected actual state {}, found {}",
1378 attempted.resource_id,
1379 attempted.after.as_str(),
1380 transaction.states[order].as_str()
1381 ),
1382 false,
1383 ))
1384 .expect("core failure has resource domain");
1385 self.failure.record_recovery_failure(
1386 failure,
1387 Some(ResourceFailurePoint::new(
1388 &transaction.reservations.reservations[order],
1389 order,
1390 transaction.states[order],
1391 )),
1392 );
1393 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1394 transaction.set_pending_failure(&self.failure);
1395 return Err(self);
1396 }
1397 let result = match self.failure.action {
1398 ResourceTransactionAction::Reserve => {
1399 let reservation = transaction.reservations.reservations[order].clone();
1400 let driver = transaction
1401 .driver
1402 .as_mut()
1403 .expect("live transaction owns driver");
1404 let context = ResourceTransactionContext {
1405 runtime: &transaction
1406 .lease
1407 .as_ref()
1408 .expect("transaction owns lease")
1409 .runtime,
1410 identity: &transaction.identity,
1411 binding: &transaction.admission,
1412 reservations: &transaction.reservations,
1413 cursor: Some(ResourceActionCursor {
1414 order,
1415 action: ResourceTransactionAction::Reserve,
1416 before: transaction.states[order],
1417 allocation_authorized: false,
1418 }),
1419 allocation_authority: None,
1420 pending_allocation: None,
1421 };
1422 driver.compensate_reserve_resource(&context, &reservation)
1423 }
1424 ResourceTransactionAction::Commit => {
1425 let reservation = transaction.reservations.reservations[order].clone();
1426 let buffer = transaction
1427 .lease
1428 .as_ref()
1429 .expect("transaction owns lease ledger")
1430 .buffer(order)
1431 .expect("uncompensated commit owns its buffer");
1432 let driver = transaction
1433 .driver
1434 .as_mut()
1435 .expect("live transaction owns driver");
1436 let context = ResourceTransactionContext {
1437 runtime: &transaction
1438 .lease
1439 .as_ref()
1440 .expect("transaction owns lease")
1441 .runtime,
1442 identity: &transaction.identity,
1443 binding: &transaction.admission,
1444 reservations: &transaction.reservations,
1445 cursor: Some(ResourceActionCursor {
1446 order,
1447 action: ResourceTransactionAction::Commit,
1448 before: transaction.states[order],
1449 allocation_authorized: false,
1450 }),
1451 allocation_authority: None,
1452 pending_allocation: None,
1453 };
1454 driver.compensate_commit_resource(&context, &reservation, buffer)
1455 }
1456 _ => unreachable!("prepare recovery handles only reserve and commit"),
1457 };
1458 match result {
1459 Ok(()) => {
1460 transaction.states[order] = attempted.before;
1461 if self.failure.action == ResourceTransactionAction::Commit {
1462 transaction
1463 .lease
1464 .as_mut()
1465 .expect("transaction owns lease ledger")
1466 .clear(order);
1467 }
1468 self.failure
1469 .compensation
1470 .push(ResourceCompensationRecord::from_transition(
1471 &attempted,
1472 self.failure.compensation.len(),
1473 ));
1474 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1475 }
1476 Err(failure) => {
1477 self.failure.record_recovery_failure(
1478 failure,
1479 Some(ResourceFailurePoint::new(
1480 &transaction.reservations.reservations[order],
1481 order,
1482 transaction.states[order],
1483 )),
1484 );
1485 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1486 transaction.set_pending_failure(&self.failure);
1487 return Err(self);
1488 }
1489 }
1490 }
1491 self.failure.recovery_complete = true;
1492 let mut transaction = self
1493 .transaction
1494 .take()
1495 .expect("prepare recovery owns transaction");
1496 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1497 transaction.clear_pending_failure();
1498 if self.failure.action == ResourceTransactionAction::Commit {
1499 for authority in &transaction.allocation_issued {
1500 authority.store(false, Ordering::Release);
1501 }
1502 }
1503 transaction.recovery_history.push(self.failure);
1504 Ok(transaction)
1505 }
1506
1507 pub fn quarantine(mut self) -> Result<ResourceTransaction<D, TransactionQuarantined>, Self> {
1508 let transaction = self
1509 .transaction
1510 .as_mut()
1511 .expect("prepare recovery owns transaction");
1512 let ledger_before = transaction.ledger_snapshot_entries();
1513 match transaction.quarantine_live() {
1514 Ok(records) => {
1515 self.failure.recovery_complete = true;
1516 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1517 transaction.clear_pending_failure();
1518 transaction.recovery_history.push(self.failure.clone());
1519 let receipt = transaction
1520 .transition_receipt(
1521 ResourceTransactionAction::Quarantine,
1522 ledger_before,
1523 records,
1524 )
1525 .expect("core quarantine journal must validate");
1526 let transaction = self
1527 .transaction
1528 .take()
1529 .expect("prepare recovery owns transaction");
1530 Ok(transaction.advance(Some(receipt)))
1531 }
1532 Err(failure) => {
1533 self.failure.record_recovery_failure(failure, None);
1534 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1535 transaction.set_pending_failure(&self.failure);
1536 Err(self)
1537 }
1538 }
1539 }
1540}
1541
1542#[must_use = "poisoned commit ownership must be reconciled or quarantined"]
1543pub struct ResourcePoisonedTransaction<D: ResourceTransactionDriver> {
1544 transaction: Option<ResourceTransaction<D, TransactionReserved>>,
1545 failure: ResourceFailureReceipt,
1546 poisoned_order: usize,
1547}
1548
1549impl<D: ResourceTransactionDriver> fmt::Debug for ResourcePoisonedTransaction<D> {
1550 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1551 formatter
1552 .debug_struct("ResourcePoisonedTransaction")
1553 .field("failure", &self.failure)
1554 .field("poisoned_order", &self.poisoned_order)
1555 .finish_non_exhaustive()
1556 }
1557}
1558
1559impl<D: ResourceTransactionDriver> ResourcePoisonedTransaction<D> {
1560 pub fn failure(&self) -> &ResourceFailureReceipt {
1561 &self.failure
1562 }
1563
1564 pub fn reconcile(
1565 mut self,
1566 ) -> Result<ResourcePrepareTransitionError<D, TransactionReserved>, Self> {
1567 let transaction = self
1568 .transaction
1569 .as_mut()
1570 .expect("poisoned owner owns transaction");
1571 let reservation = transaction.reservations.reservations[self.poisoned_order].clone();
1572 let result = {
1573 let lease = transaction
1574 .lease
1575 .as_ref()
1576 .expect("poisoned transaction owns lease ledger");
1577 let slot = &lease.slots[self.poisoned_order];
1578 let actual = ResourceCommitView {
1579 resource_id: slot
1580 .actual_resource_id
1581 .as_ref()
1582 .expect("poisoned allocation records actual resource identity"),
1583 generation: slot
1584 .actual_generation
1585 .expect("poisoned allocation records actual generation"),
1586 descriptor: slot
1587 .descriptor
1588 .as_ref()
1589 .expect("poisoned allocation records actual descriptor"),
1590 buffer: slot
1591 .buffer
1592 .as_ref()
1593 .expect("poisoned allocation remains core-owned"),
1594 };
1595 let driver = transaction
1596 .driver
1597 .as_mut()
1598 .expect("live transaction owns driver");
1599 let context = ResourceTransactionContext {
1600 runtime: &lease.runtime,
1601 identity: &transaction.identity,
1602 binding: &transaction.admission,
1603 reservations: &transaction.reservations,
1604 cursor: Some(ResourceActionCursor {
1605 order: self.poisoned_order,
1606 action: ResourceTransactionAction::Commit,
1607 before: transaction.states[self.poisoned_order],
1608 allocation_authorized: false,
1609 }),
1610 allocation_authority: None,
1611 pending_allocation: None,
1612 };
1613 driver.reconcile_commit_outcome(&context, &reservation, actual)
1614 };
1615 match result {
1616 Ok(()) => {
1617 transaction
1618 .lease
1619 .as_mut()
1620 .expect("poisoned transaction owns lease ledger")
1621 .clear(self.poisoned_order);
1622 self.failure.recovery_strategy = ResourceRecoveryStrategy::ReverseCompensation;
1623 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1624 let transaction = self
1625 .transaction
1626 .take()
1627 .expect("poisoned owner owns transaction");
1628 Ok(ResourcePrepareTransitionError {
1629 transaction: Some(transaction),
1630 failure: self.failure,
1631 })
1632 }
1633 Err(failure) => {
1634 self.failure.record_recovery_failure(
1635 failure,
1636 Some(ResourceFailurePoint::new(
1637 &reservation,
1638 self.poisoned_order,
1639 transaction.states[self.poisoned_order],
1640 )),
1641 );
1642 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1643 transaction.set_pending_failure(&self.failure);
1644 Err(self)
1645 }
1646 }
1647 }
1648
1649 pub fn quarantine(mut self) -> Result<ResourceTransaction<D, TransactionQuarantined>, Self> {
1650 let transaction = self
1651 .transaction
1652 .as_mut()
1653 .expect("poisoned owner owns transaction");
1654 let ledger_before = transaction.ledger_snapshot_entries();
1655 match transaction.quarantine_live() {
1656 Ok(records) => {
1657 self.failure.recovery_complete = true;
1658 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1659 transaction.clear_pending_failure();
1660 transaction.recovery_history.push(self.failure.clone());
1661 let receipt = transaction
1662 .transition_receipt(
1663 ResourceTransactionAction::Quarantine,
1664 ledger_before,
1665 records,
1666 )
1667 .expect("poison quarantine journal must validate");
1668 let transaction = self
1669 .transaction
1670 .take()
1671 .expect("poisoned owner owns transaction");
1672 Ok(transaction.advance(Some(receipt)))
1673 }
1674 Err(failure) => {
1675 self.failure.record_recovery_failure(failure, None);
1676 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1677 transaction.set_pending_failure(&self.failure);
1678 Err(self)
1679 }
1680 }
1681 }
1682}
1683
1684#[must_use = "rollback forward completion owns the transaction"]
1685pub struct ResourceRollbackTransitionError<D: ResourceTransactionDriver> {
1686 transaction: Option<ResourceTransaction<D, TransactionReserved>>,
1687 failure: ResourceFailureReceipt,
1688}
1689
1690impl<D: ResourceTransactionDriver> fmt::Debug for ResourceRollbackTransitionError<D> {
1691 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1692 formatter
1693 .debug_struct("ResourceRollbackTransitionError")
1694 .field("failure", &self.failure)
1695 .finish_non_exhaustive()
1696 }
1697}
1698
1699impl<D: ResourceTransactionDriver> ResourceRollbackTransitionError<D> {
1700 pub fn failure(&self) -> &ResourceFailureReceipt {
1701 &self.failure
1702 }
1703
1704 pub fn complete(mut self) -> Result<ResourceTransaction<D, TransactionRolledBack>, Self> {
1705 let transaction = self
1706 .transaction
1707 .as_mut()
1708 .expect("rollback recovery owns transaction");
1709 for order in 0..transaction.states.len() {
1710 match transaction.states[order] {
1711 ResourceTransactionState::RolledBack => continue,
1712 ResourceTransactionState::Reserved => {}
1713 actual => {
1714 let failure = ResourceDriverFailure::new(core_resource_failure(
1715 "rollback_ledger_diverged",
1716 format!("rollback found actual state {}", actual.as_str()),
1717 false,
1718 ))
1719 .expect("core failure has resource domain");
1720 self.failure.record_recovery_failure(
1721 failure,
1722 Some(ResourceFailurePoint::new(
1723 &transaction.reservations.reservations[order],
1724 order,
1725 actual,
1726 )),
1727 );
1728 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1729 transaction.set_pending_failure(&self.failure);
1730 return Err(self);
1731 }
1732 }
1733 let result = {
1734 let reservation = transaction.reservations.reservations[order].clone();
1735 let driver = transaction
1736 .driver
1737 .as_mut()
1738 .expect("live transaction owns driver");
1739 let context = ResourceTransactionContext {
1740 runtime: &transaction
1741 .lease
1742 .as_ref()
1743 .expect("transaction owns lease")
1744 .runtime,
1745 identity: &transaction.identity,
1746 binding: &transaction.admission,
1747 reservations: &transaction.reservations,
1748 cursor: Some(ResourceActionCursor {
1749 order,
1750 action: ResourceTransactionAction::Rollback,
1751 before: transaction.states[order],
1752 allocation_authorized: false,
1753 }),
1754 allocation_authority: None,
1755 pending_allocation: None,
1756 };
1757 driver.rollback_resource(&context, &reservation)
1758 };
1759 match result {
1760 Ok(()) => {
1761 let before = transaction.states[order];
1762 transaction.states[order] = ResourceTransactionState::RolledBack;
1763 self.failure.completed.push(transaction.make_record(
1764 order,
1765 ResourceTransactionAction::Rollback,
1766 before,
1767 ResourceTransactionState::RolledBack,
1768 ));
1769 }
1770 Err(failure) => {
1771 self.failure.record_recovery_failure(
1772 failure,
1773 Some(ResourceFailurePoint::new(
1774 &transaction.reservations.reservations[order],
1775 order,
1776 transaction.states[order],
1777 )),
1778 );
1779 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1780 transaction.set_pending_failure(&self.failure);
1781 return Err(self);
1782 }
1783 }
1784 }
1785 self.failure.recovery_complete = true;
1786 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1787 let receipt = transaction
1788 .transition_receipt(
1789 ResourceTransactionAction::Rollback,
1790 self.failure.ledger_before.clone(),
1791 self.failure.completed.clone(),
1792 )
1793 .expect("forward-completed rollback journal must validate");
1794 transaction.clear_pending_failure();
1795 transaction.recovery_history.push(self.failure);
1796 let transaction = self
1797 .transaction
1798 .take()
1799 .expect("rollback recovery owns transaction");
1800 Ok(transaction.advance(Some(receipt)))
1801 }
1802
1803 pub fn quarantine(mut self) -> Result<ResourceTransaction<D, TransactionQuarantined>, Self> {
1804 let transaction = self
1805 .transaction
1806 .as_mut()
1807 .expect("rollback recovery owns transaction");
1808 let ledger_before = transaction.ledger_snapshot_entries();
1809 match transaction.quarantine_live() {
1810 Ok(records) => {
1811 self.failure.recovery_complete = true;
1812 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1813 transaction.clear_pending_failure();
1814 transaction.recovery_history.push(self.failure.clone());
1815 let receipt = transaction
1816 .transition_receipt(
1817 ResourceTransactionAction::Quarantine,
1818 ledger_before,
1819 records,
1820 )
1821 .expect("rollback quarantine journal must validate");
1822 let transaction = self
1823 .transaction
1824 .take()
1825 .expect("rollback recovery owns transaction");
1826 Ok(transaction.advance(Some(receipt)))
1827 }
1828 Err(failure) => {
1829 self.failure.record_recovery_failure(failure, None);
1830 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1831 transaction.set_pending_failure(&self.failure);
1832 Err(self)
1833 }
1834 }
1835 }
1836}
1837
1838#[must_use = "release forward completion owns the transaction"]
1839pub struct ResourceReleaseTransitionError<D: ResourceTransactionDriver> {
1840 transaction: Option<ResourceTransaction<D, TransactionCommitted>>,
1841 failure: ResourceFailureReceipt,
1842}
1843
1844impl<D: ResourceTransactionDriver> fmt::Debug for ResourceReleaseTransitionError<D> {
1845 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1846 formatter
1847 .debug_struct("ResourceReleaseTransitionError")
1848 .field("failure", &self.failure)
1849 .finish_non_exhaustive()
1850 }
1851}
1852
1853impl<D: ResourceTransactionDriver> ResourceReleaseTransitionError<D> {
1854 pub fn failure(&self) -> &ResourceFailureReceipt {
1855 &self.failure
1856 }
1857
1858 pub fn complete(mut self) -> Result<ResourceTransaction<D, TransactionReleased>, Self> {
1859 let transaction = self
1860 .transaction
1861 .as_mut()
1862 .expect("release recovery owns transaction");
1863 match transaction.complete_pending_release() {
1864 Ok(_) => {
1865 self.failure = transaction
1866 .recovery_history
1867 .last()
1868 .cloned()
1869 .unwrap_or_else(|| self.failure.clone());
1870 let transaction = self
1871 .transaction
1872 .take()
1873 .expect("release recovery owns transaction");
1874 transaction.release().map_err(|next| Self {
1875 transaction: next.transaction,
1876 failure: next.failure,
1877 })
1878 }
1879 Err(failure) => {
1880 self.failure = failure;
1881 Err(self)
1882 }
1883 }
1884 }
1885
1886 pub fn quarantine(mut self) -> Result<ResourceTransaction<D, TransactionQuarantined>, Self> {
1887 let transaction = self
1888 .transaction
1889 .as_mut()
1890 .expect("release recovery owns transaction");
1891 let ledger_before = transaction.ledger_snapshot_entries();
1892 match transaction.quarantine_live() {
1893 Ok(records) => {
1894 self.failure.recovery_complete = true;
1895 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1896 transaction.clear_pending_failure();
1897 transaction.pending_subset_release = None;
1898 transaction.recovery_history.push(self.failure.clone());
1899 let receipt = transaction
1900 .transition_receipt(
1901 ResourceTransactionAction::Quarantine,
1902 ledger_before,
1903 records,
1904 )
1905 .expect("release quarantine journal must validate");
1906 let transaction = self
1907 .transaction
1908 .take()
1909 .expect("release recovery owns transaction");
1910 Ok(transaction.advance(Some(receipt)))
1911 }
1912 Err(failure) => {
1913 self.failure.record_recovery_failure(failure, None);
1914 self.failure.ledger_after = transaction.ledger_snapshot_entries();
1915 transaction.set_pending_failure(&self.failure);
1916 Err(self)
1917 }
1918 }
1919 }
1920}