1use super::{
2 expected_transition, invalid_resource, BTreeSet, BufferDescriptor, BufferUsage, Deserialize,
3 DeviceId, ElementType, FailureEnvelope, NodeId, PlanHash, PlanId, RequestIdentity,
4 ResourceCompensationAction, ResourceDriverFailure, ResourceId, ResourcePoolId,
5 ResourceRecoveryStrategy, ResourceReservation, ResourceRetentionDecision,
6 ResourceRetentionPolicy, ResourceTransactionAction, ResourceTransactionIdentity,
7 ResourceTransactionState, RunId, Serialize, StaticProvisioningBinding, TransactionId,
8 VNextError, MAX_RESOURCE_LEASE_RECEIPT_WIRE_BYTES, MAX_RESOURCE_TRANSITION_RECEIPT_WIRE_BYTES,
9};
10
11#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
12#[serde(rename_all = "snake_case")]
13pub enum ResourceLeaseState {
14 Active,
15 Deferred,
16 Cancelled,
17 Mixed,
18}
19
20impl ResourceLeaseState {
21 pub const fn as_str(self) -> &'static str {
22 match self {
23 Self::Active => "active",
24 Self::Deferred => "deferred",
25 Self::Cancelled => "cancelled",
26 Self::Mixed => "mixed",
27 }
28 }
29}
30
31#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
32#[serde(rename_all = "snake_case")]
33pub enum ResourceLeaseAction {
34 Defer,
35 Resume,
36 Cancel,
37}
38
39impl ResourceLeaseAction {
40 pub const fn as_str(self) -> &'static str {
41 match self {
42 Self::Defer => "defer",
43 Self::Resume => "resume",
44 Self::Cancel => "cancel",
45 }
46 }
47
48 const fn decision(self) -> ResourceRetentionDecision {
49 match self {
50 Self::Defer | Self::Resume => ResourceRetentionDecision::Retain,
51 Self::Cancel => ResourceRetentionDecision::ReturnRequested,
52 }
53 }
54}
55
56pub(super) const fn expected_lease_transition(
57 action: ResourceLeaseAction,
58 before: ResourceLeaseState,
59) -> Option<ResourceLeaseState> {
60 match (before, action) {
61 (ResourceLeaseState::Active, ResourceLeaseAction::Defer) => {
62 Some(ResourceLeaseState::Deferred)
63 }
64 (ResourceLeaseState::Deferred, ResourceLeaseAction::Resume) => {
65 Some(ResourceLeaseState::Active)
66 }
67 (
68 ResourceLeaseState::Active | ResourceLeaseState::Deferred,
69 ResourceLeaseAction::Cancel,
70 ) => Some(ResourceLeaseState::Cancelled),
71 _ => None,
72 }
73}
74
75#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize)]
76pub struct ResourceLeaseEntry {
77 pub(super) owner_node_id: Option<NodeId>,
78 pub(super) resource_id: ResourceId,
79 pub(super) size_bytes: u64,
80 pub(super) alignment_bytes: u64,
81 pub(super) usage: BufferUsage,
82 pub(super) element_type: ElementType,
83 pub(super) retention_policy: ResourceRetentionPolicy,
84 pub(super) generation: u64,
85 pub(super) state: ResourceLeaseState,
86}
87
88impl ResourceLeaseEntry {
89 pub(super) fn from_reservation(
90 reservation: &ResourceReservation,
91 state: ResourceLeaseState,
92 ) -> Self {
93 Self {
94 owner_node_id: reservation.owner_node_id.clone(),
95 resource_id: reservation.resource_id.clone(),
96 size_bytes: reservation.size_bytes,
97 alignment_bytes: reservation.alignment_bytes,
98 usage: reservation.usage,
99 element_type: reservation.element_type,
100 retention_policy: reservation.retention_policy,
101 generation: reservation.generation,
102 state,
103 }
104 }
105
106 pub fn owner_node_id(&self) -> Option<&NodeId> {
107 self.owner_node_id.as_ref()
108 }
109
110 pub fn resource_id(&self) -> &ResourceId {
111 &self.resource_id
112 }
113
114 pub const fn size_bytes(&self) -> u64 {
115 self.size_bytes
116 }
117
118 pub const fn alignment_bytes(&self) -> u64 {
119 self.alignment_bytes
120 }
121
122 pub const fn usage(&self) -> BufferUsage {
123 self.usage
124 }
125
126 pub const fn element_type(&self) -> ElementType {
127 self.element_type
128 }
129
130 pub const fn retention_policy(&self) -> ResourceRetentionPolicy {
131 self.retention_policy
132 }
133
134 pub const fn generation(&self) -> u64 {
135 self.generation
136 }
137
138 pub const fn state(&self) -> ResourceLeaseState {
139 self.state
140 }
141}
142
143#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
144#[serde(deny_unknown_fields)]
145pub struct UnvalidatedResourceLeaseEntry {
146 pub owner_node_id: Option<NodeId>,
147 pub resource_id: ResourceId,
148 pub size_bytes: u64,
149 pub alignment_bytes: u64,
150 pub usage: BufferUsage,
151 pub element_type: ElementType,
152 pub retention_policy: ResourceRetentionPolicy,
153 pub generation: u64,
154 pub state: ResourceLeaseState,
155}
156
157impl UnvalidatedResourceLeaseEntry {
158 fn into_trusted(self) -> Result<ResourceLeaseEntry, VNextError> {
159 if self.size_bytes == 0
160 || self.alignment_bytes == 0
161 || !self.alignment_bytes.is_power_of_two()
162 || self.generation == 0
163 || self.state == ResourceLeaseState::Mixed
164 {
165 return Err(invalid_resource("untrusted lease entry is malformed"));
166 }
167 Ok(ResourceLeaseEntry {
168 owner_node_id: self.owner_node_id,
169 resource_id: self.resource_id,
170 size_bytes: self.size_bytes,
171 alignment_bytes: self.alignment_bytes,
172 usage: self.usage,
173 element_type: self.element_type,
174 retention_policy: self.retention_policy,
175 generation: self.generation,
176 state: self.state,
177 })
178 }
179}
180
181impl From<&ResourceLeaseEntry> for UnvalidatedResourceLeaseEntry {
182 fn from(entry: &ResourceLeaseEntry) -> Self {
183 Self {
184 owner_node_id: entry.owner_node_id.clone(),
185 resource_id: entry.resource_id.clone(),
186 size_bytes: entry.size_bytes,
187 alignment_bytes: entry.alignment_bytes,
188 usage: entry.usage,
189 element_type: entry.element_type,
190 retention_policy: entry.retention_policy,
191 generation: entry.generation,
192 state: entry.state,
193 }
194 }
195}
196
197#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
198pub struct ResourceLedgerEntrySnapshot {
199 pub(super) entry: ResourceLeaseEntry,
200 pub(super) transaction_state: ResourceTransactionState,
201 pub(super) buffer_present: bool,
202 pub(super) actual_resource_id: Option<ResourceId>,
203 pub(super) actual_generation: Option<u64>,
204 pub(super) actual_descriptor: Option<BufferDescriptor>,
205}
206
207impl ResourceLedgerEntrySnapshot {
208 pub fn entry(&self) -> &ResourceLeaseEntry {
209 &self.entry
210 }
211
212 pub const fn transaction_state(&self) -> ResourceTransactionState {
213 self.transaction_state
214 }
215
216 pub const fn buffer_present(&self) -> bool {
217 self.buffer_present
218 }
219
220 pub fn actual_resource_id(&self) -> Option<&ResourceId> {
221 self.actual_resource_id.as_ref()
222 }
223
224 pub const fn actual_generation(&self) -> Option<u64> {
225 self.actual_generation
226 }
227
228 pub fn actual_descriptor(&self) -> Option<&BufferDescriptor> {
229 self.actual_descriptor.as_ref()
230 }
231}
232
233#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
234pub struct ResourceLedgerSnapshot {
235 pub(super) identity: ResourceTransactionIdentity,
236 pub(super) admission: StaticProvisioningBinding,
237 pub(super) entries: Vec<ResourceLedgerEntrySnapshot>,
238}
239
240impl ResourceLedgerSnapshot {
241 pub fn identity(&self) -> &ResourceTransactionIdentity {
242 &self.identity
243 }
244
245 pub fn admission(&self) -> &StaticProvisioningBinding {
246 &self.admission
247 }
248
249 pub fn entries(&self) -> &[ResourceLedgerEntrySnapshot] {
250 &self.entries
251 }
252}
253
254#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
257pub struct ResourceTransitionValidationContext {
258 pub(super) identity: ResourceTransactionIdentity,
259 pub(super) admission: StaticProvisioningBinding,
260 pub(super) action: ResourceTransactionAction,
261 pub(super) before: Vec<ResourceLedgerEntrySnapshot>,
262 pub(super) after: Vec<ResourceLedgerEntrySnapshot>,
263}
264
265impl ResourceTransitionValidationContext {
266 pub fn identity(&self) -> &ResourceTransactionIdentity {
267 &self.identity
268 }
269
270 pub fn admission(&self) -> &StaticProvisioningBinding {
271 &self.admission
272 }
273
274 pub const fn action(&self) -> ResourceTransactionAction {
275 self.action
276 }
277
278 pub fn before(&self) -> &[ResourceLedgerEntrySnapshot] {
279 &self.before
280 }
281
282 pub fn after(&self) -> &[ResourceLedgerEntrySnapshot] {
283 &self.after
284 }
285}
286
287#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
288pub struct ResourceLeaseValidationContext {
289 pub(super) identity: ResourceTransactionIdentity,
290 pub(super) admission: StaticProvisioningBinding,
291 pub(super) action: ResourceLeaseAction,
292 pub(super) before: Vec<ResourceLedgerEntrySnapshot>,
293 pub(super) after: Vec<ResourceLedgerEntrySnapshot>,
294}
295
296impl ResourceLeaseValidationContext {
297 pub fn identity(&self) -> &ResourceTransactionIdentity {
298 &self.identity
299 }
300
301 pub fn admission(&self) -> &StaticProvisioningBinding {
302 &self.admission
303 }
304
305 pub const fn action(&self) -> ResourceLeaseAction {
306 self.action
307 }
308
309 pub fn before(&self) -> &[ResourceLedgerEntrySnapshot] {
310 &self.before
311 }
312
313 pub fn after(&self) -> &[ResourceLedgerEntrySnapshot] {
314 &self.after
315 }
316}
317
318#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
319pub struct ResourceTransitionRecord {
320 run_id: RunId,
321 transaction_id: TransactionId,
322 request_id: RequestIdentity,
323 plan_id: PlanId,
324 plan_hash: PlanHash,
325 device_id: DeviceId,
326 admission_generation: u64,
327 owner_node_id: Option<NodeId>,
328 pub(super) resource_id: ResourceId,
329 generation: u64,
330 retention_policy: ResourceRetentionPolicy,
331 pub(super) action: ResourceTransactionAction,
332 pub(super) before: ResourceTransactionState,
333 pub(super) after: ResourceTransactionState,
334 pub(super) order: u32,
335}
336
337impl ResourceTransitionRecord {
338 pub(super) fn from_reservation(
339 identity: &ResourceTransactionIdentity,
340 admission: &StaticProvisioningBinding,
341 reservation: &ResourceReservation,
342 action: ResourceTransactionAction,
343 before: ResourceTransactionState,
344 after: ResourceTransactionState,
345 order: usize,
346 ) -> Self {
347 Self {
348 run_id: identity.run_id.clone(),
349 transaction_id: identity.transaction_id.clone(),
350 request_id: identity.request_id.clone(),
351 plan_id: admission.plan_id.clone(),
352 plan_hash: admission.plan_hash.clone(),
353 device_id: admission.device_id.clone(),
354 admission_generation: admission.admission_generation,
355 owner_node_id: reservation.owner_node_id.clone(),
356 resource_id: reservation.resource_id.clone(),
357 generation: reservation.generation,
358 retention_policy: reservation.retention_policy,
359 action,
360 before,
361 after,
362 order: order as u32,
363 }
364 }
365
366 pub fn validate(&self) -> Result<(), VNextError> {
367 if self.generation == 0
368 || self.admission_generation == 0
369 || self.generation != self.admission_generation
370 || expected_transition(self.action, self.before) != Some(self.after)
371 {
372 return Err(VNextError::InvalidResourceTransition {
373 resource_id: self.resource_id.to_string(),
374 from: self.before.as_str(),
375 action: self.action.as_str(),
376 });
377 }
378 Ok(())
379 }
380
381 pub fn run_id(&self) -> &RunId {
382 &self.run_id
383 }
384
385 pub fn transaction_id(&self) -> &TransactionId {
386 &self.transaction_id
387 }
388
389 pub fn request_id(&self) -> &RequestIdentity {
390 &self.request_id
391 }
392
393 pub fn plan_id(&self) -> &PlanId {
394 &self.plan_id
395 }
396
397 pub fn plan_hash(&self) -> &PlanHash {
398 &self.plan_hash
399 }
400
401 pub fn device_id(&self) -> &DeviceId {
402 &self.device_id
403 }
404
405 pub const fn admission_generation(&self) -> u64 {
406 self.admission_generation
407 }
408
409 pub fn owner_node_id(&self) -> Option<&NodeId> {
410 self.owner_node_id.as_ref()
411 }
412
413 pub fn resource_id(&self) -> &ResourceId {
414 &self.resource_id
415 }
416
417 pub const fn generation(&self) -> u64 {
418 self.generation
419 }
420
421 pub const fn retention_policy(&self) -> ResourceRetentionPolicy {
422 self.retention_policy
423 }
424
425 pub const fn action(&self) -> ResourceTransactionAction {
426 self.action
427 }
428
429 pub const fn before(&self) -> ResourceTransactionState {
430 self.before
431 }
432
433 pub const fn after(&self) -> ResourceTransactionState {
434 self.after
435 }
436
437 pub const fn order(&self) -> u32 {
438 self.order
439 }
440
441 pub(super) fn matches_identity_and_admission(
442 &self,
443 identity: &ResourceTransactionIdentity,
444 admission: &StaticProvisioningBinding,
445 ) -> bool {
446 self.run_id == identity.run_id
447 && self.transaction_id == identity.transaction_id
448 && self.request_id == identity.request_id
449 && self.plan_id == admission.plan_id
450 && self.plan_hash == admission.plan_hash
451 && self.device_id == admission.device_id
452 && self.admission_generation == admission.admission_generation
453 }
454
455 pub(super) fn matches_snapshot(&self, snapshot: &ResourceLedgerEntrySnapshot) -> bool {
456 self.resource_id == *snapshot.entry.resource_id()
457 && self.owner_node_id.as_ref() == snapshot.entry.owner_node_id()
458 && self.generation == snapshot.entry.generation()
459 && self.retention_policy == snapshot.entry.retention_policy()
460 }
461}
462
463#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
464#[serde(deny_unknown_fields)]
465pub struct UnvalidatedResourceTransitionRecord {
466 pub run_id: RunId,
467 pub transaction_id: TransactionId,
468 pub request_id: RequestIdentity,
469 pub plan_id: PlanId,
470 pub plan_hash: PlanHash,
471 pub device_id: DeviceId,
472 pub admission_generation: u64,
473 pub owner_node_id: Option<NodeId>,
474 pub resource_id: ResourceId,
475 pub generation: u64,
476 pub retention_policy: ResourceRetentionPolicy,
477 pub action: ResourceTransactionAction,
478 pub before: ResourceTransactionState,
479 pub after: ResourceTransactionState,
480 pub order: u32,
481}
482
483impl UnvalidatedResourceTransitionRecord {
484 pub fn try_validate(self) -> Result<ResourceTransitionRecord, VNextError> {
487 let record = ResourceTransitionRecord {
488 run_id: self.run_id,
489 transaction_id: self.transaction_id,
490 request_id: self.request_id,
491 plan_id: self.plan_id,
492 plan_hash: self.plan_hash,
493 device_id: self.device_id,
494 admission_generation: self.admission_generation,
495 owner_node_id: self.owner_node_id,
496 resource_id: self.resource_id,
497 generation: self.generation,
498 retention_policy: self.retention_policy,
499 action: self.action,
500 before: self.before,
501 after: self.after,
502 order: self.order,
503 };
504 record.validate()?;
505 Ok(record)
506 }
507}
508
509impl From<&ResourceTransitionRecord> for UnvalidatedResourceTransitionRecord {
510 fn from(record: &ResourceTransitionRecord) -> Self {
511 Self {
512 run_id: record.run_id.clone(),
513 transaction_id: record.transaction_id.clone(),
514 request_id: record.request_id.clone(),
515 plan_id: record.plan_id.clone(),
516 plan_hash: record.plan_hash.clone(),
517 device_id: record.device_id.clone(),
518 admission_generation: record.admission_generation,
519 owner_node_id: record.owner_node_id.clone(),
520 resource_id: record.resource_id.clone(),
521 generation: record.generation,
522 retention_policy: record.retention_policy,
523 action: record.action,
524 before: record.before,
525 after: record.after,
526 order: record.order,
527 }
528 }
529}
530
531#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
532pub struct ResourceCompensationRecord {
533 run_id: RunId,
534 transaction_id: TransactionId,
535 request_id: RequestIdentity,
536 plan_id: PlanId,
537 plan_hash: PlanHash,
538 owner_node_id: Option<NodeId>,
539 resource_id: ResourceId,
540 generation: u64,
541 failed_action: ResourceTransactionAction,
542 compensation_action: ResourceCompensationAction,
543 before: ResourceTransactionState,
544 after: ResourceTransactionState,
545 compensation_order: u32,
546}
547
548impl ResourceCompensationRecord {
549 pub(super) fn from_transition(
550 attempted: &ResourceTransitionRecord,
551 compensation_order: usize,
552 ) -> Self {
553 Self {
554 run_id: attempted.run_id.clone(),
555 transaction_id: attempted.transaction_id.clone(),
556 request_id: attempted.request_id.clone(),
557 plan_id: attempted.plan_id.clone(),
558 plan_hash: attempted.plan_hash.clone(),
559 owner_node_id: attempted.owner_node_id.clone(),
560 resource_id: attempted.resource_id.clone(),
561 generation: attempted.generation,
562 failed_action: attempted.action,
563 compensation_action: ResourceCompensationAction::for_prepare_action(attempted.action)
564 .expect("prepare transition has a compensation action"),
565 before: attempted.after,
566 after: attempted.before,
567 compensation_order: compensation_order as u32,
568 }
569 }
570
571 pub fn validate(&self) -> Result<(), VNextError> {
572 if self.generation == 0
573 || ResourceCompensationAction::for_prepare_action(self.failed_action)
574 != Some(self.compensation_action)
575 || expected_transition(self.failed_action, self.after) != Some(self.before)
576 {
577 return Err(VNextError::InvalidResourceTransition {
578 resource_id: self.resource_id.to_string(),
579 from: self.before.as_str(),
580 action: "compensate",
581 });
582 }
583 Ok(())
584 }
585
586 pub fn run_id(&self) -> &RunId {
587 &self.run_id
588 }
589
590 pub fn transaction_id(&self) -> &TransactionId {
591 &self.transaction_id
592 }
593
594 pub fn request_id(&self) -> &RequestIdentity {
595 &self.request_id
596 }
597
598 pub fn plan_id(&self) -> &PlanId {
599 &self.plan_id
600 }
601
602 pub fn plan_hash(&self) -> &PlanHash {
603 &self.plan_hash
604 }
605
606 pub fn owner_node_id(&self) -> Option<&NodeId> {
607 self.owner_node_id.as_ref()
608 }
609
610 pub fn resource_id(&self) -> &ResourceId {
611 &self.resource_id
612 }
613
614 pub const fn generation(&self) -> u64 {
615 self.generation
616 }
617
618 pub const fn failed_action(&self) -> ResourceTransactionAction {
619 self.failed_action
620 }
621
622 pub const fn compensation_action(&self) -> ResourceCompensationAction {
623 self.compensation_action
624 }
625
626 pub const fn before(&self) -> ResourceTransactionState {
627 self.before
628 }
629
630 pub const fn after(&self) -> ResourceTransactionState {
631 self.after
632 }
633
634 pub const fn compensation_order(&self) -> u32 {
635 self.compensation_order
636 }
637}
638
639#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
640#[serde(deny_unknown_fields)]
641pub struct UnvalidatedResourceCompensationRecord {
642 pub run_id: RunId,
643 pub transaction_id: TransactionId,
644 pub request_id: RequestIdentity,
645 pub plan_id: PlanId,
646 pub plan_hash: PlanHash,
647 pub owner_node_id: Option<NodeId>,
648 pub resource_id: ResourceId,
649 pub generation: u64,
650 pub failed_action: ResourceTransactionAction,
651 pub compensation_action: ResourceCompensationAction,
652 pub before: ResourceTransactionState,
653 pub after: ResourceTransactionState,
654 pub compensation_order: u32,
655}
656
657impl UnvalidatedResourceCompensationRecord {
658 pub fn try_validate(self) -> Result<ResourceCompensationRecord, VNextError> {
659 let record = ResourceCompensationRecord {
660 run_id: self.run_id,
661 transaction_id: self.transaction_id,
662 request_id: self.request_id,
663 plan_id: self.plan_id,
664 plan_hash: self.plan_hash,
665 owner_node_id: self.owner_node_id,
666 resource_id: self.resource_id,
667 generation: self.generation,
668 failed_action: self.failed_action,
669 compensation_action: self.compensation_action,
670 before: self.before,
671 after: self.after,
672 compensation_order: self.compensation_order,
673 };
674 record.validate()?;
675 Ok(record)
676 }
677}
678
679#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
680pub struct ResourceTransitionReceipt {
681 identity: ResourceTransactionIdentity,
682 admission: StaticProvisioningBinding,
683 action: ResourceTransactionAction,
684 records: Vec<ResourceTransitionRecord>,
685}
686
687impl ResourceTransitionReceipt {
688 pub(super) fn from_context(
689 context: &ResourceTransitionValidationContext,
690 records: Vec<ResourceTransitionRecord>,
691 ) -> Result<Self, VNextError> {
692 validate_transition_records_against_context(&records, context)?;
693 Ok(Self {
694 identity: context.identity.clone(),
695 admission: context.admission.clone(),
696 action: context.action,
697 records,
698 })
699 }
700
701 pub fn validate(&self) -> Result<(), VNextError> {
702 if self.records.is_empty() {
703 return Err(invalid_resource(
704 "resource transition receipt must not be empty",
705 ));
706 }
707 let mut previous_order = None;
708 let mut resources = BTreeSet::new();
709 for record in &self.records {
710 record.validate()?;
711 if !record.matches_identity_and_admission(&self.identity, &self.admission)
712 || record.action != self.action
713 || !resources.insert(record.resource_id.clone())
714 || previous_order.is_some_and(|previous| record.order <= previous)
715 {
716 return Err(invalid_resource(
717 "resource receipt identity, admission, action, order, or set is invalid",
718 ));
719 }
720 previous_order = Some(record.order);
721 }
722 Ok(())
723 }
724
725 pub fn identity(&self) -> &ResourceTransactionIdentity {
726 &self.identity
727 }
728
729 pub fn admission(&self) -> &StaticProvisioningBinding {
730 &self.admission
731 }
732
733 pub const fn action(&self) -> ResourceTransactionAction {
734 self.action
735 }
736
737 pub fn records(&self) -> &[ResourceTransitionRecord] {
738 &self.records
739 }
740}
741
742#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
743#[serde(deny_unknown_fields)]
744pub struct UnvalidatedResourceTransactionIdentity {
745 pub pool_id: ResourcePoolId,
746 pub run_id: RunId,
747 pub transaction_id: TransactionId,
748 pub request_id: RequestIdentity,
749}
750
751#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
752#[serde(deny_unknown_fields)]
753pub struct UnvalidatedStaticProvisioningBinding {
754 pub pool_identity: UnvalidatedResourcePoolIdentity,
755 pub plan_id: PlanId,
756 pub plan_hash: PlanHash,
757 pub request_id: RequestIdentity,
758 pub device_id: DeviceId,
759 pub device_runtime_implementation_fingerprint: String,
760 pub device_capacity_bytes: u64,
761 pub usable_capacity_bytes: u64,
762 pub plan_static_bytes: u64,
763 pub admitted_bytes: u64,
764 pub maximum_active_sequences: u32,
765 pub admission_generation: u64,
766}
767
768#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
769#[serde(deny_unknown_fields)]
770pub struct UnvalidatedResourcePoolIdentity {
771 pub pool_id: ResourcePoolId,
772 pub plan_id: PlanId,
773 pub plan_hash: PlanHash,
774 pub device_id: DeviceId,
775 pub device_runtime_implementation_fingerprint: String,
776 pub admission_generation: u64,
777}
778
779#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
780pub struct UnvalidatedResourceTransitionReceipt {
781 pub identity: UnvalidatedResourceTransactionIdentity,
782 pub admission: UnvalidatedStaticProvisioningBinding,
783 pub action: ResourceTransactionAction,
784 pub records: Vec<UnvalidatedResourceTransitionRecord>,
785}
786
787#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
788#[serde(deny_unknown_fields)]
789pub(crate) struct UnvalidatedResourceTransitionReceiptWire {
790 identity: UnvalidatedResourceTransactionIdentity,
791 admission: UnvalidatedStaticProvisioningBinding,
792 action: ResourceTransactionAction,
793 records: Vec<UnvalidatedResourceTransitionRecord>,
794}
795
796impl From<UnvalidatedResourceTransitionReceiptWire> for UnvalidatedResourceTransitionReceipt {
797 fn from(wire: UnvalidatedResourceTransitionReceiptWire) -> Self {
798 Self {
799 identity: wire.identity,
800 admission: wire.admission,
801 action: wire.action,
802 records: wire.records,
803 }
804 }
805}
806
807impl UnvalidatedResourceTransitionReceipt {
808 pub fn decode_untrusted(bytes: &[u8]) -> Result<Self, VNextError> {
809 if bytes.len() > MAX_RESOURCE_TRANSITION_RECEIPT_WIRE_BYTES {
810 return Err(VNextError::Serialization {
811 context: "decode untrusted resource transition receipt",
812 message: format!(
813 "resource transition receipt wire size {} exceeds limit {}",
814 bytes.len(),
815 MAX_RESOURCE_TRANSITION_RECEIPT_WIRE_BYTES
816 ),
817 });
818 }
819 serde_json::from_slice::<UnvalidatedResourceTransitionReceiptWire>(bytes)
820 .map(Self::from)
821 .map_err(|error| VNextError::Serialization {
822 context: "decode untrusted resource transition receipt",
823 message: error.to_string(),
824 })
825 }
826
827 pub fn try_validate(self) -> Result<ResourceTransitionReceipt, VNextError> {
830 let _ = self;
831 Err(invalid_resource(
832 "untrusted resource receipt requires a trusted ledger context",
833 ))
834 }
835
836 pub fn try_validate_against(
837 self,
838 expected: &ResourceTransitionValidationContext,
839 ) -> Result<ResourceTransitionReceipt, VNextError> {
840 if self.identity.pool_id != expected.identity.pool_id
841 || self.identity.run_id != expected.identity.run_id
842 || self.identity.transaction_id != expected.identity.transaction_id
843 || self.identity.request_id != expected.identity.request_id
844 || self.admission.pool_identity.pool_id != expected.admission.pool_identity.pool_id
845 || self.admission.pool_identity.plan_id != expected.admission.pool_identity.plan_id
846 || self.admission.pool_identity.plan_hash != expected.admission.pool_identity.plan_hash
847 || self.admission.pool_identity.device_id != expected.admission.pool_identity.device_id
848 || self
849 .admission
850 .pool_identity
851 .device_runtime_implementation_fingerprint
852 != expected
853 .admission
854 .pool_identity
855 .device_runtime_implementation_fingerprint
856 || self.admission.pool_identity.admission_generation
857 != expected.admission.pool_identity.admission_generation
858 || self.admission.plan_id != expected.admission.plan_id
859 || self.admission.plan_hash != expected.admission.plan_hash
860 || self.admission.request_id != expected.admission.request_id
861 || self.admission.device_id != expected.admission.device_id
862 || self.admission.device_runtime_implementation_fingerprint
863 != expected.admission.device_runtime_implementation_fingerprint
864 || self.admission.device_capacity_bytes != expected.admission.device_capacity_bytes
865 || self.admission.usable_capacity_bytes != expected.admission.usable_capacity_bytes
866 || self.admission.plan_static_bytes != expected.admission.plan_static_bytes
867 || self.admission.admitted_bytes != expected.admission.admitted_bytes
868 || self.admission.maximum_active_sequences
869 != expected.admission.maximum_active_sequences
870 || self.admission.admission_generation != expected.admission.admission_generation
871 || self.action != expected.action
872 {
873 return Err(invalid_resource(
874 "untrusted resource receipt does not match expected identity or admission",
875 ));
876 }
877 let records = self
878 .records
879 .into_iter()
880 .map(UnvalidatedResourceTransitionRecord::try_validate)
881 .collect::<Result<Vec<_>, _>>()?;
882 ResourceTransitionReceipt::from_context(expected, records)
883 }
884}
885
886impl From<&ResourceTransitionReceipt> for UnvalidatedResourceTransitionReceipt {
887 fn from(receipt: &ResourceTransitionReceipt) -> Self {
888 Self {
889 identity: UnvalidatedResourceTransactionIdentity {
890 pool_id: receipt.identity.pool_id,
891 run_id: receipt.identity.run_id.clone(),
892 transaction_id: receipt.identity.transaction_id.clone(),
893 request_id: receipt.identity.request_id.clone(),
894 },
895 admission: UnvalidatedStaticProvisioningBinding {
896 pool_identity: UnvalidatedResourcePoolIdentity {
897 pool_id: receipt.admission.pool_identity.pool_id,
898 plan_id: receipt.admission.pool_identity.plan_id.clone(),
899 plan_hash: receipt.admission.pool_identity.plan_hash.clone(),
900 device_id: receipt.admission.pool_identity.device_id.clone(),
901 device_runtime_implementation_fingerprint: receipt
902 .admission
903 .pool_identity
904 .device_runtime_implementation_fingerprint
905 .clone(),
906 admission_generation: receipt.admission.pool_identity.admission_generation,
907 },
908 plan_id: receipt.admission.plan_id.clone(),
909 plan_hash: receipt.admission.plan_hash.clone(),
910 request_id: receipt.admission.request_id.clone(),
911 device_id: receipt.admission.device_id.clone(),
912 device_runtime_implementation_fingerprint: receipt
913 .admission
914 .device_runtime_implementation_fingerprint
915 .clone(),
916 device_capacity_bytes: receipt.admission.device_capacity_bytes,
917 usable_capacity_bytes: receipt.admission.usable_capacity_bytes,
918 plan_static_bytes: receipt.admission.plan_static_bytes,
919 admitted_bytes: receipt.admission.admitted_bytes,
920 maximum_active_sequences: receipt.admission.maximum_active_sequences,
921 admission_generation: receipt.admission.admission_generation,
922 },
923 action: receipt.action,
924 records: receipt
925 .records
926 .iter()
927 .map(UnvalidatedResourceTransitionRecord::from)
928 .collect(),
929 }
930 }
931}
932
933#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
934pub struct ResourceLeaseTransitionReceipt {
935 run_id: RunId,
936 transaction_id: TransactionId,
937 request_id: RequestIdentity,
938 admission: StaticProvisioningBinding,
939 action: ResourceLeaseAction,
940 decision: ResourceRetentionDecision,
941 before: ResourceLeaseState,
942 after: ResourceLeaseState,
943 entries: Vec<ResourceLeaseEntry>,
944}
945
946impl ResourceLeaseTransitionReceipt {
947 pub(super) fn from_context(
948 context: &ResourceLeaseValidationContext,
949 before: ResourceLeaseState,
950 after: ResourceLeaseState,
951 entries: Vec<ResourceLeaseEntry>,
952 ) -> Result<Self, VNextError> {
953 validate_lease_entries_against_context(&entries, before, after, context)?;
954 Ok(Self {
955 run_id: context.identity.run_id.clone(),
956 transaction_id: context.identity.transaction_id.clone(),
957 request_id: context.identity.request_id.clone(),
958 admission: context.admission.clone(),
959 action: context.action,
960 decision: context.action.decision(),
961 before,
962 after,
963 entries,
964 })
965 }
966
967 pub fn validate(&self) -> Result<(), VNextError> {
968 if expected_lease_transition(self.action, self.before) != Some(self.after)
969 || self.entries.is_empty()
970 || self.decision != self.action.decision()
971 {
972 return Err(VNextError::InvalidLeaseTransition {
973 lease_id: self.transaction_id.to_string(),
974 from: self.before.as_str(),
975 action: self.action.as_str(),
976 });
977 }
978 let mut resources = BTreeSet::new();
979 if self.entries.iter().any(|entry| {
980 entry.generation == 0
981 || entry.state != self.after
982 || !resources.insert(entry.resource_id.clone())
983 }) {
984 return Err(invalid_resource(
985 "lease receipt contains an invalid or duplicate entry",
986 ));
987 }
988 Ok(())
989 }
990
991 pub fn run_id(&self) -> &RunId {
992 &self.run_id
993 }
994
995 pub fn transaction_id(&self) -> &TransactionId {
996 &self.transaction_id
997 }
998
999 pub fn request_id(&self) -> &RequestIdentity {
1000 &self.request_id
1001 }
1002
1003 pub fn admission(&self) -> &StaticProvisioningBinding {
1004 &self.admission
1005 }
1006
1007 pub const fn action(&self) -> ResourceLeaseAction {
1008 self.action
1009 }
1010
1011 pub const fn decision(&self) -> ResourceRetentionDecision {
1012 self.decision
1013 }
1014
1015 pub const fn before(&self) -> ResourceLeaseState {
1016 self.before
1017 }
1018
1019 pub const fn after(&self) -> ResourceLeaseState {
1020 self.after
1021 }
1022
1023 pub fn entries(&self) -> &[ResourceLeaseEntry] {
1024 &self.entries
1025 }
1026}
1027
1028#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1029pub struct UnvalidatedResourceLeaseTransitionReceipt {
1030 pub run_id: RunId,
1031 pub transaction_id: TransactionId,
1032 pub request_id: RequestIdentity,
1033 pub admission: UnvalidatedStaticProvisioningBinding,
1034 pub action: ResourceLeaseAction,
1035 pub decision: ResourceRetentionDecision,
1036 pub before: ResourceLeaseState,
1037 pub after: ResourceLeaseState,
1038 pub entries: Vec<UnvalidatedResourceLeaseEntry>,
1039}
1040
1041#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
1042#[serde(deny_unknown_fields)]
1043pub(crate) struct UnvalidatedResourceLeaseTransitionReceiptWire {
1044 run_id: RunId,
1045 transaction_id: TransactionId,
1046 request_id: RequestIdentity,
1047 admission: UnvalidatedStaticProvisioningBinding,
1048 action: ResourceLeaseAction,
1049 decision: ResourceRetentionDecision,
1050 before: ResourceLeaseState,
1051 after: ResourceLeaseState,
1052 entries: Vec<UnvalidatedResourceLeaseEntry>,
1053}
1054
1055impl From<UnvalidatedResourceLeaseTransitionReceiptWire>
1056 for UnvalidatedResourceLeaseTransitionReceipt
1057{
1058 fn from(wire: UnvalidatedResourceLeaseTransitionReceiptWire) -> Self {
1059 Self {
1060 run_id: wire.run_id,
1061 transaction_id: wire.transaction_id,
1062 request_id: wire.request_id,
1063 admission: wire.admission,
1064 action: wire.action,
1065 decision: wire.decision,
1066 before: wire.before,
1067 after: wire.after,
1068 entries: wire.entries,
1069 }
1070 }
1071}
1072
1073impl UnvalidatedResourceLeaseTransitionReceipt {
1074 pub fn decode_untrusted(bytes: &[u8]) -> Result<Self, VNextError> {
1075 if bytes.len() > MAX_RESOURCE_LEASE_RECEIPT_WIRE_BYTES {
1076 return Err(VNextError::Serialization {
1077 context: "decode untrusted resource lease transition receipt",
1078 message: format!(
1079 "resource lease transition receipt wire size {} exceeds limit {}",
1080 bytes.len(),
1081 MAX_RESOURCE_LEASE_RECEIPT_WIRE_BYTES
1082 ),
1083 });
1084 }
1085 serde_json::from_slice::<UnvalidatedResourceLeaseTransitionReceiptWire>(bytes)
1086 .map(Self::from)
1087 .map_err(|error| VNextError::Serialization {
1088 context: "decode untrusted resource lease transition receipt",
1089 message: error.to_string(),
1090 })
1091 }
1092
1093 pub fn try_validate(self) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
1094 let _ = self;
1095 Err(invalid_resource(
1096 "untrusted lease receipt requires a trusted ledger context",
1097 ))
1098 }
1099
1100 pub fn try_validate_against(
1101 self,
1102 expected: &ResourceLeaseValidationContext,
1103 ) -> Result<ResourceLeaseTransitionReceipt, VNextError> {
1104 if self.run_id != expected.identity.run_id
1105 || self.transaction_id != expected.identity.transaction_id
1106 || self.request_id != expected.identity.request_id
1107 || self.admission.pool_identity.pool_id != expected.admission.pool_identity.pool_id
1108 || self.admission.pool_identity.plan_id != expected.admission.pool_identity.plan_id
1109 || self.admission.pool_identity.plan_hash != expected.admission.pool_identity.plan_hash
1110 || self.admission.pool_identity.device_id != expected.admission.pool_identity.device_id
1111 || self
1112 .admission
1113 .pool_identity
1114 .device_runtime_implementation_fingerprint
1115 != expected
1116 .admission
1117 .pool_identity
1118 .device_runtime_implementation_fingerprint
1119 || self.admission.pool_identity.admission_generation
1120 != expected.admission.pool_identity.admission_generation
1121 || self.admission.plan_id != expected.admission.plan_id
1122 || self.admission.plan_hash != expected.admission.plan_hash
1123 || self.admission.request_id != expected.admission.request_id
1124 || self.admission.device_id != expected.admission.device_id
1125 || self.admission.device_runtime_implementation_fingerprint
1126 != expected.admission.device_runtime_implementation_fingerprint
1127 || self.admission.device_capacity_bytes != expected.admission.device_capacity_bytes
1128 || self.admission.usable_capacity_bytes != expected.admission.usable_capacity_bytes
1129 || self.admission.plan_static_bytes != expected.admission.plan_static_bytes
1130 || self.admission.admitted_bytes != expected.admission.admitted_bytes
1131 || self.admission.maximum_active_sequences
1132 != expected.admission.maximum_active_sequences
1133 || self.admission.admission_generation != expected.admission.admission_generation
1134 || self.action != expected.action
1135 || self.decision != self.action.decision()
1136 {
1137 return Err(invalid_resource(
1138 "untrusted lease receipt does not match expected identity or admission",
1139 ));
1140 }
1141 let entries = self
1142 .entries
1143 .into_iter()
1144 .map(UnvalidatedResourceLeaseEntry::into_trusted)
1145 .collect::<Result<Vec<_>, _>>()?;
1146 ResourceLeaseTransitionReceipt::from_context(expected, self.before, self.after, entries)
1147 }
1148}
1149
1150impl From<&ResourceLeaseTransitionReceipt> for UnvalidatedResourceLeaseTransitionReceipt {
1151 fn from(receipt: &ResourceLeaseTransitionReceipt) -> Self {
1152 Self {
1153 run_id: receipt.run_id.clone(),
1154 transaction_id: receipt.transaction_id.clone(),
1155 request_id: receipt.request_id.clone(),
1156 admission: UnvalidatedStaticProvisioningBinding {
1157 pool_identity: UnvalidatedResourcePoolIdentity {
1158 pool_id: receipt.admission.pool_identity.pool_id,
1159 plan_id: receipt.admission.pool_identity.plan_id.clone(),
1160 plan_hash: receipt.admission.pool_identity.plan_hash.clone(),
1161 device_id: receipt.admission.pool_identity.device_id.clone(),
1162 device_runtime_implementation_fingerprint: receipt
1163 .admission
1164 .pool_identity
1165 .device_runtime_implementation_fingerprint
1166 .clone(),
1167 admission_generation: receipt.admission.pool_identity.admission_generation,
1168 },
1169 plan_id: receipt.admission.plan_id.clone(),
1170 plan_hash: receipt.admission.plan_hash.clone(),
1171 request_id: receipt.admission.request_id.clone(),
1172 device_id: receipt.admission.device_id.clone(),
1173 device_runtime_implementation_fingerprint: receipt
1174 .admission
1175 .device_runtime_implementation_fingerprint
1176 .clone(),
1177 device_capacity_bytes: receipt.admission.device_capacity_bytes,
1178 usable_capacity_bytes: receipt.admission.usable_capacity_bytes,
1179 plan_static_bytes: receipt.admission.plan_static_bytes,
1180 admitted_bytes: receipt.admission.admitted_bytes,
1181 maximum_active_sequences: receipt.admission.maximum_active_sequences,
1182 admission_generation: receipt.admission.admission_generation,
1183 },
1184 action: receipt.action,
1185 decision: receipt.decision,
1186 before: receipt.before,
1187 after: receipt.after,
1188 entries: receipt
1189 .entries
1190 .iter()
1191 .map(UnvalidatedResourceLeaseEntry::from)
1192 .collect(),
1193 }
1194 }
1195}
1196
1197#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1198pub struct ResourceFailurePoint {
1199 owner_node_id: Option<NodeId>,
1200 resource_id: ResourceId,
1201 generation: u64,
1202 order: u32,
1203 actual_before: ResourceTransactionState,
1204}
1205
1206impl ResourceFailurePoint {
1207 pub(super) fn new(
1208 reservation: &ResourceReservation,
1209 order: usize,
1210 actual_before: ResourceTransactionState,
1211 ) -> Self {
1212 Self {
1213 owner_node_id: reservation.owner_node_id.clone(),
1214 resource_id: reservation.resource_id.clone(),
1215 generation: reservation.generation,
1216 order: order as u32,
1217 actual_before,
1218 }
1219 }
1220
1221 pub fn owner_node_id(&self) -> Option<&NodeId> {
1222 self.owner_node_id.as_ref()
1223 }
1224
1225 pub fn resource_id(&self) -> &ResourceId {
1226 &self.resource_id
1227 }
1228
1229 pub const fn generation(&self) -> u64 {
1230 self.generation
1231 }
1232
1233 pub const fn order(&self) -> u32 {
1234 self.order
1235 }
1236
1237 pub const fn actual_before(&self) -> ResourceTransactionState {
1238 self.actual_before
1239 }
1240}
1241
1242#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1243pub struct ResourceRecoveryFailure {
1244 failure: FailureEnvelope,
1245 resource: Option<ResourceFailurePoint>,
1246 attempt: u32,
1247}
1248
1249#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
1250#[serde(try_from = "u64", into = "u64")]
1251pub struct ResourceFailureId(u64);
1252
1253impl ResourceFailureId {
1254 pub const fn get(self) -> u64 {
1255 self.0
1256 }
1257}
1258
1259impl TryFrom<u64> for ResourceFailureId {
1260 type Error = VNextError;
1261
1262 fn try_from(value: u64) -> Result<Self, Self::Error> {
1263 if value == 0 {
1264 return Err(invalid_resource("resource failure id must be non-zero"));
1265 }
1266 Ok(Self(value))
1267 }
1268}
1269
1270impl From<ResourceFailureId> for u64 {
1271 fn from(value: ResourceFailureId) -> Self {
1272 value.0
1273 }
1274}
1275
1276impl ResourceRecoveryFailure {
1277 pub fn failure(&self) -> &FailureEnvelope {
1278 &self.failure
1279 }
1280
1281 pub fn resource(&self) -> Option<&ResourceFailurePoint> {
1282 self.resource.as_ref()
1283 }
1284
1285 pub const fn attempt(&self) -> u32 {
1286 self.attempt
1287 }
1288}
1289
1290#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1291pub struct ResourceFailureReceipt {
1292 failure_id: ResourceFailureId,
1293 identity: ResourceTransactionIdentity,
1294 admission: StaticProvisioningBinding,
1295 pub(super) action: ResourceTransactionAction,
1296 pub(super) failure: FailureEnvelope,
1297 failure_point: Option<ResourceFailurePoint>,
1298 pub(super) completed: Vec<ResourceTransitionRecord>,
1299 pub(super) compensation: Vec<ResourceCompensationRecord>,
1300 recovery_failures: Vec<ResourceRecoveryFailure>,
1301 pub(super) recovery_strategy: ResourceRecoveryStrategy,
1302 pub(super) recovery_complete: bool,
1303 pub(super) ledger_before: Vec<ResourceLedgerEntrySnapshot>,
1304 pub(super) ledger_after: Vec<ResourceLedgerEntrySnapshot>,
1305}
1306
1307impl ResourceFailureReceipt {
1308 pub(super) fn new(
1309 failure_id: ResourceFailureId,
1310 identity: &ResourceTransactionIdentity,
1311 admission: &StaticProvisioningBinding,
1312 action: ResourceTransactionAction,
1313 failure: FailureEnvelope,
1314 failure_point: Option<ResourceFailurePoint>,
1315 completed: Vec<ResourceTransitionRecord>,
1316 recovery_strategy: ResourceRecoveryStrategy,
1317 ledger_before: Vec<ResourceLedgerEntrySnapshot>,
1318 ledger_after: Vec<ResourceLedgerEntrySnapshot>,
1319 ) -> Self {
1320 Self {
1321 failure_id,
1322 identity: identity.clone(),
1323 admission: admission.clone(),
1324 action,
1325 failure,
1326 failure_point,
1327 completed,
1328 compensation: Vec::new(),
1329 recovery_failures: Vec::new(),
1330 recovery_strategy,
1331 recovery_complete: false,
1332 ledger_before,
1333 ledger_after,
1334 }
1335 }
1336
1337 pub const fn failure_id(&self) -> ResourceFailureId {
1338 self.failure_id
1339 }
1340
1341 pub fn identity(&self) -> &ResourceTransactionIdentity {
1342 &self.identity
1343 }
1344
1345 pub fn admission(&self) -> &StaticProvisioningBinding {
1346 &self.admission
1347 }
1348
1349 pub const fn action(&self) -> ResourceTransactionAction {
1350 self.action
1351 }
1352
1353 pub fn failure(&self) -> &FailureEnvelope {
1354 &self.failure
1355 }
1356
1357 pub fn failure_point(&self) -> Option<&ResourceFailurePoint> {
1358 self.failure_point.as_ref()
1359 }
1360
1361 pub fn completed(&self) -> &[ResourceTransitionRecord] {
1362 &self.completed
1363 }
1364
1365 pub fn compensation(&self) -> &[ResourceCompensationRecord] {
1366 &self.compensation
1367 }
1368
1369 pub fn recovery_failures(&self) -> &[ResourceRecoveryFailure] {
1370 &self.recovery_failures
1371 }
1372
1373 pub const fn recovery_strategy(&self) -> ResourceRecoveryStrategy {
1374 self.recovery_strategy
1375 }
1376
1377 pub const fn recovery_complete(&self) -> bool {
1378 self.recovery_complete
1379 }
1380
1381 pub fn ledger_before(&self) -> &[ResourceLedgerEntrySnapshot] {
1382 &self.ledger_before
1383 }
1384
1385 pub fn ledger_after(&self) -> &[ResourceLedgerEntrySnapshot] {
1386 &self.ledger_after
1387 }
1388
1389 pub fn validate_recovery_continuation(&self, anchor: &Self) -> Result<(), VNextError> {
1390 let strategy_continues = self.recovery_strategy == anchor.recovery_strategy
1391 || (anchor.recovery_strategy == ResourceRecoveryStrategy::ReconcileOrQuarantine
1392 && self.recovery_strategy == ResourceRecoveryStrategy::ReverseCompensation);
1393 if self.failure_id != anchor.failure_id
1394 || self.identity != anchor.identity
1395 || self.admission != anchor.admission
1396 || self.action != anchor.action
1397 || self.failure != anchor.failure
1398 || self.failure_point != anchor.failure_point
1399 || !self.completed.starts_with(&anchor.completed)
1400 || !self.compensation.starts_with(&anchor.compensation)
1401 || !self
1402 .recovery_failures
1403 .starts_with(&anchor.recovery_failures)
1404 || !strategy_continues
1405 || (anchor.recovery_complete && !self.recovery_complete)
1406 || self.ledger_before != anchor.ledger_before
1407 {
1408 return Err(invalid_resource(
1409 "resource recovery receipt does not continue its exact failure anchor",
1410 ));
1411 }
1412 Ok(())
1413 }
1414
1415 pub(super) fn record_recovery_failure(
1416 &mut self,
1417 failure: ResourceDriverFailure,
1418 resource: Option<ResourceFailurePoint>,
1419 ) {
1420 self.recovery_failures.push(ResourceRecoveryFailure {
1421 failure: failure.into_failure(),
1422 resource,
1423 attempt: (self.recovery_failures.len() + 1) as u32,
1424 });
1425 }
1426}
1427
1428#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1429pub struct ResourceAbandonSignal {
1430 pub(super) identity: ResourceTransactionIdentity,
1431 pub(super) admission: StaticProvisioningBinding,
1432 pub(super) state: ResourceTransactionState,
1433 pub(super) pending_action: Option<ResourceTransactionAction>,
1434 pub(super) ledger: Vec<ResourceLedgerEntrySnapshot>,
1435 pub(super) active_sequence_slots: Vec<u32>,
1436 pub(super) poisoned_sequence_slots: Vec<u32>,
1437 pub(super) undrained_sequence_slots: Vec<u32>,
1438 pub(super) failure: Option<FailureEnvelope>,
1439}
1440
1441impl ResourceAbandonSignal {
1442 pub fn identity(&self) -> &ResourceTransactionIdentity {
1443 &self.identity
1444 }
1445
1446 pub fn admission(&self) -> &StaticProvisioningBinding {
1447 &self.admission
1448 }
1449
1450 pub const fn state(&self) -> ResourceTransactionState {
1451 self.state
1452 }
1453
1454 pub const fn pending_action(&self) -> Option<ResourceTransactionAction> {
1455 self.pending_action
1456 }
1457
1458 pub fn ledger(&self) -> &[ResourceLedgerEntrySnapshot] {
1459 &self.ledger
1460 }
1461
1462 pub fn active_sequence_slots(&self) -> &[u32] {
1463 &self.active_sequence_slots
1464 }
1465
1466 pub fn poisoned_sequence_slots(&self) -> &[u32] {
1467 &self.poisoned_sequence_slots
1468 }
1469
1470 pub fn undrained_sequence_slots(&self) -> &[u32] {
1471 &self.undrained_sequence_slots
1472 }
1473
1474 pub fn live_resources(&self) -> impl Iterator<Item = &ResourceLeaseEntry> {
1475 self.ledger
1476 .iter()
1477 .filter(|entry| entry.transaction_state.is_live())
1478 .map(ResourceLedgerEntrySnapshot::entry)
1479 }
1480
1481 pub fn failure(&self) -> Option<&FailureEnvelope> {
1482 self.failure.as_ref()
1483 }
1484}
1485
1486fn same_allocation_entry(left: &ResourceLeaseEntry, right: &ResourceLeaseEntry) -> bool {
1487 left.owner_node_id == right.owner_node_id
1488 && left.resource_id == right.resource_id
1489 && left.size_bytes == right.size_bytes
1490 && left.alignment_bytes == right.alignment_bytes
1491 && left.usage == right.usage
1492 && left.element_type == right.element_type
1493 && left.retention_policy == right.retention_policy
1494 && left.generation == right.generation
1495}
1496
1497fn validate_context_envelope(
1498 identity: &ResourceTransactionIdentity,
1499 admission: &StaticProvisioningBinding,
1500 before: &[ResourceLedgerEntrySnapshot],
1501 after: &[ResourceLedgerEntrySnapshot],
1502) -> Result<(), VNextError> {
1503 if identity.pool_id != admission.pool_id()
1504 || identity.request_id != admission.request_id
1505 || admission.pool_identity.plan_id != admission.plan_id
1506 || admission.pool_identity.plan_hash != admission.plan_hash
1507 || admission.pool_identity.device_id != admission.device_id
1508 || admission
1509 .pool_identity
1510 .device_runtime_implementation_fingerprint
1511 != admission.device_runtime_implementation_fingerprint
1512 || admission.pool_identity.admission_generation != admission.admission_generation
1513 || admission.admission_generation == 0
1514 || admission.maximum_active_sequences == 0
1515 || admission.device_capacity_bytes == 0
1516 || admission.usable_capacity_bytes == 0
1517 || admission.usable_capacity_bytes > admission.device_capacity_bytes
1518 || admission.plan_static_bytes > admission.usable_capacity_bytes
1519 || admission.admitted_bytes == 0
1520 || admission.admitted_bytes != admission.plan_static_bytes
1521 || admission.admitted_bytes > admission.usable_capacity_bytes
1522 || before.is_empty()
1523 || before.len() != after.len()
1524 {
1525 return Err(invalid_resource(
1526 "trusted resource validation context has an invalid envelope",
1527 ));
1528 }
1529 let mut resources = BTreeSet::new();
1530 let mut total = 0_u64;
1531 for (before_entry, after_entry) in before.iter().zip(after) {
1532 if !same_allocation_entry(&before_entry.entry, &after_entry.entry)
1533 || !resources.insert(before_entry.entry.resource_id.clone())
1534 || before_entry.entry.generation != admission.admission_generation
1535 {
1536 return Err(invalid_resource(
1537 "trusted resource validation context changes allocation identity",
1538 ));
1539 }
1540 total = total
1541 .checked_add(before_entry.entry.size_bytes)
1542 .ok_or_else(|| invalid_resource("validation context bytes overflow u64"))?;
1543 }
1544 if total != admission.admitted_bytes {
1545 return Err(invalid_resource(
1546 "validation context does not cover the complete admitted allocation set",
1547 ));
1548 }
1549 Ok(())
1550}
1551
1552fn validate_transition_records_against_context(
1553 records: &[ResourceTransitionRecord],
1554 context: &ResourceTransitionValidationContext,
1555) -> Result<(), VNextError> {
1556 validate_context_envelope(
1557 &context.identity,
1558 &context.admission,
1559 &context.before,
1560 &context.after,
1561 )?;
1562 let changed = context
1563 .before
1564 .iter()
1565 .zip(&context.after)
1566 .enumerate()
1567 .filter(|(_, (before, after))| {
1568 before.transaction_state != after.transaction_state
1569 || before.buffer_present != after.buffer_present
1570 })
1571 .collect::<Vec<_>>();
1572 if changed.is_empty() || changed.len() != records.len() {
1573 return Err(invalid_resource(
1574 "resource receipt does not cover the exact ledger delta",
1575 ));
1576 }
1577 for (record, (order, (before, after))) in records.iter().zip(changed) {
1578 record.validate()?;
1579 if !record.matches_identity_and_admission(&context.identity, &context.admission)
1580 || record.action != context.action
1581 || record.order as usize != order
1582 || !record.matches_snapshot(before)
1583 || !record.matches_snapshot(after)
1584 || record.before != before.transaction_state
1585 || record.after != after.transaction_state
1586 || expected_transition(context.action, record.before) != Some(record.after)
1587 {
1588 return Err(invalid_resource(
1589 "resource receipt differs from trusted allocation or ledger state",
1590 ));
1591 }
1592 match context.action {
1593 ResourceTransactionAction::Commit => {
1594 if before.buffer_present || !after.buffer_present {
1595 return Err(invalid_resource(
1596 "commit receipt does not prove buffer acquisition",
1597 ));
1598 }
1599 }
1600 ResourceTransactionAction::Release | ResourceTransactionAction::Quarantine => {
1601 if before.transaction_state == ResourceTransactionState::Committed
1602 && (!before.buffer_present || after.buffer_present)
1603 {
1604 return Err(invalid_resource(
1605 "cleanup receipt does not prove committed buffer return",
1606 ));
1607 }
1608 }
1609 ResourceTransactionAction::Reserve | ResourceTransactionAction::Rollback => {
1610 if before.buffer_present != after.buffer_present {
1611 return Err(invalid_resource(
1612 "non-buffer transition changed buffer ownership",
1613 ));
1614 }
1615 }
1616 }
1617 }
1618 Ok(())
1619}
1620
1621fn validate_lease_entries_against_context(
1622 entries: &[ResourceLeaseEntry],
1623 before_state: ResourceLeaseState,
1624 after_state: ResourceLeaseState,
1625 context: &ResourceLeaseValidationContext,
1626) -> Result<(), VNextError> {
1627 validate_context_envelope(
1628 &context.identity,
1629 &context.admission,
1630 &context.before,
1631 &context.after,
1632 )?;
1633 if expected_lease_transition(context.action, before_state) != Some(after_state) {
1634 return Err(VNextError::InvalidLeaseTransition {
1635 lease_id: context.identity.transaction_id.to_string(),
1636 from: before_state.as_str(),
1637 action: context.action.as_str(),
1638 });
1639 }
1640 let changed = context
1641 .before
1642 .iter()
1643 .zip(&context.after)
1644 .filter(|(before, after)| before.entry.state != after.entry.state)
1645 .collect::<Vec<_>>();
1646 if changed.is_empty() || changed.len() != entries.len() {
1647 return Err(invalid_resource(
1648 "lease receipt does not cover the exact lease-state delta",
1649 ));
1650 }
1651 for (entry, (before, after)) in entries.iter().zip(changed) {
1652 if !same_allocation_entry(entry, &before.entry)
1653 || !same_allocation_entry(entry, &after.entry)
1654 || before.entry.state != before_state
1655 || after.entry.state != after_state
1656 || entry.state != after_state
1657 || before.transaction_state != after.transaction_state
1658 || before.transaction_state != ResourceTransactionState::Committed
1659 || before.buffer_present != after.buffer_present
1660 || !before.buffer_present
1661 {
1662 return Err(invalid_resource(
1663 "lease receipt differs from trusted allocation or ledger state",
1664 ));
1665 }
1666 }
1667 Ok(())
1668}