Skip to main content

ferrum_interfaces/vnext/event/
sequence_binding.rs

1use serde::Serialize;
2use std::collections::BTreeSet;
3
4use super::{
5    canonical_fingerprint, invalid_event, validate_sha256, ActiveSequenceAbortDisposition,
6    ActiveSequenceAbortReceipt, ActiveSequenceCompletionReceipt, ActiveSequencePermit,
7    DeviceRuntime, LogicalAdmissionCoordinatorId, LogicalBackingSliceEvidence, RequestIdentity,
8    ResourceLeaseEntry, ResourceLeaseState, ResourcePoolId, ResourceTransactionIdentity, RunId,
9    SequenceAuthorityId, SequenceSession, SequenceSessionEpoch, SequenceSessionFingerprint,
10    SequenceSessionLiveWitness, SequenceSessionTerminalDisposition, SequenceSessionTerminalReceipt,
11    TrustedPlanRuntimeEvidence, VNextError,
12};
13
14#[derive(Clone, Serialize)]
15enum TrustedActiveSequenceAuthority {
16    StreamActivation,
17    SequenceSession {
18        fingerprint: SequenceSessionFingerprint,
19        #[serde(skip)]
20        live_witness: SequenceSessionLiveWitness,
21    },
22}
23
24#[derive(Clone, Serialize)]
25pub struct TrustedActiveSequenceBinding {
26    plan: TrustedPlanRuntimeEvidence,
27    coordinator_id: LogicalAdmissionCoordinatorId,
28    sequence_authority: SequenceAuthorityId,
29    run_id: RunId,
30    request_id: RequestIdentity,
31    activation_epoch: u64,
32    authority: TrustedActiveSequenceAuthority,
33    runtime_implementation_fingerprint: String,
34    static_entries: Vec<ResourceLeaseEntry>,
35    #[serde(rename = "backing_slices")]
36    legacy_backing_slices: Option<Vec<LogicalBackingSliceEvidence>>,
37    #[serde(skip)]
38    static_pool_identity_fingerprint: Option<String>,
39    #[serde(skip)]
40    fingerprint: String,
41}
42
43impl TrustedActiveSequenceBinding {
44    pub fn from_permit<R>(permit: &ActiveSequencePermit<'_, '_, R>) -> Result<Self, VNextError>
45    where
46        R: DeviceRuntime,
47    {
48        let resources = permit.resources();
49        let plan = resources.plan_evidence();
50        let mut static_entries = resources
51            .static_provisioning()
52            .map(|lease| lease.plan_static_entries().cloned().collect::<Vec<_>>())
53            .unwrap_or_default();
54        static_entries.sort_by(|left, right| left.resource_id().cmp(right.resource_id()));
55        let mut resources = BTreeSet::new();
56        if permit.activation_epoch() == 0
57            || plan.coordinator_id() != permit.coordinator_id()
58            || plan.runtime_implementation_fingerprint()
59                != permit.runtime_implementation_fingerprint()
60            || static_entries.iter().any(|entry| {
61                entry.state() != ResourceLeaseState::Active
62                    || !resources.insert(entry.resource_id().clone())
63            })
64        {
65            return Err(invalid_event(
66                "active permit pool, slot, epoch, admission, or lease entries are invalid",
67            ));
68        }
69        validate_sha256(
70            permit.runtime_implementation_fingerprint(),
71            "runtime implementation fingerprint",
72        )?;
73        let static_pool_identity_fingerprint =
74            plan.static_pool_identity().map(canonical_fingerprint);
75        let mut binding = Self {
76            plan,
77            coordinator_id: permit.coordinator_id(),
78            sequence_authority: permit.sequence_authority(),
79            run_id: permit.run_id().clone(),
80            request_id: permit.request_id().clone(),
81            activation_epoch: permit.activation_epoch(),
82            authority: TrustedActiveSequenceAuthority::StreamActivation,
83            runtime_implementation_fingerprint: permit
84                .runtime_implementation_fingerprint()
85                .to_owned(),
86            static_entries,
87            legacy_backing_slices: Some(
88                permit
89                    .backing_slices()
90                    .iter()
91                    .map(|slice| slice.evidence().clone())
92                    .collect(),
93            ),
94            static_pool_identity_fingerprint,
95            fingerprint: String::new(),
96        };
97        binding.fingerprint = canonical_fingerprint(&binding);
98        Ok(binding)
99    }
100
101    pub fn from_session<R>(session: &SequenceSession<R>) -> Result<Self, VNextError>
102    where
103        R: DeviceRuntime,
104    {
105        session.ensure_open_identity()?;
106        let resources = session.resources();
107        let plan = resources.plan_evidence();
108        let runtime_fingerprint = plan.runtime_implementation_fingerprint().to_owned();
109        validate_sha256(&runtime_fingerprint, "runtime implementation fingerprint")?;
110        let static_pool_identity_fingerprint =
111            plan.static_pool_identity().map(canonical_fingerprint);
112        let mut static_entries = resources
113            .static_provisioning()
114            .map(|lease| lease.plan_static_entries().cloned().collect::<Vec<_>>())
115            .unwrap_or_default();
116        static_entries.sort_by(|left, right| left.resource_id().cmp(right.resource_id()));
117        let mut resource_ids = BTreeSet::new();
118        if plan.coordinator_id() != resources.coordinator_id()
119            || static_entries.iter().any(|entry| {
120                entry.state() != ResourceLeaseState::Active
121                    || !resource_ids.insert(entry.resource_id().clone())
122            })
123        {
124            return Err(invalid_event(
125                "active session plan, coordinator, or lease entries are invalid",
126            ));
127        }
128        let mut binding = Self {
129            plan,
130            coordinator_id: resources.coordinator_id(),
131            sequence_authority: resources.sequence_authority(),
132            run_id: resources.run_id().clone(),
133            request_id: resources.request_id().clone(),
134            activation_epoch: session.epoch().get(),
135            authority: TrustedActiveSequenceAuthority::SequenceSession {
136                fingerprint: session.fingerprint().clone(),
137                live_witness: session.live_witness()?,
138            },
139            runtime_implementation_fingerprint: runtime_fingerprint,
140            static_entries,
141            // A session identity is stable across dynamic backing generations.
142            // The exact physical authority belongs to each captured Step.
143            legacy_backing_slices: None,
144            static_pool_identity_fingerprint,
145            fingerprint: String::new(),
146        };
147        binding.fingerprint = canonical_fingerprint(&binding);
148        Ok(binding)
149    }
150
151    pub fn plan(&self) -> &TrustedPlanRuntimeEvidence {
152        &self.plan
153    }
154
155    pub const fn coordinator_id(&self) -> LogicalAdmissionCoordinatorId {
156        self.coordinator_id
157    }
158
159    pub const fn sequence_authority(&self) -> SequenceAuthorityId {
160        self.sequence_authority
161    }
162
163    pub fn static_pool_id(&self) -> Option<ResourcePoolId> {
164        self.plan.static_pool_identity().map(|pool| pool.pool_id())
165    }
166
167    pub fn static_pool_identity_fingerprint(&self) -> Option<String> {
168        self.static_pool_identity_fingerprint.clone()
169    }
170
171    pub(crate) fn static_pool_identity_fingerprint_ref(&self) -> Option<&str> {
172        self.static_pool_identity_fingerprint.as_deref()
173    }
174
175    pub fn static_provisioning_identity(&self) -> Option<&ResourceTransactionIdentity> {
176        self.plan.static_provisioning_identity()
177    }
178
179    pub fn run_id(&self) -> &RunId {
180        &self.run_id
181    }
182
183    pub fn request_id(&self) -> &RequestIdentity {
184        &self.request_id
185    }
186
187    pub const fn activation_epoch(&self) -> u64 {
188        self.activation_epoch
189    }
190
191    pub(crate) fn matches_sequence_session(
192        &self,
193        epoch: SequenceSessionEpoch,
194        fingerprint: &SequenceSessionFingerprint,
195    ) -> bool {
196        self.activation_epoch == epoch.get()
197            && matches!(
198                &self.authority,
199                TrustedActiveSequenceAuthority::SequenceSession {
200                    fingerprint: bound_fingerprint,
201                    ..
202                } if bound_fingerprint == fingerprint
203            )
204    }
205
206    fn is_stream_activation(&self) -> bool {
207        matches!(
208            &self.authority,
209            TrustedActiveSequenceAuthority::StreamActivation
210        )
211    }
212
213    pub(crate) fn ensure_open_for_emission(&self) -> Result<(), VNextError> {
214        match &self.authority {
215            TrustedActiveSequenceAuthority::StreamActivation => Err(invalid_event(
216                "node execution emission requires typed sequence-session authority",
217            )),
218            TrustedActiveSequenceAuthority::SequenceSession { live_witness, .. } => {
219                live_witness.ensure_open()
220            }
221        }
222    }
223
224    pub(super) fn ensure_live_for_emission(&self) -> Result<(), VNextError> {
225        match &self.authority {
226            TrustedActiveSequenceAuthority::StreamActivation => Err(invalid_event(
227                "operation progress emission requires typed sequence-session authority",
228            )),
229            TrustedActiveSequenceAuthority::SequenceSession { live_witness, .. } => {
230                live_witness.ensure_live()
231            }
232        }
233    }
234
235    pub(super) fn matches_abort_disposition(
236        &self,
237        disposition: ActiveSequenceAbortDisposition,
238    ) -> bool {
239        matches!(
240            (&self.authority, disposition),
241            (
242                TrustedActiveSequenceAuthority::StreamActivation,
243                ActiveSequenceAbortDisposition::SynchronizedAndPoisoned,
244            ) | (
245                TrustedActiveSequenceAuthority::SequenceSession { .. },
246                ActiveSequenceAbortDisposition::SequenceSessionTerminalized,
247            )
248        )
249    }
250
251    pub fn runtime_implementation_fingerprint(&self) -> &str {
252        &self.runtime_implementation_fingerprint
253    }
254
255    pub fn static_entries(&self) -> &[ResourceLeaseEntry] {
256        &self.static_entries
257    }
258
259    pub fn legacy_backing_slices(&self) -> Option<&[LogicalBackingSliceEvidence]> {
260        self.legacy_backing_slices.as_deref()
261    }
262
263    pub fn backing_slices(&self) -> &[LogicalBackingSliceEvidence] {
264        self.legacy_backing_slices.as_deref().unwrap_or_default()
265    }
266
267    pub fn fingerprint(&self) -> &str {
268        &self.fingerprint
269    }
270}
271
272#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
273pub struct TrustedCompletedSequenceBinding {
274    plan: TrustedPlanRuntimeEvidence,
275    coordinator_id: LogicalAdmissionCoordinatorId,
276    sequence_authority: SequenceAuthorityId,
277    run_id: RunId,
278    request_id: RequestIdentity,
279    activation_epoch: u64,
280    runtime_implementation_fingerprint: String,
281    active_sequence_fingerprint: String,
282    #[serde(skip)]
283    fingerprint: String,
284}
285
286impl TrustedCompletedSequenceBinding {
287    pub fn from_receipt(
288        receipt: &ActiveSequenceCompletionReceipt,
289        active: &TrustedActiveSequenceBinding,
290    ) -> Result<Self, VNextError> {
291        if !active.is_stream_activation()
292            || receipt.plan() != active.plan()
293            || receipt.sequence_authority() != active.sequence_authority()
294            || receipt.run_id() != active.run_id()
295            || receipt.request_id() != active.request_id()
296            || receipt.activation_epoch() != active.activation_epoch()
297            || receipt.runtime_implementation_fingerprint()
298                != active.runtime_implementation_fingerprint()
299        {
300            return Err(invalid_event(
301                "sequence completion receipt differs from the active pool, request, slot, epoch, or runtime",
302            ));
303        }
304        let mut binding = Self {
305            plan: receipt.plan().clone(),
306            coordinator_id: receipt.plan().coordinator_id(),
307            sequence_authority: receipt.sequence_authority(),
308            run_id: receipt.run_id().clone(),
309            request_id: receipt.request_id().clone(),
310            activation_epoch: receipt.activation_epoch(),
311            runtime_implementation_fingerprint: receipt
312                .runtime_implementation_fingerprint()
313                .to_owned(),
314            active_sequence_fingerprint: active.fingerprint().to_owned(),
315            fingerprint: String::new(),
316        };
317        binding.fingerprint = canonical_fingerprint(&binding);
318        Ok(binding)
319    }
320
321    pub fn from_session_receipt(
322        receipt: &SequenceSessionTerminalReceipt,
323        active: &TrustedActiveSequenceBinding,
324    ) -> Result<Self, VNextError> {
325        if receipt.disposition() != SequenceSessionTerminalDisposition::Completed
326            || receipt.retired_frames() == 0
327            || !active.matches_sequence_session(receipt.epoch(), receipt.fingerprint())
328        {
329            return Err(invalid_event(
330                "sequence session completion differs from the exact active session or has no retired frame",
331            ));
332        }
333        let mut binding = Self {
334            plan: active.plan().clone(),
335            coordinator_id: active.coordinator_id(),
336            sequence_authority: active.sequence_authority(),
337            run_id: active.run_id().clone(),
338            request_id: active.request_id().clone(),
339            activation_epoch: active.activation_epoch(),
340            runtime_implementation_fingerprint: active
341                .runtime_implementation_fingerprint()
342                .to_owned(),
343            active_sequence_fingerprint: active.fingerprint().to_owned(),
344            fingerprint: String::new(),
345        };
346        binding.fingerprint = canonical_fingerprint(&binding);
347        Ok(binding)
348    }
349
350    pub fn plan(&self) -> &TrustedPlanRuntimeEvidence {
351        &self.plan
352    }
353
354    pub const fn coordinator_id(&self) -> LogicalAdmissionCoordinatorId {
355        self.coordinator_id
356    }
357
358    pub const fn sequence_authority(&self) -> SequenceAuthorityId {
359        self.sequence_authority
360    }
361
362    pub fn run_id(&self) -> &RunId {
363        &self.run_id
364    }
365
366    pub fn request_id(&self) -> &RequestIdentity {
367        &self.request_id
368    }
369
370    pub const fn activation_epoch(&self) -> u64 {
371        self.activation_epoch
372    }
373
374    pub fn runtime_implementation_fingerprint(&self) -> &str {
375        &self.runtime_implementation_fingerprint
376    }
377
378    pub fn active_sequence_fingerprint(&self) -> &str {
379        &self.active_sequence_fingerprint
380    }
381
382    pub fn fingerprint(&self) -> &str {
383        &self.fingerprint
384    }
385}
386
387#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
388pub struct TrustedAbortedSequenceBinding {
389    plan: TrustedPlanRuntimeEvidence,
390    coordinator_id: LogicalAdmissionCoordinatorId,
391    sequence_authority: SequenceAuthorityId,
392    run_id: RunId,
393    request_id: RequestIdentity,
394    activation_epoch: u64,
395    runtime_implementation_fingerprint: String,
396    active_sequence_fingerprint: String,
397    disposition: ActiveSequenceAbortDisposition,
398    #[serde(skip)]
399    fingerprint: String,
400}
401
402impl TrustedAbortedSequenceBinding {
403    pub fn from_receipt(
404        receipt: &ActiveSequenceAbortReceipt,
405        active: &TrustedActiveSequenceBinding,
406    ) -> Result<Self, VNextError> {
407        if !active.is_stream_activation()
408            || receipt.disposition() != ActiveSequenceAbortDisposition::SynchronizedAndPoisoned
409            || receipt.plan() != active.plan()
410            || receipt.sequence_authority() != active.sequence_authority()
411            || receipt.run_id() != active.run_id()
412            || receipt.request_id() != active.request_id()
413            || receipt.activation_epoch() != active.activation_epoch()
414            || receipt.runtime_implementation_fingerprint()
415                != active.runtime_implementation_fingerprint()
416        {
417            return Err(invalid_event(
418                "sequence abort receipt differs from the active pool, request, slot, epoch, runtime, or poison disposition",
419            ));
420        }
421        let mut binding = Self {
422            plan: receipt.plan().clone(),
423            coordinator_id: receipt.plan().coordinator_id(),
424            sequence_authority: receipt.sequence_authority(),
425            run_id: receipt.run_id().clone(),
426            request_id: receipt.request_id().clone(),
427            activation_epoch: receipt.activation_epoch(),
428            runtime_implementation_fingerprint: receipt
429                .runtime_implementation_fingerprint()
430                .to_owned(),
431            active_sequence_fingerprint: active.fingerprint().to_owned(),
432            disposition: receipt.disposition(),
433            fingerprint: String::new(),
434        };
435        binding.fingerprint = canonical_fingerprint(&binding);
436        Ok(binding)
437    }
438
439    pub fn from_session_receipt(
440        receipt: &SequenceSessionTerminalReceipt,
441        active: &TrustedActiveSequenceBinding,
442    ) -> Result<Self, VNextError> {
443        if receipt.disposition() != SequenceSessionTerminalDisposition::Aborted
444            || !active.matches_sequence_session(receipt.epoch(), receipt.fingerprint())
445        {
446            return Err(invalid_event(
447                "sequence session abort differs from the exact active session",
448            ));
449        }
450        let mut binding = Self {
451            plan: active.plan().clone(),
452            coordinator_id: active.coordinator_id(),
453            sequence_authority: active.sequence_authority(),
454            run_id: active.run_id().clone(),
455            request_id: active.request_id().clone(),
456            activation_epoch: active.activation_epoch(),
457            runtime_implementation_fingerprint: active
458                .runtime_implementation_fingerprint()
459                .to_owned(),
460            active_sequence_fingerprint: active.fingerprint().to_owned(),
461            disposition: ActiveSequenceAbortDisposition::SequenceSessionTerminalized,
462            fingerprint: String::new(),
463        };
464        binding.fingerprint = canonical_fingerprint(&binding);
465        Ok(binding)
466    }
467
468    pub fn plan(&self) -> &TrustedPlanRuntimeEvidence {
469        &self.plan
470    }
471
472    pub const fn coordinator_id(&self) -> LogicalAdmissionCoordinatorId {
473        self.coordinator_id
474    }
475
476    pub const fn sequence_authority(&self) -> SequenceAuthorityId {
477        self.sequence_authority
478    }
479
480    pub fn run_id(&self) -> &RunId {
481        &self.run_id
482    }
483
484    pub fn request_id(&self) -> &RequestIdentity {
485        &self.request_id
486    }
487
488    pub const fn activation_epoch(&self) -> u64 {
489        self.activation_epoch
490    }
491
492    pub fn runtime_implementation_fingerprint(&self) -> &str {
493        &self.runtime_implementation_fingerprint
494    }
495
496    pub fn active_sequence_fingerprint(&self) -> &str {
497        &self.active_sequence_fingerprint
498    }
499
500    pub const fn disposition(&self) -> ActiveSequenceAbortDisposition {
501        self.disposition
502    }
503
504    pub fn fingerprint(&self) -> &str {
505        &self.fingerprint
506    }
507}