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
19pub 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 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#[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#[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}