use std::collections::hash_map::RandomState;
use std::fmt::Display;
use std::future::Future;
use std::hash::BuildHasher;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::{Duration, SystemTime};
use serde::Serialize;
use serde::de::DeserializeOwned;
use serde_json::{Value, json};
use tracing::{Instrument, Span, debug, field, info_span, warn};
use crate::approval::{ApprovalDecision, ApprovalProvider, ApprovalRequest, ErasedApproval};
use crate::clock::{Clock, SystemClock};
use crate::effect::{
EffectBuilder, EffectContext, EffectFailure, EffectOutcome, EffectSpec, Precondition,
};
use crate::error::RuntimeError;
use crate::failure::{Disposition, FailureClass};
#[cfg(feature = "fault-injection")]
use crate::fault::FaultInjector;
use crate::fault::FaultPoint;
use crate::fingerprint::fingerprint;
use crate::handler::{
CompensationSubmission, EffectHandler, Handler, Registered, Registry, Resume, Submission,
};
use crate::id::{EffectId, EffectName, WorkerId};
use crate::kind::EffectKind;
use crate::observer::{EffectObserver, Observation};
use crate::policy::{RiskPolicy, UnknownPlan};
use crate::redaction::{Field, Redactor};
use crate::retention::RetentionPolicy;
use crate::retry::RetryPolicy;
use crate::state::{EffectStatus, Transition};
use crate::store::{
EffectRecord, EffectStore, ErrorRecord, Lease, NewEffect, StoreError, TransitionRequest,
};
use crate::verification::{NotFoundReading, Verification, VerificationMode, Verifier};
const MAX_ROUNDS: usize = 4;
pub struct Runtime<S> {
inner: Arc<Inner<S>>,
}
struct Inner<S> {
store: S,
clock: Arc<dyn Clock>,
worker: WorkerId,
lease_ttl: Duration,
retry: RetryPolicy,
handlers: Registry<S>,
approval: Option<Arc<dyn ErasedApproval>>,
policy: RiskPolicy,
redactor: Option<Arc<dyn Redactor>>,
observers: Vec<Arc<dyn EffectObserver>>,
retention: RetentionPolicy,
#[cfg(feature = "fault-injection")]
faults: Option<Arc<FaultInjector>>,
}
impl<S> Clone for Runtime<S> {
fn clone(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
}
}
}
#[must_use]
pub struct RuntimeBuilder<S> {
store: S,
clock: Arc<dyn Clock>,
worker: Option<WorkerId>,
lease_ttl: Duration,
retry: RetryPolicy,
handlers: Registry<S>,
approval: Option<Arc<dyn ErasedApproval>>,
policy: RiskPolicy,
redactor: Option<Arc<dyn Redactor>>,
observers: Vec<Arc<dyn EffectObserver>>,
retention: RetentionPolicy,
#[cfg(feature = "fault-injection")]
faults: Option<Arc<FaultInjector>>,
}
impl<S: EffectStore> RuntimeBuilder<S> {
pub fn retention(mut self, policy: RetentionPolicy) -> Self {
self.retention = policy;
self
}
pub fn observer(mut self, observer: impl EffectObserver) -> Self {
self.observers.push(Arc::new(observer));
self
}
pub fn redactor(mut self, redactor: impl Redactor) -> Self {
self.redactor = Some(Arc::new(redactor));
self
}
pub fn risk_policy(mut self, policy: RiskPolicy) -> Self {
self.policy = policy;
self
}
pub fn approval_provider(mut self, provider: impl ApprovalProvider) -> Self {
self.approval = Some(Arc::new(provider));
self
}
pub fn register<H: EffectHandler>(mut self, handler: Handler<H>) -> Self {
assert!(
!self.handlers.contains_key(H::NAME),
"a handler is already registered for effect `{}`",
H::NAME
);
self.handlers.insert(H::NAME, Registered::new(handler));
self
}
pub fn clock(mut self, clock: impl Clock) -> Self {
self.clock = Arc::new(clock);
self
}
pub fn worker_id(mut self, worker: WorkerId) -> Self {
self.worker = Some(worker);
self
}
pub fn lease_ttl(mut self, ttl: Duration) -> Self {
self.lease_ttl = ttl.max(Duration::from_millis(3));
self
}
pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
self.retry = policy;
self
}
#[cfg(feature = "fault-injection")]
pub fn fault_injector(mut self, injector: Arc<FaultInjector>) -> Self {
self.faults = Some(injector);
self
}
pub fn build(self) -> Runtime<S> {
Runtime {
inner: Arc::new(Inner {
store: self.store,
clock: self.clock,
worker: self.worker.unwrap_or_else(WorkerId::random),
lease_ttl: self.lease_ttl,
retry: self.retry,
handlers: self.handlers,
approval: self.approval,
policy: self.policy,
redactor: self.redactor,
observers: self.observers,
retention: self.retention,
#[cfg(feature = "fault-injection")]
faults: self.faults,
}),
}
}
}
pub(crate) enum Interrupt {
LeaseLost,
Error(RuntimeError),
}
impl From<StoreError> for Interrupt {
fn from(error: StoreError) -> Self {
match error {
StoreError::LeaseLost => Self::LeaseLost,
other => Self::Error(other.into()),
}
}
}
impl<S: EffectStore> Runtime<S> {
pub fn new(store: S) -> Self {
Self::builder(store).build()
}
pub fn builder(store: S) -> RuntimeBuilder<S> {
RuntimeBuilder {
store,
clock: Arc::new(SystemClock),
worker: None,
lease_ttl: Duration::from_secs(30),
retry: RetryPolicy::default(),
handlers: Registry::new(),
approval: None,
policy: RiskPolicy::default(),
redactor: None,
observers: Vec::new(),
retention: RetentionPolicy::KEEP_ALL,
#[cfg(feature = "fault-injection")]
faults: None,
}
}
pub fn effect(&self, name: impl Into<String>, key: impl Display) -> EffectBuilder<S> {
EffectBuilder::new(self.clone(), name.into(), key.to_string())
}
pub fn submit<H: EffectHandler>(
&self,
key: impl Display,
input: H::Input,
) -> Submission<'_, S, H> {
Submission::new(self, key, input)
}
pub fn compensate<H: EffectHandler>(
&self,
key: impl Display,
) -> CompensationSubmission<'_, S, H> {
CompensationSubmission::new(self, key)
}
pub(crate) fn handler<H: EffectHandler>(&self) -> Option<Handler<H>> {
self.inner.handlers.get(H::NAME)?.typed::<H>()
}
pub(crate) fn resumer(&self, name: &str) -> Option<Resume<S>> {
self.inner.handlers.get(name).map(|r| Arc::clone(&r.resume))
}
pub fn store(&self) -> &S {
&self.inner.store
}
pub fn worker_id(&self) -> &WorkerId {
&self.inner.worker
}
pub async fn wait<T: DeserializeOwned>(
&self,
id: EffectId,
timeout: Duration,
) -> Result<EffectOutcome<T>, RuntimeError> {
let deadline = tokio::time::Instant::now() + timeout;
let mut pause = Duration::from_millis(10);
loop {
let record = self
.store()
.get(id)
.await?
.ok_or(StoreError::NotFound(id))?;
let busy = !settled(record.status) && record.live_lease_owner(self.now()).is_some();
let now = tokio::time::Instant::now();
if !busy {
return report(&record, None);
}
if now >= deadline {
return Ok(EffectOutcome::InProgress { id });
}
tokio::time::sleep(pause.min(deadline - now)).await;
pause = (pause * 2).min(Duration::from_millis(250));
}
}
pub(crate) fn default_retry(&self) -> RetryPolicy {
self.inner.retry
}
pub(crate) fn now(&self) -> SystemTime {
self.inner.clock.now()
}
pub(crate) fn retention(&self) -> RetentionPolicy {
self.inner.retention
}
pub(crate) fn lease_ttl(&self) -> Duration {
self.inner.lease_ttl
}
pub(crate) async fn with_lease<Fut: Future>(
&self,
lease: &Lease,
future: Fut,
) -> Result<Fut::Output, Interrupt> {
let ttl = self.inner.lease_ttl;
tokio::pin!(future);
loop {
tokio::select! {
output = &mut future => return Ok(output),
() = tokio::time::sleep(ttl / 3) => {
match self.store().renew_lease(lease, self.now(), ttl).await {
Ok(_) => {}
Err(StoreError::LeaseLost) => return Err(Interrupt::LeaseLost),
Err(e) => warn!(error = %e, "lease renewal failed; will retry"),
}
}
}
}
}
pub(crate) async fn transition_leased(
&self,
record: &EffectRecord,
lease: &Lease,
actor: Option<&str>,
transition: Transition,
customize: impl FnOnce(&mut TransitionRequest),
) -> Result<EffectRecord, Interrupt> {
let mut request = TransitionRequest::new(record, Some(lease), transition, self.now());
request.actor = actor.map(str::to_owned);
customize(&mut request);
let record = self.commit_transition(record, request).await?;
debug!(%transition, status = %record.status, "effect transition");
Ok(record)
}
pub(crate) async fn commit_transition(
&self,
before: &EffectRecord,
mut request: TransitionRequest,
) -> Result<EffectRecord, StoreError> {
self.redact_request(&mut request, &before.key.name);
let transition = request.transition;
let after = self.store().transition(request).await?;
if !self.inner.observers.is_empty() {
let observation = Observation {
record: &after,
transition,
from: before.status,
to: after.status,
in_previous_status: after
.updated_at
.duration_since(before.updated_at)
.unwrap_or_default(),
since_created: after
.updated_at
.duration_since(after.created_at)
.unwrap_or_default(),
};
self.notify(|observer| observer.on_transition(&observation));
}
Ok(after)
}
fn notify(&self, call: impl Fn(&dyn EffectObserver)) {
for observer in &self.inner.observers {
let observer: &dyn EffectObserver = observer.as_ref();
if std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| call(observer))).is_err() {
warn!("an effect observer panicked; ignoring it");
}
}
}
pub(crate) fn redact(&self, field: Field, effect: &EffectName, value: &mut Value) {
if let Some(redactor) = &self.inner.redactor {
redactor.redact(field, effect, value);
}
}
pub(crate) fn redact_request(&self, request: &mut TransitionRequest, effect: &EffectName) {
if self.inner.redactor.is_none() {
return;
}
if let Some(output) = request.output.as_mut() {
self.redact(Field::Output, effect, output);
}
if let Some(payload) = request.payload.as_mut() {
self.redact(Field::AuditPayload, effect, payload);
}
if let Some(error) = request.error.as_mut() {
let mut message = Value::String(std::mem::take(&mut error.message));
self.redact(Field::ErrorMessage, effect, &mut message);
error.message = match message {
Value::String(text) => text,
other => other.to_string(),
};
}
}
#[cfg_attr(not(feature = "fault-injection"), allow(clippy::unused_self))]
pub(crate) fn checkpoint(&self, point: FaultPoint) {
#[cfg(feature = "fault-injection")]
if let Some(faults) = &self.inner.faults {
faults.reach(point);
}
#[cfg(not(feature = "fault-injection"))]
let _ = point;
}
pub(crate) async fn execute<T, F, Fut, V>(
&self,
mut spec: EffectSpec,
action: F,
verifier: V,
) -> Result<EffectOutcome<T>, RuntimeError>
where
T: Serialize + DeserializeOwned + Send + 'static,
F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
V: Verifier<T>,
{
let required = self
.inner
.policy
.requirements(spec.risk, spec.capabilities.kind);
if required.verification && spec.capabilities.verification == VerificationMode::None {
return Err(RuntimeError::PolicyViolation {
key: spec.key.to_string(),
requirement: "verification",
});
}
spec.require_approval |= required.approval;
spec.automatic_retry &= !required.no_automatic_retry;
if !spec.input_stored
&& let Some(input) = spec.input.as_mut()
{
self.redact(Field::Input, &spec.key.name, input);
spec.fingerprint = Some(fingerprint(input));
}
let span = info_span!(
"agent_effect.execute",
effect.name = %spec.key.name,
effect.logical_key = %spec.key.key,
effect.kind = spec.capabilities.kind.as_str(),
effect.risk_level = %spec.risk,
effect.id = field::Empty,
effect.status = field::Empty,
effect.attempt = field::Empty,
);
if spec.capabilities.unknown_always_escalates() {
span.in_scope(|| {
warn!(
"effect is neither idempotent nor verifiable: \
any unknown outcome will need an operator"
);
});
}
let runtime = self.clone();
let task = tokio::spawn(
async move { runtime.drive(spec, action, verifier).await }.instrument(span),
);
task.await
.unwrap_or_else(|e| Err(RuntimeError::Internal(e.to_string())))
}
async fn drive<T, F, Fut, V>(
&self,
spec: EffectSpec,
action: F,
verifier: V,
) -> Result<EffectOutcome<T>, RuntimeError>
where
T: Serialize + DeserializeOwned + Send + 'static,
F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
V: Verifier<T>,
{
self.checkpoint(FaultPoint::BeforeInsert);
let store = self.store();
let inserted = store
.insert_or_get(NewEffect {
id: EffectId::new(),
key: spec.key.clone(),
kind: spec.capabilities.kind,
input: spec.input.clone(),
input_fingerprint: spec.fingerprint.clone(),
created_by: spec.actor.clone(),
now: self.now(),
})
.await?;
let mut record = inserted.record;
if inserted.inserted {
self.notify(|observer| observer.on_created(&record));
}
Span::current().record("effect.id", field::display(record.id));
self.checkpoint(FaultPoint::AfterInsert);
if !inserted.inserted {
check_matches(&record, &spec)?;
}
let checks_made = AtomicU32::new(0);
for _ in 0..MAX_ROUNDS {
if let Some(outcome) = observe(&record)? {
return Ok(outcome);
}
let lease = match store
.acquire_lease(
record.id,
self.worker_id(),
self.now(),
self.inner.lease_ttl,
)
.await
{
Ok(lease) => lease,
Err(StoreError::LeaseHeld { .. }) => {
return Ok(EffectOutcome::InProgress { id: record.id });
}
Err(e) => return Err(e.into()),
};
let current = store
.get(record.id)
.await?
.ok_or(StoreError::NotFound(record.id))?;
let driver = Driver {
rt: self,
spec: &spec,
action: &action,
verifier: &verifier,
lease: &lease,
checks_made: &checks_made,
};
let advanced = driver.advance(current).await;
if let Err(e) = store.release_lease(&lease).await {
warn!(error = %e, "could not release lease; it will expire");
}
match advanced {
Ok((settled, output)) => {
Span::current().record("effect.status", settled.status.as_str());
return report(&settled, output);
}
Err(Interrupt::LeaseLost) => {
warn!("lease lost mid-effect; re-reading the record");
record = store
.get(record.id)
.await?
.ok_or(StoreError::NotFound(record.id))?;
}
Err(Interrupt::Error(e)) => return Err(e),
}
}
Ok(EffectOutcome::InProgress { id: record.id })
}
}
fn observe<T: DeserializeOwned>(
record: &EffectRecord,
) -> Result<Option<EffectOutcome<T>>, RuntimeError> {
match record.status {
status if settled(status) => report(record, None).map(Some),
EffectStatus::Pending
| EffectStatus::AwaitingApproval
| EffectStatus::Executing
| EffectStatus::Verifying
| EffectStatus::Unknown => Ok(None),
_ => Ok(Some(EffectOutcome::InProgress { id: record.id })),
}
}
struct Driver<'a, S, F, V> {
rt: &'a Runtime<S>,
spec: &'a EffectSpec,
action: &'a F,
verifier: &'a V,
lease: &'a Lease,
checks_made: &'a AtomicU32,
}
impl<S: EffectStore, F, V> Driver<'_, S, F, V> {
async fn advance<T, Fut>(
&self,
mut record: EffectRecord,
) -> Result<(EffectRecord, Option<T>), Interrupt>
where
T: Serialize + Send + 'static,
F: Fn(EffectContext) -> Fut,
Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
V: Verifier<T>,
{
let mut output = None;
let mut verification_exhausted = false;
loop {
match record.status {
EffectStatus::Pending => {
if let Some(at) = record.next_attempt_at {
self.sleep_until(at).await?;
}
if record.attempt_count == 0
&& let Some(reason) = self.check_precondition(&record).await?
{
record = self
.transition(&record, Transition::PreconditionRejected, |r| {
r.error = Some(reason);
})
.await?;
continue;
}
if record.attempt_count == 0 && self.spec.require_approval && !record.approved {
record = self
.transition(&record, Transition::RequestApproval, |_| {})
.await?;
self.rt.checkpoint(FaultPoint::AfterApprovalRequested);
continue;
}
record = self
.transition(&record, Transition::StartAttempt, |r| {
r.payload = Some(json!({ "worker": self.rt.worker_id() }));
})
.await?;
Span::current().record("effect.attempt", record.attempt_count);
self.rt.checkpoint(FaultPoint::AfterAttemptPersisted);
let (next, produced, exhausted) = self.attempt(record).await?;
record = next;
if produced.is_some() {
output = produced;
}
verification_exhausted = exhausted;
}
EffectStatus::Executing | EffectStatus::Verifying => {
record = self
.transition(&record, Transition::LeaseExpired, |_| {})
.await?;
}
EffectStatus::AwaitingApproval => match self.ask_approval(&record).await? {
ApprovalDecision::Approved { by } => {
record = self
.transition(&record, Transition::Approve, |r| r.actor = Some(by))
.await?;
}
ApprovalDecision::Denied { by, reason } => {
record = self
.transition(&record, Transition::Deny, |r| {
r.actor = Some(by);
r.error = Some(ErrorRecord {
class: None,
message: reason,
});
})
.await?;
}
ApprovalDecision::Deferred => break,
},
EffectStatus::Unknown => match self.spec.capabilities.unknown_plan() {
UnknownPlan::Verify if !verification_exhausted => {
let (next, verified, exhausted) = self.verify(record, None).await?;
record = next;
output = verified.or(output);
verification_exhausted = exhausted;
}
UnknownPlan::Verify => break,
UnknownPlan::Reexecute if self.may_retry(record.attempt_count) => {
record = self
.schedule_retry(&record, FailureClass::Ambiguous, None)
.await?;
}
UnknownPlan::Reexecute | UnknownPlan::Escalate => {
record = self
.transition(&record, Transition::Escalate, |_| {})
.await?;
}
},
_ => break,
}
}
Ok((record, output))
}
async fn attempt<T, Fut>(
&self,
record: EffectRecord,
) -> Result<(EffectRecord, Option<T>, bool), Interrupt>
where
T: Serialize + Send + 'static,
F: Fn(EffectContext) -> Fut,
Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
V: Verifier<T>,
{
let mut task = tokio::spawn((self.action)(context(&record)));
self.rt.checkpoint(FaultPoint::AfterActionStarted);
let joined = match self.spec.attempt_timeout {
None => self.leased(&mut task).await?,
Some(limit) => match self.leased(tokio::time::timeout(limit, &mut task)).await? {
Ok(joined) => joined,
Err(_elapsed) => {
task.abort();
Ok(Err(EffectFailure::ambiguous(format!(
"attempt timed out after {limit:?}"
))))
}
},
};
self.rt.checkpoint(FaultPoint::AfterActionReturned);
match joined {
Ok(Ok(value)) if self.spec.capabilities.verification != VerificationMode::None => {
let (record, verified, exhausted) =
self.verify(record, Some(output_json(&value))).await?;
let output = match record.status {
EffectStatus::Committed => verified.or(Some(value)),
_ => None,
};
Ok((record, output, exhausted))
}
Ok(Ok(value)) => {
let (output, payload) = output_json(&value);
let record = self
.transition(&record, Transition::Succeeded, |r| {
r.output = output;
r.payload = payload;
})
.await?;
Ok((record, Some(value), false))
}
Ok(Err(failure)) => {
debug!(%failure, "action failed");
let class = failure.class();
let error = failure.to_record();
let record = match class.disposition() {
Disposition::Retry if self.may_retry(record.attempt_count) => {
self.schedule_retry(&record, class, Some(error)).await?
}
Disposition::Retry | Disposition::Fail
if record.may_have_applied && record.kind != EffectKind::Read =>
{
let unknown = self
.transition(&record, Transition::OutcomeUnknown, |r| {
r.error = Some(error);
r.payload =
Some(json!({ "earlier_attempt_may_have_applied": true }));
})
.await?;
if self.spec.capabilities.unknown_plan() == UnknownPlan::Verify {
unknown
} else {
self.transition(&unknown, Transition::Escalate, |_| {})
.await?
}
}
Disposition::Retry | Disposition::Fail => {
self.transition(&record, Transition::FailedDefinitively, |r| {
r.error = Some(error);
})
.await?
}
Disposition::Unknown => {
self.transition(&record, Transition::OutcomeUnknown, |r| {
r.error = Some(error);
})
.await?
}
};
Ok((record, None, false))
}
Err(join_error) => {
let record = self
.transition(&record, Transition::OutcomeUnknown, |r| {
r.error = Some(ErrorRecord {
class: Some(FailureClass::Ambiguous),
message: format!("action did not complete: {join_error}"),
});
})
.await?;
Ok((record, None, false))
}
}
}
async fn verify<T>(
&self,
record: EffectRecord,
succeeded: Option<(Option<Value>, Option<Value>)>,
) -> Result<(EffectRecord, Option<T>, bool), Interrupt>
where
T: Serialize + Send + 'static,
V: Verifier<T>,
{
let mut record = self
.transition(&record, Transition::StartVerification, |r| {
if let Some((output, payload)) = succeeded {
(r.output, r.payload) = (output, payload);
}
})
.await?;
self.rt.checkpoint(FaultPoint::AfterVerificationStarted);
let mode = self.spec.capabilities.verification;
let mut checks = 0;
let mut last_problem = String::from("no check completed");
let verifying_since = record.updated_at;
let ended_before = record.attempt_ended_at.map_or(Duration::ZERO, |ended| {
verifying_since.duration_since(ended).unwrap_or_default()
});
let started = tokio::time::Instant::now();
while self.take_check() {
let check = checks;
checks += 1;
let Some(future) = self.verifier.check(context(&record)) else {
break;
};
let found = match self.leased(tokio::spawn(future)).await? {
Ok(Ok(found)) => found,
Ok(Err(failure)) => {
last_problem = format!("check failed: {failure}");
Verification::Inconclusive
}
Err(join_error) => {
last_problem = format!("check did not complete: {join_error}");
Verification::Inconclusive
}
};
match found {
Verification::Confirmed(value) => {
let (output, payload) = output_json(&value);
record = self
.transition(&record, Transition::VerificationConfirmed, |r| {
r.output = output;
r.payload = payload;
})
.await?;
return Ok((record, Some(value), false));
}
Verification::Conflict { details } => {
record = self
.transition(&record, Transition::VerificationConflict, |r| {
r.error = Some(ErrorRecord {
class: None,
message: details,
});
})
.await?;
return Ok((record, None, false));
}
Verification::NotApplied => {
match mode.read_not_found(ended_before + started.elapsed()) {
NotFoundReading::NotApplied => {
return Ok((self.not_applied(&record).await?, None, false));
}
NotFoundReading::TooEarly { wait } => {
last_problem = "not visible yet within the settle delay".into();
if self.has_checks_left() {
self.sleep(wait).await?;
}
}
}
}
Verification::Inconclusive => {
if last_problem == "no check completed" {
last_problem = "remote system could not tell".into();
}
if self.has_checks_left() {
let delay =
self.retry()
.delay(check, FailureClass::Transient, jitter_sample());
self.sleep(delay).await?;
}
}
}
}
if checks == 0 {
last_problem = "this call's checks were already spent".into();
}
record = self
.transition(&record, Transition::OutcomeUnknown, |r| {
r.error = Some(ErrorRecord {
class: Some(FailureClass::Ambiguous),
message: format!(
"verification inconclusive after {checks} checks: {last_problem}"
),
});
})
.await?;
Ok((record, None, true))
}
fn take_check(&self) -> bool {
self.checks_made.fetch_add(1, Ordering::Relaxed) < self.retry().max_attempts.max(1)
}
fn has_checks_left(&self) -> bool {
self.checks_made.load(Ordering::Relaxed) < self.retry().max_attempts.max(1)
}
async fn not_applied(&self, record: &EffectRecord) -> Result<EffectRecord, Interrupt> {
let error = ErrorRecord {
class: None,
message: "verification found that the effect did not apply".into(),
};
if self.may_retry(record.attempt_count) {
self.schedule_retry(record, FailureClass::Transient, Some(error))
.await
} else {
self.transition(record, Transition::VerificationNotApplied, |r| {
r.error = Some(error);
})
.await
}
}
async fn check_precondition(
&self,
record: &EffectRecord,
) -> Result<Option<ErrorRecord>, Interrupt> {
let Some(precondition) = &self.spec.precondition else {
return Ok(None);
};
let max_checks = self.retry().max_attempts.max(1);
let mut last_reason = String::new();
for check in 1..=max_checks {
let rejection = |message: String| {
Some(ErrorRecord {
class: None,
message,
})
};
match self
.leased(tokio::spawn(precondition(context(record))))
.await?
{
Ok(Precondition::Satisfied) => return Ok(None),
Ok(Precondition::Rejected { reason }) => return Ok(rejection(reason)),
Ok(Precondition::RetryLater { after, reason }) => {
last_reason = reason;
if check < max_checks {
self.sleep(after).await?;
}
}
Err(join_error) => {
return Ok(rejection(format!(
"precondition check did not complete: {join_error}"
)));
}
}
}
Ok(Some(ErrorRecord {
class: None,
message: format!("precondition not satisfied after {max_checks} checks: {last_reason}"),
}))
}
async fn schedule_retry(
&self,
record: &EffectRecord,
class: FailureClass,
error: Option<ErrorRecord>,
) -> Result<EffectRecord, Interrupt> {
let retry = record.attempt_count.saturating_sub(1);
let delay = self.retry().delay(retry, class, jitter_sample());
let at = self.rt.now() + delay;
debug!(?delay, ?class, "retry scheduled");
self.transition(record, Transition::ScheduleRetry, |r| {
r.next_attempt_at = Some(at);
r.error = error;
r.payload =
Some(json!({ "delay_ms": u64::try_from(delay.as_millis()).unwrap_or(u64::MAX) }));
})
.await
}
async fn ask_approval(&self, record: &EffectRecord) -> Result<ApprovalDecision, Interrupt> {
let Some(provider) = self.rt.inner.approval.clone() else {
return Ok(ApprovalDecision::Deferred);
};
let request = ApprovalRequest {
effect_id: record.id,
key: record.key.clone(),
kind: record.kind,
risk: self.spec.risk,
input: record.input.clone(),
requested_by: record.created_by.clone(),
};
let task = tokio::spawn(async move { provider.request_boxed(request).await });
Ok(self.leased(task).await?.unwrap_or_else(|e| {
warn!(error = %e, "approval provider failed; deferring");
ApprovalDecision::Deferred
}))
}
async fn sleep_until(&self, at: SystemTime) -> Result<(), Interrupt> {
let wait = at.duration_since(self.rt.now()).unwrap_or_default();
self.sleep(wait).await
}
async fn sleep(&self, duration: Duration) -> Result<(), Interrupt> {
if duration.is_zero() {
return Ok(());
}
self.leased(tokio::time::sleep(duration)).await
}
async fn leased<Fut: Future>(&self, future: Fut) -> Result<Fut::Output, Interrupt> {
self.rt.with_lease(self.lease, future).await
}
async fn transition(
&self,
record: &EffectRecord,
transition: Transition,
customize: impl FnOnce(&mut TransitionRequest),
) -> Result<EffectRecord, Interrupt> {
self.rt
.transition_leased(
record,
self.lease,
self.spec.actor.as_deref(),
transition,
customize,
)
.await
}
fn retry(&self) -> &RetryPolicy {
&self.spec.retry
}
fn may_retry(&self, attempts: u32) -> bool {
self.spec.automatic_retry && self.retry().allows_another(attempts)
}
}
fn context(record: &EffectRecord) -> EffectContext {
EffectContext {
id: record.id,
key: record.key.clone(),
attempt: record.attempt_count,
}
}
fn output_json<T: Serialize>(value: &T) -> (Option<Value>, Option<Value>) {
match serde_json::to_value(value) {
Ok(output) => (Some(output), None),
Err(e) => (None, Some(json!({ "output_not_stored": e.to_string() }))),
}
}
pub(crate) fn jitter_sample() -> f64 {
let bits = RandomState::new().hash_one(0_u8) >> 11;
#[allow(clippy::cast_precision_loss)]
let sample = bits as f64 / (1_u64 << 53) as f64;
sample
}
fn settled(status: EffectStatus) -> bool {
matches!(
status,
EffectStatus::Committed
| EffectStatus::Failed
| EffectStatus::Rejected
| EffectStatus::NeedsIntervention
| EffectStatus::Compensated
| EffectStatus::CompensationFailed
)
}
fn check_matches(record: &EffectRecord, spec: &EffectSpec) -> Result<(), RuntimeError> {
if record.kind != spec.capabilities.kind {
return Err(RuntimeError::KindMismatch {
id: record.id,
stored: record.kind,
requested: spec.capabilities.kind,
});
}
if record.input_fingerprint != spec.fingerprint {
return Err(RuntimeError::InputMismatch { id: record.id });
}
Ok(())
}
fn report<T: DeserializeOwned>(
record: &EffectRecord,
fresh: Option<T>,
) -> Result<EffectOutcome<T>, RuntimeError> {
let id = record.id;
Ok(match record.status {
EffectStatus::Committed => match fresh {
Some(value) => EffectOutcome::Committed(value),
None => EffectOutcome::Committed(
serde_json::from_value(record.output.clone().unwrap_or(Value::Null))
.map_err(|source| RuntimeError::Output { id, source })?,
),
},
EffectStatus::Failed => EffectOutcome::Failed(last_error(record)),
EffectStatus::Rejected => EffectOutcome::Rejected(last_error(record)),
EffectStatus::NeedsIntervention | EffectStatus::CompensationFailed => {
EffectOutcome::NeedsIntervention { id }
}
EffectStatus::Unknown => EffectOutcome::Unknown { id },
EffectStatus::Compensated => EffectOutcome::Compensated { id },
EffectStatus::AwaitingApproval => EffectOutcome::AwaitingApproval { id },
_ => EffectOutcome::InProgress { id },
})
}
pub(crate) fn last_error(record: &EffectRecord) -> ErrorRecord {
record.last_error.clone().unwrap_or_else(|| ErrorRecord {
class: None,
message: "no error was recorded".into(),
})
}