1use std::collections::{BTreeMap, VecDeque};
35use std::fmt;
36
37use serde::Serialize;
38
39use super::checkpoint::{
40 AcceptedCancellationState, AcceptedInputState, CanonicalInput, CheckpointCandidate,
41 CheckpointDraft, KernelCheckpoint, LaunchTokenState, LogicalKernelState,
42 LogicalStateProjection, ResolvedEffectState, TransitionState,
43};
44use super::command::{CancelCommand, HostCommand};
45use super::config::{ConfigDefaults, ResolvedOperationConfig, TailBounds};
46use super::effect::{Digest, EffectKind, EffectKindTag, EffectOutcome, KernelEffect, LaunchToken};
47use super::envelope::{OperationLifecycle, WireEnvelope, WireRejection, WireRejectionKind};
48use super::fault::{
49 KernelFault, KernelFaultCode, KernelPreparation, PrepareToken, PreparedTransition,
50 RejectedTransition, ReplayedTransition,
51};
52use super::record::{
53 ChainAnchor, KernelRecord, NormalizedInput, NormalizedPayload, RecordError, RecordPreparation,
54 canonical_bytes, canonical_digest, verify_record_chain,
55};
56use super::root::{ExecutionFocus, RootKind};
57use super::scalar::{EffectId, InputId, OperationId, WireU64};
58use super::terminal::{KernelTerminal, StepDisposition, TerminalSlot};
59
60pub trait TransitionStep: Serialize + Clone {
71 fn disposition(&self) -> &StepDisposition;
72
73 fn effects(&self) -> &[KernelEffect] {
74 self.disposition().effects()
75 }
76
77 fn terminal(&self) -> Option<&KernelTerminal> {
78 self.disposition().terminal()
79 }
80}
81
82pub trait RecordIndex {
93 fn record_for_input(
95 &self,
96 operation_id: &OperationId,
97 input_id: &InputId,
98 ) -> Option<KernelRecord>;
99
100 fn note_committed(&mut self, record: &KernelRecord) {
106 let _ = record;
107 }
108}
109
110#[derive(Debug, Clone, Default, PartialEq)]
112pub struct InMemoryRecordIndex {
113 records: BTreeMap<(OperationId, InputId), KernelRecord>,
114}
115
116impl InMemoryRecordIndex {
117 pub fn new() -> Self {
118 Self::default()
119 }
120
121 pub fn from_records(records: &[KernelRecord]) -> Self {
123 let mut index = Self::new();
124 for record in records {
125 index.note_committed(record);
126 }
127 index
128 }
129
130 pub fn len(&self) -> usize {
131 self.records.len()
132 }
133
134 pub fn is_empty(&self) -> bool {
135 self.records.is_empty()
136 }
137}
138
139impl RecordIndex for InMemoryRecordIndex {
140 fn record_for_input(
141 &self,
142 operation_id: &OperationId,
143 input_id: &InputId,
144 ) -> Option<KernelRecord> {
145 self.records
146 .get(&(operation_id.clone(), input_id.clone()))
147 .cloned()
148 }
149
150 fn note_committed(&mut self, record: &KernelRecord) {
151 self.records.insert(
152 (record.operation_id().clone(), record.input_id().clone()),
153 record.clone(),
154 );
155 }
156}
157
158#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
164pub struct TailUsage {
165 pub records: u64,
166 pub bytes: u64,
167}
168
169#[derive(Debug, Clone, Copy, PartialEq, Eq)]
172pub enum TailPressure {
173 Nominal,
174 Watermark,
176 Full,
178}
179
180#[derive(Debug, Clone, PartialEq)]
187struct TailEntry {
188 step_seq: WireU64,
189 record_digest: Digest,
190 bytes: u64,
191 input: NormalizedInput,
192}
193
194#[derive(Debug, Clone, PartialEq, Eq)]
200pub struct DurableHead {
201 pub digest: Digest,
202 pub step_seq: WireU64,
203}
204
205#[derive(Debug, Clone, PartialEq, Eq)]
211pub struct CheckpointBoundary {
212 pub through_step_seq: WireU64,
213 pub covered_head: Digest,
214}
215
216#[derive(Debug)]
220pub struct PlanContext<'a> {
221 pub input: &'a NormalizedInput,
222 pub step_seq: WireU64,
223 pub previous_head: Option<&'a Digest>,
224 pub config: &'a ResolvedOperationConfig,
225 pub resolving: Option<&'a KernelEffect>,
233 pub pending: &'a [KernelEffect],
238}
239
240#[derive(Debug, Clone, PartialEq)]
242pub struct CommittedTransition<Step> {
243 pub record: KernelRecord,
244 pub step: Step,
245 pub step_seq: WireU64,
246 pub checkpoint_advice: Option<CheckpointAdvice>,
256}
257
258#[derive(Debug, Clone, Copy, PartialEq, Eq)]
265pub struct CheckpointAdvice {
266 pub through_step_seq: WireU64,
267 pub usage: TailUsage,
268 pub bounds: TailBounds,
269}
270
271impl<Step: TransitionStep> CommittedTransition<Step> {
272 pub fn published_effects(&self) -> &[KernelEffect] {
276 self.step.effects()
277 }
278
279 pub fn terminal(&self) -> Option<&KernelTerminal> {
280 self.step.terminal()
281 }
282}
283
284#[derive(Debug, Clone, PartialEq)]
285struct Candidate<Step> {
286 token: PrepareToken,
287 record: KernelRecord,
288 step: Step,
289 input: NormalizedInput,
290 record_bytes: u64,
291}
292
293struct InputFacts {
299 operation_id: OperationId,
300 input_id: InputId,
301 observed_at_ms: WireU64,
302 admissible_lifecycles: &'static [OperationLifecycle],
303}
304
305#[derive(Debug, Clone, PartialEq)]
306struct ResolvedEffectRecord {
307 outcome_digest: Digest,
308 input_id: InputId,
309 step_seq: WireU64,
310}
311
312#[derive(Debug, Clone, PartialEq)]
323pub struct KernelTransaction<Step, Index> {
324 defaults: ConfigDefaults,
325 bounds: TailBounds,
326 index: Index,
327 operation_id: Option<OperationId>,
328 genesis_digest: Option<Digest>,
331 config: Option<ResolvedOperationConfig>,
332 head: Option<ChainAnchor>,
338 lifecycle: OperationLifecycle,
339 terminal: TerminalSlot,
340 candidate: Option<Candidate<Step>>,
341 prepare_epoch: u64,
342 last_observed_at_ms: WireU64,
343 steps: BTreeMap<InputId, Step>,
355 accepted: BTreeMap<InputId, AcceptedInputState>,
363 replay_floor: Option<WireU64>,
369 pending_effects: BTreeMap<EffectId, KernelEffect>,
370 resolved_effects: BTreeMap<EffectId, ResolvedEffectRecord>,
371 launch_tokens: BTreeMap<LaunchToken, WireU64>,
372 accepted_cancellation: Option<AcceptedCancellation>,
380 tail: VecDeque<TailEntry>,
381 poison: Option<KernelFault>,
382}
383
384#[derive(Debug, Clone, PartialEq)]
386struct AcceptedCancellation {
387 command_digest: Digest,
390 input_id: InputId,
391 step_seq: WireU64,
392}
393
394impl<Step, Index> KernelTransaction<Step, Index>
395where
396 Step: TransitionStep,
397 Index: RecordIndex,
398{
399 pub fn new(defaults: ConfigDefaults, index: Index) -> Self {
406 let bounds = defaults.baseline.recovery_policy.tail_bounds;
407 Self {
408 defaults,
409 bounds,
410 index,
411 operation_id: None,
412 genesis_digest: None,
413 config: None,
414 head: None,
415 lifecycle: OperationLifecycle::Created,
416 terminal: TerminalSlot::empty(),
417 candidate: None,
418 prepare_epoch: 0,
419 last_observed_at_ms: WireU64::ZERO,
420 steps: BTreeMap::new(),
421 accepted: BTreeMap::new(),
422 replay_floor: None,
423 pending_effects: BTreeMap::new(),
424 resolved_effects: BTreeMap::new(),
425 launch_tokens: BTreeMap::new(),
426 accepted_cancellation: None,
427 tail: VecDeque::new(),
428 poison: None,
429 }
430 }
431
432 pub fn prepare<F>(&mut self, envelope: &WireEnvelope, plan: F) -> RecordPreparation<Step>
443 where
444 F: FnOnce(&PlanContext<'_>) -> Result<Step, KernelFault>,
445 {
446 match self.prepare_inner(envelope, plan) {
447 Ok(preparation) => preparation,
448 Err(fault) => KernelPreparation::Rejected(RejectedTransition { fault }),
449 }
450 }
451
452 fn prepare_inner<F>(
453 &mut self,
454 envelope: &WireEnvelope,
455 plan: F,
456 ) -> Result<RecordPreparation<Step>, KernelFault>
457 where
458 F: FnOnce(&PlanContext<'_>) -> Result<Step, KernelFault>,
459 {
460 self.check_preparable(&envelope.operation_id)?;
461 let input =
462 NormalizedInput::normalize(envelope, &self.defaults).map_err(rejection_fault)?;
463 self.prepare_normalized_inner(input, plan)
464 }
465
466 fn check_preparable(&self, operation_id: &OperationId) -> Result<(), KernelFault> {
468 if let Some(fault) = &self.poison {
469 return Err(fault.clone());
470 }
471
472 if let Some(candidate) = &self.candidate {
476 return Err(KernelFault::new(
477 KernelFaultCode::TransactionConflict,
478 format!(
479 "transaction candidate {} is still outstanding at step {}; commit or abort it \
480 before preparing another input",
481 candidate.token,
482 candidate.record.step_seq()
483 ),
484 ));
485 }
486
487 if let Some(bound) = &self.operation_id
488 && bound != operation_id
489 {
490 return Err(KernelFault::new(
491 KernelFaultCode::OperationMismatch,
492 format!(
493 "input belongs to operation {operation_id}, but this runtime is bound to \
494 {bound}"
495 ),
496 ));
497 }
498 Ok(())
499 }
500
501 fn prepare_normalized_inner<F>(
502 &mut self,
503 input: NormalizedInput,
504 plan: F,
505 ) -> Result<RecordPreparation<Step>, KernelFault>
506 where
507 F: FnOnce(&PlanContext<'_>) -> Result<Step, KernelFault>,
508 {
509 self.check_preparable(&input.operation_id)?;
510 let envelope = InputFacts {
511 operation_id: input.operation_id.clone(),
512 input_id: input.input_id.clone(),
513 observed_at_ms: input.observed_at_ms,
514 admissible_lifecycles: input.input.admissible_lifecycles(),
515 };
516 let envelope = &envelope;
517 let canonical_input = canonical_bytes(&input).map_err(record_fault)?;
518
519 if let Some(existing) = self
522 .index
523 .record_for_input(&envelope.operation_id, &envelope.input_id)
524 {
525 if existing.canonical_input() != &canonical_input {
526 return Err(KernelFault::new(
527 KernelFaultCode::DuplicateInputConflict,
528 format!(
529 "input {} was already accepted at step {} with a different payload",
530 envelope.input_id,
531 existing.step_seq()
532 ),
533 ));
534 }
535 return self.replay_of(existing);
536 }
537
538 if let Some(entry) = self.accepted.get(&envelope.input_id)
544 && self.below_replay_floor(entry.step_seq)
545 {
546 return Ok(self.replay_by_reference(entry));
547 }
548
549 if let Some(config) = &self.config
550 && canonical_input.len() > config.kernel_limits.max_input_bytes as usize
551 {
552 return Err(KernelFault::new(
553 KernelFaultCode::ResourceLimitExceeded,
554 format!(
555 "canonical input carries {} bytes; the operation limit is {}",
556 canonical_input.len(),
557 config.kernel_limits.max_input_bytes
558 ),
559 ));
560 }
561
562 if envelope.observed_at_ms.get() < self.last_observed_at_ms.get() {
563 return Err(KernelFault::new(
564 KernelFaultCode::ClockRegression,
565 format!(
566 "input observed at {} precedes the last accepted input at {}",
567 envelope.observed_at_ms, self.last_observed_at_ms
568 ),
569 ));
570 }
571
572 if let NormalizedPayload::HostControl(control) = &input.input
576 && let HostCommand::Cancel(cancel) = &control.command
577 && let Some(preparation) = self.cancel_guard(cancel)?
578 {
579 return Ok(preparation);
580 }
581
582 if !envelope.admissible_lifecycles.contains(&self.lifecycle) {
583 return Err(KernelFault::new(
584 KernelFaultCode::InvalidLifecycle,
585 format!(
586 "a {} input is not admissible while the operation is {:?}",
587 input.input.kind(),
588 self.lifecycle
589 ),
590 ));
591 }
592
593 if let NormalizedPayload::ResolveEffect(resolve) = &input.input
595 && let Some(preparation) = self.resolve_effect_guard(resolve)?
596 {
597 return Ok(preparation);
598 }
599
600 if self.tail_records() + 1 > self.bounds.hard_records.get() {
601 return Err(self.checkpoint_required(format!(
602 "the journal tail already holds {} records; its hard limit is {}",
603 self.tail_records(),
604 self.bounds.hard_records
605 )));
606 }
607
608 let step_seq = self.next_step_seq()?;
609 let genesis_config = input.resolved_config().cloned();
610 let config = match (&genesis_config, &self.config) {
611 (Some(config), _) => config,
612 (None, Some(config)) => config,
613 (None, None) => {
614 return Err(KernelFault::new(
615 KernelFaultCode::InvalidLifecycle,
616 "the operation has no genesis record, so it has no configuration to plan \
617 against"
618 .to_string(),
619 ));
620 }
621 };
622
623 let settled = match &input.input {
624 NormalizedPayload::ResolveEffect(resolve) => Some(&resolve.effect_id),
625 _ => None,
626 };
627 let resolving = settled.and_then(|effect_id| self.pending_effects.get(effect_id));
628 let outstanding: Vec<KernelEffect> = self
631 .pending_effects
632 .values()
633 .filter(|effect| Some(&effect.effect_id) != settled)
634 .cloned()
635 .collect();
636 let step = plan(&PlanContext {
637 input: &input,
638 step_seq,
639 previous_head: self.head.as_ref().map(|anchor| &anchor.record_digest),
640 config,
641 resolving,
642 pending: &outstanding,
643 })?;
644
645 self.screen_planned_effects(&step, config, settled)?;
646
647 let record =
648 KernelRecord::chain_after(self.head.as_ref(), &input, &step).map_err(record_fault)?;
649 let record_bytes = record.record_bytes().len() as u64;
650 if self.tail_bytes() + record_bytes > self.bounds.hard_bytes.get() {
651 return Err(self.checkpoint_required(format!(
652 "the journal tail holds {} bytes and this record adds {record_bytes}; \
653 the hard limit is {}",
654 self.tail_bytes(),
655 self.bounds.hard_bytes
656 )));
657 }
658
659 self.prepare_epoch += 1;
661 let token = PrepareToken::new(format!(
662 "{}:prepare:{step_seq}:{}",
663 envelope.operation_id, self.prepare_epoch
664 ))
665 .expect("an operation-scoped prepare token is always a legal branded ref");
666 self.candidate = Some(Candidate {
667 token: token.clone(),
668 record: record.clone(),
669 step: step.clone(),
670 input,
671 record_bytes,
672 });
673
674 Ok(KernelPreparation::Prepared(PreparedTransition {
675 token,
676 record,
677 planned_step: step,
678 }))
679 }
680
681 pub fn commit(
692 &mut self,
693 token: &PrepareToken,
694 appended_head: &Digest,
695 ) -> Result<CommittedTransition<Step>, KernelFault> {
696 if let Some(fault) = &self.poison {
697 return Err(fault.clone());
698 }
699
700 let Some(candidate) = self.candidate.take() else {
701 return Err(self.poison_with(KernelFault::new(
702 KernelFaultCode::TransactionConflict,
703 format!(
704 "commit reports a durable append for token {token}, but no candidate is \
705 outstanding; a committed record is never re-committed and never aborted"
706 ),
707 )));
708 };
709
710 if &candidate.token != token {
711 let outstanding = candidate.token.clone();
712 return Err(self.poison_with(KernelFault::new(
713 KernelFaultCode::TransactionConflict,
714 format!(
715 "commit names token {token}, but the outstanding candidate is {outstanding}; \
716 the runtime no longer describes what the journal holds"
717 ),
718 )));
719 }
720
721 if appended_head != candidate.record.record_digest() {
722 let expected = candidate.record.record_digest().clone();
723 return Err(self.poison_with(KernelFault::new(
724 KernelFaultCode::TransactionConflict,
725 format!(
726 "the journal head after the append is {appended_head}, but this candidate is \
727 {expected}; the append did not place this record, so the runtime must be \
728 rebuilt from the journal"
729 ),
730 )));
731 }
732
733 let Candidate {
734 record,
735 step,
736 input,
737 record_bytes,
738 ..
739 } = candidate;
740 match self.integrate(record, step, &input, record_bytes) {
741 Ok(committed) => Ok(committed),
742 Err(fault) => Err(self.poison_with(fault)),
743 }
744 }
745
746 pub fn abort(&mut self, token: &PrepareToken) -> Result<KernelRecord, KernelFault> {
754 if let Some(fault) = &self.poison {
755 return Err(fault.clone());
756 }
757
758 let Some(candidate) = &self.candidate else {
759 return Err(KernelFault::new(
760 KernelFaultCode::TransactionConflict,
761 format!(
762 "no candidate is outstanding for token {token}; a record that reached the \
763 journal is never abortable"
764 ),
765 ));
766 };
767 if &candidate.token != token {
768 return Err(KernelFault::new(
769 KernelFaultCode::TransactionConflict,
770 format!(
771 "token {token} does not name the outstanding candidate {}",
772 candidate.token
773 ),
774 ));
775 }
776
777 let candidate = self.candidate.take().expect("checked just above");
778 Ok(candidate.record)
779 }
780
781 pub fn note_append_conflict(
788 &mut self,
789 token: &PrepareToken,
790 observed_head: Option<&Digest>,
791 ) -> KernelFault {
792 let expected = self
793 .candidate
794 .as_ref()
795 .and_then(|candidate| candidate.record.expected_head().cloned());
796 self.candidate = None;
797 let observed =
798 observed_head.map_or_else(|| "an empty journal".to_string(), Digest::to_string);
799 let expected =
800 expected.map_or_else(|| "an empty journal".to_string(), |head| head.to_string());
801 let fault = KernelFault::new(
802 KernelFaultCode::TransactionConflict,
803 format!(
804 "the CAS append for token {token} expected head {expected} but the journal holds \
805 {observed}; the candidate is discarded and this runtime must be rebuilt from the \
806 journal before the input is replayed"
807 ),
808 );
809 self.poison_with(fault)
810 }
811
812 pub fn rebuild_from_records<F>(
822 records: &[KernelRecord],
823 defaults: ConfigDefaults,
824 index: Index,
825 mut plan: F,
826 ) -> Result<Self, KernelFault>
827 where
828 F: FnMut(&PlanContext<'_>) -> Result<Step, KernelFault>,
829 {
830 let mut transaction = Self::new(defaults, index);
831 if records.is_empty() {
832 return Ok(transaction);
833 }
834 verify_record_chain(records).map_err(corrupt_chain_fault)?;
835
836 for record in records {
837 let input = record.normalized_input().map_err(record_fault)?;
838 transaction.replay_committed(&input, record.record_digest(), &mut plan)?;
839 }
840 Ok(transaction)
841 }
842
843 pub fn replay_committed<F>(
857 &mut self,
858 input: &NormalizedInput,
859 expected_record_digest: &Digest,
860 plan: &mut F,
861 ) -> Result<KernelRecord, KernelFault>
862 where
863 F: FnMut(&PlanContext<'_>) -> Result<Step, KernelFault>,
864 {
865 if let Some(fault) = &self.poison {
866 return Err(fault.clone());
867 }
868 let step_seq = self.next_step_seq()?;
869 let genesis_config = input.resolved_config().cloned();
870 let config = match (&genesis_config, &self.config) {
871 (Some(config), _) => config,
872 (None, Some(config)) => config,
873 (None, None) => {
874 return Err(KernelFault::new(
875 KernelFaultCode::RecordCorrupted,
876 format!(
877 "the transition at step {step_seq} has no genesis configuration before it"
878 ),
879 ));
880 }
881 };
882
883 let settled = match &input.input {
887 NormalizedPayload::ResolveEffect(resolve) => Some(&resolve.effect_id),
888 _ => None,
889 };
890 let resolving = settled.and_then(|effect_id| self.pending_effects.get(effect_id));
891 let outstanding: Vec<KernelEffect> = self
892 .pending_effects
893 .values()
894 .filter(|effect| Some(&effect.effect_id) != settled)
895 .cloned()
896 .collect();
897 let step = plan(&PlanContext {
898 input,
899 step_seq,
900 previous_head: self.head.as_ref().map(|anchor| &anchor.record_digest),
901 config,
902 resolving,
903 pending: &outstanding,
904 })?;
905
906 let rebuilt = KernelRecord::chain_after(self.head.as_ref(), input, &step)
907 .map_err(corrupt_chain_fault)?;
908 if rebuilt.record_digest() != expected_record_digest {
909 return Err(KernelFault::new(
910 KernelFaultCode::RecordCorrupted,
911 format!(
912 "replaying the transition at step {step_seq} produced record digest {} against \
913 the durable {expected_record_digest}; this binary does not reproduce the \
914 history it is resuming",
915 rebuilt.record_digest(),
916 ),
917 ));
918 }
919
920 let bytes = rebuilt.record_bytes().len() as u64;
921 self.integrate(rebuilt.clone(), step, input, bytes)?;
922 Ok(rebuilt)
923 }
924
925 pub fn restore_from_checkpoint(
940 checkpoint: &KernelCheckpoint,
941 defaults: ConfigDefaults,
942 index: Index,
943 ) -> Result<Self, KernelFault> {
944 let state = checkpoint.logical_state();
945 let transition = &state.transition;
946 let mut transaction = Self::new(defaults, index);
947
948 transaction.operation_id = Some(checkpoint.operation_id().clone());
949 transaction.genesis_digest = Some(checkpoint.genesis_digest().clone());
950 transaction.bounds = transition.resolved_config.recovery_policy.tail_bounds;
951 transaction.config = Some(transition.resolved_config.clone());
952 transaction.head = Some(ChainAnchor {
953 operation_id: checkpoint.operation_id().clone(),
954 step_seq: checkpoint.base_step_seq(),
955 record_digest: checkpoint.base_record_digest().clone(),
956 });
957 transaction.lifecycle = transition.lifecycle;
958 transaction.last_observed_at_ms = transition.last_observed_at_ms;
959 transaction.replay_floor = Some(checkpoint.base_step_seq());
960
961 if let Some(terminal) = &transition.terminal {
962 transaction
963 .terminal
964 .commit(terminal.clone())
965 .map_err(|error| {
966 KernelFault::new(KernelFaultCode::CheckpointCorrupted, error.to_string())
967 })?;
968 }
969 transaction.pending_effects = transition
970 .pending_effects
971 .iter()
972 .map(|effect| (effect.effect_id.clone(), effect.clone()))
973 .collect();
974 transaction.resolved_effects = transition
975 .resolved_effects
976 .iter()
977 .map(|resolved| {
978 (
979 resolved.effect_id.clone(),
980 ResolvedEffectRecord {
981 outcome_digest: resolved.outcome_digest.clone(),
982 input_id: resolved.input_id.clone(),
983 step_seq: resolved.step_seq,
984 },
985 )
986 })
987 .collect();
988 transaction.launch_tokens = transition
989 .launch_tokens
990 .iter()
991 .map(|token| (token.launch_token.clone(), token.step_seq))
992 .collect();
993 transaction.accepted = transition
994 .accepted_inputs
995 .iter()
996 .map(|entry| (entry.input_id.clone(), entry.clone()))
997 .collect();
998 transaction.accepted_cancellation =
999 transition
1000 .accepted_cancellation
1001 .as_ref()
1002 .map(|cancellation| AcceptedCancellation {
1003 command_digest: cancellation.command_digest.clone(),
1004 input_id: cancellation.input_id.clone(),
1005 step_seq: cancellation.step_seq,
1006 });
1007 Ok(transaction)
1008 }
1009
1010 pub fn checkpoint_boundary(&self) -> Option<CheckpointBoundary> {
1017 self.head.as_ref().map(|head| CheckpointBoundary {
1018 through_step_seq: head.step_seq,
1019 covered_head: head.record_digest.clone(),
1020 })
1021 }
1022
1023 pub fn tail_inputs(&self) -> Vec<CanonicalInput> {
1029 self.tail
1030 .iter()
1031 .map(|entry| CanonicalInput {
1032 step_seq: entry.step_seq,
1033 record_digest: entry.record_digest.clone(),
1034 input: entry.input.clone(),
1035 })
1036 .collect()
1037 }
1038
1039 pub fn checkpoint_candidate(
1050 &self,
1051 projection: LogicalStateProjection,
1052 ) -> Result<CheckpointCandidate, KernelFault> {
1053 let head = self.require_head()?.clone();
1054 self.assemble_checkpoint(
1055 projection,
1056 head.step_seq,
1057 head.record_digest.clone(),
1058 head.step_seq,
1059 head.record_digest,
1060 Vec::new(),
1061 )
1062 }
1063
1064 pub fn checkpoint_rebase(
1078 &self,
1079 base: &CheckpointBoundary,
1080 base_state: LogicalKernelState,
1081 ) -> Result<CheckpointCandidate, KernelFault> {
1082 let head = self.require_head()?.clone();
1083 if base.through_step_seq > head.step_seq {
1084 return Err(KernelFault::new(
1085 KernelFaultCode::CheckpointIncompatible,
1086 format!(
1087 "a rebase bases at step {} but this journal's head is step {}",
1088 base.through_step_seq, head.step_seq
1089 ),
1090 ));
1091 }
1092 let tail: Vec<CanonicalInput> = self
1093 .tail_inputs()
1094 .into_iter()
1095 .filter(|entry| entry.step_seq.get() > base.through_step_seq.get())
1096 .collect();
1097 if tail.len() as u64 != head.step_seq.get() - base.through_step_seq.get() {
1098 return Err(KernelFault::new(
1099 KernelFaultCode::CheckpointIncompatible,
1100 format!(
1101 "a rebase over ({}, {}] needs {} tail inputs, but this runtime's tail holds \
1102 {} of them — the prefix it would rebase onto was already reclaimed",
1103 base.through_step_seq,
1104 head.step_seq,
1105 head.step_seq.get() - base.through_step_seq.get(),
1106 tail.len()
1107 ),
1108 ));
1109 }
1110
1111 let operation_id = self.require_operation()?;
1112 let genesis_digest = self.require_genesis()?;
1113 let checkpoint = KernelCheckpoint::assemble(CheckpointDraft {
1114 operation_id: operation_id.clone(),
1115 genesis_digest: genesis_digest.clone(),
1116 base_step_seq: base.through_step_seq,
1117 base_record_digest: base.covered_head.clone(),
1118 through_step_seq: head.step_seq,
1119 covered_transaction_head_digest: head.record_digest,
1120 logical_state: base_state,
1121 tail_inputs: tail,
1122 })
1123 .map_err(|error| error.fault())?;
1124 Ok(checkpoint.into_candidate())
1125 }
1126
1127 #[allow(clippy::too_many_arguments)]
1128 fn assemble_checkpoint(
1129 &self,
1130 projection: LogicalStateProjection,
1131 base_step_seq: WireU64,
1132 base_record_digest: Digest,
1133 through_step_seq: WireU64,
1134 covered_transaction_head_digest: Digest,
1135 tail_inputs: Vec<CanonicalInput>,
1136 ) -> Result<CheckpointCandidate, KernelFault> {
1137 let operation_id = self.require_operation()?.clone();
1138 let genesis_digest = self.require_genesis()?.clone();
1139 let LogicalStateProjection {
1140 root_kind,
1141 focus,
1142 syscall,
1143 scheduler,
1144 context_vm,
1145 } = projection;
1146 let logical_state = LogicalKernelState {
1147 transition: self.transition_state(root_kind, focus)?,
1148 syscall,
1149 scheduler,
1150 context_vm,
1151 };
1152 let checkpoint = KernelCheckpoint::assemble(CheckpointDraft {
1153 operation_id,
1154 genesis_digest,
1155 base_step_seq,
1156 base_record_digest,
1157 through_step_seq,
1158 covered_transaction_head_digest,
1159 logical_state,
1160 tail_inputs,
1161 })
1162 .map_err(|error| error.fault())?;
1163 Ok(checkpoint.into_candidate())
1164 }
1165
1166 fn require_head(&self) -> Result<&ChainAnchor, KernelFault> {
1167 if let Some(fault) = &self.poison {
1168 return Err(fault.clone());
1169 }
1170 self.head.as_ref().ok_or_else(|| {
1171 KernelFault::new(
1172 KernelFaultCode::InvalidLifecycle,
1173 "an operation with no genesis record has no logical state to checkpoint"
1174 .to_string(),
1175 )
1176 })
1177 }
1178
1179 fn require_operation(&self) -> Result<&OperationId, KernelFault> {
1180 self.operation_id.as_ref().ok_or_else(|| {
1181 KernelFault::new(
1182 KernelFaultCode::InvalidLifecycle,
1183 "an unbound operation has no logical state to checkpoint".to_string(),
1184 )
1185 })
1186 }
1187
1188 fn require_genesis(&self) -> Result<&Digest, KernelFault> {
1189 self.genesis_digest.as_ref().ok_or_else(|| {
1190 KernelFault::new(
1191 KernelFaultCode::InvalidLifecycle,
1192 "an operation with no genesis record has no identity to bind a checkpoint to"
1193 .to_string(),
1194 )
1195 })
1196 }
1197
1198 pub fn transition_state_for_restore(
1204 &self,
1205 root_kind: Option<RootKind>,
1206 focus: Option<ExecutionFocus>,
1207 ) -> Result<TransitionState, KernelFault> {
1208 self.transition_state(root_kind, focus)
1209 }
1210
1211 fn transition_state(
1217 &self,
1218 root_kind: Option<RootKind>,
1219 focus: Option<ExecutionFocus>,
1220 ) -> Result<TransitionState, KernelFault> {
1221 let config = self.config.as_ref().ok_or_else(|| {
1222 KernelFault::new(
1223 KernelFaultCode::InvalidLifecycle,
1224 "an operation with no genesis record has no resolved configuration to checkpoint"
1225 .to_string(),
1226 )
1227 })?;
1228
1229 Ok(TransitionState {
1230 lifecycle: self.lifecycle,
1231 resolved_config: config.clone(),
1232 root_kind,
1233 focus,
1234 last_observed_at_ms: self.last_observed_at_ms,
1235 pending_effects: self.pending_effects.values().cloned().collect(),
1236 resolved_effects: self
1237 .resolved_effects
1238 .iter()
1239 .map(|(effect_id, resolved)| ResolvedEffectState {
1240 effect_id: effect_id.clone(),
1241 outcome_digest: resolved.outcome_digest.clone(),
1242 input_id: resolved.input_id.clone(),
1243 step_seq: resolved.step_seq,
1244 })
1245 .collect(),
1246 launch_tokens: self
1247 .launch_tokens
1248 .iter()
1249 .map(|(launch_token, step_seq)| LaunchTokenState {
1250 launch_token: launch_token.clone(),
1251 step_seq: *step_seq,
1252 })
1253 .collect(),
1254 accepted_inputs: self.accepted.values().cloned().collect(),
1255 accepted_cancellation: self.accepted_cancellation.as_ref().map(|cancellation| {
1256 AcceptedCancellationState {
1257 command_digest: cancellation.command_digest.clone(),
1258 input_id: cancellation.input_id.clone(),
1259 step_seq: cancellation.step_seq,
1260 }
1261 }),
1262 terminal: self.terminal.get().cloned(),
1263 })
1264 }
1265
1266 pub fn note_checkpoint_acked(
1276 &mut self,
1277 boundary: &CheckpointBoundary,
1278 ) -> Result<TailUsage, KernelFault> {
1279 if let Some(fault) = &self.poison {
1280 return Err(fault.clone());
1281 }
1282 let matches = self.tail.iter().any(|entry| {
1283 entry.step_seq == boundary.through_step_seq
1284 && entry.record_digest == boundary.covered_head
1285 });
1286 if !matches {
1287 return Err(KernelFault::new(
1288 KernelFaultCode::CheckpointIncompatible,
1289 format!(
1290 "no tail record at step {} has digest {}; this checkpoint does not cover a \
1291 prefix of this journal",
1292 boundary.through_step_seq, boundary.covered_head
1293 ),
1294 ));
1295 }
1296 while let Some(entry) = self.tail.front() {
1297 if entry.step_seq.get() <= boundary.through_step_seq.get() {
1298 self.tail.pop_front();
1299 } else {
1300 break;
1301 }
1302 }
1303 Ok(self.tail_usage())
1304 }
1305
1306 pub fn operation_id(&self) -> Option<&OperationId> {
1309 self.operation_id.as_ref()
1310 }
1311
1312 pub fn config(&self) -> Option<&ResolvedOperationConfig> {
1313 self.config.as_ref()
1314 }
1315
1316 pub fn head(&self) -> Option<DurableHead> {
1317 self.head.as_ref().map(|anchor| DurableHead {
1318 digest: anchor.record_digest.clone(),
1319 step_seq: anchor.step_seq,
1320 })
1321 }
1322
1323 pub fn lifecycle(&self) -> OperationLifecycle {
1324 self.lifecycle
1325 }
1326
1327 pub fn terminal(&self) -> Option<&KernelTerminal> {
1328 self.terminal.get()
1329 }
1330
1331 pub fn pending_effects(&self) -> impl Iterator<Item = &KernelEffect> {
1338 self.pending_effects.values()
1339 }
1340
1341 pub fn pending_effects_in_order(&self) -> Vec<&KernelEffect> {
1345 let mut effects: Vec<&KernelEffect> = self.pending_effects.values().collect();
1346 effects.sort_by_key(|effect| effect_position(effect.effect_id.as_str()));
1347 effects
1348 }
1349
1350 pub fn is_effect_resolved(&self, effect_id: &EffectId) -> bool {
1351 self.resolved_effects.contains_key(effect_id)
1352 }
1353
1354 pub fn knows_launch_token(&self, token: &LaunchToken) -> bool {
1355 self.launch_tokens.contains_key(token)
1356 }
1357
1358 pub fn committed_step(&self, input_id: &InputId) -> Option<&Step> {
1359 self.steps.get(input_id)
1360 }
1361
1362 pub fn outstanding_token(&self) -> Option<&PrepareToken> {
1363 self.candidate.as_ref().map(|candidate| &candidate.token)
1364 }
1365
1366 pub fn has_candidate(&self) -> bool {
1367 self.candidate.is_some()
1368 }
1369
1370 pub fn poison(&self) -> Option<&KernelFault> {
1373 self.poison.as_ref()
1374 }
1375
1376 pub fn is_poisoned(&self) -> bool {
1377 self.poison.is_some()
1378 }
1379
1380 pub fn bounds(&self) -> TailBounds {
1381 self.bounds
1382 }
1383
1384 pub fn index(&self) -> &Index {
1385 &self.index
1386 }
1387
1388 pub fn tail_usage(&self) -> TailUsage {
1389 TailUsage {
1390 records: self.tail_records(),
1391 bytes: self.tail_bytes(),
1392 }
1393 }
1394
1395 pub fn tail_pressure(&self) -> TailPressure {
1396 let usage = self.tail_usage();
1397 if usage.records >= self.bounds.hard_records.get()
1398 || usage.bytes >= self.bounds.hard_bytes.get()
1399 {
1400 TailPressure::Full
1401 } else if usage.records >= self.bounds.soft_records.get()
1402 || usage.bytes >= self.bounds.soft_bytes.get()
1403 {
1404 TailPressure::Watermark
1405 } else {
1406 TailPressure::Nominal
1407 }
1408 }
1409
1410 fn tail_records(&self) -> u64 {
1413 self.tail.len() as u64
1414 }
1415
1416 fn tail_bytes(&self) -> u64 {
1417 self.tail.iter().map(|entry| entry.bytes).sum()
1418 }
1419
1420 fn next_step_seq(&self) -> Result<WireU64, KernelFault> {
1421 match &self.head {
1422 None => Ok(WireU64::ZERO),
1423 Some(head) => head
1424 .step_seq
1425 .get()
1426 .checked_add(1)
1427 .map(WireU64::new)
1428 .ok_or_else(|| {
1429 KernelFault::new(
1430 KernelFaultCode::ResourceLimitExceeded,
1431 "step sequence overflowed u64".to_string(),
1432 )
1433 }),
1434 }
1435 }
1436
1437 fn checkpoint_required(&self, detail: String) -> KernelFault {
1438 KernelFault::new(
1439 KernelFaultCode::CheckpointRequired,
1440 format!(
1441 "{detail}; take a checkpoint candidate, install and ack it, then retry this input \
1442 unchanged — it was never accepted"
1443 ),
1444 )
1445 }
1446
1447 fn poison_with(&mut self, fault: KernelFault) -> KernelFault {
1448 self.candidate = None;
1449 self.poison.get_or_insert(fault).clone()
1450 }
1451
1452 fn replay_of(&self, existing: KernelRecord) -> Result<RecordPreparation<Step>, KernelFault> {
1460 let step_seq = existing.step_seq();
1461 let record_digest = existing.record_digest().clone();
1462 match self.steps.get(existing.input_id()) {
1463 Some(step) => {
1464 existing.verify_step(step).map_err(record_fault)?;
1465 Ok(KernelPreparation::Replayed(ReplayedTransition {
1466 record: Some(existing),
1467 record_digest,
1468 committed_step: Some(step.clone()),
1469 step_seq,
1470 }))
1471 }
1472 None if self.below_replay_floor(step_seq) => {
1473 Ok(KernelPreparation::Replayed(ReplayedTransition {
1474 record: Some(existing),
1475 record_digest,
1476 committed_step: None,
1477 step_seq,
1478 }))
1479 }
1480 None => Err(KernelFault::new(
1481 KernelFaultCode::RecordCorrupted,
1482 format!(
1483 "the journal holds record {record_digest} at step {step_seq} for input {}, but \
1484 this runtime has not replayed it and cannot reproduce its step; rebuild from \
1485 the journal first",
1486 existing.input_id()
1487 ),
1488 )),
1489 }
1490 }
1491
1492 fn replay_by_reference(&self, entry: &AcceptedInputState) -> RecordPreparation<Step> {
1495 KernelPreparation::Replayed(ReplayedTransition {
1496 record: None,
1497 record_digest: entry.record_digest.clone(),
1498 committed_step: None,
1499 step_seq: entry.step_seq,
1500 })
1501 }
1502
1503 fn below_replay_floor(&self, step_seq: WireU64) -> bool {
1504 self.replay_floor
1505 .is_some_and(|floor| step_seq.get() <= floor.get())
1506 }
1507
1508 fn resolve_effect_guard(
1513 &self,
1514 resolve: &super::envelope::ResolveEffect,
1515 ) -> Result<Option<RecordPreparation<Step>>, KernelFault> {
1517 let outcome_digest = outcome_digest(&resolve.outcome)?;
1518 if let Some(resolved) = self.resolved_effects.get(&resolve.effect_id) {
1519 if resolved.outcome_digest != outcome_digest {
1520 return Err(KernelFault::new(
1521 KernelFaultCode::UnexpectedEffectOutcome,
1522 format!(
1523 "effect {} was already resolved at step {} with a different outcome",
1524 resolve.effect_id, resolved.step_seq
1525 ),
1526 ));
1527 }
1528 let Some(existing) = self.index.record_for_input(
1531 self.operation_id.as_ref().expect("bound by the genesis"),
1532 &resolved.input_id,
1533 ) else {
1534 return Err(KernelFault::new(
1535 KernelFaultCode::RecordCorrupted,
1536 format!(
1537 "effect {} is resolved by input {} at step {}, but the journal has no such \
1538 record",
1539 resolve.effect_id, resolved.input_id, resolved.step_seq
1540 ),
1541 ));
1542 };
1543 return self.replay_of(existing).map(Some);
1544 }
1545
1546 let Some(pending) = self.pending_effects.get(&resolve.effect_id) else {
1547 return Err(KernelFault::new(
1548 KernelFaultCode::UnexpectedEffectOutcome,
1549 format!(
1550 "effect {} is not pending; the kernel is not waiting on it",
1551 resolve.effect_id
1552 ),
1553 ));
1554 };
1555 pending.accept_outcome(&resolve.outcome).map_err(|error| {
1556 KernelFault::new(KernelFaultCode::UnexpectedEffectOutcome, error.to_string())
1557 })?;
1558 Ok(None)
1559 }
1560
1561 fn cancel_guard(
1569 &self,
1570 cancel: &CancelCommand,
1571 ) -> Result<Option<RecordPreparation<Step>>, KernelFault> {
1572 let Some(accepted) = &self.accepted_cancellation else {
1573 return Ok(None);
1574 };
1575 if accepted.command_digest != cancel_digest(cancel)? {
1576 return Err(KernelFault::new(
1577 KernelFaultCode::DuplicateInputConflict,
1578 format!(
1579 "this operation committed a different cancellation at step {}; a cancel is not \
1580 re-decided once the terminal it produced exists",
1581 accepted.step_seq
1582 ),
1583 ));
1584 }
1585 let Some(existing) = self.index.record_for_input(
1586 self.operation_id.as_ref().expect("bound by the genesis"),
1587 &accepted.input_id,
1588 ) else {
1589 return Err(KernelFault::new(
1590 KernelFaultCode::RecordCorrupted,
1591 format!(
1592 "this operation was cancelled by input {} at step {}, but the journal has no \
1593 such record",
1594 accepted.input_id, accepted.step_seq
1595 ),
1596 ));
1597 };
1598 self.replay_of(existing).map(Some)
1599 }
1600
1601 fn screen_planned_effects(
1604 &self,
1605 step: &Step,
1606 config: &ResolvedOperationConfig,
1607 settled: Option<&EffectId>,
1608 ) -> Result<(), KernelFault> {
1609 let mut kinds_in_step: Vec<EffectKindTag> = Vec::new();
1610 let mut tokens_in_step: Vec<&LaunchToken> = Vec::new();
1611 let still_pending = |effect_id: &EffectId| {
1615 self.pending_effects.contains_key(effect_id) && Some(effect_id) != settled
1616 };
1617
1618 for effect in step.effects() {
1619 let tag = effect.tag();
1620
1621 if !config.host_effect_support.supports(tag) {
1623 return Err(KernelFault::new(
1624 KernelFaultCode::UnsupportedEffect,
1625 format!(
1626 "this operation's host does not declare support for {tag} effects, so \
1627 effect {} is refused before emission",
1628 effect.effect_id
1629 ),
1630 ));
1631 }
1632
1633 if self.pending_effects.contains_key(&effect.effect_id)
1636 || self.resolved_effects.contains_key(&effect.effect_id)
1637 || Some(&effect.effect_id) == settled
1638 {
1639 return Err(KernelFault::new(
1640 KernelFaultCode::TransactionConflict,
1641 format!(
1642 "effect id {} was already published; effect identity is minted once",
1643 effect.effect_id
1644 ),
1645 ));
1646 }
1647
1648 if kinds_in_step.contains(&tag)
1650 || self
1651 .pending_effects
1652 .iter()
1653 .any(|(id, pending)| pending.tag() == tag && still_pending(id))
1654 {
1655 return Err(KernelFault::new(
1656 KernelFaultCode::ResourceLimitExceeded,
1657 format!(
1658 "a {tag} effect is already pending; resolve it before emitting another \
1659 (§15.3 admits at most one pending effect per kind)"
1660 ),
1661 ));
1662 }
1663 kinds_in_step.push(tag);
1664
1665 if let EffectKind::SpawnTasks(spawn) = &effect.effect {
1666 for launch in &spawn.tasks {
1667 if self.launch_tokens.contains_key(&launch.launch_token)
1668 || tokens_in_step.contains(&&launch.launch_token)
1669 {
1670 return Err(KernelFault::new(
1671 KernelFaultCode::TransactionConflict,
1672 format!(
1673 "launch token {} was already published; a re-launch reuses the \
1674 committed token so the host's launch dedup stays exact",
1675 launch.launch_token
1676 ),
1677 ));
1678 }
1679 tokens_in_step.push(&launch.launch_token);
1680 }
1681 }
1682 }
1683 Ok(())
1684 }
1685
1686 fn integrate(
1689 &mut self,
1690 record: KernelRecord,
1691 step: Step,
1692 input: &NormalizedInput,
1693 record_bytes: u64,
1694 ) -> Result<CommittedTransition<Step>, KernelFault> {
1695 let step_seq = record.step_seq();
1696 let was_nominal = self.tail_pressure() == TailPressure::Nominal;
1697
1698 if self.operation_id.is_none() {
1699 self.operation_id = Some(record.operation_id().clone());
1700 }
1701 if record.is_genesis() {
1702 self.genesis_digest = Some(record.record_digest().clone());
1703 }
1704 if let Some(config) = input.resolved_config() {
1705 self.bounds = config.recovery_policy.tail_bounds;
1710 self.config = Some(config.clone());
1711 self.lifecycle = OperationLifecycle::Configured;
1712 } else if matches!(input.input, NormalizedPayload::StartOperation(_)) {
1713 self.lifecycle = OperationLifecycle::Running;
1714 }
1715
1716 if let NormalizedPayload::HostControl(control) = &input.input
1717 && let HostCommand::Cancel(cancel) = &control.command
1718 {
1719 self.accepted_cancellation = Some(AcceptedCancellation {
1720 command_digest: cancel_digest(cancel)?,
1721 input_id: record.input_id().clone(),
1722 step_seq,
1723 });
1724 }
1725
1726 if let NormalizedPayload::ResolveEffect(resolve) = &input.input {
1727 self.pending_effects.remove(&resolve.effect_id);
1728 self.resolved_effects.insert(
1729 resolve.effect_id.clone(),
1730 ResolvedEffectRecord {
1731 outcome_digest: outcome_digest(&resolve.outcome)?,
1732 input_id: record.input_id().clone(),
1733 step_seq,
1734 },
1735 );
1736 }
1737
1738 for effect in step.effects() {
1739 self.pending_effects
1740 .insert(effect.effect_id.clone(), effect.clone());
1741 if let EffectKind::SpawnTasks(spawn) = &effect.effect {
1742 for launch in &spawn.tasks {
1743 self.launch_tokens
1744 .insert(launch.launch_token.clone(), step_seq);
1745 }
1746 }
1747 }
1748
1749 if let Some(terminal) = step.terminal() {
1750 self.terminal.commit(terminal.clone()).map_err(|error| {
1751 KernelFault::new(KernelFaultCode::InvalidLifecycle, error.to_string())
1752 })?;
1753 self.lifecycle = terminal_lifecycle(terminal);
1754 self.pending_effects.clear();
1756 }
1757
1758 self.last_observed_at_ms = input.observed_at_ms;
1759 self.steps.insert(record.input_id().clone(), step.clone());
1760 self.accepted.insert(
1761 record.input_id().clone(),
1762 AcceptedInputState {
1763 input_id: record.input_id().clone(),
1764 step_seq,
1765 record_digest: record.record_digest().clone(),
1766 },
1767 );
1768 self.tail.push_back(TailEntry {
1769 step_seq,
1770 record_digest: record.record_digest().clone(),
1771 bytes: record_bytes,
1772 input: input.clone(),
1773 });
1774 self.index.note_committed(&record);
1775 self.head = Some(record.anchor());
1776
1777 let checkpoint_advice = (was_nominal && self.tail_pressure() != TailPressure::Nominal)
1780 .then(|| CheckpointAdvice {
1781 through_step_seq: step_seq,
1782 usage: self.tail_usage(),
1783 bounds: self.bounds,
1784 });
1785
1786 Ok(CommittedTransition {
1787 record,
1788 step,
1789 step_seq,
1790 checkpoint_advice,
1791 })
1792 }
1793}
1794
1795fn outcome_digest(outcome: &EffectOutcome) -> Result<Digest, KernelFault> {
1800 canonical_bytes(outcome)
1801 .map(|bytes| canonical_digest(bytes.as_slice()))
1802 .map_err(record_fault)
1803}
1804
1805fn effect_position(id: &str) -> (u64, u64, &str) {
1809 let step = id
1810 .rsplit_once(":step:")
1811 .and_then(|(_, rest)| rest.split_once(':'))
1812 .and_then(|(digits, _)| digits.parse::<u64>().ok())
1813 .unwrap_or(u64::MAX);
1814 let index = id
1815 .rsplit_once(":effect:")
1816 .and_then(|(_, digits)| digits.parse::<u64>().ok())
1817 .unwrap_or(u64::MAX);
1818 (step, index, id)
1819}
1820
1821fn cancel_digest(cancel: &CancelCommand) -> Result<Digest, KernelFault> {
1822 canonical_bytes(cancel)
1823 .map(|bytes| canonical_digest(bytes.as_slice()))
1824 .map_err(record_fault)
1825}
1826
1827fn record_fault(error: RecordError) -> KernelFault {
1828 KernelFault::new(error.code(), error.message().to_string())
1829}
1830
1831fn corrupt_chain_fault(error: RecordError) -> KernelFault {
1834 match error {
1835 RecordError::ChainBroken(message) => {
1836 KernelFault::new(KernelFaultCode::RecordCorrupted, message)
1837 }
1838 other => record_fault(other),
1839 }
1840}
1841
1842fn rejection_fault(rejection: WireRejection) -> KernelFault {
1843 let code = match rejection.kind {
1844 WireRejectionKind::PolicyViolation => KernelFaultCode::InvalidConfig,
1845 _ => KernelFaultCode::MalformedEnvelope,
1846 };
1847 KernelFault::new(code, rejection.message)
1848}
1849
1850fn terminal_lifecycle(terminal: &KernelTerminal) -> OperationLifecycle {
1851 match terminal {
1852 KernelTerminal::Agent(_) | KernelTerminal::Workflow(_) => OperationLifecycle::Completed,
1853 KernelTerminal::Cancelled(_) => OperationLifecycle::Cancelled,
1854 KernelTerminal::Failed(_) => OperationLifecycle::Failed,
1855 }
1856}
1857
1858impl fmt::Display for CheckpointBoundary {
1859 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1860 write!(
1861 f,
1862 "through step {} at head {}",
1863 self.through_step_seq, self.covered_head
1864 )
1865 }
1866}
1867
1868#[cfg(test)]
1869mod tests {
1870 use serde::Serialize;
1871
1872 use super::super::*;
1873
1874 #[test]
1875 fn effect_position_parses_numeric_step_order_not_lexicographic() {
1876 let mut ids = [
1878 "op-1:step:10:effect:0",
1879 "op-1:step:9:effect:2",
1880 "op-1:step:9:effect:10",
1881 "op-1:step:9:effect:1",
1882 ];
1883 ids.sort_by_key(|id| super::effect_position(id));
1884 assert_eq!(
1885 ids,
1886 [
1887 "op-1:step:9:effect:1",
1888 "op-1:step:9:effect:2",
1889 "op-1:step:9:effect:10",
1890 "op-1:step:10:effect:0",
1891 ]
1892 );
1893 assert!(
1895 super::effect_position("not-a-kernel-id")
1896 > super::effect_position("op-1:step:1:effect:0")
1897 );
1898 }
1899
1900 const OPERATION: &str = "op-tx-1";
1905
1906 fn operation() -> OperationId {
1907 OperationId::new(OPERATION).unwrap()
1908 }
1909
1910 fn input_id(name: &str) -> InputId {
1911 InputId::new(name).unwrap()
1912 }
1913
1914 fn boot_config(supported: impl IntoIterator<Item = EffectKindTag>) -> OperationConfig {
1915 OperationConfig {
1916 execution_policy: Some(ExecutionPolicy {
1917 max_turns: Some(12),
1918 ..ExecutionPolicy::default()
1919 }),
1920 host_effect_support: HostEffectSupport::new(supported),
1921 ..OperationConfig::default()
1922 }
1923 }
1924
1925 fn envelope(id: &str, observed_at_ms: u64, input: KernelInput) -> WireEnvelope {
1926 WireEnvelope::new(
1927 operation(),
1928 input_id(id),
1929 WireU64::new(observed_at_ms),
1930 input,
1931 )
1932 }
1933
1934 fn configure_at(id: &str, supported: impl IntoIterator<Item = EffectKindTag>) -> WireEnvelope {
1935 envelope(
1936 id,
1937 1_700_000_000_000,
1938 KernelInput::ConfigureOperation(ConfigureOperation {
1939 config: boot_config(supported),
1940 }),
1941 )
1942 }
1943
1944 fn configure() -> WireEnvelope {
1945 configure_at(
1946 "in-configure",
1947 [EffectKindTag::CallProvider, EffectKindTag::SpawnTasks],
1948 )
1949 }
1950
1951 fn start_at(id: &str, observed_at_ms: u64) -> WireEnvelope {
1952 envelope(
1953 id,
1954 observed_at_ms,
1955 KernelInput::StartOperation(StartOperation {
1956 entry: RootEntry::Agent(RootAgentEntry {
1957 task: LogicalTask::new("write the brief"),
1958 run_spec: None,
1959 }),
1960 initial_context: InitialContext::default(),
1961 }),
1962 )
1963 }
1964
1965 fn start() -> WireEnvelope {
1966 start_at("in-start", 1_700_000_001_000)
1967 }
1968
1969 fn provider_outcome() -> EffectOutcome {
1970 EffectOutcome::Succeeded(EffectSucceeded {
1971 result: EffectSuccess::Provider(ProviderSuccess {
1972 outcome: ProviderOutcome::ContextOverflow(ProviderContextOverflow::default()),
1973 }),
1974 })
1975 }
1976
1977 fn failure_outcome() -> EffectOutcome {
1978 EffectOutcome::Failed(EffectFailed {
1979 failure: HostEffectFailure {
1980 kind: HostEffectFailureKind::TransportExhausted,
1981 message: "the vendor gave up".to_string(),
1982 retryable: Some(false),
1983 },
1984 })
1985 }
1986
1987 fn resolve_at(
1988 id: &str,
1989 observed_at_ms: u64,
1990 effect_id: &EffectId,
1991 outcome: EffectOutcome,
1992 ) -> WireEnvelope {
1993 envelope(
1994 id,
1995 observed_at_ms,
1996 KernelInput::ResolveEffect(ResolveEffect {
1997 effect_id: effect_id.clone(),
1998 outcome,
1999 }),
2000 )
2001 }
2002
2003 fn cancel_at(id: &str, observed_at_ms: u64) -> WireEnvelope {
2004 envelope(
2005 id,
2006 observed_at_ms,
2007 KernelInput::HostControl(HostControl {
2008 command: HostCommand::Cancel(CancelCommand {
2009 reason: CancellationReason::User,
2010 pending_call_ids: vec![],
2011 }),
2012 }),
2013 )
2014 }
2015
2016 fn signal_at(id: &str, observed_at_ms: u64) -> WireEnvelope {
2017 use super::super::event::{DeliverSignal, ExternalEvent, LogicalSignal};
2018 use super::super::scalar::{DeliveryId, SignalId};
2019
2020 envelope(
2021 id,
2022 observed_at_ms,
2023 KernelInput::DeliverExternalEvent(DeliverExternalEvent {
2024 event: ExternalEvent::DeliverSignal(DeliverSignal {
2025 delivery_id: DeliveryId::new(format!("delivery-{id}")).unwrap(),
2026 attempt: 1,
2027 signal: LogicalSignal::new(SignalId::new("sig-late").unwrap()),
2028 }),
2029 }),
2030 )
2031 }
2032
2033 #[derive(Debug, Clone, PartialEq, Serialize)]
2036 struct TestStep {
2037 plan: String,
2038 disposition: StepDisposition,
2039 }
2040
2041 impl TransitionStep for TestStep {
2042 fn disposition(&self) -> &StepDisposition {
2043 &self.disposition
2044 }
2045 }
2046
2047 fn nothing(plan: &str) -> TestStep {
2048 TestStep {
2049 plan: plan.to_string(),
2050 disposition: StepDisposition::Effects(EffectsDisposition::default()),
2051 }
2052 }
2053
2054 fn publishing(plan: &str, effects: Vec<KernelEffect>) -> TestStep {
2055 TestStep {
2056 plan: plan.to_string(),
2057 disposition: StepDisposition::Effects(EffectsDisposition { effects }),
2058 }
2059 }
2060
2061 fn effect(id: &str, causation: &InputId, kind: EffectKind) -> KernelEffect {
2062 KernelEffect {
2063 effect_id: EffectId::new(id).unwrap(),
2064 causation_input_id: causation.clone(),
2065 effect: kind,
2066 }
2067 }
2068
2069 fn provider_payload() -> CallProviderEffect {
2070 let manager = crate::context::manager::ContextManager::new(128_000);
2071 let (mut candidate, _) = manager
2072 .prepare_candidate(
2073 OPERATION.to_string(),
2074 "step-1".to_string(),
2075 1,
2076 crate::evolution::ContentDigest::from_bytes(b"test-policy"),
2077 )
2078 .unwrap();
2079 let context = super::super::effect::RenderedContext::default();
2080 let tools: Vec<super::super::effect::ToolSchema> = vec![];
2081 candidate.rendered_snapshot = crate::evolution::ContentDigest::from_bytes(
2082 super::super::record::canonical_bytes(&(&context, &tools))
2083 .unwrap()
2084 .as_slice(),
2085 );
2086 CallProviderEffect {
2087 context,
2088 tools,
2089 context_candidate: Box::new(candidate),
2090 }
2091 }
2092
2093 fn provider_effect_id(step_seq: WireU64) -> EffectId {
2094 EffectId::new(format!("{OPERATION}:step:{step_seq}:effect:0")).unwrap()
2095 }
2096
2097 fn plan(context: &PlanContext<'_>) -> Result<TestStep, KernelFault> {
2103 let label = format!("{}@{}", context.input.input.kind(), context.step_seq);
2104 Ok(match &context.input.input {
2105 NormalizedPayload::StartOperation(_) => publishing(
2106 &label,
2107 vec![effect(
2108 provider_effect_id(context.step_seq).as_str(),
2109 &context.input.input_id,
2110 EffectKind::CallProvider(provider_payload()),
2111 )],
2112 ),
2113 NormalizedPayload::HostControl(_) => TestStep {
2114 plan: label,
2115 disposition: StepDisposition::Terminal(TerminalDisposition {
2116 terminal: KernelTerminal::Cancelled(CancelledTerminal {
2117 reason: CancellationReason::User,
2118 usage: UsageReport::default(),
2119 }),
2120 }),
2121 },
2122 _ => nothing(&label),
2123 })
2124 }
2125
2126 type Tx = KernelTransaction<TestStep, InMemoryRecordIndex>;
2127
2128 fn transaction() -> Tx {
2129 KernelTransaction::new(ConfigDefaults::default(), InMemoryRecordIndex::new())
2130 }
2131
2132 fn bounded(bounds: TailBounds) -> Tx {
2140 let mut defaults = ConfigDefaults::default();
2141 defaults.baseline.recovery_policy.tail_bounds = bounds;
2142 KernelTransaction::new(defaults, InMemoryRecordIndex::new())
2143 }
2144
2145 fn run(tx: &mut Tx, envelope: &WireEnvelope) -> CommittedTransition<TestStep> {
2147 run_with(tx, envelope, plan)
2148 }
2149
2150 fn run_with<F>(
2151 tx: &mut Tx,
2152 envelope: &WireEnvelope,
2153 planner: F,
2154 ) -> CommittedTransition<TestStep>
2155 where
2156 F: FnOnce(&PlanContext<'_>) -> Result<TestStep, KernelFault>,
2157 {
2158 let preparation = tx.prepare(envelope, planner);
2159 let token = preparation
2160 .token()
2161 .unwrap_or_else(|| {
2162 panic!(
2163 "expected a prepared transition, got {:?}",
2164 preparation.fault()
2165 )
2166 })
2167 .clone();
2168 let head = preparation.record().unwrap().record_digest().clone();
2169 tx.commit(&token, &head).expect("the commit must succeed")
2170 }
2171
2172 fn fault_of(preparation: &RecordPreparation<TestStep>) -> KernelFaultCode {
2173 preparation
2174 .fault()
2175 .unwrap_or_else(|| panic!("expected a rejection, got a success"))
2176 .code
2177 }
2178
2179 fn observable(
2182 tx: &Tx,
2183 ) -> (
2184 Option<DurableHead>,
2185 OperationLifecycle,
2186 Vec<String>,
2187 TailUsage,
2188 bool,
2189 ) {
2190 (
2191 tx.head(),
2192 tx.lifecycle(),
2193 tx.pending_effects()
2194 .map(|effect| effect.effect_id.to_string())
2195 .collect(),
2196 tx.tail_usage(),
2197 tx.has_candidate(),
2198 )
2199 }
2200
2201 fn started() -> (Tx, Vec<KernelRecord>, EffectId) {
2202 let mut tx = transaction();
2203 let genesis = run(&mut tx, &configure());
2204 let started = run(&mut tx, &start());
2205 let effect_id = provider_effect_id(started.step_seq);
2206 (tx, vec![genesis.record, started.record], effect_id)
2207 }
2208
2209 #[test]
2215 fn matrix_row1_before_a_prepare_nothing_exists() {
2216 let tx = transaction();
2217 assert!(!tx.has_candidate());
2218 assert_eq!(tx.head(), None);
2219 assert_eq!(tx.lifecycle(), OperationLifecycle::Created);
2220 assert_eq!(tx.pending_effects().count(), 0);
2221 assert_eq!(tx.tail_usage(), TailUsage::default());
2222 assert!(tx.index().is_empty(), "no record exists yet");
2223 assert_eq!(tx.checkpoint_boundary(), None);
2224 assert!(tx.outstanding_token().is_none());
2225 }
2226
2227 #[test]
2232 fn matrix_row2_a_rejected_prepare_hands_out_no_token_and_moves_nothing() {
2233 let (mut tx, _, effect_id) = started();
2234
2235 let rejections: Vec<(&str, RecordPreparation<TestStep>)> = vec![
2236 (
2238 "second configure",
2239 tx.prepare(
2240 &configure_at("in-again", [EffectKindTag::CallProvider]),
2241 plan,
2242 ),
2243 ),
2244 (
2246 "unknown effect",
2247 tx.prepare(
2248 &resolve_at(
2249 "in-unknown",
2250 1_700_000_002_000,
2251 &EffectId::new("op-tx-1:step:9:effect:0").unwrap(),
2252 failure_outcome(),
2253 ),
2254 plan,
2255 ),
2256 ),
2257 (
2259 "planner fault",
2260 tx.prepare(
2261 &resolve_at(
2262 "in-planned",
2263 1_700_000_002_000,
2264 &effect_id,
2265 failure_outcome(),
2266 ),
2267 |_| {
2268 Err(KernelFault::new(
2269 KernelFaultCode::ResourceLimitExceeded,
2270 "the planner refused",
2271 ))
2272 },
2273 ),
2274 ),
2275 ];
2276
2277 let before = tx.clone();
2278 for (label, preparation) in rejections {
2279 assert!(preparation.is_zero_mutation(), "{label}");
2280 assert!(preparation.token().is_none(), "{label}");
2281 assert!(preparation.record().is_none(), "{label}");
2282 assert!(preparation.step().is_none(), "{label}");
2283 assert!(preparation.step_seq().is_none(), "{label}");
2284 }
2285 assert_eq!(tx, before, "a rejected prepare must not move one byte");
2286 }
2287
2288 #[test]
2290 fn matrix_row3_a_crash_before_the_append_is_undone_by_abort() {
2291 let (mut tx, records, _) = started();
2292 let before = observable(&tx);
2293
2294 let preparation = tx.prepare(&cancel_at("in-cancel", 1_700_000_003_000), plan);
2295 let token = preparation.token().expect("prepared").clone();
2296 assert!(tx.has_candidate());
2297 assert_eq!(
2298 tx.head().map(|head| head.digest),
2299 Some(records[1].record_digest().clone()),
2300 "a candidate does not move the durable head"
2301 );
2302
2303 let discarded = tx.abort(&token).expect("the candidate was never appended");
2304 assert_eq!(discarded.step_seq(), WireU64::new(2));
2305 assert_eq!(
2306 observable(&tx),
2307 before,
2308 "an aborted candidate leaves no trace"
2309 );
2310
2311 assert_eq!(
2313 tx.abort(&token).unwrap_err().code,
2314 KernelFaultCode::TransactionConflict
2315 );
2316 assert!(
2317 !tx.is_poisoned(),
2318 "an abort is a normal, non-poisoning path"
2319 );
2320 }
2321
2322 #[test]
2324 fn matrix_row4_a_cas_conflict_discards_the_candidate_and_demands_a_rebuild() {
2325 let (mut tx, records, _) = started();
2326 let preparation = tx.prepare(&cancel_at("in-cancel", 1_700_000_003_000), plan);
2327 let token = preparation.token().expect("prepared").clone();
2328
2329 let forked = canonical_digest(b"another writer's record");
2331 let fault = tx.note_append_conflict(&token, Some(&forked));
2332
2333 assert_eq!(fault.code, KernelFaultCode::TransactionConflict);
2334 assert!(!fault.is_retryable(), "a conflict is not a bare retry");
2335 assert!(
2336 !tx.has_candidate(),
2337 "the candidate is discarded, not appended"
2338 );
2339 assert!(tx.is_poisoned(), "this layer never rebuilds itself");
2340
2341 assert_eq!(
2343 fault_of(&tx.prepare(&cancel_at("in-cancel-2", 1_700_000_004_000), plan)),
2344 KernelFaultCode::TransactionConflict
2345 );
2346 assert_eq!(
2347 tx.abort(&token).unwrap_err().code,
2348 KernelFaultCode::TransactionConflict
2349 );
2350 assert_eq!(
2351 tx.commit(&token, &forked).unwrap_err().code,
2352 KernelFaultCode::TransactionConflict
2353 );
2354
2355 let rebuilt = Tx::rebuild_from_records(
2357 &records,
2358 ConfigDefaults::default(),
2359 InMemoryRecordIndex::from_records(&records),
2360 plan,
2361 )
2362 .expect("the journal still verifies");
2363 assert_eq!(
2364 rebuilt.head().map(|head| head.digest),
2365 Some(records[1].record_digest().clone())
2366 );
2367 assert!(!rebuilt.is_poisoned());
2368 }
2369
2370 #[test]
2373 fn matrix_row5_a_crash_between_append_and_commit_rebuilds_and_publishes_once() {
2374 let mut tx = transaction();
2375 let genesis = run(&mut tx, &configure());
2376 let mut journal = vec![genesis.record];
2377
2378 let preparation = tx.prepare(&start(), plan);
2379 let appended = preparation.record().expect("prepared").clone();
2380 journal.push(appended.clone());
2382 assert_eq!(
2384 tx.pending_effects().count(),
2385 0,
2386 "a record that has not been committed publishes no effect (§15.2)"
2387 );
2388 drop(tx);
2389
2390 let rebuilt = Tx::rebuild_from_records(
2391 &journal,
2392 ConfigDefaults::default(),
2393 InMemoryRecordIndex::from_records(&journal),
2394 plan,
2395 )
2396 .expect("the journal rebuilds");
2397
2398 assert_eq!(
2399 rebuilt.head().map(|head| head.digest),
2400 Some(appended.record_digest().clone())
2401 );
2402 assert_eq!(rebuilt.lifecycle(), OperationLifecycle::Running);
2403 let pending: Vec<&EffectId> = rebuilt
2404 .pending_effects()
2405 .map(|effect| &effect.effect_id)
2406 .collect();
2407 assert_eq!(
2408 pending,
2409 vec![&provider_effect_id(WireU64::new(1))],
2410 "the rebuilt runtime re-exposes the one pending effect, with the same identity"
2411 );
2412 }
2413
2414 #[test]
2417 fn matrix_row6_a_failed_commit_never_becomes_an_abort() {
2418 let mut tx = transaction();
2419 let genesis = run(&mut tx, &configure());
2420 let journal = vec![genesis.record];
2421
2422 let preparation = tx.prepare(&start(), plan);
2423 let token = preparation.token().expect("prepared").clone();
2424 let wrong_head = canonical_digest(b"some other record");
2427 let fault = tx
2428 .commit(&token, &wrong_head)
2429 .expect_err("a commit that cannot be attributed must fail");
2430 assert_eq!(fault.code, KernelFaultCode::TransactionConflict);
2431
2432 assert!(tx.is_poisoned());
2435 assert!(!tx.has_candidate());
2436 assert_eq!(
2437 tx.abort(&token).unwrap_err().code,
2438 KernelFaultCode::TransactionConflict,
2439 "append-then-abort is not expressible"
2440 );
2441
2442 let rebuilt = Tx::rebuild_from_records(
2444 &journal,
2445 ConfigDefaults::default(),
2446 InMemoryRecordIndex::from_records(&journal),
2447 plan,
2448 )
2449 .expect("the durable record survives a failed commit");
2450 assert_eq!(rebuilt.lifecycle(), OperationLifecycle::Configured);
2451 }
2452
2453 #[test]
2455 fn matrix_row7_a_lost_commit_response_replays_the_same_record_and_step() {
2456 let mut tx = transaction();
2457 run(&mut tx, &configure());
2458 let committed = run(&mut tx, &start());
2459 let before = observable(&tx);
2460
2461 let replay = tx.prepare(&start(), plan);
2463
2464 assert!(replay.token().is_none(), "a replay has nothing to commit");
2465 assert_eq!(replay.step_seq(), Some(committed.step_seq));
2466 assert_eq!(replay.record(), Some(&committed.record));
2467 assert_eq!(replay.step(), Some(&committed.step));
2468 assert_eq!(
2469 observable(&tx),
2470 before,
2471 "a replay creates no record and re-publishes no effect"
2472 );
2473 assert_eq!(tx.index().len(), 2, "no second record was minted");
2474 }
2475
2476 #[test]
2479 fn matrix_row8_a_lost_resolution_response_replays_idempotently() {
2480 let (mut tx, _, effect_id) = started();
2481 let resolved = run(
2482 &mut tx,
2483 &resolve_at(
2484 "in-resolve",
2485 1_700_000_002_000,
2486 &effect_id,
2487 provider_outcome(),
2488 ),
2489 );
2490 assert_eq!(tx.pending_effects().count(), 0, "the effect is settled");
2491 let before = observable(&tx);
2492
2493 let same_input = tx.prepare(
2495 &resolve_at(
2496 "in-resolve",
2497 1_700_000_002_000,
2498 &effect_id,
2499 provider_outcome(),
2500 ),
2501 plan,
2502 );
2503 assert_eq!(same_input.step_seq(), Some(resolved.step_seq));
2504 assert_eq!(same_input.record(), Some(&resolved.record));
2505
2506 assert_eq!(
2510 fault_of(&tx.prepare(
2511 &resolve_at(
2512 "in-resolve",
2513 1_700_000_002_500,
2514 &effect_id,
2515 provider_outcome()
2516 ),
2517 plan
2518 )),
2519 KernelFaultCode::DuplicateInputConflict
2520 );
2521
2522 let new_input = tx.prepare(
2525 &resolve_at(
2526 "in-resolve-again",
2527 1_700_000_003_000,
2528 &effect_id,
2529 provider_outcome(),
2530 ),
2531 plan,
2532 );
2533 assert_eq!(
2534 new_input.step_seq(),
2535 Some(resolved.step_seq),
2536 "effect-level dedup points at the existing record's step_seq"
2537 );
2538 assert_eq!(new_input.record(), Some(&resolved.record));
2539 assert!(
2540 new_input.token().is_none(),
2541 "reporting this as Prepared with an old step_seq is the dead end this replaces"
2542 );
2543 assert_eq!(observable(&tx), before);
2544 assert_eq!(tx.index().len(), 3, "still three records");
2545 }
2546
2547 #[test]
2552 fn a_committed_record_can_never_be_aborted() {
2553 let mut tx = transaction();
2554 let preparation = tx.prepare(&configure(), plan);
2555 let token = preparation.token().expect("prepared").clone();
2556 let head = preparation.record().unwrap().record_digest().clone();
2557 tx.commit(&token, &head).expect("commit");
2558
2559 let error = tx
2560 .abort(&token)
2561 .expect_err("a committed record has no candidate to abort");
2562 assert_eq!(error.code, KernelFaultCode::TransactionConflict);
2563 assert!(!tx.is_poisoned(), "asking is not itself a corruption");
2564 assert_eq!(
2565 tx.head().map(|head| head.step_seq),
2566 Some(WireU64::ZERO),
2567 "the committed record stands"
2568 );
2569 }
2570
2571 #[test]
2572 fn a_second_prepare_is_refused_while_a_candidate_is_outstanding() {
2573 let mut tx = transaction();
2574 let first = tx.prepare(&configure(), plan);
2575 let token = first.token().expect("prepared").clone();
2576
2577 let second = tx.prepare(&start(), plan);
2578 assert_eq!(fault_of(&second), KernelFaultCode::TransactionConflict);
2579 assert_eq!(
2580 tx.outstanding_token(),
2581 Some(&token),
2582 "the first candidate is not displaced by the second attempt"
2583 );
2584
2585 let head = first.record().unwrap().record_digest().clone();
2587 tx.commit(&token, &head).expect("the first candidate wins");
2588 }
2589
2590 #[test]
2591 fn effects_become_visible_only_when_their_record_is_durable() {
2592 let mut tx = transaction();
2593 run(&mut tx, &configure());
2594
2595 let preparation = tx.prepare(&start(), plan);
2596 assert_eq!(
2597 preparation.step().unwrap().effects().len(),
2598 1,
2599 "the planned step carries the effect"
2600 );
2601 assert_eq!(
2602 tx.pending_effects().count(),
2603 0,
2604 "but nothing is pending until the record is durable"
2605 );
2606
2607 let token = preparation.token().unwrap().clone();
2608 let head = preparation.record().unwrap().record_digest().clone();
2609 let committed = tx.commit(&token, &head).unwrap();
2610 assert_eq!(committed.published_effects().len(), 1);
2611 assert_eq!(tx.pending_effects().count(), 1);
2612 }
2613
2614 #[test]
2615 fn an_aborted_candidate_publishes_nothing() {
2616 let mut tx = transaction();
2617 run(&mut tx, &configure());
2618 let preparation = tx.prepare(&start(), plan);
2619 let token = preparation.token().unwrap().clone();
2620 tx.abort(&token).unwrap();
2621 assert_eq!(tx.pending_effects().count(), 0);
2622 assert_eq!(tx.lifecycle(), OperationLifecycle::Configured);
2623 }
2624
2625 #[test]
2626 fn the_operation_id_is_bound_by_the_genesis_record() {
2627 let mut tx = transaction();
2628 run(&mut tx, &configure());
2629 assert_eq!(tx.operation_id(), Some(&operation()));
2630
2631 let foreign = WireEnvelope::new(
2632 OperationId::new("op-other").unwrap(),
2633 input_id("in-foreign"),
2634 WireU64::new(1_700_000_002_000),
2635 KernelInput::HostControl(HostControl {
2636 command: HostCommand::Cancel(CancelCommand {
2637 reason: CancellationReason::User,
2638 pending_call_ids: vec![],
2639 }),
2640 }),
2641 );
2642 assert_eq!(
2643 fault_of(&tx.prepare(&foreign, plan)),
2644 KernelFaultCode::OperationMismatch
2645 );
2646 }
2647
2648 #[test]
2649 fn a_duplicate_input_id_with_a_different_payload_is_a_conflict() {
2650 let mut tx = transaction();
2651 run(&mut tx, &configure());
2652 run(&mut tx, &start());
2653
2654 let mut divergent = start();
2656 divergent.input = KernelInput::StartOperation(StartOperation {
2657 entry: RootEntry::Agent(RootAgentEntry {
2658 task: LogicalTask::new("write something else"),
2659 run_spec: None,
2660 }),
2661 initial_context: InitialContext::default(),
2662 });
2663 assert_eq!(
2664 fault_of(&tx.prepare(&divergent, plan)),
2665 KernelFaultCode::DuplicateInputConflict
2666 );
2667 }
2668
2669 #[test]
2670 fn an_exact_replay_is_answered_before_the_clock_check() {
2671 let mut tx = transaction();
2672 run(&mut tx, &configure());
2673 let started = run(&mut tx, &start_at("in-start", 1_700_000_005_000));
2674
2675 let replay = tx.prepare(&start_at("in-start", 1_700_000_005_000), plan);
2677 assert_eq!(replay.step_seq(), Some(started.step_seq));
2678
2679 assert_eq!(
2681 fault_of(&tx.prepare(&cancel_at("in-past", 1_700_000_004_000), plan)),
2682 KernelFaultCode::ClockRegression
2683 );
2684 }
2685
2686 #[test]
2687 fn a_terminal_closes_the_operation_to_every_later_input() {
2688 let (mut tx, _, effect_id) = started();
2689 assert_eq!(tx.pending_effects().count(), 1);
2690
2691 let terminal = run(&mut tx, &cancel_at("in-cancel", 1_700_000_003_000));
2692 assert!(matches!(
2693 terminal.terminal(),
2694 Some(KernelTerminal::Cancelled(_))
2695 ));
2696 assert_eq!(tx.lifecycle(), OperationLifecycle::Cancelled);
2697 assert!(tx.lifecycle().is_terminal());
2698 assert_eq!(
2699 tx.pending_effects().count(),
2700 0,
2701 "a terminal leaves nothing waiting on the host"
2702 );
2703
2704 let before = tx.clone();
2707 for envelope in [
2708 resolve_at("in-late", 1_700_000_004_000, &effect_id, provider_outcome()),
2709 start_at("in-restart", 1_700_000_004_000),
2710 signal_at("in-late-signal", 1_700_000_004_000),
2711 ] {
2712 assert_eq!(
2713 fault_of(&tx.prepare(&envelope, plan)),
2714 KernelFaultCode::InvalidLifecycle,
2715 "{} must be refused after a terminal",
2716 envelope.input_id
2717 );
2718 }
2719 assert_eq!(tx, before, "a refused input leaves the transaction alone");
2720 }
2721
2722 #[test]
2724 fn a_re_issued_cancellation_replays_the_terminal_it_already_committed() {
2725 let (mut tx, _, _) = started();
2726 let cancelled = run(&mut tx, &cancel_at("in-cancel", 1_700_000_003_000));
2727 let before = tx.clone();
2728
2729 let replay = tx.prepare(&cancel_at("in-cancel-again", 1_700_000_004_000), plan);
2732 assert!(
2733 matches!(replay, KernelPreparation::Replayed(_)),
2734 "a re-issued cancellation must replay, not be refused by the latch it created"
2735 );
2736 assert_eq!(replay.step_seq(), Some(cancelled.step_seq));
2737 assert_eq!(
2738 replay.record().unwrap().record_digest(),
2739 cancelled.record.record_digest(),
2740 "the replay points at the existing record, so no second record exists"
2741 );
2742 assert_eq!(tx, before, "a replay moves nothing");
2743
2744 let exact = tx.prepare(&cancel_at("in-cancel", 1_700_000_003_000), plan);
2746 assert_eq!(exact.step_seq(), Some(cancelled.step_seq));
2747 assert_eq!(tx, before);
2748 }
2749
2750 #[test]
2751 fn configured_input_byte_limit_rejects_later_inputs_without_mutation() {
2752 let mut tx = transaction();
2753 let mut config = boot_config([EffectKindTag::CallProvider]);
2754 config.kernel_limits = Some(KernelLimits {
2755 max_input_bytes: Some(1_024),
2756 ..KernelLimits::default()
2757 });
2758 let configure = envelope(
2759 "in-configure-limited",
2760 1_700_000_000_000,
2761 KernelInput::ConfigureOperation(ConfigureOperation { config }),
2762 );
2763 run(&mut tx, &configure);
2764
2765 let oversized = envelope(
2766 "in-oversized",
2767 1_700_000_001_000,
2768 KernelInput::StartOperation(StartOperation {
2769 entry: RootEntry::Agent(RootAgentEntry {
2770 task: LogicalTask::new("x".repeat(4_096)),
2771 run_spec: None,
2772 }),
2773 initial_context: InitialContext::default(),
2774 }),
2775 );
2776 let before = tx.clone();
2777 let rejected = tx.prepare(&oversized, plan);
2778
2779 assert_eq!(fault_of(&rejected), KernelFaultCode::ResourceLimitExceeded);
2780 assert!(rejected.is_zero_mutation());
2781 assert_eq!(tx, before);
2782 }
2783
2784 #[test]
2785 fn a_differing_cancellation_after_the_terminal_is_a_conflict_not_an_overwrite() {
2786 let (mut tx, _, _) = started();
2787 run(&mut tx, &cancel_at("in-cancel", 1_700_000_003_000));
2788 let before = tx.clone();
2789
2790 let divergent = envelope(
2791 "in-cancel-other",
2792 1_700_000_004_000,
2793 KernelInput::HostControl(HostControl {
2794 command: HostCommand::Cancel(CancelCommand {
2795 reason: CancellationReason::HostShutdown,
2796 pending_call_ids: vec![],
2797 }),
2798 }),
2799 );
2800 assert_eq!(
2801 fault_of(&tx.prepare(&divergent, plan)),
2802 KernelFaultCode::DuplicateInputConflict,
2803 "the operation already ended for the first reason; the second cannot re-decide it"
2804 );
2805 assert_eq!(tx, before);
2806 }
2807
2808 #[test]
2813 fn a_conflicting_resolution_of_a_settled_effect_fails_closed() {
2814 let (mut tx, _, effect_id) = started();
2815 run(
2816 &mut tx,
2817 &resolve_at(
2818 "in-resolve",
2819 1_700_000_002_000,
2820 &effect_id,
2821 provider_outcome(),
2822 ),
2823 );
2824 let before = tx.clone();
2825
2826 let conflicting = tx.prepare(
2827 &resolve_at(
2828 "in-resolve-conflict",
2829 1_700_000_003_000,
2830 &effect_id,
2831 failure_outcome(),
2832 ),
2833 plan,
2834 );
2835 assert_eq!(
2836 fault_of(&conflicting),
2837 KernelFaultCode::UnexpectedEffectOutcome
2838 );
2839 assert_eq!(tx, before);
2840 }
2841
2842 #[test]
2843 fn an_unknown_effect_id_fails_closed() {
2844 let (mut tx, _, _) = started();
2845 let unknown = EffectId::new("op-tx-1:step:99:effect:0").unwrap();
2846 assert_eq!(
2847 fault_of(&tx.prepare(
2848 &resolve_at(
2849 "in-unknown",
2850 1_700_000_002_000,
2851 &unknown,
2852 provider_outcome()
2853 ),
2854 plan
2855 )),
2856 KernelFaultCode::UnexpectedEffectOutcome
2857 );
2858 }
2859
2860 #[test]
2861 fn a_resolution_carrying_another_effect_kinds_payload_fails_closed() {
2862 let (mut tx, _, effect_id) = started();
2863 let wrong_shape = EffectOutcome::Succeeded(EffectSucceeded {
2864 result: EffectSuccess::Tools(ToolsSuccess::default()),
2865 });
2866 assert_eq!(
2867 fault_of(&tx.prepare(
2868 &resolve_at("in-wrong", 1_700_000_002_000, &effect_id, wrong_shape),
2869 plan
2870 )),
2871 KernelFaultCode::UnexpectedEffectOutcome
2872 );
2873 }
2874
2875 #[test]
2876 fn at_most_one_pending_effect_per_kind() {
2877 let (mut tx, _, _) = started();
2878 let before = tx.clone();
2879
2880 let greedy = tx.prepare(&cancel_at("in-second", 1_700_000_002_000), |context| {
2882 Ok(publishing(
2883 "greedy",
2884 vec![effect(
2885 "op-tx-1:step:2:effect:0",
2886 &context.input.input_id,
2887 EffectKind::CallProvider(provider_payload()),
2888 )],
2889 ))
2890 });
2891 assert_eq!(fault_of(&greedy), KernelFaultCode::ResourceLimitExceeded);
2892 assert_eq!(tx, before, "the refusal is zero mutation");
2893
2894 let doubled = tx.prepare(&cancel_at("in-double", 1_700_000_002_000), |context| {
2896 Ok(publishing(
2897 "doubled",
2898 vec![
2899 effect(
2900 "op-tx-1:step:2:effect:0",
2901 &context.input.input_id,
2902 EffectKind::SpawnTasks(SpawnTasksEffect::default()),
2903 ),
2904 effect(
2905 "op-tx-1:step:2:effect:1",
2906 &context.input.input_id,
2907 EffectKind::SpawnTasks(SpawnTasksEffect::default()),
2908 ),
2909 ],
2910 ))
2911 });
2912 assert_eq!(fault_of(&doubled), KernelFaultCode::ResourceLimitExceeded);
2913 }
2914
2915 #[test]
2916 fn an_effect_kind_the_host_did_not_declare_is_refused_before_emission() {
2917 let mut tx = transaction();
2918 run(
2920 &mut tx,
2921 &configure_at("in-configure", [EffectKindTag::CallProvider]),
2922 );
2923 let before = tx.clone();
2924
2925 let refused = tx.prepare(&start(), |context| {
2926 Ok(publishing(
2927 "tools",
2928 vec![effect(
2929 "op-tx-1:step:1:effect:0",
2930 &context.input.input_id,
2931 EffectKind::ExecuteTools(ExecuteToolsEffect::default()),
2932 )],
2933 ))
2934 });
2935 assert_eq!(fault_of(&refused), KernelFaultCode::UnsupportedEffect);
2936 assert_eq!(tx, before, "no record, no effect, no state change");
2937 }
2938
2939 #[test]
2940 fn a_launch_token_is_never_minted_twice() {
2941 let (mut tx, _, effect_id) = started();
2942 let spawn = |id: &str, token: &str| {
2943 let effect_id = id.to_string();
2944 let launch_token = token.to_string();
2945 move |context: &PlanContext<'_>| {
2946 Ok(publishing(
2947 "spawn",
2948 vec![effect(
2949 &effect_id,
2950 &context.input.input_id,
2951 EffectKind::SpawnTasks(SpawnTasksEffect {
2952 tasks: vec![TaskLaunch {
2953 task_id: TaskId::new("task-1").unwrap(),
2954 attempt_id: AttemptId::new("task-1:attempt:1").unwrap(),
2955 launch_token: LaunchToken::new(launch_token.clone()).unwrap(),
2956 node_id: NodeId::new("node-1").unwrap(),
2957 spec: LogicalAgentSpec::new("do the thing"),
2958 }],
2959 budget: None,
2960 }),
2961 )],
2962 ))
2963 }
2964 };
2965
2966 let spawn_effect_id = EffectId::new("op-tx-1:step:2:effect:0").unwrap();
2967 run_with(
2968 &mut tx,
2969 &resolve_at(
2970 "in-resolve",
2971 1_700_000_002_000,
2972 &effect_id,
2973 provider_outcome(),
2974 ),
2975 spawn(spawn_effect_id.as_str(), "op-tx-1:launch:1"),
2976 );
2977 assert!(tx.knows_launch_token(&LaunchToken::new("op-tx-1:launch:1").unwrap()));
2978
2979 run(
2981 &mut tx,
2982 &resolve_at(
2983 "in-spawned",
2984 1_700_000_003_000,
2985 &spawn_effect_id,
2986 EffectOutcome::Succeeded(EffectSucceeded {
2987 result: EffectSuccess::TasksSpawned(TasksSpawnedSuccess::default()),
2988 }),
2989 ),
2990 );
2991
2992 let reused = tx.prepare(
2995 &cancel_at("in-relaunch", 1_700_000_004_000),
2996 spawn("op-tx-1:step:4:effect:0", "op-tx-1:launch:1"),
2997 );
2998 assert_eq!(fault_of(&reused), KernelFaultCode::TransactionConflict);
2999
3000 let fresh = tx.prepare(
3002 &cancel_at("in-relaunch", 1_700_000_004_000),
3003 spawn("op-tx-1:step:4:effect:0", "op-tx-1:launch:2"),
3004 );
3005 assert!(fresh.token().is_some(), "{:?}", fresh.fault());
3006 }
3007
3008 #[test]
3009 fn an_effect_id_is_never_minted_twice() {
3010 let (mut tx, _, effect_id) = started();
3011 let collision = tx.prepare(
3012 &resolve_at(
3013 "in-resolve",
3014 1_700_000_002_000,
3015 &effect_id,
3016 provider_outcome(),
3017 ),
3018 move |context| {
3019 Ok(publishing(
3020 "collide",
3021 vec![effect(
3022 "op-tx-1:step:1:effect:0",
3024 &context.input.input_id,
3025 EffectKind::SpawnTasks(SpawnTasksEffect::default()),
3026 )],
3027 ))
3028 },
3029 );
3030 assert_eq!(fault_of(&collision), KernelFaultCode::TransactionConflict);
3031 }
3032
3033 #[test]
3038 fn a_full_tail_asks_for_a_checkpoint_and_the_retry_is_a_fresh_prepare() {
3039 let mut tx = bounded(TailBounds::new(1, 2, 1024, 1024 * 1024).unwrap());
3040 run(&mut tx, &configure());
3041 run(&mut tx, &start());
3042 assert_eq!(tx.tail_pressure(), TailPressure::Full);
3043 let before = tx.clone();
3044
3045 let retry_envelope = cancel_at("in-cancel", 1_700_000_003_000);
3046 let refused = tx.prepare(&retry_envelope, plan);
3047 let fault = refused.fault().expect("rejected").clone();
3048 assert_eq!(fault.code, KernelFaultCode::CheckpointRequired);
3049 assert!(fault.is_retryable(), "the one retryable code (GAP-2)");
3050 assert!(refused.is_zero_mutation());
3051 assert_eq!(tx, before, "the input was never accepted");
3052
3053 let boundary = tx.checkpoint_boundary().expect("a head exists");
3055 assert_eq!(boundary.through_step_seq, WireU64::new(1));
3056 let usage = tx.note_checkpoint_acked(&boundary).expect("ack");
3057 assert_eq!(
3058 usage,
3059 TailUsage::default(),
3060 "the covered prefix is reclaimed"
3061 );
3062 assert_eq!(tx.tail_pressure(), TailPressure::Nominal, "not a latch");
3063
3064 let retried = tx.prepare(&retry_envelope, plan);
3066 assert!(retried.token().is_some(), "{:?}", retried.fault());
3067 assert_eq!(retried.record().unwrap().step_seq(), WireU64::new(2));
3068 }
3069
3070 #[test]
3071 fn the_tail_bounds_the_byte_axis_too() {
3072 let mut tx = bounded(TailBounds::new(64, 128, 128, 512).unwrap());
3073 let refused = tx.prepare(&configure(), plan);
3075 assert_eq!(fault_of(&refused), KernelFaultCode::CheckpointRequired);
3076 assert!(refused.fault().unwrap().is_retryable());
3077 }
3078
3079 #[test]
3080 fn tail_bounds_refuse_an_incoherent_watermark() {
3081 assert_eq!(
3082 TailBounds::new(10, 4, 100, 100).unwrap_err().code,
3083 KernelFaultCode::InvalidConfig
3084 );
3085 assert_eq!(
3086 TailBounds::new(1, 4, 0, 0).unwrap_err().code,
3087 KernelFaultCode::InvalidConfig
3088 );
3089 assert_eq!(TailBounds::default(), TailBounds::DEFAULT);
3090 }
3091
3092 #[test]
3093 fn the_tail_reports_its_soft_watermark_before_its_hard_limit() {
3094 let mut tx = bounded(TailBounds::new(2, 8, 1024 * 1024, 4 * 1024 * 1024).unwrap());
3095 assert_eq!(tx.tail_pressure(), TailPressure::Nominal);
3096 run(&mut tx, &configure());
3097 assert_eq!(tx.tail_pressure(), TailPressure::Nominal);
3098 run(&mut tx, &start());
3099 assert_eq!(tx.tail_pressure(), TailPressure::Watermark);
3100 }
3101
3102 #[test]
3105 fn a_checkpoint_boundary_neither_blocks_nor_is_blocked_by_a_candidate() {
3106 let (mut tx, records, _) = started();
3107
3108 let preparation = tx.prepare(&cancel_at("in-cancel", 1_700_000_003_000), plan);
3109 let token = preparation.token().expect("prepared").clone();
3110 let candidate_digest = preparation.record().unwrap().record_digest().clone();
3111
3112 let boundary = tx.checkpoint_boundary().expect("a head exists");
3113 assert_eq!(
3114 boundary.covered_head,
3115 *records[1].record_digest(),
3116 "the boundary follows the durable head, not the outstanding candidate"
3117 );
3118 assert_ne!(boundary.covered_head, candidate_digest);
3119
3120 tx.note_checkpoint_acked(&boundary).expect("ack");
3122 assert_eq!(tx.tail_usage(), TailUsage::default());
3123 assert!(tx.has_candidate(), "the candidate survived the checkpoint");
3124
3125 let committed = tx.commit(&token, &candidate_digest).expect("commit");
3127 assert_eq!(committed.step_seq, WireU64::new(2));
3128 assert_eq!(
3129 tx.tail_usage().records,
3130 1,
3131 "records after the candidate are tail"
3132 );
3133 }
3134
3135 #[test]
3136 fn a_checkpoint_that_covers_no_prefix_of_this_journal_is_refused() {
3137 let (mut tx, _, _) = started();
3138 let bogus = CheckpointBoundary {
3139 through_step_seq: WireU64::new(1),
3140 covered_head: canonical_digest(b"another operation's head"),
3141 };
3142 assert_eq!(
3143 tx.note_checkpoint_acked(&bogus).unwrap_err().code,
3144 KernelFaultCode::CheckpointIncompatible
3145 );
3146 assert_eq!(tx.tail_usage().records, 2, "nothing was reclaimed");
3147 }
3148
3149 #[test]
3154 fn a_rebuild_reproduces_every_step_digest_and_the_next_transition() {
3155 let mut live = transaction();
3156 let mut journal = Vec::new();
3157 journal.push(run(&mut live, &configure()).record);
3158 journal.push(run(&mut live, &start()).record);
3159 let effect_id = provider_effect_id(WireU64::new(1));
3160 journal.push(
3161 run(
3162 &mut live,
3163 &resolve_at(
3164 "in-resolve",
3165 1_700_000_002_000,
3166 &effect_id,
3167 provider_outcome(),
3168 ),
3169 )
3170 .record,
3171 );
3172
3173 let mut rebuilt = Tx::rebuild_from_records(
3174 &journal,
3175 ConfigDefaults::default(),
3176 InMemoryRecordIndex::from_records(&journal),
3177 plan,
3178 )
3179 .expect("the journal rebuilds");
3180
3181 assert_eq!(rebuilt.head(), live.head());
3182 assert_eq!(rebuilt.lifecycle(), live.lifecycle());
3183 assert_eq!(rebuilt.config(), live.config());
3184 assert_eq!(rebuilt.tail_usage(), live.tail_usage());
3185 for record in &journal {
3186 let step = rebuilt
3187 .committed_step(record.input_id())
3188 .expect("every replayed step is recoverable");
3189 record
3190 .verify_step(step)
3191 .expect("the rebuilt step matches the frozen digest");
3192 assert_eq!(Some(step), live.committed_step(record.input_id()));
3193 }
3194
3195 let next = cancel_at("in-cancel", 1_700_000_003_000);
3197 let uninterrupted = live.prepare(&next, plan);
3198 let after_rebuild = rebuilt.prepare(&next, plan);
3199 assert_eq!(
3200 after_rebuild.record().unwrap().step_digest(),
3201 uninterrupted.record().unwrap().step_digest()
3202 );
3203 assert_eq!(
3204 after_rebuild.record().unwrap().record_digest(),
3205 uninterrupted.record().unwrap().record_digest()
3206 );
3207 }
3208
3209 #[test]
3210 fn a_rebuild_refuses_a_broken_chain() {
3211 let mut live = transaction();
3212 let genesis = run(&mut live, &configure()).record;
3213 let started = run(&mut live, &start()).record;
3214 let resolved = run(
3215 &mut live,
3216 &resolve_at(
3217 "in-resolve",
3218 1_700_000_002_000,
3219 &provider_effect_id(WireU64::new(1)),
3220 provider_outcome(),
3221 ),
3222 )
3223 .record;
3224
3225 let gapped = vec![genesis.clone(), resolved.clone()];
3227 let error = Tx::rebuild_from_records(
3228 &gapped,
3229 ConfigDefaults::default(),
3230 InMemoryRecordIndex::from_records(&gapped),
3231 plan,
3232 )
3233 .expect_err("a chain with a hole is not a journal");
3234 assert_eq!(error.code, KernelFaultCode::RecordCorrupted);
3235
3236 let intact = vec![genesis, started, resolved];
3238 let error = Tx::rebuild_from_records(
3239 &intact,
3240 ConfigDefaults::default(),
3241 InMemoryRecordIndex::from_records(&intact),
3242 |context| {
3243 let mut step = plan(context)?;
3244 step.plan.push_str(" (drifted)");
3245 Ok(step)
3246 },
3247 )
3248 .expect_err("a drifted planner must not silently resume");
3249 assert_eq!(error.code, KernelFaultCode::RecordCorrupted);
3250 assert!(
3251 error.message.contains("step"),
3252 "the fault names the digest that disagreed: {}",
3253 error.message
3254 );
3255 }
3256
3257 #[test]
3258 fn a_rebuild_of_an_empty_journal_is_a_fresh_operation() {
3259 let rebuilt = Tx::rebuild_from_records(
3260 &[],
3261 ConfigDefaults::default(),
3262 InMemoryRecordIndex::new(),
3263 plan,
3264 )
3265 .expect("an empty journal is a legal starting point");
3266 assert_eq!(rebuilt.lifecycle(), OperationLifecycle::Created);
3267 assert_eq!(rebuilt.head(), None);
3268 }
3269
3270 #[test]
3273 fn a_replay_of_a_record_this_runtime_never_saw_demands_a_rebuild() {
3274 let mut source = transaction();
3275 let genesis = run(&mut source, &configure()).record;
3276
3277 let mut cold = KernelTransaction::<TestStep, _>::new(
3278 ConfigDefaults::default(),
3279 InMemoryRecordIndex::from_records(&[genesis]),
3280 );
3281 assert_eq!(
3282 fault_of(&cold.prepare(&configure(), plan)),
3283 KernelFaultCode::RecordCorrupted
3284 );
3285 }
3286
3287 #[test]
3288 fn a_poisoned_transaction_refuses_every_call() {
3289 let (mut tx, _, _) = started();
3290 let preparation = tx.prepare(&cancel_at("in-cancel", 1_700_000_003_000), plan);
3291 let token = preparation.token().unwrap().clone();
3292 tx.note_append_conflict(&token, None);
3293
3294 assert!(tx.is_poisoned());
3295 assert_eq!(
3296 tx.poison().map(|fault| fault.code),
3297 Some(KernelFaultCode::TransactionConflict)
3298 );
3299 assert_eq!(
3300 fault_of(&tx.prepare(&start_at("in-any", 1_700_000_009_000), plan)),
3301 KernelFaultCode::TransactionConflict
3302 );
3303 let boundary = CheckpointBoundary {
3304 through_step_seq: WireU64::new(1),
3305 covered_head: canonical_digest(b"whatever"),
3306 };
3307 assert_eq!(
3308 tx.note_checkpoint_acked(&boundary).unwrap_err().code,
3309 KernelFaultCode::TransactionConflict
3310 );
3311 }
3312}