use std::collections::BTreeMap;
use crate::core::{
Disposition, EffectDescriptor, EffectKey, Phase, Recovery, Seq, StepError, StepId,
};
use super::{Record, RecordKind};
#[derive(Debug, Clone, PartialEq)]
pub enum EffectReplay {
Done {
output: serde_json::Value,
source: Option<String>,
spend: crate::core::Spend,
declared: crate::core::DeclaredOutput,
},
Failed {
error: String,
disposition: Disposition,
spend: crate::core::Spend,
permanent: bool,
},
Refused { limit: String, used: String },
Denied {
reason: String,
action: String,
resource: String,
},
Orphan {
descriptor: Box<EffectDescriptor>,
recovery: Recovery,
},
}
#[derive(Debug, Clone, Default)]
pub struct StepCursor {
effects: Vec<(EffectKey, Seq, EffectReplay)>,
pos: usize,
}
impl StepCursor {
fn settle(&mut self, key: EffectKey, state: EffectReplay) {
if let Some(slot) = self.effects.iter_mut().rev().find(|(k, _, _)| *k == key) {
slot.2 = state;
}
}
fn apply(&mut self, key: EffectKey, seq: Seq, kind: &RecordKind) {
match kind {
RecordKind::EffectStarted {
descriptor,
recovery,
..
} => {
self.effects.push((
key,
seq,
EffectReplay::Orphan {
descriptor: Box::new(descriptor.clone()),
recovery: recovery.clone(),
},
));
}
RecordKind::EffectDone {
output,
source,
spend,
declared,
} => {
self.settle(
key,
EffectReplay::Done {
output: output.clone(),
source: source.clone(),
spend: *spend,
declared: *declared,
},
);
}
RecordKind::EffectReconciled {
disposition,
output,
spend,
declared,
..
} => {
self.settle(
key,
Self::reconciled(*disposition, output.as_ref(), *spend, *declared),
);
}
RecordKind::BudgetRefused { limit, used } => {
self.effects.push((
key,
seq,
EffectReplay::Refused {
limit: limit.clone(),
used: used.clone(),
},
));
}
RecordKind::BudgetReadmitted { .. } => {
if let Some(pos) = self
.effects
.iter()
.rposition(|(k, _, r)| *k == key && matches!(r, EffectReplay::Refused { .. }))
{
self.effects.remove(pos);
}
}
RecordKind::PolicyDenied {
reason,
action,
resource,
} => {
self.effects.push((
key,
seq,
EffectReplay::Denied {
reason: reason.clone(),
action: action.clone(),
resource: resource.clone(),
},
));
}
RecordKind::Released { .. } => self.record_release(key, seq),
RecordKind::EffectFailed {
error,
disposition,
spend,
permanent,
} => {
self.settle(
key,
EffectReplay::Failed {
error: error.clone(),
disposition: *disposition,
spend: *spend,
permanent: *permanent,
},
);
}
_ => {}
}
}
fn reconciled(
disposition: Disposition,
output: Option<&serde_json::Value>,
spend: crate::core::Spend,
declared: Option<crate::core::DeclaredOutput>,
) -> EffectReplay {
match (disposition, output) {
(Disposition::Landed, Some(output)) => EffectReplay::Done {
output: output.clone(),
source: None,
spend,
declared: declared.unwrap_or_else(crate::core::DeclaredOutput::untrusted),
},
(disposition, _) => EffectReplay::Failed {
error: "resolved by reconciliation".to_owned(),
disposition,
spend: crate::core::Spend::default(),
permanent: false,
},
}
}
fn record_release(&mut self, key: EffectKey, seq: Seq) {
self.effects.push((
key,
seq,
EffectReplay::Done {
output: serde_json::Value::Null,
source: None,
spend: crate::core::Spend::default(),
declared: crate::core::DeclaredOutput::untrusted(),
},
));
}
#[must_use]
pub fn exhausted(&self) -> bool {
self.pos >= self.effects.len()
}
#[must_use]
pub fn first_unconsumed(&self) -> Option<EffectKey> {
self.effects.get(self.pos).map(|(k, _, _)| *k)
}
#[must_use]
pub fn peek_is(&self, key: EffectKey) -> bool {
self.effects
.get(self.pos)
.is_some_and(|(k, _, _)| *k == key)
}
pub fn next(&mut self, recomputed: EffectKey) -> Result<Option<EffectReplay>, StepError> {
let Some((expected, seq, replay)) = self.effects.get(self.pos) else {
return Ok(None);
};
if *expected != recomputed {
return Err(StepError::NonDeterminism {
seq: *seq,
expected: *expected,
actual: recomputed,
});
}
self.pos += 1;
Ok(Some(replay.clone()))
}
}
#[derive(Debug, Clone, Default)]
pub struct ReplayCursor {
by_step: BTreeMap<(StepId, Phase), StepCursor>,
}
impl ReplayCursor {
#[must_use]
pub fn from_records(records: &[Record]) -> Self {
let mut by_step: BTreeMap<(StepId, Phase), StepCursor> = BTreeMap::new();
for r in records {
let Some(key) = r.effect_key() else { continue };
let step = r.body.step.unwrap_or(StepId(0));
by_step
.entry((step, r.body.phase))
.or_default()
.apply(key, r.seq(), r.kind());
}
Self { by_step }
}
#[must_use]
pub fn take(&mut self, step: StepId, phase: Phase) -> StepCursor {
self.by_step.remove(&(step, phase)).unwrap_or_default()
}
pub fn restore(&mut self, step: StepId, phase: Phase, cursor: StepCursor) {
self.by_step.insert((step, phase), cursor);
}
#[must_use]
pub fn exhausted(&self, step: StepId, phase: Phase) -> bool {
self.by_step
.get(&(step, phase))
.is_none_or(StepCursor::exhausted)
}
#[must_use]
pub fn fully_exhausted(&self) -> bool {
self.by_step.values().all(StepCursor::exhausted)
}
#[must_use]
pub fn first_unconsumed(&self) -> Option<(StepId, Phase, EffectKey)> {
self.by_step.iter().find_map(|((step, phase), cursor)| {
cursor.first_unconsumed().map(|k| (*step, *phase, k))
})
}
#[must_use]
pub fn unconsumed_in(&self, step: StepId, phase: Phase) -> Option<EffectKey> {
self.by_step
.get(&(step, phase))
.and_then(StepCursor::first_unconsumed)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::{Digest, RunId};
use crate::journal::Append;
use serde_json::json;
fn key(n: u8) -> EffectKey {
EffectKey::from_hex(&Digest::of(&[n]).to_hex()).unwrap()
}
fn desc() -> EffectDescriptor {
EffectDescriptor::nullary("test.effect")
}
fn records(entries: Vec<(StepId, EffectKey, RecordKind)>) -> Vec<Record> {
let run = RunId::generate();
let mut prev = Digest::ZERO;
let mut out = Vec::new();
for (i, (step, k, kind)) in entries.into_iter().enumerate() {
let a = Append::new(run, kind).effect(k).step(step);
let r = Record::seal(a.into_body(i as u64 + 1, 1), prev).unwrap();
prev = r.hash;
out.push(r);
}
out
}
fn started() -> RecordKind {
RecordKind::EffectStarted {
descriptor: desc(),
recovery: Recovery::Retry,
mutates: false,
attempt: 1,
backoff_ms: 0,
outbound_label: None,
}
}
const S0: StepId = StepId(0);
const S1: StepId = StepId(1);
fn next(
cur: &mut ReplayCursor,
step: StepId,
key: EffectKey,
) -> Result<Option<EffectReplay>, StepError> {
let mut slice = cur.take(step, Phase::Forward);
let out = slice.next(key);
cur.restore(step, Phase::Forward, slice);
out
}
#[test]
fn replays_a_completed_effect_without_performing_it() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!("recorded"),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert_eq!(
next(&mut cur, S0, key(1)).unwrap(),
Some(EffectReplay::Done {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!("recorded"),
source: None,
spend: crate::core::Spend::default()
})
);
assert!(cur.exhausted(S0, Phase::Forward));
}
#[test]
fn divergent_key_is_rejected() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert!(matches!(
next(&mut cur, S0, key(99)).unwrap_err(),
StepError::NonDeterminism { .. }
));
}
#[test]
fn reordered_effects_within_a_step_are_rejected() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
(S0, key(2), started()),
(
S0,
key(2),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!(2),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert!(
next(&mut cur, S0, key(2)).is_err(),
"order within a step is verified"
);
}
#[test]
fn steps_replay_independently_of_journal_interleaving() {
let recs = records(vec![
(S0, key(1), started()),
(S1, key(10), started()),
(
S1,
key(10),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!("b"),
source: None,
spend: crate::core::Spend::default(),
},
),
(
S0,
key(1),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!("a"),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert_eq!(
next(&mut cur, S0, key(1)).unwrap(),
Some(EffectReplay::Done {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!("a"),
source: None,
spend: crate::core::Spend::default()
})
);
assert_eq!(
next(&mut cur, S1, key(10)).unwrap(),
Some(EffectReplay::Done {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!("b"),
source: None,
spend: crate::core::Spend::default()
})
);
assert!(cur.fully_exhausted());
}
#[test]
fn exhaustion_is_per_step() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
(S1, key(10), started()),
(
S1,
key(10),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!(2),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
next(&mut cur, S0, key(1)).unwrap();
assert!(cur.exhausted(S0, Phase::Forward));
assert!(
!cur.exhausted(S1, Phase::Forward),
"step 1 still has history to replay"
);
assert!(!cur.fully_exhausted());
}
#[test]
fn an_unknown_step_has_no_history() {
let cur = ReplayCursor::from_records(&[]);
assert!(cur.exhausted(StepId(7), Phase::Forward));
assert!(cur.fully_exhausted());
assert!(cur.first_unconsumed().is_none());
}
#[test]
fn exhausted_cursor_yields_none_so_execution_continues_live() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
next(&mut cur, S0, key(1)).unwrap();
assert_eq!(
next(&mut cur, S0, key(2)).unwrap(),
None,
"resume runs the rest live"
);
}
#[test]
fn orphaned_start_is_surfaced_with_its_recovery_mode() {
let recs = records(vec![(S0, key(1), started())]);
let mut cur = ReplayCursor::from_records(&recs);
match next(&mut cur, S0, key(1)).unwrap() {
Some(EffectReplay::Orphan { recovery, .. }) => {
assert!(matches!(recovery, Recovery::Retry));
}
other => panic!("expected orphan, got {other:?}"),
}
}
#[test]
fn failure_is_part_of_history() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectFailed {
error: "boom".into(),
disposition: Disposition::DidNotHappen,
spend: crate::core::Spend::default(),
permanent: false,
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert_eq!(
next(&mut cur, S0, key(1)).unwrap(),
Some(EffectReplay::Failed {
error: "boom".into(),
disposition: Disposition::DidNotHappen,
spend: crate::core::Spend::default(),
permanent: false,
})
);
}
}