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,
},
Failed {
error: String,
disposition: Disposition,
spend: crate::core::Spend,
},
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 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(),
},
));
}
#[must_use]
pub fn exhausted(&self) -> bool {
self.pos >= self.effects.len()
}
#[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));
let cursor = by_step.entry((step, r.body.phase)).or_default();
match r.kind() {
RecordKind::EffectStarted {
descriptor,
recovery,
..
} => {
cursor.effects.push((
key,
r.seq(),
EffectReplay::Orphan {
descriptor: Box::new(descriptor.clone()),
recovery: recovery.clone(),
},
));
}
RecordKind::EffectDone {
output,
source,
spend,
} => {
if let Some(slot) = cursor.effects.iter_mut().rev().find(|(k, _, _)| *k == key)
{
slot.2 = EffectReplay::Done {
output: output.clone(),
source: source.clone(),
spend: *spend,
};
}
}
RecordKind::EffectReconciled {
disposition,
output,
spend,
..
} => {
if let Some(slot) = cursor.effects.iter_mut().rev().find(|(k, _, _)| *k == key)
{
slot.2 = match (disposition, output) {
(Disposition::Landed, Some(output)) => EffectReplay::Done {
output: output.clone(),
source: None,
spend: *spend,
},
(d, _) => EffectReplay::Failed {
error: "resolved by reconciliation".to_owned(),
disposition: *d,
spend: crate::core::Spend::default(),
},
};
}
}
RecordKind::BudgetRefused { limit, used } => {
cursor.effects.push((
key,
r.seq(),
EffectReplay::Refused {
limit: limit.clone(),
used: used.clone(),
},
));
}
RecordKind::PolicyDenied {
reason,
action,
resource,
} => {
cursor.effects.push((
key,
r.seq(),
EffectReplay::Denied {
reason: reason.clone(),
action: action.clone(),
resource: resource.clone(),
},
));
}
RecordKind::Released { .. } => cursor.record_release(key, r.seq()),
RecordKind::EffectFailed {
error,
disposition,
spend,
} => {
if let Some(slot) = cursor.effects.iter_mut().rev().find(|(k, _, _)| *k == key)
{
slot.2 = EffectReplay::Failed {
error: error.clone(),
disposition: *disposition,
spend: *spend,
};
}
}
_ => {}
}
}
Self { by_step }
}
#[must_use]
pub fn len(&self) -> usize {
self.by_step.values().map(|c| c.effects.len()).sum()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
#[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(|c| c.pos >= c.effects.len())
}
#[must_use]
pub fn peek_is(&self, step: StepId, phase: Phase, key: EffectKey) -> bool {
self.by_step
.get(&(step, phase))
.and_then(|c| c.effects.get(c.pos))
.is_some_and(|(k, _, _)| *k == key)
}
pub fn next(
&mut self,
step: StepId,
phase: Phase,
recomputed: EffectKey,
) -> Result<Option<EffectReplay>, StepError> {
let Some(cursor) = self.by_step.get_mut(&(step, phase)) else {
return Ok(None);
};
let Some((expected, seq, replay)) = cursor.effects.get(cursor.pos) else {
return Ok(None);
};
if *expected != recomputed {
return Err(StepError::NonDeterminism {
seq: *seq,
expected: *expected,
actual: recomputed,
});
}
cursor.pos += 1;
Ok(Some(replay.clone()))
}
}
#[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);
#[test]
fn replays_a_completed_effect_without_performing_it() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectDone {
output: json!("recorded"),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert_eq!(
cur.next(S0, Phase::Forward, key(1)).unwrap(),
Some(EffectReplay::Done {
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 {
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert!(matches!(
cur.next(S0, Phase::Forward, 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 {
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
(S0, key(2), started()),
(
S0,
key(2),
RecordKind::EffectDone {
output: json!(2),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert!(
cur.next(S0, Phase::Forward, 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 {
output: json!("b"),
source: None,
spend: crate::core::Spend::default(),
},
),
(
S0,
key(1),
RecordKind::EffectDone {
output: json!("a"),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert_eq!(
cur.next(S0, Phase::Forward, key(1)).unwrap(),
Some(EffectReplay::Done {
output: json!("a"),
source: None,
spend: crate::core::Spend::default()
})
);
assert_eq!(
cur.next(S1, Phase::Forward, key(10)).unwrap(),
Some(EffectReplay::Done {
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 {
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
(S1, key(10), started()),
(
S1,
key(10),
RecordKind::EffectDone {
output: json!(2),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
cur.next(S0, Phase::Forward, 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.is_empty());
}
#[test]
fn exhausted_cursor_yields_none_so_execution_continues_live() {
let recs = records(vec![
(S0, key(1), started()),
(
S0,
key(1),
RecordKind::EffectDone {
output: json!(1),
source: None,
spend: crate::core::Spend::default(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
cur.next(S0, Phase::Forward, key(1)).unwrap();
assert_eq!(
cur.next(S0, Phase::Forward, 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 cur.next(S0, Phase::Forward, 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(),
},
),
]);
let mut cur = ReplayCursor::from_records(&recs);
assert_eq!(
cur.next(S0, Phase::Forward, key(1)).unwrap(),
Some(EffectReplay::Failed {
error: "boom".into(),
disposition: Disposition::DidNotHappen,
spend: crate::core::Spend::default(),
})
);
}
}