use std::collections::BTreeMap;
use crate::core::{
Disposition, Doubt, EffectDescriptor, EffectKey, Phase, Recovery, Seq, StepError, StepId,
Undecided,
};
use super::{Record, RecordKind};
#[derive(Debug, Clone, PartialEq)]
pub enum EffectReplay {
Done {
output: serde_json::Value,
source: Option<String>,
by: Option<crate::core::Operator>,
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 { recovery: Recovery },
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Attempted {
kind: String,
attempt: u32,
}
impl EffectReplay {
#[must_use]
pub const fn spend(&self) -> crate::core::Spend {
match self {
Self::Done { spend, .. } | Self::Failed { spend, .. } => *spend,
Self::Refused { .. } | Self::Denied { .. } | Self::Orphan { .. } => {
crate::core::Spend::ZERO
}
}
}
fn add_spend(&mut self, extra: crate::core::Spend) {
match self {
Self::Done { spend, .. } | Self::Failed { spend, .. } => *spend += extra,
Self::Refused { .. } | Self::Denied { .. } | Self::Orphan { .. } => {}
}
}
}
#[derive(Debug, Clone, Default)]
pub struct StepCursor {
effects: Vec<Journaled>,
pos: usize,
diverged: Option<Divergence>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct Divergence {
pub step: StepId,
pub phase: Phase,
pub recorded: Option<EffectKey>,
pub recomputed: Option<EffectKey>,
pub kind: Option<String>,
pub detail: String,
}
impl Divergence {
#[must_use]
pub fn of(step: StepId, phase: Phase, error: &StepError, history: &StepCursor) -> Option<Self> {
match error {
StepError::NonDeterminism {
seq,
expected,
actual,
detail,
} => Some(Self {
step,
phase,
recorded: Some(*expected),
recomputed: Some(*actual),
kind: history.kind_at(*seq),
detail: detail.clone(),
}),
StepError::ReplayOverrun { actual, kind } => Some(Self {
step,
phase,
recorded: None,
recomputed: Some(*actual),
kind: Some(kind.clone()),
detail: format!(
"this build asks for `{kind}`, past the end of the recorded history"
),
}),
_ => None,
}
}
#[must_use]
pub fn unconsumed(step: StepId, phase: Phase, missing: &Unconsumed) -> Self {
Self {
step,
phase,
recorded: Some(missing.key),
recomputed: None,
kind: missing.kind.clone(),
detail: format!("the record holds {missing}, which this build never asks for"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Unconsumed {
pub key: EffectKey,
pub kind: Option<String>,
}
impl std::fmt::Display for Unconsumed {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match &self.kind {
Some(kind) => write!(f, "`{kind}` ({})", self.key),
None => write!(f, "a refused call ({})", self.key),
}
}
}
#[derive(Debug, Clone)]
struct Journaled {
key: EffectKey,
seq: Seq,
what: Option<Attempted>,
replay: EffectReplay,
}
impl StepCursor {
fn settle(&mut self, key: EffectKey, state: EffectReplay) {
if let Some(slot) = self.effects.iter_mut().rev().find(|e| e.key == key) {
let carried = slot.replay.spend();
slot.replay = state;
slot.replay.add_spend(carried);
}
}
fn announce(
&mut self,
key: EffectKey,
seq: Seq,
descriptor: &EffectDescriptor,
recovery: &Recovery,
attempt: u32,
) {
self.effects.push(Journaled {
key,
seq,
what: Some(Attempted {
kind: descriptor.kind.clone(),
attempt,
}),
replay: EffectReplay::Orphan {
recovery: recovery.clone(),
},
});
}
fn apply(&mut self, key: EffectKey, seq: Seq, kind: &RecordKind) {
match kind {
RecordKind::EffectStarted {
descriptor,
recovery,
attempt,
..
} => self.announce(key, seq, descriptor, recovery, *attempt),
RecordKind::EffectDone {
output,
source,
by,
spend,
declared,
} => {
self.settle(
key,
EffectReplay::Done {
output: output.clone(),
source: source.clone(),
by: by.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(Journaled {
key,
seq,
what: None,
replay: EffectReplay::Refused {
limit: limit.clone(),
used: used.clone(),
},
});
}
RecordKind::BudgetReadmitted { .. } => {
if let Some(pos) = self
.effects
.iter()
.rposition(|e| e.key == key && matches!(e.replay, EffectReplay::Refused { .. }))
{
self.effects.remove(pos);
}
}
RecordKind::PolicyDenied {
reason,
action,
resource,
} => {
self.effects.push(Journaled {
key,
seq,
what: None,
replay: 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,
by: 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(Journaled {
key,
seq,
what: None,
replay: EffectReplay::Done {
output: serde_json::Value::Null,
source: None,
by: 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<Unconsumed> {
self.effects.get(self.pos).map(|e| Unconsumed {
key: e.key,
kind: e.what.as_ref().map(|w| w.kind.clone()),
})
}
#[must_use]
pub fn kind_at(&self, seq: Seq) -> Option<String> {
self.effects
.iter()
.find(|e| e.seq == seq)
.and_then(|e| e.what.as_ref().map(|w| w.kind.clone()))
}
pub fn note_divergence(&mut self, divergence: Divergence) {
self.diverged.get_or_insert(divergence);
}
#[must_use]
pub fn peek_is(&self, key: EffectKey) -> bool {
self.effects.get(self.pos).is_some_and(|e| e.key == key)
}
pub fn next(
&mut self,
recomputed: EffectKey,
asked: &EffectDescriptor,
attempt: u32,
) -> Result<Option<EffectReplay>, StepError> {
let Some(entry) = self.effects.get(self.pos) else {
return Ok(None);
};
if entry.key != recomputed {
return Err(StepError::NonDeterminism {
seq: entry.seq,
expected: entry.key,
actual: recomputed,
detail: entry.diverged_from(asked, attempt),
});
}
self.pos += 1;
Ok(Some(entry.replay.clone()))
}
}
impl Journaled {
fn diverged_from(&self, asked: &EffectDescriptor, attempt: u32) -> String {
let Some(what) = &self.what else {
return format!(
"history announced no call at this position — it recorded a refusal, which \
stops a call before it starts — and this build asks for `{}`",
asked.kind
);
};
if what.kind != asked.kind {
return format!(
"history performed `{}` here and this build asks for `{}` — the code takes a \
different path than the code that wrote this journal",
what.kind, asked.kind
);
}
if what.attempt != attempt {
return format!(
"both perform `{}`, at different attempts — history recorded attempt {} and \
this build is on attempt {}, so the retry decisions differ rather than the \
call",
what.kind, what.attempt, attempt
);
}
format!(
"both perform `{}` at attempt {}, so the arguments differ — they are sealed with \
the record, so comparing them needs a reader holding the key",
what.kind, what.attempt
)
}
}
#[derive(Debug, Clone, Default)]
pub struct ReplayCursor {
by_step: BTreeMap<(StepId, Phase), StepCursor>,
noted: Option<Divergence>,
}
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,
noted: None,
}
}
#[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, Unconsumed)> {
self.by_step.iter().find_map(|((step, phase), cursor)| {
cursor.first_unconsumed().map(|u| (*step, *phase, u))
})
}
pub fn note_divergence(&mut self, divergence: Divergence) {
self.noted.get_or_insert(divergence);
}
#[must_use]
pub fn divergence(&self) -> Option<&Divergence> {
self.noted
.as_ref()
.or_else(|| self.by_step.values().find_map(|c| c.diverged.as_ref()))
}
#[must_use]
pub fn unconsumed_in(&self, step: StepId, phase: Phase) -> Option<Unconsumed> {
self.by_step
.get(&(step, phase))
.and_then(StepCursor::first_unconsumed)
}
}
#[must_use]
pub fn undecided_effects(records: &[Record]) -> Vec<Undecided> {
let mut open: BTreeMap<EffectKey, Pending> = BTreeMap::new();
let mut doubts: BTreeMap<EffectKey, Pending> = BTreeMap::new();
for r in records {
let Some(key) = r.effect_key() else { continue };
match r.kind() {
RecordKind::EffectStarted {
mutates: true,
descriptor,
..
} => {
if let Some(step) = r.body.step {
open.insert(
key,
Pending {
step,
phase: r.body.phase,
descriptor: descriptor.clone(),
seq: r.body.seq,
},
);
}
}
RecordKind::EffectDone { .. } => {
if let Some(settled) = open.remove(&key) {
doubts.retain(|_, p| {
!(p.step == settled.step && p.descriptor == settled.descriptor)
});
}
}
RecordKind::EffectFailed { disposition, .. } => {
if let Some(pending) = open.remove(&key)
&& *disposition == Disposition::InDoubt
{
doubts.insert(key, pending);
}
}
RecordKind::EffectReconciled { disposition, .. } => {
let settled = open.remove(&key);
match disposition {
Disposition::Landed | Disposition::DidNotHappen => {
doubts.remove(&key);
}
Disposition::InDoubt => {
if let Some(pending) = settled {
doubts.insert(key, pending);
}
}
}
}
_ => {}
}
}
let mut out: Vec<(StepId, Seq, Undecided)> = open
.into_iter()
.map(|(effect, p)| p.into_undecided(effect, Doubt::Announced))
.chain(
doubts
.into_iter()
.map(|(effect, p)| p.into_undecided(effect, Doubt::Inconclusive)),
)
.collect();
out.sort_by_key(|(step, seq, _)| (*step, *seq));
out.into_iter().map(|(_, _, u)| u).collect()
}
struct Pending {
step: StepId,
phase: Phase,
descriptor: EffectDescriptor,
seq: Seq,
}
impl Pending {
fn into_undecided(self, effect: EffectKey, doubt: Doubt) -> (StepId, Seq, Undecided) {
(
self.step,
self.seq,
Undecided {
effect,
step: self.step,
phase: self.phase,
kind: self.descriptor.kind,
doubt,
},
)
}
}
#[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,
outbound_bytes: None,
}
}
const S0: StepId = StepId(0);
const S1: StepId = StepId(1);
fn next(
cur: &mut ReplayCursor,
step: StepId,
key: EffectKey,
) -> Result<Option<EffectReplay>, StepError> {
next_asked(cur, step, key, &desc())
}
fn next_asked(
cur: &mut ReplayCursor,
step: StepId,
key: EffectKey,
asked: &EffectDescriptor,
) -> Result<Option<EffectReplay>, StepError> {
let mut slice = cur.take(step, Phase::Forward);
let out = slice.next(key, asked, 1);
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,
by: 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,
by: 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,
by: 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 a_divergence_says_which_of_the_three_things_moved() {
let recs = records(vec![(S0, key(1), started())]);
let elsewhere = EffectDescriptor::nullary("test.somewhere-else");
let msg = next_asked(
&mut ReplayCursor::from_records(&recs),
S0,
key(99),
&elsewhere,
)
.expect_err("diverges")
.to_string();
assert!(
msg.contains("test.effect") && msg.contains("test.somewhere-else"),
"a changed call must name both sides: {msg}"
);
let msg = next_asked(&mut ReplayCursor::from_records(&recs), S0, key(99), &desc())
.expect_err("diverges")
.to_string();
assert!(
msg.contains("the arguments differ") && msg.contains("sealed"),
"the remaining input has to be named, and named as unquotable: {msg}"
);
let mut slice = ReplayCursor::from_records(&recs).take(S0, Phase::Forward);
let msg = slice
.next(key(99), &desc(), 2)
.expect_err("diverges")
.to_string();
assert!(
msg.contains("attempt 1") && msg.contains("attempt 2"),
"a differing attempt is a different finding from a differing argument: {msg}"
);
let refused = records(vec![(
S0,
key(1),
RecordKind::BudgetRefused {
limit: "spend".into(),
used: "1".into(),
},
)]);
let msg = next_asked(
&mut ReplayCursor::from_records(&refused),
S0,
key(99),
&desc(),
)
.expect_err("diverges")
.to_string();
assert!(
msg.contains("announced no call"),
"a refusal has no descriptor and the reason must say so: {msg}"
);
}
#[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,
by: 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,
by: 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,
by: None,
spend: crate::core::Spend::default(),
},
),
(
S0,
key(1),
RecordKind::EffectDone {
declared: crate::core::DeclaredOutput::untrusted(),
output: json!("a"),
source: None,
by: 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,
by: 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,
by: 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,
by: 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,
by: 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,
by: 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,
})
);
}
}