use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{AnyEffect, Effect, GroupOutcome, StepError, Tainted};
use crate::journal::RecordKind;
use crate::runtime::StepCtx;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Invariant {
pub what: String,
pub holds: bool,
}
impl Invariant {
pub fn new(what: impl Into<String>, holds: bool) -> Self {
Self {
what: what.into(),
holds,
}
}
}
struct Reversal {
resource: String,
kind: String,
undo: Box<dyn AnyEffect>,
}
pub(crate) struct AtomicMember {
resource: String,
inner: std::sync::Arc<dyn crate::journal::AtomicResource>,
}
struct Deferred {
resource: String,
effect: Box<dyn AnyEffect>,
}
pub(crate) struct OpenGroup {
pub(crate) name: String,
resources: Vec<String>,
reversals: Vec<Reversal>,
deferred: Vec<Deferred>,
atomic: Vec<AtomicMember>,
}
impl std::fmt::Debug for OpenGroup {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("OpenGroup")
.field("name", &self.name)
.field("resources", &self.resources)
.field("reversible", &self.reversals.len())
.field("deferred", &self.deferred.len())
.field("atomic", &self.atomic.len())
.finish()
}
}
#[derive(Debug)]
pub struct EffectGroup<'g, 'c> {
cx: &'g mut StepCtx<'c>,
}
impl<'g, 'c> EffectGroup<'g, 'c> {
pub(crate) fn new(cx: &'g mut StepCtx<'c>) -> Self {
Self { cx }
}
pub async fn reversible<E, U, F>(
&mut self,
resource: &str,
effect: E,
undo: F,
) -> Result<Tainted<E::Output>, StepError>
where
E: Effect,
U: AnyEffect + 'static,
F: FnOnce(&E::Output) -> U,
{
self.check_footprint(resource)?;
let kind = effect.descriptor().kind;
let out = self.cx.effect_as_member(effect).await?;
let undo = undo(out.peek());
if let Some(detail) = Self::check_dispatchable(&undo) {
let group = self.group_name();
self.cx
.settle_open_group(GroupOutcome::Quarantined, Some(&detail))
.await?;
return Err(StepError::GroupUnsettled { group, detail });
}
let reversal = Reversal {
resource: resource.to_owned(),
kind,
undo: Box::new(undo),
};
self.open()?.reversals.push(reversal);
Ok(out)
}
pub async fn read<E: Effect>(
&mut self,
resource: &str,
effect: E,
) -> Result<Tainted<E::Output>, StepError> {
self.check_footprint(resource)?;
if effect.mutates() {
return Err(StepError::GroupFootprint {
group: self.open()?.name.clone(),
detail: format!(
"effect '{}' mutates, so it cannot be a group read — declare it \
reversible with the call that undoes it, or deferred so an aborted \
group never runs it",
effect.descriptor().kind
),
});
}
self.cx.effect_as_member(effect).await
}
pub fn deferred<E: AnyEffect + 'static>(
&mut self,
resource: &str,
effect: E,
) -> Result<(), StepError> {
self.check_footprint(resource)?;
if let Some(detail) = Self::check_dispatchable(&effect) {
return Err(StepError::GroupFootprint {
group: self.group_name(),
detail,
});
}
let member = Deferred {
resource: resource.to_owned(),
effect: Box::new(effect),
};
self.open()?.deferred.push(member);
Ok(())
}
pub fn atomic(
&mut self,
resource: &str,
member: std::sync::Arc<dyn crate::journal::AtomicResource>,
) -> Result<(), StepError> {
self.check_footprint(resource)?;
if !self.cx.store_is_atomic() {
return Err(StepError::GroupFootprint {
group: self.group_name(),
detail: format!(
"'{}' must commit with the journal, and this store has no transaction \
a resource can join — an embedded backend has no notion of a foreign \
table, so the capability is absent rather than failing",
member.descriptor().kind
),
});
}
let member = AtomicMember {
resource: resource.to_owned(),
inner: member,
};
self.open()?.atomic.push(member);
Ok(())
}
pub async fn commit(self, invariants: &[Invariant]) -> Result<Vec<Tainted<Value>>, StepError> {
if let Some(broken) = invariants.iter().find(|i| !i.holds) {
let what = broken.what.clone();
self.cx
.abort_open_group(&format!("invariant: {what}"))
.await?;
return Err(StepError::GroupAborted { what });
}
let (name, deferred, atomic) = {
let open = self.cx.open_group_mut().ok_or_else(no_group)?;
(
open.name.clone(),
std::mem::take(&mut open.deferred),
std::mem::take(&mut open.atomic),
)
};
let touched = atomic
.iter()
.map(|m| m.resource.as_str())
.collect::<std::collections::BTreeSet<_>>()
.into_iter()
.collect::<Vec<_>>()
.join(", ");
let atomic_committed = !atomic.is_empty();
if atomic_committed && let Err(e) = self.cx.commit_atomic(&name, atomic).await {
if matches!(
&e,
StepError::Store(crate::core::StoreError::CommitUnknown { .. })
) {
let detail = format!(
"the atomic members on [{touched}] may have committed — the \
acknowledgement was lost ({e}); aborting would claim 'taken \
back whole' over a write that may stand"
);
self.cx
.settle_open_group(GroupOutcome::Quarantined, Some(&detail))
.await?;
return Err(StepError::GroupUnsettled {
group: name,
detail,
});
}
let what = format!("the atomic members on [{touched}] did not commit: {e}");
self.cx.abort_open_group(&what).await?;
return Err(StepError::GroupAborted { what });
}
let mut outputs = Vec::with_capacity(deferred.len());
for member in deferred {
let kind = AnyEffect::descriptor(&*member.effect).kind;
let resource = member.resource.clone();
match self.cx.effect_as_member(member.effect).await {
Ok(v) => outputs.push(v),
Err(e) if outputs.is_empty() && !atomic_committed && !may_have_externalised(&e) => {
let what = format!("deferred member '{kind}' on '{resource}' failed: {e}");
self.cx.abort_open_group(&what).await?;
return Err(StepError::GroupAborted { what });
}
Err(e) => {
let landed = if atomic_committed {
format!(
"{} deferred and the atomic members on [{touched}]",
outputs.len()
)
} else {
outputs.len().to_string()
};
let detail = format!(
"deferred member '{kind}' on '{resource}' failed after {landed} \
landed: {e} — reversing now would undo everything except the \
thing that actually happened",
);
self.cx
.settle_open_group(GroupOutcome::Quarantined, Some(&detail))
.await?;
return Err(StepError::GroupUnsettled {
group: name,
detail,
});
}
}
}
self.cx
.settle_open_group(GroupOutcome::Committed, None)
.await?;
Ok(outputs)
}
pub async fn abort(self, why: &str) -> Result<(), StepError> {
self.cx.abort_open_group(why).await
}
fn open(&mut self) -> Result<&mut OpenGroup, StepError> {
self.cx.open_group_mut().ok_or_else(no_group)
}
fn check_dispatchable(effect: &dyn AnyEffect) -> Option<String> {
effect.sink_arguments().is_some().then(|| {
format!(
"'{}' binds the arguments it sends, so it must be dispatched with \
StepCtx::sink and cannot be a group member — a group has no labelled \
value to bind on its behalf",
effect.descriptor().kind
)
})
}
fn group_name(&self) -> String {
self.cx
.open_group()
.map_or_else(String::new, |g| g.name.clone())
}
fn check_footprint(&self, resource: &str) -> Result<(), StepError> {
let open = self.cx.open_group().ok_or_else(no_group)?;
if open.resources.iter().any(|r| r == resource) {
return Ok(());
}
Err(StepError::GroupFootprint {
group: open.name.clone(),
detail: format!(
"member touches '{resource}', which is not among the declared resources \
[{}] — a frontier over an undeclared resource is a frontier over nothing",
open.resources.join(", ")
),
})
}
}
fn no_group() -> StepError {
StepError::GroupUnsettled {
group: String::new(),
detail: "the group was already settled".to_owned(),
}
}
pub(crate) fn may_have_externalised(e: &StepError) -> bool {
match e {
StepError::Undecidable { .. } => true,
StepError::Unrecorded { disposition, .. } => {
*disposition != crate::core::Disposition::DidNotHappen
}
StepError::Effect(inner) => inner.disposition() != crate::core::Disposition::DidNotHappen,
_ => false,
}
}
impl<'a> StepCtx<'a> {
pub async fn group<'g, R, S>(
&'g mut self,
name: impl Into<String>,
resources: R,
) -> Result<EffectGroup<'g, 'a>, StepError>
where
R: IntoIterator<Item = S>,
S: Into<String>,
{
let name = name.into();
if let Some(open) = self.open_group() {
return Err(StepError::GroupFootprint {
group: name,
detail: format!(
"group '{}' is still open — groups do not nest, because a nested abort \
would have to decide whether it takes the outer group with it",
open.name
),
});
}
let resources: Vec<String> = resources.into_iter().map(Into::into).collect();
if resources.is_empty() {
return Err(StepError::GroupFootprint {
group: name,
detail: "a group must declare the resources it touches; an empty footprint \
admits every member and refuses none"
.to_owned(),
});
}
if let Some(recorded) = self.recorded_groups.get_mut(&name)
&& recorded.opened > 0
{
recorded.opened -= 1;
} else {
self.append(RecordKind::GroupOpened {
group: name.clone(),
resources: resources.clone(),
})
.await?;
}
self.set_open_group(OpenGroup {
name,
resources,
reversals: Vec::new(),
deferred: Vec::new(),
atomic: Vec::new(),
});
Ok(EffectGroup::new(self))
}
pub(crate) async fn abort_open_group(&mut self, why: &str) -> Result<(), StepError> {
let Some(open) = self.open_group_mut() else {
return Err(no_group());
};
let name = open.name.clone();
let reversals = std::mem::take(&mut open.reversals);
open.deferred.clear();
self.set_reversing(true);
let reversed = self.reverse_each(reversals).await;
self.set_reversing(false);
match reversed {
Ok(()) => {
self.settle_open_group(GroupOutcome::Aborted, Some(why))
.await
}
Err(detail) => {
self.settle_open_group(GroupOutcome::Quarantined, Some(&detail))
.await?;
Err(StepError::GroupUnsettled {
group: name,
detail,
})
}
}
}
async fn reverse_each(&mut self, reversals: Vec<Reversal>) -> Result<(), String> {
let total = reversals.len();
for (done, member) in reversals.into_iter().rev().enumerate() {
let Reversal {
resource,
kind,
undo,
} = member;
if let Err(e) = self.effect_as_member(undo).await {
return Err(format!(
"reversing '{kind}' on '{resource}' failed after {done} of {total}: {e}"
));
}
}
Ok(())
}
pub(crate) async fn settle_open_group(
&mut self,
outcome: GroupOutcome,
detail: Option<&str>,
) -> Result<(), StepError> {
let Some(open) = self.take_open_group() else {
return Err(no_group());
};
if let Some(recorded) = self.recorded_groups.get_mut(&open.name)
&& recorded.settled > 0
{
recorded.settled -= 1;
return Ok(());
}
self.append(RecordKind::GroupSettled {
group: open.name,
outcome,
detail: detail.map(ToOwned::to_owned),
})
.await
}
}
struct GroupCommit {
run: crate::core::RunId,
step: crate::core::StepId,
phase: crate::core::Phase,
case: Option<crate::core::CaseId>,
members: Vec<(crate::core::EffectKey, AtomicMember)>,
}
impl GroupCommit {
fn stamp(&self, kind: RecordKind) -> crate::journal::Append {
let mut a = crate::journal::Append::new(self.run, kind)
.step(self.step)
.phase(self.phase);
if let Some(c) = self.case {
a = a.case(c);
}
a
}
}
#[async_trait::async_trait]
impl crate::journal::AtomicWork for GroupCommit {
async fn run(
&self,
tx: &dyn crate::journal::AtomicTx,
) -> Result<Vec<crate::journal::Append>, crate::core::EffectError> {
let mut appends = Vec::with_capacity(self.members.len() * 2 + 1);
for (key, member) in &self.members {
let descriptor = member.inner.descriptor();
appends.push(
self.stamp(RecordKind::EffectStarted {
descriptor: descriptor.clone(),
recovery: crate::core::Recovery::Retry,
mutates: true,
attempt: 1,
backoff_ms: 0,
outbound_label: None,
})
.effect(*key),
);
let output = member.inner.apply(tx).await?;
appends.push(
self.stamp(RecordKind::EffectDone {
output,
source: None,
spend: crate::core::Spend::default(),
declared: crate::core::DeclaredOutput::untrusted(),
})
.effect(*key),
);
}
Ok(appends)
}
}
impl StepCtx<'_> {
pub(crate) fn store_is_atomic(&self) -> bool {
self.journal().atomic().is_some()
}
pub(crate) async fn commit_atomic(
&mut self,
group: &str,
members: Vec<AtomicMember>,
) -> Result<(), StepError> {
let mut keyed = Vec::with_capacity(members.len());
for member in members {
let descriptor = member.inner.descriptor();
let key = self.next_effect_key(&descriptor);
if self.replaying() {
match self.cursor_next(key)? {
Some(crate::journal::EffectReplay::Done { .. }) => continue,
Some(
refusal @ (crate::journal::EffectReplay::Refused { .. }
| crate::journal::EffectReplay::Denied { .. }),
) => {
self.replayed_refusal(key, refusal).await?;
}
Some(other) => {
return Err(StepError::GroupUnsettled {
group: group.to_owned(),
detail: format!(
"atomic member '{}' replays {other:?}, which an atomic \
member cannot have written — its records commit with \
its transaction, so the recorded run dispatched this \
call through a different protocol class than this \
build does",
descriptor.kind
),
});
}
None if self.is_strict() => {
return Err(StepError::ReplayOverrun { actual: key });
}
None => {}
}
}
self.gate(key, &descriptor, true, None, None).await?;
keyed.push((key, member));
}
if keyed.is_empty() {
return Ok(());
}
let work = GroupCommit {
run: self.run_id(),
step: self.step_id(),
phase: self.phase_of(),
case: self.bound_case(),
members: keyed,
};
let atomic = self
.journal()
.atomic()
.ok_or_else(|| StepError::GroupFootprint {
group: group.to_owned(),
detail: "the store stopped offering a transaction between registration \
and commit"
.to_owned(),
})?;
atomic
.append_atomic(self.run_id(), self.epoch(), &work)
.await?;
Ok(())
}
}