Skip to main content

ferrum_interfaces/vnext/event/
resource_pool.rs

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}