Skip to main content

agent_effects/
handler.rs

1//! Durable effect handlers: effects the runtime can finish without a caller.
2//!
3//! A closure effect can only be finished by a caller holding the closure. A
4//! handler is registered with the runtime under its effect name, and its
5//! input is stored with the record. So after a crash,
6//! [`Runtime::recover`](crate::Runtime::recover) can rebuild the call from
7//! the store and finish the effect with nobody calling: verify it, re-run it
8//! if that is safe, or escalate it.
9//!
10//! ```
11//! use agent_effects::handler::{EffectHandler, Handler, VerifiableEffect};
12//! use agent_effects::{EffectContext, EffectFailure, EffectKind, EffectOutcome, Runtime, Verification};
13//! use agent_effects_memory::MemoryStore;
14//!
15//! struct SendInvoice;
16//!
17//! impl EffectHandler for SendInvoice {
18//!     const NAME: &'static str = "invoice.send";
19//!     type Input = String; // the customer's email
20//!     type Output = String; // the provider's message id
21//!     type Error = EffectFailure;
22//!
23//!     fn kind(&self) -> EffectKind {
24//!         EffectKind::IrreversibleWrite
25//!     }
26//!
27//!     async fn execute(&self, ctx: &EffectContext, to: &String) -> Result<String, EffectFailure> {
28//!         // Send through the provider, forwarding ctx.idempotency_key().
29//!         Ok(format!("msg-for-{to}"))
30//!     }
31//! }
32//!
33//! impl VerifiableEffect for SendInvoice {
34//!     async fn verify(&self, _: &EffectContext, to: &String) -> Result<Verification<String>, EffectFailure> {
35//!         // Look the message up in the provider's outbox.
36//!         Ok(Verification::Confirmed(format!("msg-for-{to}")))
37//!     }
38//! }
39//!
40//! # #[tokio::main(flavor = "current_thread")]
41//! # async fn main() -> Result<(), agent_effects::RuntimeError> {
42//! let runtime = Runtime::builder(MemoryStore::new())
43//!     .register(Handler::new(SendInvoice).verifiable())
44//!     .build();
45//!
46//! let outcome = runtime
47//!     .submit::<SendInvoice>("invoice-1001", "ada@example.com".to_string())
48//!     .actor("agent:billing")
49//!     .await?;
50//! assert_eq!(outcome, EffectOutcome::Committed("msg-for-ada@example.com".to_string()));
51//! # Ok(())
52//! # }
53//! ```
54//!
55//! Optional capabilities are separate traits ([`VerifiableEffect`]; the
56//! compensation trait follows), so a handler only implements what it can
57//! actually do. The [`Handler`] builder only offers `.verifiable()` for
58//! handlers that implement it.
59
60use std::any::Any;
61use std::collections::HashMap;
62use std::fmt::Display;
63use std::future::{Future, IntoFuture};
64use std::pin::Pin;
65use std::sync::Arc;
66use std::time::Duration;
67
68use serde::Serialize;
69use serde::de::DeserializeOwned;
70use serde_json::Value;
71
72use crate::compensation::{
73    CompensationContext, CompensationOutcome, CompensationSpec, Compensator,
74};
75use crate::effect::{EffectContext, EffectFailure, EffectOutcome, EffectSpec, Precondition};
76use crate::error::RuntimeError;
77use crate::id::{EffectKey, EffectName, LogicalKey};
78use crate::kind::EffectKind;
79use crate::policy::{Capabilities, RiskLevel};
80use crate::retry::RetryPolicy;
81use crate::runtime::Runtime;
82use crate::state::EffectStatus;
83use crate::store::{EffectRecord, EffectStore, StoreError};
84use crate::verification::{NoVerification, Verification, VerificationMode, VerifyWith};
85
86/// An effect the runtime can execute, and re-execute after a crash, from
87/// its stored input.
88///
89/// Implement `execute`; override the other methods to describe the effect.
90/// The defaults are the most cautious choices.
91pub trait EffectHandler: Send + Sync + 'static {
92    /// The effect's name, e.g. `payment.charge`. One handler per name.
93    const NAME: &'static str;
94
95    /// What the effect acts on. Stored with the record so recovery can
96    /// rebuild the call: keep credentials in the handler, not here.
97    type Input: Serialize + DeserializeOwned + Send + Sync + 'static;
98
99    /// What the effect produces. Stored and replayed to later callers.
100    type Output: Serialize + DeserializeOwned + Send + 'static;
101
102    /// The handler's error, classified for the runtime.
103    type Error: Into<EffectFailure> + Send + 'static;
104
105    /// What re-executing does to the outside world. Defaults to
106    /// [`EffectKind::IrreversibleWrite`].
107    fn kind(&self) -> EffectKind {
108        EffectKind::IrreversibleWrite
109    }
110
111    /// Whether the remote system deduplicates on
112    /// [`EffectContext::idempotency_key`], which `execute` must send.
113    fn remote_idempotency(&self) -> bool {
114        false
115    }
116
117    /// The retry policy; `None` uses the runtime's.
118    fn retry_policy(&self) -> Option<RetryPolicy> {
119        None
120    }
121
122    /// How long to wait for one attempt; `None` waits indefinitely.
123    fn attempt_timeout(&self) -> Option<Duration> {
124        None
125    }
126
127    /// How much damage the effect could do; see
128    /// [`RiskPolicy`](crate::RiskPolicy). Defaults to [`RiskLevel::Low`].
129    fn risk(&self) -> RiskLevel {
130        RiskLevel::Low
131    }
132
133    /// Whether a human must approve the effect before its first attempt;
134    /// see [`approval`](crate::approval).
135    fn requires_approval(&self) -> bool {
136        false
137    }
138
139    /// Checked before the first attempt; see
140    /// [`EffectBuilder::precondition`](crate::EffectBuilder::precondition).
141    fn precondition(
142        &self,
143        ctx: &EffectContext,
144        input: &Self::Input,
145    ) -> impl Future<Output = Precondition> + Send {
146        let _ = (ctx, input);
147        async { Precondition::Satisfied }
148    }
149
150    /// Performs the effect once.
151    fn execute(
152        &self,
153        ctx: &EffectContext,
154        input: &Self::Input,
155    ) -> impl Future<Output = Result<Self::Output, Self::Error>> + Send;
156}
157
158/// A handler whose effect can be looked up in the remote system.
159///
160/// Register it with [`Handler::verifiable`] to verify after every success and
161/// to resolve unknown outcomes; see
162/// [`EffectBuilder::verify`](crate::EffectBuilder::verify).
163pub trait VerifiableEffect: EffectHandler {
164    /// How far a lookup can be trusted. Defaults to
165    /// [`VerificationMode::Authoritative`]; use
166    /// [`VerificationMode::EventuallyConsistent`] for lookups that lag.
167    fn verification_mode(&self) -> VerificationMode {
168        VerificationMode::Authoritative
169    }
170
171    /// Asks the remote system whether the effect applied.
172    fn verify(
173        &self,
174        ctx: &EffectContext,
175        input: &Self::Input,
176    ) -> impl Future<Output = Result<Verification<Self::Output>, Self::Error>> + Send;
177}
178
179/// A handler whose effect can be undone.
180///
181/// Register it with [`Handler::compensable`]; undo an effect with
182/// [`Runtime::compensate`](crate::Runtime::compensate). The compensation must
183/// be idempotent: a crashed or ambiguous attempt is run again. Send
184/// [`CompensationContext::idempotency_key`] to the remote system.
185pub trait CompensableEffect: EffectHandler {
186    /// Undoes the effect. `output` is what `execute` returned, if it was
187    /// stored.
188    fn compensate(
189        &self,
190        ctx: &CompensationContext,
191        input: &Self::Input,
192        output: Option<&Self::Output>,
193    ) -> impl Future<Output = Result<(), Self::Error>> + Send;
194}
195
196type BoxFuture<T> = Pin<Box<dyn Future<Output = T> + Send>>;
197
198type VerifyFn<H> = Arc<
199    dyn Fn(
200            Arc<H>,
201            EffectContext,
202            Arc<<H as EffectHandler>::Input>,
203        ) -> BoxFuture<Result<Verification<<H as EffectHandler>::Output>, EffectFailure>>
204        + Send
205        + Sync,
206>;
207
208/// A handler and the capabilities it is registered with. Pass it to
209/// [`RuntimeBuilder::register`](crate::RuntimeBuilder::register).
210pub struct Handler<H: EffectHandler> {
211    effect: Arc<H>,
212    verify: Option<(VerificationMode, VerifyFn<H>)>,
213    compensate: Option<Compensator>,
214}
215
216impl<H: EffectHandler> Clone for Handler<H> {
217    fn clone(&self) -> Self {
218        Self {
219            effect: Arc::clone(&self.effect),
220            verify: self.verify.clone(),
221            compensate: self.compensate.clone(),
222        }
223    }
224}
225
226impl<H: EffectHandler> Handler<H> {
227    /// Registers `handler` with no optional capabilities.
228    pub fn new(handler: H) -> Self {
229        Self {
230            effect: Arc::new(handler),
231            verify: None,
232            compensate: None,
233        }
234    }
235}
236
237impl<H: CompensableEffect> Handler<H> {
238    /// Lets the effect be undone with
239    /// [`Runtime::compensate`](crate::Runtime::compensate), using
240    /// [`CompensableEffect::compensate`], and lets recovery finish an
241    /// interrupted compensation.
242    #[must_use]
243    pub fn compensable(mut self) -> Self {
244        let handler = Arc::clone(&self.effect);
245        self.compensate = Some(Arc::new(move |ctx, input, output| {
246            let handler = Arc::clone(&handler);
247            Box::pin(async move {
248                let input: H::Input = serde_json::from_value(input.unwrap_or(Value::Null))
249                    .map_err(|e| {
250                        EffectFailure::permanent(format!("stored input does not match: {e}"))
251                    })?;
252                let output: Option<H::Output> = output
253                    .map(serde_json::from_value)
254                    .transpose()
255                    .map_err(|e| {
256                        EffectFailure::permanent(format!("stored output does not match: {e}"))
257                    })?;
258                handler
259                    .compensate(&ctx, &input, output.as_ref())
260                    .await
261                    .map_err(Into::into)
262            })
263        }));
264        self
265    }
266}
267
268impl<H: VerifiableEffect> Handler<H> {
269    /// Verifies the effect after every success and to resolve unknown
270    /// outcomes, using [`VerifiableEffect::verify`].
271    #[must_use]
272    pub fn verifiable(mut self) -> Self {
273        let mode = self.effect.verification_mode();
274        let verify: VerifyFn<H> = Arc::new(|handler, ctx, input| {
275            Box::pin(async move { handler.verify(&ctx, &input).await.map_err(Into::into) })
276        });
277        self.verify = Some((mode, verify));
278        self
279    }
280}
281
282/// A registered handler, erased so the runtime can hold handlers of many
283/// types and resume their effects by name.
284pub(crate) struct Registered<S> {
285    /// The `Handler<H>`, for typed submission.
286    typed: Arc<dyn Any + Send + Sync>,
287    /// Finishes an effect of this handler from its stored record.
288    pub(crate) resume: Resume<S>,
289}
290
291pub(crate) type Resume<S> = Arc<
292    dyn Fn(Runtime<S>, EffectRecord) -> BoxFuture<Result<EffectStatus, RuntimeError>> + Send + Sync,
293>;
294
295pub(crate) type Registry<S> = HashMap<&'static str, Registered<S>>;
296
297impl<S: EffectStore> Registered<S> {
298    pub(crate) fn new<H: EffectHandler>(handler: Handler<H>) -> Self {
299        let typed: Arc<dyn Any + Send + Sync> = Arc::new(handler.clone());
300        let resume: Resume<S> = Arc::new(move |runtime: Runtime<S>, record: EffectRecord| {
301            let handler = handler.clone();
302            Box::pin(async move { resume(&runtime, &handler, record).await })
303        });
304        Self { typed, resume }
305    }
306
307    pub(crate) fn typed<H: EffectHandler>(&self) -> Option<Handler<H>> {
308        self.typed.downcast_ref::<Handler<H>>().cloned()
309    }
310}
311
312/// A submission of input to a registered handler. Await it to run the
313/// effect. Created by [`Runtime::submit`].
314#[must_use = "a submission does nothing until it is awaited"]
315pub struct Submission<'a, S, H: EffectHandler> {
316    runtime: &'a Runtime<S>,
317    key: String,
318    input: H::Input,
319    actor: Option<String>,
320}
321
322impl<'a, S, H: EffectHandler> Submission<'a, S, H> {
323    pub(crate) fn new(runtime: &'a Runtime<S>, key: impl Display, input: H::Input) -> Self {
324        Self {
325            runtime,
326            key: key.to_string(),
327            input,
328            actor: None,
329        }
330    }
331
332    /// Who is asking for the effect, e.g. `agent:billing`.
333    pub fn actor(mut self, actor: impl Into<String>) -> Self {
334        self.actor = Some(actor.into());
335        self
336    }
337}
338
339impl<'a, S: EffectStore, H: EffectHandler> IntoFuture for Submission<'a, S, H> {
340    type Output = Result<EffectOutcome<H::Output>, RuntimeError>;
341    type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>;
342
343    fn into_future(self) -> Self::IntoFuture {
344        Box::pin(async move {
345            let handler = self
346                .runtime
347                .handler::<H>()
348                .ok_or(RuntimeError::NotRegistered { name: H::NAME })?;
349            let key = EffectKey::new(EffectName::new(H::NAME)?, LogicalKey::new(self.key)?);
350            let json = serde_json::to_value(&self.input).map_err(RuntimeError::Input)?;
351            let stored = Stored {
352                // Computed by the runtime, after redaction.
353                fingerprint: None,
354                json,
355                actor: self.actor,
356                from_record: false,
357            };
358            run(self.runtime, &handler, key, Arc::new(self.input), stored).await
359        })
360    }
361}
362
363/// Finishes an effect of `handler` from its record, with no caller: its
364/// compensation if one is under way, else the effect itself.
365async fn resume<S: EffectStore, H: EffectHandler>(
366    runtime: &Runtime<S>,
367    handler: &Handler<H>,
368    record: EffectRecord,
369) -> Result<EffectStatus, RuntimeError> {
370    let id = record.id;
371    if record.status == EffectStatus::Compensating {
372        let compensate = handler
373            .compensate
374            .clone()
375            .ok_or(RuntimeError::NotCompensable { name: H::NAME })?;
376        let spec = CompensationSpec {
377            key: record.key.clone(),
378            reason: None,
379            actor: Some(format!("recovery:{}", runtime.worker_id())),
380            retry: handler.effect.retry_policy(),
381            attempt_timeout: handler.effect.attempt_timeout(),
382        };
383        runtime.compensate_effect(spec, compensate).await?;
384        let settled = runtime
385            .store()
386            .get(id)
387            .await?
388            .ok_or(StoreError::NotFound(id))?;
389        return Ok(settled.status);
390    }
391    let json = record.input.clone().unwrap_or(Value::Null);
392    let input: H::Input = serde_json::from_value(json.clone())
393        .map_err(|source| RuntimeError::StoredInput { id, source })?;
394    // Rebuilt from the record itself, so it is the same effect by
395    // definition: reuse its fingerprint rather than recompute it, which
396    // would refuse every older record if canonicalization ever changed.
397    let stored = Stored {
398        json,
399        fingerprint: record.input_fingerprint.clone(),
400        actor: record.created_by.clone(),
401        from_record: true,
402    };
403    run(
404        runtime,
405        handler,
406        record.key.clone(),
407        Arc::new(input),
408        stored,
409    )
410    .await?;
411    let settled = runtime
412        .store()
413        .get(id)
414        .await?
415        .ok_or(StoreError::NotFound(id))?;
416    Ok(settled.status)
417}
418
419/// What is recorded about an effect besides its key.
420struct Stored {
421    json: Value,
422    fingerprint: Option<String>,
423    actor: Option<String>,
424    /// From an existing record: already redacted, fingerprint kept as is.
425    from_record: bool,
426}
427
428/// Runs `handler`'s effect through the same machinery as a closure effect.
429async fn run<S: EffectStore, H: EffectHandler>(
430    runtime: &Runtime<S>,
431    handler: &Handler<H>,
432    key: EffectKey,
433    input: Arc<H::Input>,
434    stored: Stored,
435) -> Result<EffectOutcome<H::Output>, RuntimeError> {
436    let effect = Arc::clone(&handler.effect);
437    let precondition = {
438        let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
439        Arc::new(move |ctx: EffectContext| -> BoxFuture<Precondition> {
440            let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
441            Box::pin(async move { effect.precondition(&ctx, &input).await })
442        })
443    };
444    let spec = EffectSpec {
445        fingerprint: stored.fingerprint,
446        input: Some(stored.json),
447        key,
448        capabilities: Capabilities {
449            kind: effect.kind(),
450            remote_idempotency: effect.remote_idempotency(),
451            verification: handler
452                .verify
453                .as_ref()
454                .map_or(VerificationMode::None, |(mode, _)| *mode),
455        },
456        actor: stored.actor,
457        retry: effect
458            .retry_policy()
459            .unwrap_or_else(|| runtime.default_retry()),
460        attempt_timeout: effect.attempt_timeout(),
461        precondition: Some(precondition),
462        require_approval: effect.requires_approval(),
463        risk: effect.risk(),
464        automatic_retry: true,
465        input_stored: stored.from_record,
466    };
467    let action = {
468        let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
469        move |ctx: EffectContext| {
470            let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
471            async move { effect.execute(&ctx, &input).await.map_err(Into::into) }
472        }
473    };
474    match &handler.verify {
475        Some((_, verify)) => {
476            let verify = Arc::clone(verify);
477            let checker =
478                VerifyWith(move |ctx| verify(Arc::clone(&effect), ctx, Arc::clone(&input)));
479            runtime.execute(spec, action, checker).await
480        }
481        None => runtime.execute(spec, action, NoVerification).await,
482    }
483}
484
485/// A request to undo a registered handler's effect. Await it. Created by
486/// [`Runtime::compensate`](crate::Runtime::compensate).
487#[must_use = "a compensation does nothing until it is awaited"]
488pub struct CompensationSubmission<'a, S, H: EffectHandler> {
489    runtime: &'a Runtime<S>,
490    key: String,
491    reason: Option<String>,
492    actor: Option<String>,
493    _handler: std::marker::PhantomData<fn() -> H>,
494}
495
496impl<'a, S, H: EffectHandler> CompensationSubmission<'a, S, H> {
497    pub(crate) fn new(runtime: &'a Runtime<S>, key: impl Display) -> Self {
498        Self {
499            runtime,
500            key: key.to_string(),
501            reason: None,
502            actor: None,
503            _handler: std::marker::PhantomData,
504        }
505    }
506
507    /// Why the effect is being undone; recorded in the audit trail.
508    pub fn reason(mut self, reason: impl Into<String>) -> Self {
509        self.reason = Some(reason.into());
510        self
511    }
512
513    /// Who is undoing it.
514    pub fn actor(mut self, actor: impl Into<String>) -> Self {
515        self.actor = Some(actor.into());
516        self
517    }
518}
519
520impl<'a, S: EffectStore, H: EffectHandler> IntoFuture for CompensationSubmission<'a, S, H> {
521    type Output = Result<CompensationOutcome, RuntimeError>;
522    type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>;
523
524    fn into_future(self) -> Self::IntoFuture {
525        Box::pin(async move {
526            let handler = self
527                .runtime
528                .handler::<H>()
529                .ok_or(RuntimeError::NotRegistered { name: H::NAME })?;
530            let compensate = handler
531                .compensate
532                .clone()
533                .ok_or(RuntimeError::NotCompensable { name: H::NAME })?;
534            let spec = CompensationSpec {
535                key: EffectKey::new(EffectName::new(H::NAME)?, LogicalKey::new(self.key)?),
536                reason: self.reason,
537                actor: self.actor,
538                retry: handler.effect.retry_policy(),
539                attempt_timeout: handler.effect.attempt_timeout(),
540            };
541            self.runtime.compensate_effect(spec, compensate).await
542        })
543    }
544}