1use serde::{Deserialize, Serialize};
2use std::collections::{BTreeMap, BTreeSet};
3
4use super::{
5 canonical_fingerprint, invalid_event, validate_sha256, FailureDomain, MonotonicTimestamp,
6 RequestIdentity, ResourceFailureId, ResourceFailureReceipt, ResourceId, ResourceLeaseState,
7 ResourceLeaseTransitionReceipt, ResourceLeaseValidationContext, ResourceLedgerEntrySnapshot,
8 ResourceLedgerSnapshot, ResourcePoolId, ResourceTransactionIdentity, ResourceTransactionState,
9 ResourceTransitionReceipt, ResourceTransitionValidationContext, RunId,
10 StaticProvisioningBinding, TransactionId, TrustedActiveSequenceBinding,
11 TrustedExecutionTopology, UnvalidatedResourceLeaseTransitionReceipt,
12 UnvalidatedResourceLeaseTransitionReceiptWire, UnvalidatedResourceTransitionReceipt,
13 UnvalidatedResourceTransitionReceiptWire, VNextError, MAX_RESOURCE_POOL_EVENT_WIRE_BYTES,
14};
15
16#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
17pub struct ResourcePoolEvidence {
18 topology_fingerprint: String,
19 pool_id: ResourcePoolId,
20 pool_identity_fingerprint: String,
21 admission: StaticProvisioningBinding,
22 provisioning_identity: ResourceTransactionIdentity,
23}
24
25impl ResourcePoolEvidence {
26 pub fn from_external(
27 topology: &TrustedExecutionTopology,
28 admission: &StaticProvisioningBinding,
29 provisioning_identity: &ResourceTransactionIdentity,
30 ) -> Result<Self, VNextError> {
31 let pool_fingerprint = canonical_fingerprint(admission.pool_identity());
32 if admission.plan_id() != topology.plan_id()
33 || admission.plan_hash() != topology.plan_hash()
34 || admission.device_id() != topology.device_id()
35 || admission.device_runtime_implementation_fingerprint()
36 != topology.device_runtime_implementation_fingerprint()
37 || admission
38 .pool_identity()
39 .device_runtime_implementation_fingerprint()
40 != topology.device_runtime_implementation_fingerprint()
41 || provisioning_identity.pool_id() != admission.pool_id()
42 || provisioning_identity.request_id() != admission.request_id()
43 {
44 return Err(invalid_event(
45 "pool evidence differs from topology, admission, or provisioning identity",
46 ));
47 }
48 Ok(Self {
49 topology_fingerprint: topology.fingerprint().to_owned(),
50 pool_id: admission.pool_id(),
51 pool_identity_fingerprint: pool_fingerprint,
52 admission: admission.clone(),
53 provisioning_identity: provisioning_identity.clone(),
54 })
55 }
56
57 pub const fn pool_id(&self) -> ResourcePoolId {
58 self.pool_id
59 }
60
61 pub fn topology_fingerprint(&self) -> &str {
62 &self.topology_fingerprint
63 }
64
65 pub fn pool_identity_fingerprint(&self) -> &str {
66 &self.pool_identity_fingerprint
67 }
68
69 pub fn admission(&self) -> &StaticProvisioningBinding {
70 &self.admission
71 }
72
73 pub fn provisioning_identity(&self) -> &ResourceTransactionIdentity {
74 &self.provisioning_identity
75 }
76}
77
78#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
79#[serde(rename_all = "snake_case")]
80pub enum ResourcePoolEventKind {
81 ResourcePoolOpened,
82 ResourceTransition,
83 ResourceLeaseTransition,
84 ResourceFailed,
85 ResourceRecoveryCompleted,
86 ResourcePoolClosed,
87}
88
89#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
90pub struct ResourcePoolEventIdentity {
91 sequence: u64,
92 pool_id: ResourcePoolId,
93 pool_identity_fingerprint: String,
94 provisioning_run_id: RunId,
95 provisioning_request_id: RequestIdentity,
96 transaction_id: TransactionId,
97 resource_id: Option<ResourceId>,
98 resource_generation: Option<u64>,
99 resource_batch_fingerprint: Option<String>,
100}
101
102#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize)]
103#[serde(deny_unknown_fields)]
104pub struct UnvalidatedResourcePoolEventIdentity {
105 sequence: u64,
106 pool_id: ResourcePoolId,
107 pool_identity_fingerprint: String,
108 provisioning_run_id: RunId,
109 provisioning_request_id: RequestIdentity,
110 transaction_id: TransactionId,
111 resource_id: Option<ResourceId>,
112 resource_generation: Option<u64>,
113 resource_batch_fingerprint: Option<String>,
114}
115
116impl ResourcePoolEventIdentity {
117 fn for_evidence(
118 sequence: u64,
119 evidence: &ResourcePoolEvidence,
120 resource_id: Option<ResourceId>,
121 resource_generation: Option<u64>,
122 resource_batch_fingerprint: Option<String>,
123 ) -> Result<Self, VNextError> {
124 if sequence == 0
125 || resource_id.is_some() != resource_generation.is_some()
126 || resource_generation == Some(0)
127 || resource_id.is_some() && resource_batch_fingerprint.is_some()
128 {
129 return Err(invalid_event(
130 "resource pool event identity shape is invalid",
131 ));
132 }
133 if let Some(fingerprint) = &resource_batch_fingerprint {
134 validate_sha256(fingerprint, "resource pool event batch fingerprint")?;
135 }
136 Ok(Self {
137 sequence,
138 pool_id: evidence.pool_id,
139 pool_identity_fingerprint: evidence.pool_identity_fingerprint.clone(),
140 provisioning_run_id: evidence.provisioning_identity.run_id().clone(),
141 provisioning_request_id: evidence.provisioning_identity.request_id().clone(),
142 transaction_id: evidence.provisioning_identity.transaction_id().clone(),
143 resource_id,
144 resource_generation,
145 resource_batch_fingerprint,
146 })
147 }
148
149 pub const fn sequence(&self) -> u64 {
150 self.sequence
151 }
152
153 pub const fn pool_id(&self) -> ResourcePoolId {
154 self.pool_id
155 }
156
157 pub fn pool_identity_fingerprint(&self) -> &str {
158 &self.pool_identity_fingerprint
159 }
160
161 pub fn transaction_id(&self) -> &TransactionId {
162 &self.transaction_id
163 }
164
165 pub fn provisioning_run_id(&self) -> &RunId {
166 &self.provisioning_run_id
167 }
168
169 pub fn provisioning_request_id(&self) -> &RequestIdentity {
170 &self.provisioning_request_id
171 }
172
173 pub fn resource_id(&self) -> Option<&ResourceId> {
174 self.resource_id.as_ref()
175 }
176
177 pub const fn resource_generation(&self) -> Option<u64> {
178 self.resource_generation
179 }
180
181 pub fn resource_batch_fingerprint(&self) -> Option<&str> {
182 self.resource_batch_fingerprint.as_deref()
183 }
184}
185
186#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
187#[serde(rename_all = "snake_case")]
188pub enum ResourcePoolEventDetail {
189 Opened(StaticProvisioningBinding),
190 Transition {
191 receipt: ResourceTransitionReceipt,
192 context: ResourceTransitionValidationContext,
193 },
194 LeaseTransition {
195 receipt: ResourceLeaseTransitionReceipt,
196 context: ResourceLeaseValidationContext,
197 },
198 Failure(ResourceFailureReceipt),
199 Closed {
200 ledger: Vec<ResourceLedgerEntrySnapshot>,
201 },
202}
203
204#[derive(Debug, Clone, PartialEq, Serialize)]
205#[serde(rename_all = "snake_case")]
206pub enum UnvalidatedResourcePoolEventDetail {
207 Opened(serde_json::Value),
208 Transition {
209 receipt: UnvalidatedResourceTransitionReceipt,
210 context: serde_json::Value,
211 },
212 LeaseTransition {
213 receipt: UnvalidatedResourceLeaseTransitionReceipt,
214 context: serde_json::Value,
215 },
216 Failure(serde_json::Value),
217 Closed {
218 ledger: serde_json::Value,
219 },
220}
221
222#[derive(Debug, Clone, PartialEq, Deserialize)]
223#[serde(rename_all = "snake_case")]
224enum UnvalidatedResourcePoolEventDetailWire {
225 Opened(serde_json::Value),
226 Transition {
227 receipt: UnvalidatedResourceTransitionReceiptWire,
228 context: serde_json::Value,
229 },
230 LeaseTransition {
231 receipt: UnvalidatedResourceLeaseTransitionReceiptWire,
232 context: serde_json::Value,
233 },
234 Failure(serde_json::Value),
235 Closed {
236 ledger: serde_json::Value,
237 },
238}
239
240impl From<UnvalidatedResourcePoolEventDetailWire> for UnvalidatedResourcePoolEventDetail {
241 fn from(wire: UnvalidatedResourcePoolEventDetailWire) -> Self {
242 match wire {
243 UnvalidatedResourcePoolEventDetailWire::Opened(value) => Self::Opened(value),
244 UnvalidatedResourcePoolEventDetailWire::Transition { receipt, context } => {
245 Self::Transition {
246 receipt: receipt.into(),
247 context,
248 }
249 }
250 UnvalidatedResourcePoolEventDetailWire::LeaseTransition { receipt, context } => {
251 Self::LeaseTransition {
252 receipt: receipt.into(),
253 context,
254 }
255 }
256 UnvalidatedResourcePoolEventDetailWire::Failure(value) => Self::Failure(value),
257 UnvalidatedResourcePoolEventDetailWire::Closed { ledger } => Self::Closed { ledger },
258 }
259 }
260}
261
262#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
263pub struct ResourcePoolEvent {
264 timestamp: MonotonicTimestamp,
265 kind: ResourcePoolEventKind,
266 identity: ResourcePoolEventIdentity,
267 detail: ResourcePoolEventDetail,
268}
269
270#[derive(Debug, Clone, PartialEq, Serialize)]
271pub struct UnvalidatedResourcePoolEvent {
272 timestamp: MonotonicTimestamp,
273 kind: ResourcePoolEventKind,
274 identity: UnvalidatedResourcePoolEventIdentity,
275 detail: UnvalidatedResourcePoolEventDetail,
276}
277
278#[derive(Deserialize)]
279#[serde(deny_unknown_fields)]
280struct ResourcePoolEventWire {
281 timestamp: MonotonicTimestamp,
282 kind: ResourcePoolEventKind,
283 identity: UnvalidatedResourcePoolEventIdentity,
284 detail: UnvalidatedResourcePoolEventDetailWire,
285}
286
287impl From<ResourcePoolEventWire> for UnvalidatedResourcePoolEvent {
288 fn from(wire: ResourcePoolEventWire) -> Self {
289 Self {
290 timestamp: wire.timestamp,
291 kind: wire.kind,
292 identity: wire.identity,
293 detail: wire.detail.into(),
294 }
295 }
296}
297
298pub enum TrustedResourcePoolEventContext<'a> {
299 Opened {
300 evidence: &'a ResourcePoolEvidence,
301 },
302 Transition {
303 evidence: &'a ResourcePoolEvidence,
304 expected: &'a ResourceTransitionValidationContext,
305 },
306 LeaseTransition {
307 evidence: &'a ResourcePoolEvidence,
308 expected: &'a ResourceLeaseValidationContext,
309 },
310 Failure {
311 evidence: &'a ResourcePoolEvidence,
312 expected: &'a ResourceFailureReceipt,
313 },
314 Closed {
315 evidence: &'a ResourcePoolEvidence,
316 expected_snapshot: &'a ResourceLedgerSnapshot,
317 },
318}
319
320impl<'a> TrustedResourcePoolEventContext<'a> {
321 pub fn opened(evidence: &'a ResourcePoolEvidence) -> Self {
322 Self::Opened { evidence }
323 }
324
325 pub fn transition(
326 evidence: &'a ResourcePoolEvidence,
327 expected: &'a ResourceTransitionValidationContext,
328 ) -> Self {
329 Self::Transition { evidence, expected }
330 }
331
332 pub fn lease_transition(
333 evidence: &'a ResourcePoolEvidence,
334 expected: &'a ResourceLeaseValidationContext,
335 ) -> Self {
336 Self::LeaseTransition { evidence, expected }
337 }
338
339 pub fn failure(
340 evidence: &'a ResourcePoolEvidence,
341 expected: &'a ResourceFailureReceipt,
342 ) -> Self {
343 Self::Failure { evidence, expected }
344 }
345
346 pub fn closed(
347 evidence: &'a ResourcePoolEvidence,
348 expected_snapshot: &'a ResourceLedgerSnapshot,
349 ) -> Self {
350 Self::Closed {
351 evidence,
352 expected_snapshot,
353 }
354 }
355}
356
357fn validate_pool_binding(
358 evidence: &ResourcePoolEvidence,
359 identity: &ResourceTransactionIdentity,
360 admission: &StaticProvisioningBinding,
361 label: &str,
362) -> Result<(), VNextError> {
363 if identity != evidence.provisioning_identity() || admission != evidence.admission() {
364 return Err(invalid_event(format!(
365 "resource pool {label} identity or admission differs from its exact pool evidence"
366 )));
367 }
368 Ok(())
369}
370
371fn validate_failure_pool_binding(
372 evidence: &ResourcePoolEvidence,
373 failure: &ResourceFailureReceipt,
374) -> Result<(), VNextError> {
375 validate_pool_binding(evidence, failure.identity(), failure.admission(), "failure")?;
376 if failure.failure().domain() != FailureDomain::Resource {
377 return Err(invalid_event(
378 "resource pool failure must carry the resource failure domain",
379 ));
380 }
381 Ok(())
382}
383
384fn validate_terminal_snapshot(
385 evidence: &ResourcePoolEvidence,
386 snapshot: &ResourceLedgerSnapshot,
387) -> Result<(), VNextError> {
388 validate_pool_binding(
389 evidence,
390 snapshot.identity(),
391 snapshot.admission(),
392 "terminal snapshot",
393 )?;
394 if snapshot.entries().is_empty()
395 || snapshot.entries().iter().any(|entry| {
396 entry.entry().generation() != evidence.admission().admission_generation()
397 || !matches!(
398 entry.transaction_state(),
399 ResourceTransactionState::RolledBack
400 | ResourceTransactionState::Released
401 | ResourceTransactionState::Quarantined
402 )
403 || entry.buffer_present()
404 })
405 {
406 return Err(invalid_event(
407 "resource pool terminal snapshot must be non-empty, generation-bound, terminal, and buffer-free",
408 ));
409 }
410 Ok(())
411}
412
413fn item_or_batch_identity(
414 entries: &[ResourceLedgerEntrySnapshot],
415) -> (Option<ResourceId>, Option<u64>, Option<String>) {
416 if entries.len() == 1 {
417 (
418 Some(entries[0].entry().resource_id().clone()),
419 Some(entries[0].entry().generation()),
420 None,
421 )
422 } else {
423 (None, None, Some(canonical_fingerprint(&entries)))
424 }
425}
426
427fn failure_item_or_batch_identity(
428 failure: &ResourceFailureReceipt,
429) -> (Option<ResourceId>, Option<u64>, Option<String>) {
430 if let Some(point) = failure.failure_point() {
431 (
432 Some(point.resource_id().clone()),
433 Some(point.generation()),
434 None,
435 )
436 } else {
437 (
438 None,
439 None,
440 Some(canonical_fingerprint(&(
441 failure.failure_id(),
442 failure.completed(),
443 failure.compensation(),
444 failure.recovery_failures(),
445 failure.recovery_strategy(),
446 failure.ledger_before(),
447 failure.ledger_after(),
448 ))),
449 )
450 }
451}
452
453fn validate_transition_receipt_context(
454 receipt: &ResourceTransitionReceipt,
455 context: &ResourceTransitionValidationContext,
456) -> Result<(), VNextError> {
457 let rebuilt =
458 UnvalidatedResourceTransitionReceipt::from(receipt).try_validate_against(context)?;
459 if &rebuilt != receipt {
460 return Err(invalid_event(
461 "resource transition receipt does not match its exact ledger delta",
462 ));
463 }
464 Ok(())
465}
466
467fn validate_lease_receipt_context(
468 receipt: &ResourceLeaseTransitionReceipt,
469 context: &ResourceLeaseValidationContext,
470) -> Result<(), VNextError> {
471 let rebuilt =
472 UnvalidatedResourceLeaseTransitionReceipt::from(receipt).try_validate_against(context)?;
473 if &rebuilt != receipt {
474 return Err(invalid_event(
475 "resource lease receipt does not match its exact ledger delta",
476 ));
477 }
478 Ok(())
479}
480
481impl ResourcePoolEvent {
482 pub fn opened(
483 sequence: u64,
484 timestamp: MonotonicTimestamp,
485 evidence: &ResourcePoolEvidence,
486 ) -> Result<Self, VNextError> {
487 Ok(Self {
488 timestamp,
489 kind: ResourcePoolEventKind::ResourcePoolOpened,
490 identity: ResourcePoolEventIdentity::for_evidence(
491 sequence, evidence, None, None, None,
492 )?,
493 detail: ResourcePoolEventDetail::Opened(evidence.admission.clone()),
494 })
495 }
496
497 pub fn transition(
498 sequence: u64,
499 timestamp: MonotonicTimestamp,
500 evidence: &ResourcePoolEvidence,
501 receipt: &ResourceTransitionReceipt,
502 context: &ResourceTransitionValidationContext,
503 ) -> Result<Self, VNextError> {
504 validate_pool_binding(
505 evidence,
506 context.identity(),
507 context.admission(),
508 "transition",
509 )?;
510 validate_transition_receipt_context(receipt, context)?;
511 let (resource, generation, batch) = item_or_batch_identity(context.after());
512 Ok(Self {
513 timestamp,
514 kind: ResourcePoolEventKind::ResourceTransition,
515 identity: ResourcePoolEventIdentity::for_evidence(
516 sequence, evidence, resource, generation, batch,
517 )?,
518 detail: ResourcePoolEventDetail::Transition {
519 receipt: receipt.clone(),
520 context: context.clone(),
521 },
522 })
523 }
524
525 pub fn lease_transition(
526 sequence: u64,
527 timestamp: MonotonicTimestamp,
528 evidence: &ResourcePoolEvidence,
529 receipt: &ResourceLeaseTransitionReceipt,
530 context: &ResourceLeaseValidationContext,
531 ) -> Result<Self, VNextError> {
532 validate_pool_binding(
533 evidence,
534 context.identity(),
535 context.admission(),
536 "lease transition",
537 )?;
538 validate_lease_receipt_context(receipt, context)?;
539 let (resource, generation, batch) = item_or_batch_identity(context.after());
540 Ok(Self {
541 timestamp,
542 kind: ResourcePoolEventKind::ResourceLeaseTransition,
543 identity: ResourcePoolEventIdentity::for_evidence(
544 sequence, evidence, resource, generation, batch,
545 )?,
546 detail: ResourcePoolEventDetail::LeaseTransition {
547 receipt: receipt.clone(),
548 context: context.clone(),
549 },
550 })
551 }
552
553 pub fn failed(
554 sequence: u64,
555 timestamp: MonotonicTimestamp,
556 evidence: &ResourcePoolEvidence,
557 failure: &ResourceFailureReceipt,
558 ) -> Result<Self, VNextError> {
559 validate_failure_pool_binding(evidence, failure)?;
560 if failure.recovery_complete() {
561 return Err(invalid_event(
562 "ResourceFailed requires an incomplete recovery anchor",
563 ));
564 }
565 let (resource, generation, batch) = failure_item_or_batch_identity(failure);
566 Ok(Self {
567 timestamp,
568 kind: ResourcePoolEventKind::ResourceFailed,
569 identity: ResourcePoolEventIdentity::for_evidence(
570 sequence, evidence, resource, generation, batch,
571 )?,
572 detail: ResourcePoolEventDetail::Failure(failure.clone()),
573 })
574 }
575
576 pub fn recovery_completed(
577 sequence: u64,
578 timestamp: MonotonicTimestamp,
579 evidence: &ResourcePoolEvidence,
580 recovery: &ResourceFailureReceipt,
581 ) -> Result<Self, VNextError> {
582 validate_failure_pool_binding(evidence, recovery)?;
583 if !recovery.recovery_complete() {
584 return Err(invalid_event(
585 "ResourceRecoveryCompleted requires completed recovery evidence",
586 ));
587 }
588 let (resource, generation, batch) = failure_item_or_batch_identity(recovery);
589 Ok(Self {
590 timestamp,
591 kind: ResourcePoolEventKind::ResourceRecoveryCompleted,
592 identity: ResourcePoolEventIdentity::for_evidence(
593 sequence, evidence, resource, generation, batch,
594 )?,
595 detail: ResourcePoolEventDetail::Failure(recovery.clone()),
596 })
597 }
598
599 pub fn closed(
600 sequence: u64,
601 timestamp: MonotonicTimestamp,
602 evidence: &ResourcePoolEvidence,
603 snapshot: &ResourceLedgerSnapshot,
604 ) -> Result<Self, VNextError> {
605 validate_terminal_snapshot(evidence, snapshot)?;
606 Ok(Self {
607 timestamp,
608 kind: ResourcePoolEventKind::ResourcePoolClosed,
609 identity: ResourcePoolEventIdentity::for_evidence(
610 sequence, evidence, None, None, None,
611 )?,
612 detail: ResourcePoolEventDetail::Closed {
613 ledger: snapshot.entries().to_vec(),
614 },
615 })
616 }
617
618 pub const fn timestamp(&self) -> MonotonicTimestamp {
619 self.timestamp
620 }
621
622 pub const fn kind(&self) -> ResourcePoolEventKind {
623 self.kind
624 }
625
626 pub fn identity(&self) -> &ResourcePoolEventIdentity {
627 &self.identity
628 }
629
630 pub fn detail(&self) -> &ResourcePoolEventDetail {
631 &self.detail
632 }
633
634 pub fn decode_untrusted(bytes: &[u8]) -> Result<UnvalidatedResourcePoolEvent, VNextError> {
635 if bytes.len() > MAX_RESOURCE_POOL_EVENT_WIRE_BYTES {
636 return Err(invalid_event(
637 "untrusted resource pool event exceeds the wire byte limit",
638 ));
639 }
640 let raw = serde_json::from_slice::<serde_json::Value>(bytes).map_err(|error| {
641 VNextError::Serialization {
642 context: "decode untrusted resource pool event",
643 message: error.to_string(),
644 }
645 })?;
646 let event = serde_json::from_value::<ResourcePoolEventWire>(raw.clone())
647 .map(UnvalidatedResourcePoolEvent::from)
648 .map_err(|error| VNextError::Serialization {
649 context: "decode untrusted resource pool event",
650 message: error.to_string(),
651 })?;
652 let canonical =
653 serde_json::to_value(&event).map_err(|error| VNextError::Serialization {
654 context: "serialize untrusted resource pool event",
655 message: error.to_string(),
656 })?;
657 if canonical != raw {
658 return Err(invalid_event(
659 "resource pool event wire contains unknown or non-canonical nested fields",
660 ));
661 }
662 Ok(event)
663 }
664}
665
666impl UnvalidatedResourcePoolEventIdentity {
667 fn matches(&self, expected: &ResourcePoolEventIdentity) -> bool {
668 self.sequence == expected.sequence
669 && self.pool_id == expected.pool_id
670 && self.pool_identity_fingerprint == expected.pool_identity_fingerprint
671 && self.provisioning_run_id == expected.provisioning_run_id
672 && self.provisioning_request_id == expected.provisioning_request_id
673 && self.transaction_id == expected.transaction_id
674 && self.resource_id == expected.resource_id
675 && self.resource_generation == expected.resource_generation
676 && self.resource_batch_fingerprint == expected.resource_batch_fingerprint
677 }
678}
679
680fn require_exact_wire_value(
681 supplied: &serde_json::Value,
682 expected: &(impl Serialize + ?Sized),
683 label: &str,
684) -> Result<(), VNextError> {
685 let expected = serde_json::to_value(expected).map_err(|error| VNextError::Serialization {
686 context: "serialize trusted resource pool event evidence",
687 message: error.to_string(),
688 })?;
689 if supplied != &expected {
690 return Err(invalid_event(format!(
691 "untrusted resource pool {label} differs from independent evidence"
692 )));
693 }
694 Ok(())
695}
696
697impl UnvalidatedResourcePoolEvent {
698 pub fn revalidate(
699 self,
700 context: &TrustedResourcePoolEventContext<'_>,
701 ) -> Result<ResourcePoolEvent, VNextError> {
702 let rebuilt = match (self.kind, self.detail, context) {
703 (
704 ResourcePoolEventKind::ResourcePoolOpened,
705 UnvalidatedResourcePoolEventDetail::Opened(admission),
706 TrustedResourcePoolEventContext::Opened { evidence },
707 ) => {
708 require_exact_wire_value(&admission, evidence.admission(), "admission")?;
709 ResourcePoolEvent::opened(self.identity.sequence, self.timestamp, evidence)?
710 }
711 (
712 ResourcePoolEventKind::ResourceTransition,
713 UnvalidatedResourcePoolEventDetail::Transition {
714 receipt,
715 context: supplied_context,
716 },
717 TrustedResourcePoolEventContext::Transition { evidence, expected },
718 ) => {
719 validate_pool_binding(
720 evidence,
721 expected.identity(),
722 expected.admission(),
723 "transition",
724 )?;
725 require_exact_wire_value(&supplied_context, expected, "transition context")?;
726 let receipt = receipt.try_validate_against(expected)?;
727 ResourcePoolEvent::transition(
728 self.identity.sequence,
729 self.timestamp,
730 evidence,
731 &receipt,
732 expected,
733 )?
734 }
735 (
736 ResourcePoolEventKind::ResourceLeaseTransition,
737 UnvalidatedResourcePoolEventDetail::LeaseTransition {
738 receipt,
739 context: supplied_context,
740 },
741 TrustedResourcePoolEventContext::LeaseTransition { evidence, expected },
742 ) => {
743 validate_pool_binding(
744 evidence,
745 expected.identity(),
746 expected.admission(),
747 "lease transition",
748 )?;
749 require_exact_wire_value(&supplied_context, expected, "lease context")?;
750 let receipt = receipt.try_validate_against(expected)?;
751 ResourcePoolEvent::lease_transition(
752 self.identity.sequence,
753 self.timestamp,
754 evidence,
755 &receipt,
756 expected,
757 )?
758 }
759 (
760 kind @ (ResourcePoolEventKind::ResourceFailed
761 | ResourcePoolEventKind::ResourceRecoveryCompleted),
762 UnvalidatedResourcePoolEventDetail::Failure(supplied),
763 TrustedResourcePoolEventContext::Failure { evidence, expected },
764 ) => {
765 validate_failure_pool_binding(evidence, expected)?;
766 require_exact_wire_value(&supplied, expected, "failure receipt")?;
767 match kind {
768 ResourcePoolEventKind::ResourceFailed => ResourcePoolEvent::failed(
769 self.identity.sequence,
770 self.timestamp,
771 evidence,
772 expected,
773 )?,
774 ResourcePoolEventKind::ResourceRecoveryCompleted => {
775 ResourcePoolEvent::recovery_completed(
776 self.identity.sequence,
777 self.timestamp,
778 evidence,
779 expected,
780 )?
781 }
782 _ => unreachable!(),
783 }
784 }
785 (
786 ResourcePoolEventKind::ResourcePoolClosed,
787 UnvalidatedResourcePoolEventDetail::Closed { ledger },
788 TrustedResourcePoolEventContext::Closed {
789 evidence,
790 expected_snapshot,
791 },
792 ) => {
793 validate_terminal_snapshot(evidence, expected_snapshot)?;
794 require_exact_wire_value(&ledger, expected_snapshot.entries(), "terminal ledger")?;
795 ResourcePoolEvent::closed(
796 self.identity.sequence,
797 self.timestamp,
798 evidence,
799 expected_snapshot,
800 )?
801 }
802 _ => {
803 return Err(invalid_event(
804 "resource pool wire kind, detail, and external evidence do not match",
805 ));
806 }
807 };
808 if !self.identity.matches(&rebuilt.identity) {
809 return Err(invalid_event(
810 "resource pool wire identity differs from independently rebuilt evidence",
811 ));
812 }
813 Ok(rebuilt)
814 }
815}
816
817#[derive(Debug, Clone)]
818pub struct ResourcePoolEventCursor {
819 evidence: ResourcePoolEvidence,
820 last_sequence: u64,
821 last_timestamp: Option<MonotonicTimestamp>,
822 opened: bool,
823 closed: bool,
824 ledgers: BTreeMap<TransactionId, Vec<ResourceLedgerEntrySnapshot>>,
825 pending_failures: BTreeMap<TransactionId, ResourceFailureReceipt>,
826 seen_failure_ids: BTreeSet<ResourceFailureId>,
827 committed_resource_generations: Option<BTreeSet<(ResourceId, u64)>>,
828}
829
830impl ResourcePoolEventCursor {
831 pub fn new(evidence: ResourcePoolEvidence) -> Self {
832 Self {
833 evidence,
834 last_sequence: 0,
835 last_timestamp: None,
836 opened: false,
837 closed: false,
838 ledgers: BTreeMap::new(),
839 pending_failures: BTreeMap::new(),
840 seen_failure_ids: BTreeSet::new(),
841 committed_resource_generations: None,
842 }
843 }
844
845 pub fn observe(&mut self, event: &ResourcePoolEvent) -> Result<(), VNextError> {
846 let mut next = self.clone();
847 next.observe_inner(event)?;
848 *self = next;
849 Ok(())
850 }
851
852 pub const fn last_sequence(&self) -> u64 {
853 self.last_sequence
854 }
855
856 pub const fn is_open(&self) -> bool {
857 self.opened && !self.closed
858 }
859
860 pub const fn is_closed(&self) -> bool {
861 self.closed
862 }
863
864 pub(super) const fn has_opened(&self) -> bool {
865 self.opened
866 }
867
868 pub(super) fn proves_active_binding(&self, active: &TrustedActiveSequenceBinding) -> bool {
869 let active_resources = active
870 .static_entries()
871 .iter()
872 .map(|entry| (entry.resource_id().clone(), entry.generation()))
873 .collect::<BTreeSet<_>>();
874 self.committed_resource_generations.as_ref() == Some(&active_resources)
875 }
876
877 fn observe_inner(&mut self, event: &ResourcePoolEvent) -> Result<(), VNextError> {
878 if event.identity.sequence != self.last_sequence.saturating_add(1)
879 || self
880 .last_timestamp
881 .is_some_and(|timestamp| event.timestamp <= timestamp)
882 || self.closed
883 || event.identity.pool_id != self.evidence.pool_id
884 || event.identity.pool_identity_fingerprint != self.evidence.pool_identity_fingerprint
885 || event.identity.provisioning_run_id != *self.evidence.provisioning_identity.run_id()
886 || event.identity.provisioning_request_id
887 != *self.evidence.provisioning_identity.request_id()
888 || event.identity.transaction_id
889 != *self.evidence.provisioning_identity.transaction_id()
890 {
891 return Err(invalid_event(
892 "pool journal sequence, timestamp, lifecycle, or pool identity is invalid",
893 ));
894 }
895 match (&event.kind, &event.detail) {
896 (
897 ResourcePoolEventKind::ResourcePoolOpened,
898 ResourcePoolEventDetail::Opened(admission),
899 ) => {
900 if self.opened
901 || self.last_sequence != 0
902 || admission != &self.evidence.admission
903 || event.identity.resource_id.is_some()
904 || event.identity.resource_batch_fingerprint.is_some()
905 {
906 return Err(invalid_event(
907 "ResourcePoolOpened must be the exact first admission event",
908 ));
909 }
910 self.opened = true;
911 }
912 (
913 ResourcePoolEventKind::ResourceTransition,
914 ResourcePoolEventDetail::Transition { receipt, context },
915 ) => {
916 self.require_open()?;
917 validate_transition_receipt_context(receipt, context)?;
918 self.validate_resource_evidence(
919 context.identity(),
920 context.admission(),
921 context.after(),
922 &event.identity,
923 )?;
924 self.apply_ledger(
925 context.identity().transaction_id(),
926 context.before(),
927 context.after(),
928 )?;
929 }
930 (
931 ResourcePoolEventKind::ResourceLeaseTransition,
932 ResourcePoolEventDetail::LeaseTransition { receipt, context },
933 ) => {
934 self.require_open()?;
935 validate_lease_receipt_context(receipt, context)?;
936 self.validate_resource_evidence(
937 context.identity(),
938 context.admission(),
939 context.after(),
940 &event.identity,
941 )?;
942 self.apply_ledger(
943 context.identity().transaction_id(),
944 context.before(),
945 context.after(),
946 )?;
947 }
948 (ResourcePoolEventKind::ResourceFailed, ResourcePoolEventDetail::Failure(failure)) => {
949 self.require_open()?;
950 self.validate_failure(failure, &event.identity)?;
951 let transaction = failure.identity().transaction_id().clone();
952 if failure.recovery_complete()
953 || self.pending_failures.contains_key(&transaction)
954 || !self.seen_failure_ids.insert(failure.failure_id())
955 {
956 return Err(invalid_event(
957 "resource failure is already complete, duplicated, or has a pending anchor",
958 ));
959 }
960 if let Some(current) = self.ledgers.get(&transaction) {
961 if current != failure.ledger_before() {
962 return Err(invalid_event(
963 "resource failure ledger does not continue pool journal",
964 ));
965 }
966 }
967 self.ledgers
968 .insert(transaction.clone(), failure.ledger_after().to_vec());
969 self.pending_failures.insert(transaction, failure.clone());
970 }
971 (
972 ResourcePoolEventKind::ResourceRecoveryCompleted,
973 ResourcePoolEventDetail::Failure(recovery),
974 ) => {
975 self.require_open()?;
976 self.validate_failure(recovery, &event.identity)?;
977 let transaction = recovery.identity().transaction_id().clone();
978 let anchor = self.pending_failures.get(&transaction).ok_or_else(|| {
979 invalid_event("resource recovery has no exact pending failure anchor")
980 })?;
981 recovery.validate_recovery_continuation(anchor)?;
982 if recovery.failure_id() != anchor.failure_id() || !recovery.recovery_complete() {
983 return Err(invalid_event(
984 "resource recovery failure id or completion flag is invalid",
985 ));
986 }
987 self.ledgers
988 .insert(transaction.clone(), recovery.ledger_after().to_vec());
989 self.pending_failures.remove(&transaction);
990 }
991 (
992 ResourcePoolEventKind::ResourcePoolClosed,
993 ResourcePoolEventDetail::Closed { ledger },
994 ) => {
995 self.require_open()?;
996 let current = self
997 .ledgers
998 .get(self.evidence.provisioning_identity.transaction_id())
999 .ok_or_else(|| {
1000 invalid_event("pool closure lacks a complete transaction ledger")
1001 })?;
1002 if ledger != current
1003 || !self.pending_failures.is_empty()
1004 || current.iter().any(|entry| {
1005 !matches!(
1006 entry.transaction_state(),
1007 ResourceTransactionState::RolledBack
1008 | ResourceTransactionState::Released
1009 | ResourceTransactionState::Quarantined
1010 ) || entry.buffer_present()
1011 })
1012 {
1013 return Err(invalid_event(
1014 "pool closure requires the exact terminal ledger and no pending recovery",
1015 ));
1016 }
1017 self.closed = true;
1018 }
1019 _ => {
1020 return Err(invalid_event(
1021 "resource pool event kind and detail do not match",
1022 ));
1023 }
1024 }
1025 self.last_sequence = event.identity.sequence;
1026 self.last_timestamp = Some(event.timestamp);
1027 Ok(())
1028 }
1029
1030 fn require_open(&self) -> Result<(), VNextError> {
1031 if !self.opened || self.closed {
1032 return Err(invalid_event("resource pool is not open"));
1033 }
1034 Ok(())
1035 }
1036
1037 fn validate_resource_evidence(
1038 &self,
1039 identity: &ResourceTransactionIdentity,
1040 admission: &StaticProvisioningBinding,
1041 after: &[ResourceLedgerEntrySnapshot],
1042 event_identity: &ResourcePoolEventIdentity,
1043 ) -> Result<(), VNextError> {
1044 let expected = item_or_batch_identity(after);
1045 if identity != &self.evidence.provisioning_identity
1046 || admission != &self.evidence.admission
1047 || event_identity.resource_id != expected.0
1048 || event_identity.resource_generation != expected.1
1049 || event_identity.resource_batch_fingerprint != expected.2
1050 {
1051 return Err(invalid_event(
1052 "resource event identity differs from receipt/admission/transaction",
1053 ));
1054 }
1055 Ok(())
1056 }
1057
1058 fn validate_failure(
1059 &self,
1060 failure: &ResourceFailureReceipt,
1061 event_identity: &ResourcePoolEventIdentity,
1062 ) -> Result<(), VNextError> {
1063 let expected = failure_item_or_batch_identity(failure);
1064 if failure.failure().domain() != FailureDomain::Resource
1065 || failure.identity() != &self.evidence.provisioning_identity
1066 || failure.admission() != &self.evidence.admission
1067 || event_identity.resource_id != expected.0
1068 || event_identity.resource_generation != expected.1
1069 || event_identity.resource_batch_fingerprint != expected.2
1070 {
1071 return Err(invalid_event(
1072 "resource failure event differs from its full external anchor",
1073 ));
1074 }
1075 Ok(())
1076 }
1077
1078 fn apply_ledger(
1079 &mut self,
1080 transaction_id: &TransactionId,
1081 before: &[ResourceLedgerEntrySnapshot],
1082 after: &[ResourceLedgerEntrySnapshot],
1083 ) -> Result<(), VNextError> {
1084 if before.is_empty()
1085 || after.is_empty()
1086 || before.len() != after.len()
1087 || self
1088 .ledgers
1089 .get(transaction_id)
1090 .is_some_and(|current| current != before)
1091 {
1092 return Err(invalid_event(
1093 "resource ledger before/after does not continue exactly",
1094 ));
1095 }
1096 let mut resources = BTreeSet::new();
1097 if after.iter().any(|entry| {
1098 entry.entry().generation() != self.evidence.admission.admission_generation()
1099 || !resources.insert(entry.entry().resource_id().clone())
1100 }) {
1101 return Err(invalid_event(
1102 "resource ledger contains duplicate or wrong-generation entries",
1103 ));
1104 }
1105 if after.iter().all(|entry| {
1106 entry.transaction_state() == ResourceTransactionState::Committed
1107 && entry.buffer_present()
1108 && entry.entry().state() == ResourceLeaseState::Active
1109 }) {
1110 self.committed_resource_generations = Some(
1111 after
1112 .iter()
1113 .map(|entry| {
1114 (
1115 entry.entry().resource_id().clone(),
1116 entry.entry().generation(),
1117 )
1118 })
1119 .collect(),
1120 );
1121 }
1122 self.ledgers.insert(transaction_id.clone(), after.to_vec());
1123 Ok(())
1124 }
1125}