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