Skip to main content

ferrum_interfaces/vnext/event/
replay.rs

1use serde::{Deserialize, Serialize};
2use std::collections::BTreeSet;
3
4use super::{
5    canonical_fingerprint, has_aborted, has_active, has_completed, invalid_event,
6    same_operation_authority_except_observation, sha256_bytes, validate_sha256,
7    BatchOperationIdentity, CompletionDrainReceipt, CompletionQuarantineReceipt, CompletionSlotId,
8    ContractVersion, ExecutionEvent, ExecutionEventCursor, ExecutionEventDetail,
9    ExecutionEventKind, ExecutionIdentityEnvelope, IdentifiedFailure, OperationCompletionReceipt,
10    OperationParticipantCompletionDisposition, OperationParticipantCompletionReceipt,
11    PlanRuntimeCloseReceipt, PlanRuntimeQuarantineReceipt, ResolvedModelPlan, ResourcePoolEvent,
12    ResourcePoolEventCursor, ResourcePoolEvidence, ResourcePoolId,
13    SubmittedOperationParticipantReceipt, SubmittedOperationReceipt, TrustedAbortedSequenceBinding,
14    TrustedActiveSequenceBinding, TrustedCompletedSequenceBinding, TrustedExecutionEventContext,
15    TrustedExecutionTopology, UnvalidatedExecutionIdentityParts, VNextError,
16    EXECUTION_IDENTITY_VERSION, MAX_REPLAY_IDENTITY_WIRE_BYTES,
17};
18
19/// Independent evidence needed to rebuild a replay identity. None of these
20/// values are accepted from the serialized replay envelope itself.
21pub struct ReplayEvidence<'a> {
22    resolved_plan: &'a ResolvedModelPlan,
23    request_input: &'a [u8],
24    initial_state: &'a [u8],
25    random_seed: u64,
26    request_journal: &'a [ExecutionEvent],
27    active_binding: &'a TrustedActiveSequenceBinding,
28    completed_binding: Option<&'a TrustedCompletedSequenceBinding>,
29    aborted_binding: Option<&'a TrustedAbortedSequenceBinding>,
30    cleanup_requirement: ReplayCleanupRequirement,
31    plan_cleanup: ReplayPlanCleanupEvidence<'a>,
32    operation_completions: &'a [OperationCompletionReceipt],
33    operation_drains: &'a [CompletionDrainReceipt],
34    // Quarantine retains invocation ownership and is never terminal replay evidence.
35    operation_quarantines: &'a [CompletionQuarantineReceipt],
36    pool_evidence: Option<&'a ResourcePoolEvidence>,
37    pool_journal: &'a [ResourcePoolEvent],
38}
39
40impl<'a> ReplayEvidence<'a> {
41    pub fn new(
42        resolved_plan: &'a ResolvedModelPlan,
43        request_input: &'a [u8],
44        initial_state: &'a [u8],
45        random_seed: u64,
46        request_journal: &'a [ExecutionEvent],
47        active_binding: &'a TrustedActiveSequenceBinding,
48        completed_binding: Option<&'a TrustedCompletedSequenceBinding>,
49        aborted_binding: Option<&'a TrustedAbortedSequenceBinding>,
50        cleanup_requirement: ReplayCleanupRequirement,
51        plan_cleanup: ReplayPlanCleanupEvidence<'a>,
52        operation_completions: &'a [OperationCompletionReceipt],
53        operation_drains: &'a [CompletionDrainReceipt],
54        operation_quarantines: &'a [CompletionQuarantineReceipt],
55        pool_evidence: &'a ResourcePoolEvidence,
56        pool_journal: &'a [ResourcePoolEvent],
57    ) -> Self {
58        Self {
59            resolved_plan,
60            request_input,
61            initial_state,
62            random_seed,
63            request_journal,
64            active_binding,
65            completed_binding,
66            aborted_binding,
67            cleanup_requirement,
68            plan_cleanup,
69            operation_completions,
70            operation_drains,
71            operation_quarantines,
72            pool_evidence: Some(pool_evidence),
73            pool_journal,
74        }
75    }
76
77    #[allow(clippy::too_many_arguments)]
78    pub fn new_no_static(
79        resolved_plan: &'a ResolvedModelPlan,
80        request_input: &'a [u8],
81        initial_state: &'a [u8],
82        random_seed: u64,
83        request_journal: &'a [ExecutionEvent],
84        active_binding: &'a TrustedActiveSequenceBinding,
85        completed_binding: Option<&'a TrustedCompletedSequenceBinding>,
86        aborted_binding: Option<&'a TrustedAbortedSequenceBinding>,
87        cleanup_requirement: ReplayCleanupRequirement,
88        plan_cleanup: ReplayPlanCleanupEvidence<'a>,
89        operation_completions: &'a [OperationCompletionReceipt],
90        operation_drains: &'a [CompletionDrainReceipt],
91        operation_quarantines: &'a [CompletionQuarantineReceipt],
92    ) -> Self {
93        Self {
94            resolved_plan,
95            request_input,
96            initial_state,
97            random_seed,
98            request_journal,
99            active_binding,
100            completed_binding,
101            aborted_binding,
102            cleanup_requirement,
103            plan_cleanup,
104            operation_completions,
105            operation_drains,
106            operation_quarantines,
107            pool_evidence: None,
108            pool_journal: &[],
109        }
110    }
111
112    pub fn resolved_plan(&self) -> &ResolvedModelPlan {
113        self.resolved_plan
114    }
115
116    pub fn request_input(&self) -> &[u8] {
117        self.request_input
118    }
119
120    pub fn initial_state(&self) -> &[u8] {
121        self.initial_state
122    }
123
124    pub const fn random_seed(&self) -> u64 {
125        self.random_seed
126    }
127
128    pub fn request_journal(&self) -> &[ExecutionEvent] {
129        self.request_journal
130    }
131
132    pub fn active_binding(&self) -> &TrustedActiveSequenceBinding {
133        self.active_binding
134    }
135
136    pub fn completed_binding(&self) -> Option<&TrustedCompletedSequenceBinding> {
137        self.completed_binding
138    }
139
140    pub fn aborted_binding(&self) -> Option<&TrustedAbortedSequenceBinding> {
141        self.aborted_binding
142    }
143
144    pub const fn cleanup_requirement(&self) -> ReplayCleanupRequirement {
145        self.cleanup_requirement
146    }
147
148    pub const fn plan_cleanup(&self) -> ReplayPlanCleanupEvidence<'a> {
149        self.plan_cleanup
150    }
151
152    pub fn operation_completions(&self) -> &[OperationCompletionReceipt] {
153        self.operation_completions
154    }
155
156    pub fn operation_drains(&self) -> &[CompletionDrainReceipt] {
157        self.operation_drains
158    }
159
160    pub fn operation_quarantines(&self) -> &[CompletionQuarantineReceipt] {
161        self.operation_quarantines
162    }
163
164    pub fn pool_evidence(&self) -> Option<&ResourcePoolEvidence> {
165        self.pool_evidence
166    }
167
168    pub fn pool_journal(&self) -> &[ResourcePoolEvent] {
169        self.pool_journal
170    }
171}
172
173#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
174enum ReplayOperationTerminalKey {
175    Completion(usize),
176    Drain(usize),
177    Quarantine(usize),
178}
179
180#[derive(Clone, Copy)]
181enum ReplayOperationTerminalRef<'a> {
182    Completion(&'a OperationCompletionReceipt),
183    Drain(&'a CompletionDrainReceipt),
184    Quarantine(&'a CompletionQuarantineReceipt),
185}
186
187impl<'a> ReplayOperationTerminalRef<'a> {
188    fn slot_id(self) -> CompletionSlotId {
189        match self {
190            Self::Completion(receipt) => receipt.submission().slot_id(),
191            Self::Drain(receipt) => receipt.slot_id(),
192            Self::Quarantine(receipt) => receipt.slot_id(),
193        }
194    }
195
196    fn batch_identity(self) -> &'a BatchOperationIdentity {
197        match self {
198            Self::Completion(receipt) => receipt.submission().batch_identity(),
199            Self::Drain(receipt) => receipt.batch_identity(),
200            Self::Quarantine(receipt) => receipt.batch_identity(),
201        }
202    }
203
204    fn participant_submission(
205        self,
206        identity: &ExecutionIdentityEnvelope,
207    ) -> Option<&'a SubmittedOperationParticipantReceipt> {
208        self.submission()?
209            .participants()
210            .iter()
211            .find(|participant| participant.identity() == identity)
212    }
213
214    fn submission(self) -> Option<&'a SubmittedOperationReceipt> {
215        match self {
216            Self::Completion(receipt) => Some(receipt.submission()),
217            Self::Drain(receipt) => receipt.submission(),
218            Self::Quarantine(receipt) => receipt.submission(),
219        }
220    }
221
222    fn participant_completion(
223        self,
224        identity: &ExecutionIdentityEnvelope,
225    ) -> Option<&'a OperationParticipantCompletionReceipt> {
226        match self {
227            Self::Completion(receipt) => receipt.participants().iter().find(|participant| {
228                same_operation_authority_except_observation(
229                    identity.parts(),
230                    participant.submission().identity().parts(),
231                )
232            }),
233            Self::Drain(_) | Self::Quarantine(_) => None,
234        }
235    }
236
237    fn contains_identity(self, identity: &ExecutionIdentityEnvelope) -> bool {
238        self.batch_identity()
239            .participants()
240            .iter()
241            .any(|participant| participant.identity() == identity)
242    }
243
244    fn had_submission_fence(self) -> bool {
245        match self {
246            Self::Completion(_) => true,
247            Self::Drain(receipt) => receipt.had_submission_fence(),
248            Self::Quarantine(receipt) => receipt.had_submission_fence(),
249        }
250    }
251
252    fn exact_failed_completion(
253        self,
254        identity: &ExecutionIdentityEnvelope,
255    ) -> Option<&'a IdentifiedFailure> {
256        self.participant_completion(identity)
257            .and_then(|participant| match participant.disposition() {
258                OperationParticipantCompletionDisposition::FailedButQuiescent(failure) => {
259                    Some(failure)
260                }
261                OperationParticipantCompletionDisposition::Succeeded
262                | OperationParticipantCompletionDisposition::ContractFailedButQuiescent(_) => None,
263            })
264    }
265
266    fn participant_is_success(self, identity: &ExecutionIdentityEnvelope) -> bool {
267        self.participant_completion(identity)
268            .is_some_and(|participant| {
269                matches!(
270                    participant.disposition(),
271                    OperationParticipantCompletionDisposition::Succeeded
272                )
273            })
274    }
275}
276
277#[derive(Serialize)]
278struct ReplayOperationTerminalFingerprint<'a> {
279    completions: &'a [OperationCompletionReceipt],
280    drains: &'a [CompletionDrainReceipt],
281    quarantines: &'a [CompletionQuarantineReceipt],
282}
283
284fn replay_operation_terminals<'a>(
285    evidence: &'a ReplayEvidence<'_>,
286) -> Vec<(ReplayOperationTerminalKey, ReplayOperationTerminalRef<'a>)> {
287    evidence
288        .operation_completions
289        .iter()
290        .enumerate()
291        .map(|(index, receipt)| {
292            (
293                ReplayOperationTerminalKey::Completion(index),
294                ReplayOperationTerminalRef::Completion(receipt),
295            )
296        })
297        .chain(
298            evidence
299                .operation_drains
300                .iter()
301                .enumerate()
302                .map(|(index, receipt)| {
303                    (
304                        ReplayOperationTerminalKey::Drain(index),
305                        ReplayOperationTerminalRef::Drain(receipt),
306                    )
307                }),
308        )
309        .chain(
310            evidence
311                .operation_quarantines
312                .iter()
313                .enumerate()
314                .map(|(index, receipt)| {
315                    (
316                        ReplayOperationTerminalKey::Quarantine(index),
317                        ReplayOperationTerminalRef::Quarantine(receipt),
318                    )
319                }),
320        )
321        .collect()
322}
323
324#[derive(Debug, Clone, Copy, PartialEq, Eq)]
325pub enum ReplayCleanupRequirement {
326    RequireClean,
327    AllowPending,
328}
329
330/// Independent root-cleanup evidence supplied while rebuilding replay. Pending
331/// is explicit and is accepted only when the caller allows pending cleanup.
332/// The receipt variants are core-signed outputs and cannot be deserialized or
333/// constructed by the replay caller.
334#[derive(Clone, Copy)]
335pub enum ReplayPlanCleanupEvidence<'a> {
336    Pending,
337    Closed(&'a PlanRuntimeCloseReceipt),
338    Quarantined(&'a PlanRuntimeQuarantineReceipt),
339}
340
341#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
342#[serde(rename_all = "snake_case")]
343pub enum ReplayCleanupStatus {
344    Completed,
345    SequenceQuiescent,
346    Quarantined,
347    CleanupPending,
348}
349
350/// A replay identity is trusted output. Deserialization always goes through
351/// `UnvalidatedReplayIdentity` and reconstruction from independent evidence.
352#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
353pub struct ReplayIdentity {
354    identity_version: ContractVersion,
355    terminal_identity: ExecutionIdentityEnvelope,
356    resolved_plan_fingerprint: String,
357    execution_topology_fingerprint: String,
358    request_input_fingerprint: String,
359    initial_state_fingerprint: String,
360    random_seed: u64,
361    request_journal_event_count: u64,
362    request_journal_fingerprint: String,
363    active_sequence_fingerprint: String,
364    completed_sequence_fingerprint: Option<String>,
365    aborted_sequence_fingerprint: Option<String>,
366    cleanup_status: ReplayCleanupStatus,
367    operation_terminal_evidence_count: u64,
368    operation_terminal_evidence_fingerprint: String,
369    resource_pool_id: Option<ResourcePoolId>,
370    resource_pool_identity_fingerprint: Option<String>,
371    pool_journal_event_count: u64,
372    pool_journal_fingerprint: Option<String>,
373    plan_cleanup_fingerprint: Option<String>,
374}
375
376#[derive(Debug, Clone, PartialEq, Eq)]
377pub struct UnvalidatedReplayIdentity {
378    identity_version: ContractVersion,
379    terminal_identity: UnvalidatedExecutionIdentityParts,
380    resolved_plan_fingerprint: String,
381    execution_topology_fingerprint: String,
382    request_input_fingerprint: String,
383    initial_state_fingerprint: String,
384    random_seed: u64,
385    request_journal_event_count: u64,
386    request_journal_fingerprint: String,
387    active_sequence_fingerprint: String,
388    completed_sequence_fingerprint: Option<String>,
389    aborted_sequence_fingerprint: Option<String>,
390    cleanup_status: ReplayCleanupStatus,
391    operation_terminal_evidence_count: u64,
392    operation_terminal_evidence_fingerprint: String,
393    resource_pool_id: Option<ResourcePoolId>,
394    resource_pool_identity_fingerprint: Option<String>,
395    pool_journal_event_count: u64,
396    pool_journal_fingerprint: Option<String>,
397    plan_cleanup_fingerprint: Option<String>,
398}
399
400#[derive(Deserialize)]
401#[serde(deny_unknown_fields)]
402struct ReplayIdentityWire {
403    identity_version: ContractVersion,
404    terminal_identity: UnvalidatedExecutionIdentityParts,
405    resolved_plan_fingerprint: String,
406    execution_topology_fingerprint: String,
407    request_input_fingerprint: String,
408    initial_state_fingerprint: String,
409    random_seed: u64,
410    request_journal_event_count: u64,
411    request_journal_fingerprint: String,
412    active_sequence_fingerprint: String,
413    completed_sequence_fingerprint: Option<String>,
414    aborted_sequence_fingerprint: Option<String>,
415    cleanup_status: ReplayCleanupStatus,
416    operation_terminal_evidence_count: u64,
417    operation_terminal_evidence_fingerprint: String,
418    resource_pool_id: Option<ResourcePoolId>,
419    resource_pool_identity_fingerprint: Option<String>,
420    pool_journal_event_count: u64,
421    pool_journal_fingerprint: Option<String>,
422    plan_cleanup_fingerprint: Option<String>,
423}
424
425impl From<ReplayIdentityWire> for UnvalidatedReplayIdentity {
426    fn from(wire: ReplayIdentityWire) -> Self {
427        Self {
428            identity_version: wire.identity_version,
429            terminal_identity: wire.terminal_identity,
430            resolved_plan_fingerprint: wire.resolved_plan_fingerprint,
431            execution_topology_fingerprint: wire.execution_topology_fingerprint,
432            request_input_fingerprint: wire.request_input_fingerprint,
433            initial_state_fingerprint: wire.initial_state_fingerprint,
434            random_seed: wire.random_seed,
435            request_journal_event_count: wire.request_journal_event_count,
436            request_journal_fingerprint: wire.request_journal_fingerprint,
437            active_sequence_fingerprint: wire.active_sequence_fingerprint,
438            completed_sequence_fingerprint: wire.completed_sequence_fingerprint,
439            aborted_sequence_fingerprint: wire.aborted_sequence_fingerprint,
440            cleanup_status: wire.cleanup_status,
441            operation_terminal_evidence_count: wire.operation_terminal_evidence_count,
442            operation_terminal_evidence_fingerprint: wire.operation_terminal_evidence_fingerprint,
443            resource_pool_id: wire.resource_pool_id,
444            resource_pool_identity_fingerprint: wire.resource_pool_identity_fingerprint,
445            pool_journal_event_count: wire.pool_journal_event_count,
446            pool_journal_fingerprint: wire.pool_journal_fingerprint,
447            plan_cleanup_fingerprint: wire.plan_cleanup_fingerprint,
448        }
449    }
450}
451
452fn validate_replay_plan_cleanup(
453    active: &TrustedActiveSequenceBinding,
454    aborted: bool,
455    requirement: ReplayCleanupRequirement,
456    cleanup: ReplayPlanCleanupEvidence<'_>,
457) -> Result<(ReplayCleanupStatus, Option<String>), VNextError> {
458    let expected_static_resources = active.static_entries().len();
459    match cleanup {
460        ReplayPlanCleanupEvidence::Pending => {
461            if requirement == ReplayCleanupRequirement::RequireClean {
462                return Err(invalid_event(
463                    "clean replay requires an exact plan close or quarantine receipt",
464                ));
465            }
466            Ok((ReplayCleanupStatus::CleanupPending, None))
467        }
468        ReplayPlanCleanupEvidence::Closed(receipt) => {
469            if receipt.evidence() != active.plan()
470                || receipt.released_static_resources() != expected_static_resources
471            {
472                return Err(invalid_event(
473                    "replay plan close receipt differs from the active plan or static resource set",
474                ));
475            }
476            Ok((
477                if aborted {
478                    ReplayCleanupStatus::SequenceQuiescent
479                } else {
480                    ReplayCleanupStatus::Completed
481                },
482                Some(canonical_fingerprint(receipt)),
483            ))
484        }
485        ReplayPlanCleanupEvidence::Quarantined(receipt) => {
486            let accounted_static_resources = receipt
487                .released_static_resources()
488                .checked_add(receipt.quarantined_static_resources())
489                .ok_or_else(|| invalid_event("replay quarantine resource count overflows usize"))?;
490            if expected_static_resources == 0
491                || receipt.evidence() != active.plan()
492                || accounted_static_resources != expected_static_resources
493            {
494                return Err(invalid_event(
495                    "replay quarantine receipt differs from the active static plan resource set",
496                ));
497            }
498            Ok((
499                ReplayCleanupStatus::Quarantined,
500                Some(canonical_fingerprint(receipt)),
501            ))
502        }
503    }
504}
505
506impl ReplayIdentity {
507    pub fn from_evidence(evidence: &ReplayEvidence<'_>) -> Result<Self, VNextError> {
508        if evidence.request_journal.is_empty() {
509            return Err(invalid_event(
510                "replay evidence requires a non-empty request journal",
511            ));
512        }
513        if !evidence.operation_quarantines.is_empty() {
514            if evidence
515                .operation_quarantines
516                .iter()
517                .any(|receipt| !receipt.is_current())
518            {
519                return Err(invalid_event(
520                    "completion quarantine replay evidence was superseded by a successful drain",
521                ));
522            }
523            if evidence.cleanup_requirement != ReplayCleanupRequirement::AllowPending
524                || !matches!(evidence.plan_cleanup, ReplayPlanCleanupEvidence::Pending)
525            {
526                return Err(invalid_event(
527                    "current completion quarantine is pending ownership and requires explicitly pending replay cleanup",
528                ));
529            }
530        }
531
532        let topology =
533            TrustedExecutionTopology::from_plan(evidence.resolved_plan.execution_plan())?;
534        let active = evidence.active_binding;
535        let completed = evidence.completed_binding;
536        let aborted = evidence.aborted_binding;
537        if completed.is_some() == aborted.is_some() {
538            return Err(invalid_event(
539                "replay requires exactly one external sequence completion or abort binding",
540            ));
541        }
542        if active.plan().plan_id() != topology.plan_id()
543            || active.plan().plan_hash() != topology.plan_hash()
544            || active.plan().device_id() != topology.device_id()
545            || active.plan().runtime_implementation_fingerprint()
546                != topology.device_runtime_implementation_fingerprint()
547            || active.runtime_implementation_fingerprint()
548                != topology.device_runtime_implementation_fingerprint()
549        {
550            return Err(invalid_event(
551                "replay plan and active binding do not share one authority",
552            ));
553        }
554        if let Some(completed) = completed {
555            if completed.active_sequence_fingerprint() != active.fingerprint()
556                || completed.sequence_authority() != active.sequence_authority()
557                || completed.run_id() != active.run_id()
558                || completed.request_id() != active.request_id()
559                || completed.activation_epoch() != active.activation_epoch()
560                || completed.runtime_implementation_fingerprint()
561                    != active.runtime_implementation_fingerprint()
562            {
563                return Err(invalid_event(
564                    "replay completion evidence differs from its active sequence",
565                ));
566            }
567        }
568        if let Some(aborted) = aborted {
569            if !active.matches_abort_disposition(aborted.disposition())
570                || aborted.active_sequence_fingerprint() != active.fingerprint()
571                || aborted.sequence_authority() != active.sequence_authority()
572                || aborted.run_id() != active.run_id()
573                || aborted.request_id() != active.request_id()
574                || aborted.activation_epoch() != active.activation_epoch()
575                || aborted.runtime_implementation_fingerprint()
576                    != active.runtime_implementation_fingerprint()
577            {
578                return Err(invalid_event(
579                    "replay abort evidence differs from its active sequence",
580                ));
581            }
582        }
583
584        let first = evidence
585            .request_journal
586            .first()
587            .expect("non-empty request journal");
588        let terminal = evidence
589            .request_journal
590            .last()
591            .expect("non-empty request journal");
592        let run_id = &first.identity().parts().run_id;
593        let request_id = &first.identity().parts().request_id;
594        if run_id != active.run_id() || request_id != active.request_id() {
595            return Err(invalid_event(
596                "replay request journal differs from the active run/request binding",
597            ));
598        }
599
600        let operation_terminals = replay_operation_terminals(evidence);
601        let mut terminal_slots = BTreeSet::new();
602        if operation_terminals
603            .iter()
604            .any(|(_, terminal)| !terminal_slots.insert(terminal.slot_id()))
605        {
606            return Err(invalid_event(
607                "replay operation terminal evidence reuses a completion slot across terminal types",
608            ));
609        }
610
611        let mut request_cursor = ExecutionEventCursor::new(run_id.clone(), request_id.clone());
612        let mut observed_active_identity = false;
613        let mut used_submitted_terminals = BTreeSet::new();
614        let mut first_failure: Option<IdentifiedFailure> = None;
615        for event in evidence.request_journal {
616            observed_active_identity |= has_active(event.identity().parts());
617            let context = match event.kind() {
618                ExecutionEventKind::RequestAccepted => {
619                    TrustedExecutionEventContext::pre_plan(run_id, request_id)
620                }
621                ExecutionEventKind::PlanBuilt => {
622                    TrustedExecutionEventContext::bound(run_id, request_id, &topology)
623                }
624                ExecutionEventKind::FailureObserved => {
625                    let failure = match event.detail() {
626                        ExecutionEventDetail::Failure(failure) => failure,
627                        _ => unreachable!("trusted FailureObserved shape was validated"),
628                    };
629                    if first_failure.is_some() {
630                        return Err(invalid_event(
631                            "replay request journal contains more than one first failure",
632                        ));
633                    }
634                    first_failure = Some(failure.clone());
635                    let unsubmitted_recoveries = operation_terminals
636                        .iter()
637                        .filter(|(_, terminal)| {
638                            !terminal.had_submission_fence()
639                                && terminal.contains_identity(failure.identity())
640                        })
641                        .collect::<Vec<_>>();
642                    if unsubmitted_recoveries.len() > 1 {
643                        return Err(invalid_event(
644                            "operation failure matches multiple unsubmitted recovery receipts",
645                        ));
646                    }
647                    TrustedExecutionEventContext::replay_failure(
648                        run_id,
649                        request_id,
650                        event
651                            .identity()
652                            .parts()
653                            .plan_id
654                            .is_some()
655                            .then_some(&topology),
656                        has_active(event.identity().parts()).then_some(active),
657                        failure,
658                        unsubmitted_recoveries.first().map(|_| failure.identity()),
659                    )
660                }
661                ExecutionEventKind::RequestFailed => match event.detail() {
662                    ExecutionEventDetail::Failure(failure) => {
663                        TrustedExecutionEventContext::failure(
664                            run_id,
665                            request_id,
666                            event
667                                .identity()
668                                .parts()
669                                .plan_id
670                                .is_some()
671                                .then_some(&topology),
672                            has_active(event.identity().parts()).then_some(active),
673                            failure,
674                        )
675                    }
676                    ExecutionEventDetail::FailureTerminal { .. } => {
677                        let failure = first_failure.as_ref().ok_or_else(|| {
678                            invalid_event(
679                                "terminal replay failure lacks its first FailureObserved evidence",
680                            )
681                        })?;
682                        TrustedExecutionEventContext::failure_with_disposition(
683                            run_id,
684                            request_id,
685                            &topology,
686                            active,
687                            has_completed(event.identity().parts())
688                                .then_some(completed)
689                                .flatten(),
690                            has_aborted(event.identity().parts())
691                                .then_some(aborted)
692                                .flatten(),
693                            failure,
694                        )
695                    }
696                    _ => unreachable!("trusted RequestFailed shape was validated"),
697                },
698                ExecutionEventKind::OperationSubmitted => {
699                    let matches = operation_terminals
700                        .iter()
701                        .filter_map(|(key, terminal)| {
702                            terminal
703                                .had_submission_fence()
704                                .then(|| {
705                                    terminal
706                                        .participant_submission(event.identity())
707                                        .map(|participant| (*key, participant))
708                                })
709                                .flatten()
710                        })
711                        .collect::<Vec<_>>();
712                    if matches.len() != 1 || !used_submitted_terminals.insert(matches[0].0) {
713                        return Err(invalid_event(
714                            "replay operation event lacks one exact terminal receipt proving submission",
715                        ));
716                    }
717                    TrustedExecutionEventContext::replay_operation_submitted(
718                        run_id,
719                        request_id,
720                        &topology,
721                        active,
722                        operation_terminals
723                            .iter()
724                            .find(|(key, _)| *key == matches[0].0)
725                            .and_then(|(_, terminal)| terminal.submission())
726                            .expect("matched submitted participant has a batch receipt"),
727                    )
728                }
729                ExecutionEventKind::NodeRetired => {
730                    let matches = operation_terminals
731                        .iter()
732                        .filter_map(|(key, terminal)| {
733                            terminal
734                                .participant_completion(event.identity())
735                                .map(|participant| (*key, participant))
736                        })
737                        .collect::<Vec<_>>();
738                    if matches.len() != 1 || !used_submitted_terminals.contains(&matches[0].0) {
739                        return Err(invalid_event(
740                            "replay NodeRetired lacks the exact submitted batch completion projection",
741                        ));
742                    }
743                    TrustedExecutionEventContext::replay_node_retired(
744                        run_id,
745                        request_id,
746                        &topology,
747                        active,
748                        matches[0].1,
749                    )
750                }
751                ExecutionEventKind::SequenceCompleted | ExecutionEventKind::RequestCompleted => {
752                    let completed = completed.ok_or_else(|| {
753                        invalid_event(
754                            "successful replay journal lacks external sequence completion evidence",
755                        )
756                    })?;
757                    TrustedExecutionEventContext::completed(
758                        run_id, request_id, &topology, active, completed,
759                    )
760                }
761                ExecutionEventKind::SequenceAborted => {
762                    let aborted = aborted.ok_or_else(|| {
763                        invalid_event(
764                            "failed replay journal lacks external sequence abort evidence",
765                        )
766                    })?;
767                    TrustedExecutionEventContext::aborted(
768                        run_id, request_id, &topology, active, aborted,
769                    )
770                }
771                _ => TrustedExecutionEventContext::active(run_id, request_id, &topology, active),
772            };
773            request_cursor.observe_against(event, &context)?;
774        }
775        if !request_cursor.is_terminal()
776            || !observed_active_identity
777            || !matches!(
778                terminal.kind(),
779                ExecutionEventKind::RequestCompleted | ExecutionEventKind::RequestFailed
780            )
781        {
782            return Err(invalid_event(
783                "replay request journal is incomplete or lacks an exact terminal event",
784            ));
785        }
786        let submitted_terminal_count = operation_terminals
787            .iter()
788            .filter(|(_, terminal)| terminal.had_submission_fence())
789            .count();
790        if used_submitted_terminals.len() != submitted_terminal_count {
791            return Err(invalid_event(
792                "replay contains unused or missing operation terminal evidence for submitted work",
793            ));
794        }
795
796        let operation_failures = first_failure
797            .iter()
798            .filter(|failure| has_active(failure.identity().parts()))
799            .collect::<Vec<_>>();
800        let submitted_identities = evidence
801            .request_journal
802            .iter()
803            .filter(|event| event.kind() == ExecutionEventKind::OperationSubmitted)
804            .map(ExecutionEvent::identity)
805            .collect::<Vec<_>>();
806        let mut used_operation_failures = BTreeSet::new();
807        let request_failed = terminal.kind() == ExecutionEventKind::RequestFailed;
808        for (_, operation_terminal) in &operation_terminals {
809            let mut relevant_identities = submitted_identities
810                .iter()
811                .copied()
812                .filter(|identity| operation_terminal.contains_identity(identity))
813                .collect::<Vec<_>>();
814            if relevant_identities.is_empty() {
815                relevant_identities.extend(
816                    operation_failures
817                        .iter()
818                        .map(|failure| failure.identity())
819                        .filter(|identity| operation_terminal.contains_identity(identity)),
820                );
821            }
822            if relevant_identities.len() != 1 {
823                return Err(invalid_event(
824                    "operation terminal evidence has no unique participant projection in this request journal",
825                ));
826            }
827            let operation_identity = relevant_identities[0];
828            let matching_failures = operation_failures
829                .iter()
830                .enumerate()
831                .filter(|(_, failure)| failure.identity() == operation_identity)
832                .collect::<Vec<_>>();
833            if operation_terminal.participant_is_success(operation_identity) {
834                if !matching_failures.is_empty() {
835                    return Err(invalid_event(
836                        "successful operation completion coexists with a failure for the same submitted operation",
837                    ));
838                }
839                continue;
840            }
841            if !request_failed || matching_failures.len() != 1 {
842                return Err(invalid_event(
843                    "non-success operation terminal evidence requires one exact operation failure and RequestFailed",
844                ));
845            }
846            if let Some(expected_failure) =
847                operation_terminal.exact_failed_completion(operation_identity)
848            {
849                if expected_failure.identity() != operation_identity
850                    || *matching_failures[0].1 != expected_failure
851                {
852                    return Err(invalid_event(
853                        "failed-but-quiescent completion differs from the exact observed operation failure",
854                    ));
855                }
856            }
857            if !used_operation_failures.insert(matching_failures[0].0) {
858                return Err(invalid_event(
859                    "one operation failure was reused by multiple terminal evidence receipts",
860                ));
861            }
862        }
863        if used_operation_failures.len() != operation_failures.len() {
864            return Err(invalid_event(
865                "operation FailureObserved lacks one exact non-success terminal evidence receipt",
866            ));
867        }
868
869        let (
870            resource_pool_id,
871            resource_pool_identity_fingerprint,
872            pool_journal_event_count,
873            pool_journal_fingerprint,
874        ) = match evidence.pool_evidence {
875            Some(pool_evidence) => {
876                if evidence.pool_journal.is_empty() {
877                    return Err(invalid_event(
878                        "static replay evidence requires a non-empty pool journal",
879                    ));
880                }
881                let active_static_pool_id = active.static_pool_id().ok_or_else(|| {
882                    invalid_event("static pool replay evidence was supplied for a no-static plan")
883                })?;
884                let active_static_fingerprint = active
885                    .static_pool_identity_fingerprint()
886                    .ok_or_else(|| invalid_event("static replay pool lacks identity evidence"))?;
887                let active_static_provisioning =
888                    active.static_provisioning_identity().ok_or_else(|| {
889                        invalid_event("static replay pool lacks provisioning identity evidence")
890                    })?;
891                if active.static_entries().is_empty()
892                    || active_static_pool_id != pool_evidence.pool_id()
893                    || active_static_fingerprint != pool_evidence.pool_identity_fingerprint()
894                    || active_static_provisioning != pool_evidence.provisioning_identity()
895                    || pool_evidence.topology_fingerprint() != topology.fingerprint()
896                {
897                    return Err(invalid_event(
898                        "replay active binding and static pool evidence do not share one authority",
899                    ));
900                }
901                let mut pool_cursor = ResourcePoolEventCursor::new(pool_evidence.clone());
902                for event in evidence.pool_journal {
903                    pool_cursor.observe(event)?;
904                }
905                if !pool_cursor.has_opened() || !pool_cursor.proves_active_binding(active) {
906                    return Err(invalid_event(
907                        "replay pool journal does not prove the complete committed active lease",
908                    ));
909                }
910                (
911                    Some(active_static_pool_id),
912                    Some(active_static_fingerprint),
913                    u64::try_from(evidence.pool_journal.len())
914                        .map_err(|_| invalid_event("pool journal length exceeds u64"))?,
915                    Some(canonical_fingerprint(&evidence.pool_journal)),
916                )
917            }
918            None => {
919                if !evidence.pool_journal.is_empty()
920                    || active.static_pool_id().is_some()
921                    || active.static_provisioning_identity().is_some()
922                    || active.plan().static_provisioning_binding().is_some()
923                    || active.plan().static_pool_identity().is_some()
924                    || !active.static_entries().is_empty()
925                {
926                    return Err(invalid_event(
927                        "no-static replay evidence contains static pool identity or journal state",
928                    ));
929                }
930                (None, None, 0, None)
931            }
932        };
933        let (cleanup_status, plan_cleanup_fingerprint) = validate_replay_plan_cleanup(
934            active,
935            aborted.is_some(),
936            evidence.cleanup_requirement,
937            evidence.plan_cleanup,
938        )?;
939
940        let request_journal_event_count = u64::try_from(evidence.request_journal.len())
941            .map_err(|_| invalid_event("request journal length exceeds u64"))?;
942        let operation_terminal_evidence_len = evidence
943            .operation_completions
944            .len()
945            .checked_add(evidence.operation_drains.len())
946            .and_then(|count| count.checked_add(evidence.operation_quarantines.len()))
947            .ok_or_else(|| invalid_event("operation terminal evidence count exceeds usize"))?;
948        let operation_terminal_evidence_count = u64::try_from(operation_terminal_evidence_len)
949            .map_err(|_| invalid_event("operation terminal evidence count exceeds u64"))?;
950        let identity = Self {
951            identity_version: EXECUTION_IDENTITY_VERSION,
952            terminal_identity: terminal.identity().clone(),
953            resolved_plan_fingerprint: evidence.resolved_plan.fingerprint().to_owned(),
954            execution_topology_fingerprint: topology.fingerprint().to_owned(),
955            request_input_fingerprint: sha256_bytes(evidence.request_input),
956            initial_state_fingerprint: sha256_bytes(evidence.initial_state),
957            random_seed: evidence.random_seed,
958            request_journal_event_count,
959            request_journal_fingerprint: canonical_fingerprint(&evidence.request_journal),
960            active_sequence_fingerprint: active.fingerprint().to_owned(),
961            completed_sequence_fingerprint: completed
962                .map(|binding| binding.fingerprint().to_owned()),
963            aborted_sequence_fingerprint: aborted.map(|binding| binding.fingerprint().to_owned()),
964            cleanup_status,
965            operation_terminal_evidence_count,
966            operation_terminal_evidence_fingerprint: canonical_fingerprint(
967                &ReplayOperationTerminalFingerprint {
968                    completions: evidence.operation_completions,
969                    drains: evidence.operation_drains,
970                    quarantines: evidence.operation_quarantines,
971                },
972            ),
973            resource_pool_id,
974            resource_pool_identity_fingerprint,
975            pool_journal_event_count,
976            pool_journal_fingerprint,
977            plan_cleanup_fingerprint,
978        };
979        identity.validate_fingerprint_shape()?;
980        Ok(identity)
981    }
982
983    fn validate_fingerprint_shape(&self) -> Result<(), VNextError> {
984        for (value, label) in [
985            (&self.resolved_plan_fingerprint, "resolved plan fingerprint"),
986            (
987                &self.execution_topology_fingerprint,
988                "execution topology fingerprint",
989            ),
990            (&self.request_input_fingerprint, "request input fingerprint"),
991            (&self.initial_state_fingerprint, "initial state fingerprint"),
992            (
993                &self.request_journal_fingerprint,
994                "request journal fingerprint",
995            ),
996            (
997                &self.active_sequence_fingerprint,
998                "active sequence fingerprint",
999            ),
1000            (
1001                &self.operation_terminal_evidence_fingerprint,
1002                "operation terminal evidence fingerprint",
1003            ),
1004        ] {
1005            validate_sha256(value, label)?;
1006        }
1007        if let Some(fingerprint) = &self.completed_sequence_fingerprint {
1008            validate_sha256(fingerprint, "completed sequence fingerprint")?;
1009        }
1010        if let Some(fingerprint) = &self.aborted_sequence_fingerprint {
1011            validate_sha256(fingerprint, "aborted sequence fingerprint")?;
1012        }
1013        if let Some(fingerprint) = &self.resource_pool_identity_fingerprint {
1014            validate_sha256(fingerprint, "resource pool identity fingerprint")?;
1015        }
1016        if let Some(fingerprint) = &self.pool_journal_fingerprint {
1017            validate_sha256(fingerprint, "pool journal fingerprint")?;
1018        }
1019        if let Some(fingerprint) = &self.plan_cleanup_fingerprint {
1020            validate_sha256(fingerprint, "plan cleanup fingerprint")?;
1021        }
1022        let has_static_pool = self.resource_pool_id.is_some();
1023        if self.completed_sequence_fingerprint.is_some()
1024            == self.aborted_sequence_fingerprint.is_some()
1025            || self.cleanup_status == ReplayCleanupStatus::Completed
1026                && self.completed_sequence_fingerprint.is_none()
1027            || self.cleanup_status == ReplayCleanupStatus::SequenceQuiescent
1028                && self.aborted_sequence_fingerprint.is_none()
1029            || self.cleanup_status == ReplayCleanupStatus::CleanupPending
1030                && self.plan_cleanup_fingerprint.is_some()
1031            || self.cleanup_status != ReplayCleanupStatus::CleanupPending
1032                && self.plan_cleanup_fingerprint.is_none()
1033            || has_static_pool != self.resource_pool_identity_fingerprint.is_some()
1034            || has_static_pool != self.pool_journal_fingerprint.is_some()
1035            || has_static_pool != (self.pool_journal_event_count > 0)
1036        {
1037            return Err(invalid_event(
1038                "replay cleanup status differs from its exact sequence disposition",
1039            ));
1040        }
1041        Ok(())
1042    }
1043
1044    pub fn terminal_identity(&self) -> &ExecutionIdentityEnvelope {
1045        &self.terminal_identity
1046    }
1047
1048    pub fn resolved_plan_fingerprint(&self) -> &str {
1049        &self.resolved_plan_fingerprint
1050    }
1051
1052    pub fn request_input_fingerprint(&self) -> &str {
1053        &self.request_input_fingerprint
1054    }
1055
1056    pub fn initial_state_fingerprint(&self) -> &str {
1057        &self.initial_state_fingerprint
1058    }
1059
1060    pub const fn random_seed(&self) -> u64 {
1061        self.random_seed
1062    }
1063
1064    pub fn request_journal_fingerprint(&self) -> &str {
1065        &self.request_journal_fingerprint
1066    }
1067
1068    pub fn pool_journal_fingerprint(&self) -> Option<&str> {
1069        self.pool_journal_fingerprint.as_deref()
1070    }
1071
1072    pub fn plan_cleanup_fingerprint(&self) -> Option<&str> {
1073        self.plan_cleanup_fingerprint.as_deref()
1074    }
1075
1076    pub const fn cleanup_status(&self) -> ReplayCleanupStatus {
1077        self.cleanup_status
1078    }
1079
1080    pub fn decode_untrusted(bytes: &[u8]) -> Result<UnvalidatedReplayIdentity, VNextError> {
1081        if bytes.len() > MAX_REPLAY_IDENTITY_WIRE_BYTES {
1082            return Err(invalid_event(
1083                "untrusted replay identity exceeds the wire byte limit",
1084            ));
1085        }
1086        serde_json::from_slice::<ReplayIdentityWire>(bytes)
1087            .map(Into::into)
1088            .map_err(|error| VNextError::Serialization {
1089                context: "decode untrusted replay identity",
1090                message: error.to_string(),
1091            })
1092    }
1093}
1094
1095impl UnvalidatedReplayIdentity {
1096    pub fn revalidate(self, evidence: &ReplayEvidence<'_>) -> Result<ReplayIdentity, VNextError> {
1097        let rebuilt = ReplayIdentity::from_evidence(evidence)?;
1098        let supplied_terminal = ExecutionIdentityEnvelope::new(self.terminal_identity.into())?;
1099        if self.identity_version != rebuilt.identity_version
1100            || supplied_terminal != rebuilt.terminal_identity
1101            || self.resolved_plan_fingerprint != rebuilt.resolved_plan_fingerprint
1102            || self.execution_topology_fingerprint != rebuilt.execution_topology_fingerprint
1103            || self.request_input_fingerprint != rebuilt.request_input_fingerprint
1104            || self.initial_state_fingerprint != rebuilt.initial_state_fingerprint
1105            || self.random_seed != rebuilt.random_seed
1106            || self.request_journal_event_count != rebuilt.request_journal_event_count
1107            || self.request_journal_fingerprint != rebuilt.request_journal_fingerprint
1108            || self.active_sequence_fingerprint != rebuilt.active_sequence_fingerprint
1109            || self.completed_sequence_fingerprint != rebuilt.completed_sequence_fingerprint
1110            || self.aborted_sequence_fingerprint != rebuilt.aborted_sequence_fingerprint
1111            || self.cleanup_status != rebuilt.cleanup_status
1112            || self.operation_terminal_evidence_count != rebuilt.operation_terminal_evidence_count
1113            || self.operation_terminal_evidence_fingerprint
1114                != rebuilt.operation_terminal_evidence_fingerprint
1115            || self.resource_pool_id != rebuilt.resource_pool_id
1116            || self.resource_pool_identity_fingerprint != rebuilt.resource_pool_identity_fingerprint
1117            || self.pool_journal_event_count != rebuilt.pool_journal_event_count
1118            || self.pool_journal_fingerprint != rebuilt.pool_journal_fingerprint
1119            || self.plan_cleanup_fingerprint != rebuilt.plan_cleanup_fingerprint
1120        {
1121            return Err(invalid_event(
1122                "serialized replay identity differs from independently rebuilt evidence",
1123            ));
1124        }
1125        rebuilt.validate_fingerprint_shape()?;
1126        Ok(rebuilt)
1127    }
1128}