Skip to main content

ferrum_interfaces/vnext/resource/
transaction.rs

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            // Driver cleanup is an observation/transfer hook, not permission
535            // to turn an existing unwind into a process abort. Ownership that
536            // drops inside a panicking callback retains its backend lifetimes.
537            let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
538                driver.abandon_transaction(ownership);
539            }));
540        } else {
541            // A live transaction always owns its driver. Leaking is safer than
542            // returning capacity while backend ownership is unknown.
543            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    /// Terminal release always targets the complete pool. On partial backend
1037    /// failure, the actual prefix remains released in the ledger and
1038    /// `complete_pending_release` resumes the outstanding suffix.
1039    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}