Skip to main content

agent_effects/
effect.rs

1//! Describing an effect: the builder, what its action sees, and how it ends.
2
3use std::error::Error;
4use std::fmt;
5use std::future::Future;
6use std::pin::Pin;
7use std::sync::Arc;
8use std::time::Duration;
9
10use serde::Serialize;
11use serde::de::DeserializeOwned;
12use serde_json::Value;
13
14use crate::error::RuntimeError;
15use crate::failure::FailureClass;
16use crate::id::{EffectId, EffectKey, IdempotencyKey};
17use crate::kind::EffectKind;
18use crate::policy::{Capabilities, RiskLevel};
19use crate::retry::RetryPolicy;
20use crate::runtime::Runtime;
21use crate::store::{EffectStore, ErrorRecord};
22use crate::verification::{NoVerification, Verification, VerificationMode, Verifier, VerifyWith};
23
24/// What an effect's action knows about the attempt it is running.
25#[derive(Clone, Debug)]
26pub struct EffectContext {
27    pub(crate) id: EffectId,
28    pub(crate) key: EffectKey,
29    pub(crate) attempt: u32,
30}
31
32impl EffectContext {
33    /// The effect's record id.
34    pub fn effect_id(&self) -> EffectId {
35        self.id
36    }
37
38    /// The effect's logical identity.
39    pub fn key(&self) -> &EffectKey {
40        &self.key
41    }
42
43    /// The key to forward to the remote system, e.g. as an HTTP
44    /// `Idempotency-Key` header. The same for every attempt of the effect.
45    pub fn idempotency_key(&self) -> IdempotencyKey {
46        self.key.idempotency_key()
47    }
48
49    /// The attempt number, starting at 1.
50    pub fn attempt(&self) -> u32 {
51        self.attempt
52    }
53}
54
55/// A failed action, classified so the runtime knows whether it applied.
56///
57/// ```
58/// # use agent_effects::EffectFailure;
59/// // The connection was refused: nothing reached the remote system.
60/// let failure = EffectFailure::ambiguous("connection refused").request_sent(false);
61/// ```
62#[derive(Debug)]
63pub struct EffectFailure {
64    class: FailureClass,
65    request_sent: Option<bool>,
66    source: Box<dyn Error + Send + Sync>,
67}
68
69impl EffectFailure {
70    /// A failure of class `class`.
71    pub fn new(class: FailureClass, source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
72        Self {
73            class,
74            request_sent: None,
75            source: source.into(),
76        }
77    }
78
79    /// The request did not apply and may succeed if repeated.
80    pub fn transient(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
81        Self::new(FailureClass::Transient, source)
82    }
83
84    /// The request did not apply and will not succeed if repeated.
85    pub fn permanent(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
86        Self::new(FailureClass::Permanent, source)
87    }
88
89    /// The request may or may not have applied, e.g. a timeout after sending.
90    pub fn ambiguous(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
91        Self::new(FailureClass::Ambiguous, source)
92    }
93
94    /// The remote system asked the caller to slow down.
95    pub fn rate_limited(
96        retry_after: Option<Duration>,
97        source: impl Into<Box<dyn Error + Send + Sync>>,
98    ) -> Self {
99        Self::new(FailureClass::RateLimited { retry_after }, source)
100    }
101
102    /// The credentials were missing or invalid.
103    pub fn authentication(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
104        Self::new(FailureClass::Authentication, source)
105    }
106
107    /// The credentials were not allowed to perform the action.
108    pub fn authorization(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
109        Self::new(FailureClass::Authorization, source)
110    }
111
112    /// The remote system rejected the request as malformed.
113    pub fn validation(source: impl Into<Box<dyn Error + Send + Sync>>) -> Self {
114        Self::new(FailureClass::Validation, source)
115    }
116
117    /// Records whether the request reached the remote system. `false` turns
118    /// an ambiguous failure into a transient one: nothing can have applied.
119    #[must_use]
120    pub fn request_sent(mut self, sent: bool) -> Self {
121        self.request_sent = Some(sent);
122        self
123    }
124
125    /// The classification after taking [`Self::request_sent`] into account.
126    pub fn class(&self) -> FailureClass {
127        self.class.with_request_sent(self.request_sent)
128    }
129
130    pub(crate) fn to_record(&self) -> ErrorRecord {
131        ErrorRecord {
132            class: Some(self.class()),
133            message: self.source.to_string(),
134        }
135    }
136}
137
138impl fmt::Display for EffectFailure {
139    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
140        write!(f, "{:?} failure: {}", self.class(), self.source)
141    }
142}
143
144impl Error for EffectFailure {
145    fn source(&self) -> Option<&(dyn Error + 'static)> {
146        Some(self.source.as_ref())
147    }
148}
149
150/// How an effect ended, or where it stands.
151///
152/// `Unknown` and `NeedsIntervention` are results, not errors: the effect may
153/// have changed the outside world, and the caller must not treat it as
154/// failed.
155#[derive(Debug, Clone, PartialEq)]
156#[non_exhaustive]
157pub enum EffectOutcome<T> {
158    /// The effect applied. Carries the action's output, from this call or
159    /// from the record of an earlier one.
160    Committed(T),
161    /// The effect definitely did not apply and will not be retried.
162    Failed(ErrorRecord),
163    /// The effect was refused before it ran.
164    Rejected(ErrorRecord),
165    /// The effect may or may not have applied. Calling again with the same
166    /// key resolves it if that is safe.
167    Unknown {
168        /// The effect.
169        id: EffectId,
170    },
171    /// The runtime cannot resolve the outcome safely; an operator must.
172    NeedsIntervention {
173        /// The effect.
174        id: EffectId,
175    },
176    /// Another caller or worker is running the effect right now.
177    InProgress {
178        /// The effect.
179        id: EffectId,
180    },
181    /// The effect waits for a human decision. Call again later, or have an
182    /// operator use [`Runtime::approve`](crate::Runtime::approve) /
183    /// [`Runtime::deny`](crate::Runtime::deny).
184    AwaitingApproval {
185        /// The effect.
186        id: EffectId,
187    },
188    /// The effect applied and was later undone by compensation.
189    Compensated {
190        /// The effect.
191        id: EffectId,
192    },
193}
194
195/// Whether an effect may still run, checked before its first attempt.
196#[derive(Clone, Debug, PartialEq, Eq)]
197pub enum Precondition {
198    /// Go ahead.
199    Satisfied,
200    /// Never run this effect, e.g. "order is no longer pending". The effect
201    /// ends [`EffectOutcome::Rejected`].
202    Rejected {
203        /// Why, recorded on the effect.
204        reason: String,
205    },
206    /// Not yet; check again after `after`, e.g. while a dependency settles.
207    /// After as many checks as the retry policy allows attempts, the effect
208    /// is rejected.
209    RetryLater {
210        /// How long to wait before the next check.
211        after: Duration,
212        /// Why, recorded if the effect is eventually rejected.
213        reason: String,
214    },
215}
216
217impl Precondition {
218    /// A rejection with `reason`.
219    pub fn reject(reason: impl Into<String>) -> Self {
220        Self::Rejected {
221            reason: reason.into(),
222        }
223    }
224
225    /// A deferral of `after`, with `reason`.
226    pub fn retry_later(after: Duration, reason: impl Into<String>) -> Self {
227        Self::RetryLater {
228            after,
229            reason: reason.into(),
230        }
231    }
232}
233
234pub(crate) type PreconditionFn =
235    Arc<dyn Fn(EffectContext) -> Pin<Box<dyn Future<Output = Precondition> + Send>> + Send + Sync>;
236
237/// Builds and runs one effect. Created by [`Runtime::effect`].
238///
239/// Invalid names, keys or inputs are reported when the effect is run, so the
240/// whole chain needs a single `?`.
241#[must_use = "an effect does nothing until `run` is awaited"]
242pub struct EffectBuilder<S, V = NoVerification> {
243    runtime: Runtime<S>,
244    name: String,
245    key: String,
246    kind: EffectKind,
247    remote_idempotency: bool,
248    input: Option<Result<Value, serde_json::Error>>,
249    actor: Option<String>,
250    retry: RetryPolicy,
251    attempt_timeout: Option<Duration>,
252    precondition: Option<PreconditionFn>,
253    require_approval: bool,
254    risk: RiskLevel,
255    verification: VerificationMode,
256    verifier: V,
257}
258
259impl<S: EffectStore> EffectBuilder<S> {
260    pub(crate) fn new(runtime: Runtime<S>, name: String, key: String) -> Self {
261        let retry = runtime.default_retry();
262        Self {
263            runtime,
264            name,
265            key,
266            kind: EffectKind::IrreversibleWrite,
267            remote_idempotency: false,
268            input: None,
269            actor: None,
270            retry,
271            attempt_timeout: None,
272            precondition: None,
273            require_approval: false,
274            risk: RiskLevel::Low,
275            verification: VerificationMode::None,
276            verifier: NoVerification,
277        }
278    }
279
280    /// Verifies the effect against a remote system that reads its own
281    /// writes: [`Verification::NotApplied`] is trusted at once.
282    ///
283    /// The check runs after every successful attempt (a postcondition) and
284    /// to reconcile an attempt whose outcome is unknown. It makes re-running
285    /// safe even for irreversible effects, because the runtime only re-runs
286    /// an effect the remote system says did not apply.
287    pub fn verify<T, F, Fut>(self, check: F) -> EffectBuilder<S, VerifyWith<F>>
288    where
289        F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
290        Fut: Future<Output = Result<Verification<T>, EffectFailure>> + Send + 'static,
291    {
292        self.with_verifier(VerificationMode::Authoritative, VerifyWith(check))
293    }
294
295    /// Like [`Self::verify`], for a remote lookup that lags behind its writes
296    /// by up to `settle`. [`Verification::NotApplied`] is trusted only once
297    /// `settle` has passed since the attempt ended; until then the runtime
298    /// waits and checks again.
299    pub fn verify_eventually<T, F, Fut>(
300        self,
301        settle: Duration,
302        check: F,
303    ) -> EffectBuilder<S, VerifyWith<F>>
304    where
305        F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
306        Fut: Future<Output = Result<Verification<T>, EffectFailure>> + Send + 'static,
307    {
308        self.with_verifier(
309            VerificationMode::EventuallyConsistent { settle },
310            VerifyWith(check),
311        )
312    }
313
314    fn with_verifier<V>(self, mode: VerificationMode, verifier: V) -> EffectBuilder<S, V> {
315        EffectBuilder {
316            runtime: self.runtime,
317            name: self.name,
318            key: self.key,
319            kind: self.kind,
320            remote_idempotency: self.remote_idempotency,
321            input: self.input,
322            actor: self.actor,
323            retry: self.retry,
324            attempt_timeout: self.attempt_timeout,
325            precondition: self.precondition,
326            require_approval: self.require_approval,
327            risk: self.risk,
328            verification: mode,
329            verifier,
330        }
331    }
332}
333
334impl<S: EffectStore, V> EffectBuilder<S, V> {
335    /// The effect's kind. Defaults to [`EffectKind::IrreversibleWrite`], the
336    /// most cautious choice.
337    pub fn kind(mut self, kind: EffectKind) -> Self {
338        self.kind = kind;
339        self
340    }
341
342    /// Declares that the remote system deduplicates on
343    /// [`EffectContext::idempotency_key`], which the action must send. This
344    /// makes re-running after an unknown outcome safe.
345    pub fn remote_idempotency(mut self, supported: bool) -> Self {
346        self.remote_idempotency = supported;
347        self
348    }
349
350    /// The input the action acts on. It is stored for audit and fingerprinted:
351    /// reusing the key with a different input fails with
352    /// [`RuntimeError::InputMismatch`].
353    ///
354    /// The input is persisted as-is. Keep credentials in the action's
355    /// captured state, not in the input.
356    pub fn input<I: Serialize + ?Sized>(mut self, input: &I) -> Self {
357        self.input = Some(serde_json::to_value(input));
358        self
359    }
360
361    /// Who is asking for the effect, e.g. `agent:refund-agent`. Recorded on
362    /// the effect and its audit events.
363    pub fn actor(mut self, actor: impl Into<String>) -> Self {
364        self.actor = Some(actor.into());
365        self
366    }
367
368    /// The retry policy. Defaults to the runtime's
369    /// ([`RuntimeBuilder::retry_policy`](crate::RuntimeBuilder::retry_policy)).
370    ///
371    /// `max_attempts` bounds the attempts over the effect's whole life,
372    /// across calls and restarts. It also bounds verification and
373    /// precondition checks per call.
374    pub fn retry(mut self, policy: RetryPolicy) -> Self {
375        self.retry = policy;
376        self
377    }
378
379    /// Gives up waiting for an attempt after `timeout`. The request may have
380    /// been sent, so a timeout is an ambiguous failure.
381    pub fn attempt_timeout(mut self, timeout: Duration) -> Self {
382        self.attempt_timeout = Some(timeout);
383        self
384    }
385
386    /// A check that must pass before the effect's first attempt, so a stale
387    /// decision is not carried out, e.g. "refund only if the order is still
388    /// unrefunded".
389    ///
390    /// It runs only before the first attempt. Once an attempt may have
391    /// applied, the effect's own success could falsify the check: a refund
392    /// that went through makes "not yet refunded" false.
393    pub fn precondition<F, Fut>(mut self, check: F) -> Self
394    where
395        F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
396        Fut: Future<Output = Precondition> + Send + 'static,
397    {
398        self.precondition = Some(Arc::new(move |ctx| Box::pin(check(ctx))));
399        self
400    }
401
402    /// How much damage the effect could do; the runtime's
403    /// [`RiskPolicy`](crate::RiskPolicy) adds requirements by risk. Defaults
404    /// to [`RiskLevel::Low`].
405    pub fn risk(mut self, risk: RiskLevel) -> Self {
406        self.risk = risk;
407        self
408    }
409
410    /// Requires a human decision before the first attempt; see
411    /// [`approval`](crate::approval). The effect waits in
412    /// `AwaitingApproval`, durably, until the runtime's approval provider or
413    /// an operator decides.
414    pub fn require_approval(mut self) -> Self {
415        self.require_approval = true;
416        self
417    }
418
419    /// Runs the effect, or attaches to an earlier run with the same key.
420    ///
421    /// The action may be called more than once over the effect's life (for
422    /// example to re-run an idempotent effect whose outcome was unknown), so
423    /// it is an `Fn`. It runs on a spawned task: dropping the returned future
424    /// does not abort an attempt that has started, and its result is still
425    /// recorded.
426    ///
427    /// # Errors
428    ///
429    /// Infrastructure failures only; see [`RuntimeError`]. Every effect
430    /// result, including an unknown one, is an [`EffectOutcome`].
431    pub async fn run<T, F, Fut>(self, action: F) -> Result<EffectOutcome<T>, RuntimeError>
432    where
433        T: Serialize + DeserializeOwned + Send + 'static,
434        F: Fn(EffectContext) -> Fut + Send + Sync + 'static,
435        Fut: Future<Output = Result<T, EffectFailure>> + Send + 'static,
436        V: Verifier<T>,
437    {
438        let key = EffectKey::new(
439            crate::id::EffectName::new(self.name)?,
440            crate::id::LogicalKey::new(self.key)?,
441        );
442        let input = self.input.transpose().map_err(RuntimeError::Input)?;
443        let spec = EffectSpec {
444            // Computed by the runtime, after redaction.
445            fingerprint: None,
446            input,
447            key,
448            capabilities: Capabilities {
449                kind: self.kind,
450                remote_idempotency: self.remote_idempotency,
451                verification: self.verification,
452            },
453            actor: self.actor,
454            retry: self.retry,
455            attempt_timeout: self.attempt_timeout,
456            precondition: self.precondition,
457            require_approval: self.require_approval,
458            risk: self.risk,
459            automatic_retry: true,
460            input_stored: false,
461        };
462        self.runtime.execute(spec, action, self.verifier).await
463    }
464}
465
466/// Everything about an effect except its action and verifier.
467pub(crate) struct EffectSpec {
468    pub(crate) key: EffectKey,
469    pub(crate) capabilities: Capabilities,
470    pub(crate) input: Option<Value>,
471    pub(crate) fingerprint: Option<String>,
472    pub(crate) actor: Option<String>,
473    pub(crate) retry: RetryPolicy,
474    pub(crate) attempt_timeout: Option<Duration>,
475    pub(crate) precondition: Option<PreconditionFn>,
476    pub(crate) require_approval: bool,
477    pub(crate) risk: RiskLevel,
478    /// Whether the runtime may retry or re-run on its own; a risk policy
479    /// can turn it off.
480    pub(crate) automatic_retry: bool,
481    /// The input comes from the store (a resumed effect): it is already
482    /// redacted, and `fingerprint` is the stored one, used as is.
483    pub(crate) input_stored: bool,
484}