Skip to main content

agent_effects/
runtime.rs

1//! The effect runtime.
2//!
3//! [`Runtime::effect`] starts a builder; running it inserts or finds the
4//! effect's record, then either reports its settled outcome or takes the
5//! execution lease and moves it forward (design ยง7, "Re-attaching").
6
7use std::collections::hash_map::RandomState;
8use std::fmt::Display;
9use std::future::Future;
10use std::hash::BuildHasher;
11use std::sync::Arc;
12use std::sync::atomic::{AtomicU32, Ordering};
13use std::time::{Duration, SystemTime};
14
15use serde::Serialize;
16use serde::de::DeserializeOwned;
17use serde_json::{Value, json};
18use tracing::{Instrument, Span, debug, field, info_span, warn};
19
20use crate::approval::{ApprovalDecision, ApprovalProvider, ApprovalRequest, ErasedApproval};
21use crate::clock::{Clock, SystemClock};
22use crate::effect::{
23    EffectBuilder, EffectContext, EffectFailure, EffectOutcome, EffectSpec, Precondition,
24};
25use crate::error::RuntimeError;
26use crate::failure::{Disposition, FailureClass};
27#[cfg(feature = "fault-injection")]
28use crate::fault::FaultInjector;
29use crate::fault::FaultPoint;
30use crate::fingerprint::fingerprint;
31use crate::handler::{
32    CompensationSubmission, EffectHandler, Handler, Registered, Registry, Resume, Submission,
33};
34use crate::id::{EffectId, EffectName, WorkerId};
35use crate::kind::EffectKind;
36use crate::observer::{EffectObserver, Observation};
37use crate::policy::{RiskPolicy, UnknownPlan};
38use crate::redaction::{Field, Redactor};
39use crate::retention::RetentionPolicy;
40use crate::retry::RetryPolicy;
41use crate::state::{EffectStatus, Transition};
42use crate::store::{
43    EffectRecord, EffectStore, ErrorRecord, Lease, NewEffect, StoreError, TransitionRequest,
44};
45use crate::verification::{NotFoundReading, Verification, VerificationMode, Verifier};
46
47/// How many times one call re-reads the record after losing a lease race
48/// before reporting the effect as in progress.
49const MAX_ROUNDS: usize = 4;
50
51/// Executes effects against a store. Cheap to clone; clones share state.
52pub struct Runtime<S> {
53    inner: Arc<Inner<S>>,
54}
55
56struct Inner<S> {
57    store: S,
58    clock: Arc<dyn Clock>,
59    worker: WorkerId,
60    lease_ttl: Duration,
61    retry: RetryPolicy,
62    handlers: Registry<S>,
63    approval: Option<Arc<dyn ErasedApproval>>,
64    policy: RiskPolicy,
65    redactor: Option<Arc<dyn Redactor>>,
66    observers: Vec<Arc<dyn EffectObserver>>,
67    retention: RetentionPolicy,
68    #[cfg(feature = "fault-injection")]
69    faults: Option<Arc<FaultInjector>>,
70}
71
72impl<S> Clone for Runtime<S> {
73    fn clone(&self) -> Self {
74        Self {
75            inner: Arc::clone(&self.inner),
76        }
77    }
78}
79
80/// Configures a [`Runtime`].
81#[must_use]
82pub struct RuntimeBuilder<S> {
83    store: S,
84    clock: Arc<dyn Clock>,
85    worker: Option<WorkerId>,
86    lease_ttl: Duration,
87    retry: RetryPolicy,
88    handlers: Registry<S>,
89    approval: Option<Arc<dyn ErasedApproval>>,
90    policy: RiskPolicy,
91    redactor: Option<Arc<dyn Redactor>>,
92    observers: Vec<Arc<dyn EffectObserver>>,
93    retention: RetentionPolicy,
94    #[cfg(feature = "fault-injection")]
95    faults: Option<Arc<FaultInjector>>,
96}
97
98impl<S: EffectStore> RuntimeBuilder<S> {
99    /// How long settled records are kept before [`Runtime::prune`] (and
100    /// every round of [`Runtime::run_recovery`]) deletes them; see
101    /// [`retention`](crate::retention). Defaults to keeping them forever.
102    pub fn retention(mut self, policy: RetentionPolicy) -> Self {
103        self.retention = policy;
104        self
105    }
106
107    /// Tells `observer` about every effect recorded and every transition
108    /// stored, for metrics; see [`observer`](crate::observer). May be called
109    /// more than once.
110    pub fn observer(mut self, observer: impl EffectObserver) -> Self {
111        self.observers.push(Arc::new(observer));
112        self
113    }
114
115    /// Rewrites every input, output, audit payload and error message before
116    /// it is stored; see [`redaction`](crate::redaction).
117    pub fn redactor(mut self, redactor: impl Redactor) -> Self {
118        self.redactor = Some(Arc::new(redactor));
119        self
120    }
121
122    /// Adds requirements by risk level and effect kind to every effect; see
123    /// [`RiskPolicy`]. Requirements only accumulate on top of each effect's
124    /// own settings.
125    pub fn risk_policy(mut self, policy: RiskPolicy) -> Self {
126        self.policy = policy;
127        self
128    }
129
130    /// Asks `provider` to decide on effects that require approval; see
131    /// [`approval`](crate::approval). Without one, such effects wait for an
132    /// operator's [`Runtime::approve`] or [`Runtime::deny`].
133    pub fn approval_provider(mut self, provider: impl ApprovalProvider) -> Self {
134        self.approval = Some(Arc::new(provider));
135        self
136    }
137
138    /// Registers a durable handler under its [`EffectHandler::NAME`], so
139    /// [`Runtime::submit`] can run it and [`Runtime::recover`] can finish
140    /// its effects without a caller.
141    ///
142    /// # Panics
143    ///
144    /// If a handler is already registered under the same name.
145    pub fn register<H: EffectHandler>(mut self, handler: Handler<H>) -> Self {
146        assert!(
147            !self.handlers.contains_key(H::NAME),
148            "a handler is already registered for effect `{}`",
149            H::NAME
150        );
151        self.handlers.insert(H::NAME, Registered::new(handler));
152        self
153    }
154
155    /// The time source for leases and schedules. Defaults to the system
156    /// clock. Tests with paused Tokio time want
157    /// [`TokioClock`](crate::clock::TokioClock).
158    pub fn clock(mut self, clock: impl Clock) -> Self {
159        self.clock = Arc::new(clock);
160        self
161    }
162
163    /// This runtime's identity as a lease holder. Defaults to a random id.
164    /// Must be unique among runtimes sharing a store.
165    pub fn worker_id(mut self, worker: WorkerId) -> Self {
166        self.worker = Some(worker);
167        self
168    }
169
170    /// How long a lease lasts without renewal. The runtime renews it every
171    /// third of this while it works on an effect, including while it waits
172    /// between retries. If a worker dies, others wait this long before
173    /// treating its in-flight effect as unknown. Defaults to 30 seconds;
174    /// values under 3 ms are raised to 3 ms.
175    pub fn lease_ttl(mut self, ttl: Duration) -> Self {
176        self.lease_ttl = ttl.max(Duration::from_millis(3));
177        self
178    }
179
180    /// The retry policy for effects that do not set their own. Defaults to
181    /// [`RetryPolicy::default`]: 5 attempts, 1 s to 30 s exponential backoff
182    /// with jitter.
183    pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
184        self.retry = policy;
185        self
186    }
187
188    /// Simulates crashes at the injector's armed points. For tests; see
189    /// [`fault`](crate::fault).
190    #[cfg(feature = "fault-injection")]
191    pub fn fault_injector(mut self, injector: Arc<FaultInjector>) -> Self {
192        self.faults = Some(injector);
193        self
194    }
195
196    /// Builds the runtime.
197    pub fn build(self) -> Runtime<S> {
198        Runtime {
199            inner: Arc::new(Inner {
200                store: self.store,
201                clock: self.clock,
202                worker: self.worker.unwrap_or_else(WorkerId::random),
203                lease_ttl: self.lease_ttl,
204                retry: self.retry,
205                handlers: self.handlers,
206                approval: self.approval,
207                policy: self.policy,
208                redactor: self.redactor,
209                observers: self.observers,
210                retention: self.retention,
211                #[cfg(feature = "fault-injection")]
212                faults: self.faults,
213            }),
214        }
215    }
216}
217
218/// Why moving an effect forward stopped early.
219pub(crate) enum Interrupt {
220    /// Our lease expired or was taken over; re-read the record.
221    LeaseLost,
222    Error(RuntimeError),
223}
224
225impl From<StoreError> for Interrupt {
226    fn from(error: StoreError) -> Self {
227        match error {
228            StoreError::LeaseLost => Self::LeaseLost,
229            other => Self::Error(other.into()),
230        }
231    }
232}
233
234impl<S: EffectStore> Runtime<S> {
235    /// A runtime with default settings.
236    pub fn new(store: S) -> Self {
237        Self::builder(store).build()
238    }
239
240    /// Starts configuring a runtime.
241    pub fn builder(store: S) -> RuntimeBuilder<S> {
242        RuntimeBuilder {
243            store,
244            clock: Arc::new(SystemClock),
245            worker: None,
246            lease_ttl: Duration::from_secs(30),
247            retry: RetryPolicy::default(),
248            handlers: Registry::new(),
249            approval: None,
250            policy: RiskPolicy::default(),
251            redactor: None,
252            observers: Vec::new(),
253            retention: RetentionPolicy::KEEP_ALL,
254            #[cfg(feature = "fault-injection")]
255            faults: None,
256        }
257    }
258
259    /// Describes an effect: `name` is its type (e.g. `payment.charge`),
260    /// `key` identifies this occurrence (e.g. the order id). Every call with
261    /// the same name and key refers to the same effect.
262    pub fn effect(&self, name: impl Into<String>, key: impl Display) -> EffectBuilder<S> {
263        EffectBuilder::new(self.clone(), name.into(), key.to_string())
264    }
265
266    /// Runs a registered handler's effect with `input`, or attaches to an
267    /// earlier run with the same `key`. Await the returned [`Submission`].
268    ///
269    /// Behaves like [`EffectBuilder::run`], with the handler's properties.
270    /// The input is stored in full, so recovery can finish the effect if
271    /// this process dies.
272    pub fn submit<H: EffectHandler>(
273        &self,
274        key: impl Display,
275        input: H::Input,
276    ) -> Submission<'_, S, H> {
277        Submission::new(self, key, input)
278    }
279
280    /// Undoes a registered, [compensable](Handler::compensable) handler's
281    /// effect. Await the returned submission. See
282    /// [`compensation`](crate::compensation).
283    pub fn compensate<H: EffectHandler>(
284        &self,
285        key: impl Display,
286    ) -> CompensationSubmission<'_, S, H> {
287        CompensationSubmission::new(self, key)
288    }
289
290    pub(crate) fn handler<H: EffectHandler>(&self) -> Option<Handler<H>> {
291        self.inner.handlers.get(H::NAME)?.typed::<H>()
292    }
293
294    pub(crate) fn resumer(&self, name: &str) -> Option<Resume<S>> {
295        self.inner.handlers.get(name).map(|r| Arc::clone(&r.resume))
296    }
297
298    /// The underlying store.
299    pub fn store(&self) -> &S {
300        &self.inner.store
301    }
302
303    /// This runtime's lease-holder identity.
304    pub fn worker_id(&self) -> &WorkerId {
305        &self.inner.worker
306    }
307
308    /// Waits until nobody is working on effect `id`, or `timeout` passes, and
309    /// reports where it stands. For a caller that got
310    /// [`EffectOutcome::InProgress`].
311    ///
312    /// Returns `InProgress` if the effect is still being worked on at the
313    /// deadline, or if it is unsettled and nobody holds it (for example a
314    /// worker crashed while waiting to retry). Running the effect again then
315    /// takes it over.
316    ///
317    /// # Errors
318    ///
319    /// [`RuntimeError::Store`] if the store fails or has no such effect, and
320    /// [`RuntimeError::Output`] if a committed output does not deserialize
321    /// into `T`.
322    pub async fn wait<T: DeserializeOwned>(
323        &self,
324        id: EffectId,
325        timeout: Duration,
326    ) -> Result<EffectOutcome<T>, RuntimeError> {
327        let deadline = tokio::time::Instant::now() + timeout;
328        let mut pause = Duration::from_millis(10);
329        loop {
330            let record = self
331                .store()
332                .get(id)
333                .await?
334                .ok_or(StoreError::NotFound(id))?;
335            let busy = !settled(record.status) && record.live_lease_owner(self.now()).is_some();
336            let now = tokio::time::Instant::now();
337            if !busy {
338                return report(&record, None);
339            }
340            if now >= deadline {
341                return Ok(EffectOutcome::InProgress { id });
342            }
343            tokio::time::sleep(pause.min(deadline - now)).await;
344            pause = (pause * 2).min(Duration::from_millis(250));
345        }
346    }
347
348    pub(crate) fn default_retry(&self) -> RetryPolicy {
349        self.inner.retry
350    }
351
352    pub(crate) fn now(&self) -> SystemTime {
353        self.inner.clock.now()
354    }
355
356    pub(crate) fn retention(&self) -> RetentionPolicy {
357        self.inner.retention
358    }
359
360    pub(crate) fn lease_ttl(&self) -> Duration {
361        self.inner.lease_ttl
362    }
363
364    /// Awaits `future`, renewing `lease` every third of its TTL. Stops with
365    /// [`Interrupt::LeaseLost`] if the lease is lost meanwhile.
366    pub(crate) async fn with_lease<Fut: Future>(
367        &self,
368        lease: &Lease,
369        future: Fut,
370    ) -> Result<Fut::Output, Interrupt> {
371        let ttl = self.inner.lease_ttl;
372        tokio::pin!(future);
373        loop {
374            tokio::select! {
375                output = &mut future => return Ok(output),
376                () = tokio::time::sleep(ttl / 3) => {
377                    match self.store().renew_lease(lease, self.now(), ttl).await {
378                        Ok(_) => {}
379                        Err(StoreError::LeaseLost) => return Err(Interrupt::LeaseLost),
380                        Err(e) => warn!(error = %e, "lease renewal failed; will retry"),
381                    }
382                }
383            }
384        }
385    }
386
387    /// Applies `transition` under `lease`, attributing it to `actor`.
388    pub(crate) async fn transition_leased(
389        &self,
390        record: &EffectRecord,
391        lease: &Lease,
392        actor: Option<&str>,
393        transition: Transition,
394        customize: impl FnOnce(&mut TransitionRequest),
395    ) -> Result<EffectRecord, Interrupt> {
396        let mut request = TransitionRequest::new(record, Some(lease), transition, self.now());
397        request.actor = actor.map(str::to_owned);
398        customize(&mut request);
399        let record = self.commit_transition(record, request).await?;
400        debug!(%transition, status = %record.status, "effect transition");
401        Ok(record)
402    }
403
404    /// The one way a transition is written: redact, store, then tell the
405    /// observers. `before` is the record the request was built from.
406    pub(crate) async fn commit_transition(
407        &self,
408        before: &EffectRecord,
409        mut request: TransitionRequest,
410    ) -> Result<EffectRecord, StoreError> {
411        self.redact_request(&mut request, &before.key.name);
412        let transition = request.transition;
413        let after = self.store().transition(request).await?;
414        if !self.inner.observers.is_empty() {
415            let observation = Observation {
416                record: &after,
417                transition,
418                from: before.status,
419                to: after.status,
420                in_previous_status: after
421                    .updated_at
422                    .duration_since(before.updated_at)
423                    .unwrap_or_default(),
424                since_created: after
425                    .updated_at
426                    .duration_since(after.created_at)
427                    .unwrap_or_default(),
428            };
429            self.notify(|observer| observer.on_transition(&observation));
430        }
431        Ok(after)
432    }
433
434    /// Calls every observer, containing panics: the write already happened.
435    fn notify(&self, call: impl Fn(&dyn EffectObserver)) {
436        for observer in &self.inner.observers {
437            let observer: &dyn EffectObserver = observer.as_ref();
438            if std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| call(observer))).is_err() {
439                warn!("an effect observer panicked; ignoring it");
440            }
441        }
442    }
443
444    /// Applies the redactor, if any, to one value about to be stored.
445    pub(crate) fn redact(&self, field: Field, effect: &EffectName, value: &mut Value) {
446        if let Some(redactor) = &self.inner.redactor {
447            redactor.redact(field, effect, value);
448        }
449    }
450
451    /// Applies the redactor to everything a transition is about to store.
452    pub(crate) fn redact_request(&self, request: &mut TransitionRequest, effect: &EffectName) {
453        if self.inner.redactor.is_none() {
454            return;
455        }
456        if let Some(output) = request.output.as_mut() {
457            self.redact(Field::Output, effect, output);
458        }
459        if let Some(payload) = request.payload.as_mut() {
460            self.redact(Field::AuditPayload, effect, payload);
461        }
462        if let Some(error) = request.error.as_mut() {
463            let mut message = Value::String(std::mem::take(&mut error.message));
464            self.redact(Field::ErrorMessage, effect, &mut message);
465            error.message = match message {
466                Value::String(text) => text,
467                other => other.to_string(),
468            };
469        }
470    }
471
472    /// A point where a crash can be injected; a no-op without the
473    /// `fault-injection` feature.
474    #[cfg_attr(not(feature = "fault-injection"), allow(clippy::unused_self))]
475    pub(crate) fn checkpoint(&self, point: FaultPoint) {
476        #[cfg(feature = "fault-injection")]
477        if let Some(faults) = &self.inner.faults {
478            faults.reach(point);
479        }
480        #[cfg(not(feature = "fault-injection"))]
481        let _ = point;
482    }
483
484    pub(crate) async fn execute<T, F, Fut, V>(
485        &self,
486        mut spec: EffectSpec,
487        action: F,
488        verifier: V,
489    ) -> Result<EffectOutcome<T>, RuntimeError>
490    where
491        T: Serialize + DeserializeOwned + Send + 'static,
492        F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
493        Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
494        V: Verifier<T>,
495    {
496        // The risk policy only adds requirements; see `RiskPolicy`.
497        let required = self
498            .inner
499            .policy
500            .requirements(spec.risk, spec.capabilities.kind);
501        if required.verification && spec.capabilities.verification == VerificationMode::None {
502            return Err(RuntimeError::PolicyViolation {
503                key: spec.key.to_string(),
504                requirement: "verification",
505            });
506        }
507        spec.require_approval |= required.approval;
508        spec.automatic_retry &= !required.no_automatic_retry;
509
510        // Redact before fingerprinting: secrets are not part of an effect's
511        // identity, and never stored, not even hashed. A resumed effect
512        // brings its stored fingerprint, and its input is already redacted.
513        if !spec.input_stored
514            && let Some(input) = spec.input.as_mut()
515        {
516            self.redact(Field::Input, &spec.key.name, input);
517            spec.fingerprint = Some(fingerprint(input));
518        }
519
520        let span = info_span!(
521            "agent_effect.execute",
522            effect.name = %spec.key.name,
523            effect.logical_key = %spec.key.key,
524            effect.kind = spec.capabilities.kind.as_str(),
525            effect.risk_level = %spec.risk,
526            effect.id = field::Empty,
527            effect.status = field::Empty,
528            effect.attempt = field::Empty,
529        );
530        if spec.capabilities.unknown_always_escalates() {
531            span.in_scope(|| {
532                warn!(
533                    "effect is neither idempotent nor verifiable: \
534                     any unknown outcome will need an operator"
535                );
536            });
537        }
538        // Run on a spawned task so that dropping the caller's future cannot
539        // abort an attempt between invoking the action and recording it.
540        let runtime = self.clone();
541        let task = tokio::spawn(
542            async move { runtime.drive(spec, action, verifier).await }.instrument(span),
543        );
544        task.await
545            .unwrap_or_else(|e| Err(RuntimeError::Internal(e.to_string())))
546    }
547
548    async fn drive<T, F, Fut, V>(
549        &self,
550        spec: EffectSpec,
551        action: F,
552        verifier: V,
553    ) -> Result<EffectOutcome<T>, RuntimeError>
554    where
555        T: Serialize + DeserializeOwned + Send + 'static,
556        F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
557        Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
558        V: Verifier<T>,
559    {
560        self.checkpoint(FaultPoint::BeforeInsert);
561        let store = self.store();
562        let inserted = store
563            .insert_or_get(NewEffect {
564                id: EffectId::new(),
565                key: spec.key.clone(),
566                kind: spec.capabilities.kind,
567                input: spec.input.clone(),
568                input_fingerprint: spec.fingerprint.clone(),
569                created_by: spec.actor.clone(),
570                now: self.now(),
571            })
572            .await?;
573        let mut record = inserted.record;
574        if inserted.inserted {
575            self.notify(|observer| observer.on_created(&record));
576        }
577        Span::current().record("effect.id", field::display(record.id));
578        self.checkpoint(FaultPoint::AfterInsert);
579        if !inserted.inserted {
580            check_matches(&record, &spec)?;
581        }
582
583        // Verification checks this call has made, across every attempt and
584        // lease it goes through: the budget is per call.
585        let checks_made = AtomicU32::new(0);
586        for _ in 0..MAX_ROUNDS {
587            if let Some(outcome) = observe(&record)? {
588                return Ok(outcome);
589            }
590            let lease = match store
591                .acquire_lease(
592                    record.id,
593                    self.worker_id(),
594                    self.now(),
595                    self.inner.lease_ttl,
596                )
597                .await
598            {
599                Ok(lease) => lease,
600                Err(StoreError::LeaseHeld { .. }) => {
601                    return Ok(EffectOutcome::InProgress { id: record.id });
602                }
603                Err(e) => return Err(e.into()),
604            };
605            // Re-read under the lease: the record may have moved on since.
606            let current = store
607                .get(record.id)
608                .await?
609                .ok_or(StoreError::NotFound(record.id))?;
610            let driver = Driver {
611                rt: self,
612                spec: &spec,
613                action: &action,
614                verifier: &verifier,
615                lease: &lease,
616                checks_made: &checks_made,
617            };
618            let advanced = driver.advance(current).await;
619            if let Err(e) = store.release_lease(&lease).await {
620                warn!(error = %e, "could not release lease; it will expire");
621            }
622            match advanced {
623                Ok((settled, output)) => {
624                    Span::current().record("effect.status", settled.status.as_str());
625                    return report(&settled, output);
626                }
627                Err(Interrupt::LeaseLost) => {
628                    warn!("lease lost mid-effect; re-reading the record");
629                    record = store
630                        .get(record.id)
631                        .await?
632                        .ok_or(StoreError::NotFound(record.id))?;
633                }
634                Err(Interrupt::Error(e)) => return Err(e),
635            }
636        }
637        Ok(EffectOutcome::InProgress { id: record.id })
638    }
639}
640
641/// The outcome to report without acting, or `None` if this call should try
642/// to take the lease and move the effect forward.
643///
644/// Whether a lease is still live is left to `acquire_lease`: a store may
645/// judge it by its own clock (`PostgresStore` does), and this worker's
646/// clock could disagree.
647fn observe<T: DeserializeOwned>(
648    record: &EffectRecord,
649) -> Result<Option<EffectOutcome<T>>, RuntimeError> {
650    match record.status {
651        status if settled(status) => report(record, None).map(Some),
652        EffectStatus::Pending
653        | EffectStatus::AwaitingApproval
654        | EffectStatus::Executing
655        | EffectStatus::Verifying
656        | EffectStatus::Unknown => Ok(None),
657        _ => Ok(Some(EffectOutcome::InProgress { id: record.id })),
658    }
659}
660
661/// One call's work on one leased effect.
662struct Driver<'a, S, F, V> {
663    rt: &'a Runtime<S>,
664    spec: &'a EffectSpec,
665    action: &'a F,
666    verifier: &'a V,
667    lease: &'a Lease,
668    /// Verification checks made so far in this call.
669    checks_made: &'a AtomicU32,
670}
671
672impl<S: EffectStore, F, V> Driver<'_, S, F, V> {
673    /// Moves the effect forward until it settles or needs something this
674    /// call cannot do. Returns the record and, if this call produced one, the
675    /// output to report.
676    async fn advance<T, Fut>(
677        &self,
678        mut record: EffectRecord,
679    ) -> Result<(EffectRecord, Option<T>), Interrupt>
680    where
681        T: Serialize + Send + 'static,
682        F: Fn(EffectContext) -> Fut,
683        Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
684        V: Verifier<T>,
685    {
686        let mut output = None;
687        let mut verification_exhausted = false;
688        loop {
689            match record.status {
690                EffectStatus::Pending => {
691                    if let Some(at) = record.next_attempt_at {
692                        self.sleep_until(at).await?;
693                    }
694                    if record.attempt_count == 0
695                        && let Some(reason) = self.check_precondition(&record).await?
696                    {
697                        record = self
698                            .transition(&record, Transition::PreconditionRejected, |r| {
699                                r.error = Some(reason);
700                            })
701                            .await?;
702                        continue;
703                    }
704                    if record.attempt_count == 0 && self.spec.require_approval && !record.approved {
705                        record = self
706                            .transition(&record, Transition::RequestApproval, |_| {})
707                            .await?;
708                        self.rt.checkpoint(FaultPoint::AfterApprovalRequested);
709                        continue;
710                    }
711                    record = self
712                        .transition(&record, Transition::StartAttempt, |r| {
713                            r.payload = Some(json!({ "worker": self.rt.worker_id() }));
714                        })
715                        .await?;
716                    Span::current().record("effect.attempt", record.attempt_count);
717                    self.rt.checkpoint(FaultPoint::AfterAttemptPersisted);
718                    let (next, produced, exhausted) = self.attempt(record).await?;
719                    record = next;
720                    if produced.is_some() {
721                        output = produced;
722                    }
723                    verification_exhausted = exhausted;
724                }
725                // The previous holder's lease expired mid-attempt or
726                // mid-verification. Whatever it was doing may have happened.
727                EffectStatus::Executing | EffectStatus::Verifying => {
728                    record = self
729                        .transition(&record, Transition::LeaseExpired, |_| {})
730                        .await?;
731                }
732                // Ask the provider; an approval loops back to `Pending`, where
733                // the precondition is checked again before the attempt.
734                EffectStatus::AwaitingApproval => match self.ask_approval(&record).await? {
735                    ApprovalDecision::Approved { by } => {
736                        record = self
737                            .transition(&record, Transition::Approve, |r| r.actor = Some(by))
738                            .await?;
739                    }
740                    ApprovalDecision::Denied { by, reason } => {
741                        record = self
742                            .transition(&record, Transition::Deny, |r| {
743                                r.actor = Some(by);
744                                r.error = Some(ErrorRecord {
745                                    class: None,
746                                    message: reason,
747                                });
748                            })
749                            .await?;
750                    }
751                    ApprovalDecision::Deferred => break,
752                },
753                EffectStatus::Unknown => match self.spec.capabilities.unknown_plan() {
754                    UnknownPlan::Verify if !verification_exhausted => {
755                        let (next, verified, exhausted) = self.verify(record, None).await?;
756                        record = next;
757                        output = verified.or(output);
758                        verification_exhausted = exhausted;
759                    }
760                    // Still unknown after every check this call may make;
761                    // a later call or recovery tries again.
762                    UnknownPlan::Verify => break,
763                    UnknownPlan::Reexecute if self.may_retry(record.attempt_count) => {
764                        record = self
765                            .schedule_retry(&record, FailureClass::Ambiguous, None)
766                            .await?;
767                    }
768                    UnknownPlan::Reexecute | UnknownPlan::Escalate => {
769                        record = self
770                            .transition(&record, Transition::Escalate, |_| {})
771                            .await?;
772                    }
773                },
774                _ => break,
775            }
776        }
777        Ok((record, output))
778    }
779
780    /// Runs the action once and records what happened. The flag reports
781    /// that the postcondition checks ran out while inconclusive.
782    async fn attempt<T, Fut>(
783        &self,
784        record: EffectRecord,
785    ) -> Result<(EffectRecord, Option<T>, bool), Interrupt>
786    where
787        T: Serialize + Send + 'static,
788        F: Fn(EffectContext) -> Fut,
789        Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
790        V: Verifier<T>,
791    {
792        // A separate task, so a panicking action is caught as a JoinError.
793        // If the lease is lost while waiting, the handle is dropped and the
794        // task finishes detached: the request is already in flight.
795        let mut task = tokio::spawn((self.action)(context(&record)));
796        self.rt.checkpoint(FaultPoint::AfterActionStarted);
797        let joined = match self.spec.attempt_timeout {
798            None => self.leased(&mut task).await?,
799            Some(limit) => match self.leased(tokio::time::timeout(limit, &mut task)).await? {
800                Ok(joined) => joined,
801                Err(_elapsed) => {
802                    task.abort();
803                    Ok(Err(EffectFailure::ambiguous(format!(
804                        "attempt timed out after {limit:?}"
805                    ))))
806                }
807            },
808        };
809        self.rt.checkpoint(FaultPoint::AfterActionReturned);
810
811        match joined {
812            Ok(Ok(value)) if self.spec.capabilities.verification != VerificationMode::None => {
813                let (record, verified, exhausted) =
814                    self.verify(record, Some(output_json(&value))).await?;
815                let output = match record.status {
816                    EffectStatus::Committed => verified.or(Some(value)),
817                    _ => None,
818                };
819                Ok((record, output, exhausted))
820            }
821            Ok(Ok(value)) => {
822                let (output, payload) = output_json(&value);
823                let record = self
824                    .transition(&record, Transition::Succeeded, |r| {
825                        r.output = output;
826                        r.payload = payload;
827                    })
828                    .await?;
829                Ok((record, Some(value), false))
830            }
831            Ok(Err(failure)) => {
832                debug!(%failure, "action failed");
833                let class = failure.class();
834                let error = failure.to_record();
835                let record = match class.disposition() {
836                    Disposition::Retry if self.may_retry(record.attempt_count) => {
837                        self.schedule_retry(&record, class, Some(error)).await?
838                    }
839                    // This attempt definitely failed, but an earlier one may
840                    // have applied the effect: that is not "Failed". Look
841                    // again if the effect can be verified, else escalate.
842                    Disposition::Retry | Disposition::Fail
843                        if record.may_have_applied && record.kind != EffectKind::Read =>
844                    {
845                        let unknown = self
846                            .transition(&record, Transition::OutcomeUnknown, |r| {
847                                r.error = Some(error);
848                                r.payload =
849                                    Some(json!({ "earlier_attempt_may_have_applied": true }));
850                            })
851                            .await?;
852                        if self.spec.capabilities.unknown_plan() == UnknownPlan::Verify {
853                            unknown
854                        } else {
855                            self.transition(&unknown, Transition::Escalate, |_| {})
856                                .await?
857                        }
858                    }
859                    Disposition::Retry | Disposition::Fail => {
860                        self.transition(&record, Transition::FailedDefinitively, |r| {
861                            r.error = Some(error);
862                        })
863                        .await?
864                    }
865                    Disposition::Unknown => {
866                        self.transition(&record, Transition::OutcomeUnknown, |r| {
867                            r.error = Some(error);
868                        })
869                        .await?
870                    }
871                };
872                Ok((record, None, false))
873            }
874            Err(join_error) => {
875                // The action may have sent its request before panicking.
876                let record = self
877                    .transition(&record, Transition::OutcomeUnknown, |r| {
878                        r.error = Some(ErrorRecord {
879                            class: Some(FailureClass::Ambiguous),
880                            message: format!("action did not complete: {join_error}"),
881                        });
882                    })
883                    .await?;
884                Ok((record, None, false))
885            }
886        }
887    }
888
889    /// Asks the remote system what happened, starting from `Executing`
890    /// (after a success, whose output JSON and audit note are `succeeded`)
891    /// or `Unknown`.
892    ///
893    /// Ends in `Committed`, `Failed`, `Pending` (re-run scheduled: the effect
894    /// verifiably did not apply), `NeedsIntervention` (conflict) or `Unknown`.
895    /// The flag reports that the checks ran out while inconclusive.
896    async fn verify<T>(
897        &self,
898        record: EffectRecord,
899        succeeded: Option<(Option<Value>, Option<Value>)>,
900    ) -> Result<(EffectRecord, Option<T>, bool), Interrupt>
901    where
902        T: Serialize + Send + 'static,
903        V: Verifier<T>,
904    {
905        let mut record = self
906            .transition(&record, Transition::StartVerification, |r| {
907                if let Some((output, payload)) = succeeded {
908                    (r.output, r.payload) = (output, payload);
909                }
910            })
911            .await?;
912        self.rt.checkpoint(FaultPoint::AfterVerificationStarted);
913        let mode = self.spec.capabilities.verification;
914        let mut checks = 0;
915        let mut last_problem = String::from("no check completed");
916        // Settle delays count from the end of the attempt: a slow request
917        // may write just before it returns. Both times come from the store's
918        // clock (`updated_at` was just stamped by `StartVerification`), and
919        // the time since is measured locally, so this worker's clock never
920        // meets the store's.
921        let verifying_since = record.updated_at;
922        let ended_before = record.attempt_ended_at.map_or(Duration::ZERO, |ended| {
923            verifying_since.duration_since(ended).unwrap_or_default()
924        });
925        let started = tokio::time::Instant::now();
926
927        while self.take_check() {
928            let check = checks;
929            checks += 1;
930            let Some(future) = self.verifier.check(context(&record)) else {
931                break;
932            };
933            let found = match self.leased(tokio::spawn(future)).await? {
934                Ok(Ok(found)) => found,
935                Ok(Err(failure)) => {
936                    last_problem = format!("check failed: {failure}");
937                    Verification::Inconclusive
938                }
939                Err(join_error) => {
940                    last_problem = format!("check did not complete: {join_error}");
941                    Verification::Inconclusive
942                }
943            };
944            match found {
945                Verification::Confirmed(value) => {
946                    let (output, payload) = output_json(&value);
947                    record = self
948                        .transition(&record, Transition::VerificationConfirmed, |r| {
949                            r.output = output;
950                            r.payload = payload;
951                        })
952                        .await?;
953                    return Ok((record, Some(value), false));
954                }
955                Verification::Conflict { details } => {
956                    record = self
957                        .transition(&record, Transition::VerificationConflict, |r| {
958                            r.error = Some(ErrorRecord {
959                                class: None,
960                                message: details,
961                            });
962                        })
963                        .await?;
964                    return Ok((record, None, false));
965                }
966                Verification::NotApplied => {
967                    match mode.read_not_found(ended_before + started.elapsed()) {
968                        NotFoundReading::NotApplied => {
969                            return Ok((self.not_applied(&record).await?, None, false));
970                        }
971                        NotFoundReading::TooEarly { wait } => {
972                            last_problem = "not visible yet within the settle delay".into();
973                            if self.has_checks_left() {
974                                self.sleep(wait).await?;
975                            }
976                        }
977                    }
978                }
979                Verification::Inconclusive => {
980                    if last_problem == "no check completed" {
981                        last_problem = "remote system could not tell".into();
982                    }
983                    if self.has_checks_left() {
984                        let delay =
985                            self.retry()
986                                .delay(check, FailureClass::Transient, jitter_sample());
987                        self.sleep(delay).await?;
988                    }
989                }
990            }
991        }
992
993        if checks == 0 {
994            last_problem = "this call's checks were already spent".into();
995        }
996        record = self
997            .transition(&record, Transition::OutcomeUnknown, |r| {
998                r.error = Some(ErrorRecord {
999                    class: Some(FailureClass::Ambiguous),
1000                    message: format!(
1001                        "verification inconclusive after {checks} checks: {last_problem}"
1002                    ),
1003                });
1004            })
1005            .await?;
1006        Ok((record, None, true))
1007    }
1008
1009    /// Takes one of this call's verification checks (`max_attempts` per
1010    /// call, across every attempt it makes), if any are left.
1011    fn take_check(&self) -> bool {
1012        self.checks_made.fetch_add(1, Ordering::Relaxed) < self.retry().max_attempts.max(1)
1013    }
1014
1015    fn has_checks_left(&self) -> bool {
1016        self.checks_made.load(Ordering::Relaxed) < self.retry().max_attempts.max(1)
1017    }
1018
1019    /// Records a trusted "not applied": re-run if the budget allows, else
1020    /// `Failed`.
1021    async fn not_applied(&self, record: &EffectRecord) -> Result<EffectRecord, Interrupt> {
1022        let error = ErrorRecord {
1023            class: None,
1024            message: "verification found that the effect did not apply".into(),
1025        };
1026        if self.may_retry(record.attempt_count) {
1027            self.schedule_retry(record, FailureClass::Transient, Some(error))
1028                .await
1029        } else {
1030            self.transition(record, Transition::VerificationNotApplied, |r| {
1031                r.error = Some(error);
1032            })
1033            .await
1034        }
1035    }
1036
1037    /// Evaluates the precondition. Returns the reason to reject, if any.
1038    async fn check_precondition(
1039        &self,
1040        record: &EffectRecord,
1041    ) -> Result<Option<ErrorRecord>, Interrupt> {
1042        let Some(precondition) = &self.spec.precondition else {
1043            return Ok(None);
1044        };
1045        let max_checks = self.retry().max_attempts.max(1);
1046        let mut last_reason = String::new();
1047        for check in 1..=max_checks {
1048            let rejection = |message: String| {
1049                Some(ErrorRecord {
1050                    class: None,
1051                    message,
1052                })
1053            };
1054            match self
1055                .leased(tokio::spawn(precondition(context(record))))
1056                .await?
1057            {
1058                Ok(Precondition::Satisfied) => return Ok(None),
1059                Ok(Precondition::Rejected { reason }) => return Ok(rejection(reason)),
1060                Ok(Precondition::RetryLater { after, reason }) => {
1061                    last_reason = reason;
1062                    if check < max_checks {
1063                        self.sleep(after).await?;
1064                    }
1065                }
1066                // A broken check must not let the effect through.
1067                Err(join_error) => {
1068                    return Ok(rejection(format!(
1069                        "precondition check did not complete: {join_error}"
1070                    )));
1071                }
1072            }
1073        }
1074        Ok(Some(ErrorRecord {
1075            class: None,
1076            message: format!("precondition not satisfied after {max_checks} checks: {last_reason}"),
1077        }))
1078    }
1079
1080    /// Records that the next attempt waits for a backoff delay. The wait
1081    /// itself happens when the loop next sees the record as `Pending`, so a
1082    /// crash during it leaves a record that is safe to resume.
1083    async fn schedule_retry(
1084        &self,
1085        record: &EffectRecord,
1086        class: FailureClass,
1087        error: Option<ErrorRecord>,
1088    ) -> Result<EffectRecord, Interrupt> {
1089        let retry = record.attempt_count.saturating_sub(1);
1090        let delay = self.retry().delay(retry, class, jitter_sample());
1091        let at = self.rt.now() + delay;
1092        debug!(?delay, ?class, "retry scheduled");
1093        self.transition(record, Transition::ScheduleRetry, |r| {
1094            r.next_attempt_at = Some(at);
1095            r.error = error;
1096            r.payload =
1097                Some(json!({ "delay_ms": u64::try_from(delay.as_millis()).unwrap_or(u64::MAX) }));
1098        })
1099        .await
1100    }
1101
1102    /// Asks the runtime's approval provider, on its own task and under the
1103    /// lease. No provider, or a provider that panics, defers.
1104    async fn ask_approval(&self, record: &EffectRecord) -> Result<ApprovalDecision, Interrupt> {
1105        let Some(provider) = self.rt.inner.approval.clone() else {
1106            return Ok(ApprovalDecision::Deferred);
1107        };
1108        let request = ApprovalRequest {
1109            effect_id: record.id,
1110            key: record.key.clone(),
1111            kind: record.kind,
1112            risk: self.spec.risk,
1113            input: record.input.clone(),
1114            requested_by: record.created_by.clone(),
1115        };
1116        let task = tokio::spawn(async move { provider.request_boxed(request).await });
1117        Ok(self.leased(task).await?.unwrap_or_else(|e| {
1118            warn!(error = %e, "approval provider failed; deferring");
1119            ApprovalDecision::Deferred
1120        }))
1121    }
1122
1123    async fn sleep_until(&self, at: SystemTime) -> Result<(), Interrupt> {
1124        let wait = at.duration_since(self.rt.now()).unwrap_or_default();
1125        self.sleep(wait).await
1126    }
1127
1128    async fn sleep(&self, duration: Duration) -> Result<(), Interrupt> {
1129        if duration.is_zero() {
1130            return Ok(());
1131        }
1132        self.leased(tokio::time::sleep(duration)).await
1133    }
1134
1135    async fn leased<Fut: Future>(&self, future: Fut) -> Result<Fut::Output, Interrupt> {
1136        self.rt.with_lease(self.lease, future).await
1137    }
1138
1139    async fn transition(
1140        &self,
1141        record: &EffectRecord,
1142        transition: Transition,
1143        customize: impl FnOnce(&mut TransitionRequest),
1144    ) -> Result<EffectRecord, Interrupt> {
1145        self.rt
1146            .transition_leased(
1147                record,
1148                self.lease,
1149                self.spec.actor.as_deref(),
1150                transition,
1151                customize,
1152            )
1153            .await
1154    }
1155
1156    fn retry(&self) -> &RetryPolicy {
1157        &self.spec.retry
1158    }
1159
1160    /// Whether the runtime may run the effect again on its own after
1161    /// `attempts`: within budget, and not forbidden by the risk policy.
1162    fn may_retry(&self, attempts: u32) -> bool {
1163        self.spec.automatic_retry && self.retry().allows_another(attempts)
1164    }
1165}
1166
1167fn context(record: &EffectRecord) -> EffectContext {
1168    EffectContext {
1169        id: record.id,
1170        key: record.key.clone(),
1171        attempt: record.attempt_count,
1172    }
1173}
1174
1175/// An output as JSON, or a note for the audit trail if it cannot be stored.
1176/// The effect applied either way: failing to store its output must not make
1177/// it look failed, and the caller still gets the value.
1178fn output_json<T: Serialize>(value: &T) -> (Option<Value>, Option<Value>) {
1179    match serde_json::to_value(value) {
1180        Ok(output) => (Some(output), None),
1181        Err(e) => (None, Some(json!({ "output_not_stored": e.to_string() }))),
1182    }
1183}
1184
1185/// A uniform sample in `[0, 1)` for jitter, from the standard library's
1186/// per-instance random hash keys.
1187pub(crate) fn jitter_sample() -> f64 {
1188    let bits = RandomState::new().hash_one(0_u8) >> 11;
1189    #[allow(clippy::cast_precision_loss)]
1190    let sample = bits as f64 / (1_u64 << 53) as f64;
1191    sample
1192}
1193
1194/// Statuses a new call reports as they are, without acting.
1195fn settled(status: EffectStatus) -> bool {
1196    matches!(
1197        status,
1198        EffectStatus::Committed
1199            | EffectStatus::Failed
1200            | EffectStatus::Rejected
1201            | EffectStatus::NeedsIntervention
1202            | EffectStatus::Compensated
1203            | EffectStatus::CompensationFailed
1204    )
1205}
1206
1207fn check_matches(record: &EffectRecord, spec: &EffectSpec) -> Result<(), RuntimeError> {
1208    if record.kind != spec.capabilities.kind {
1209        return Err(RuntimeError::KindMismatch {
1210            id: record.id,
1211            stored: record.kind,
1212            requested: spec.capabilities.kind,
1213        });
1214    }
1215    if record.input_fingerprint != spec.fingerprint {
1216        return Err(RuntimeError::InputMismatch { id: record.id });
1217    }
1218    Ok(())
1219}
1220
1221/// The outcome a record represents. `fresh` is this call's output,
1222/// preferred over the stored copy.
1223fn report<T: DeserializeOwned>(
1224    record: &EffectRecord,
1225    fresh: Option<T>,
1226) -> Result<EffectOutcome<T>, RuntimeError> {
1227    let id = record.id;
1228    Ok(match record.status {
1229        EffectStatus::Committed => match fresh {
1230            Some(value) => EffectOutcome::Committed(value),
1231            None => EffectOutcome::Committed(
1232                serde_json::from_value(record.output.clone().unwrap_or(Value::Null))
1233                    .map_err(|source| RuntimeError::Output { id, source })?,
1234            ),
1235        },
1236        EffectStatus::Failed => EffectOutcome::Failed(last_error(record)),
1237        EffectStatus::Rejected => EffectOutcome::Rejected(last_error(record)),
1238        // A failed compensation also waits for an operator.
1239        EffectStatus::NeedsIntervention | EffectStatus::CompensationFailed => {
1240            EffectOutcome::NeedsIntervention { id }
1241        }
1242        EffectStatus::Unknown => EffectOutcome::Unknown { id },
1243        EffectStatus::Compensated => EffectOutcome::Compensated { id },
1244        EffectStatus::AwaitingApproval => EffectOutcome::AwaitingApproval { id },
1245        _ => EffectOutcome::InProgress { id },
1246    })
1247}
1248
1249pub(crate) fn last_error(record: &EffectRecord) -> ErrorRecord {
1250    record.last_error.clone().unwrap_or_else(|| ErrorRecord {
1251        class: None,
1252        message: "no error was recorded".into(),
1253    })
1254}