agent-effects 0.1.1

Reliable side-effect execution for AI agents and autonomous applications
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
//! Durable effect handlers: effects the runtime can finish without a caller.
//!
//! A closure effect can only be finished by a caller holding the closure. A
//! handler is registered with the runtime under its effect name, and its
//! input is stored with the record. So after a crash,
//! [`Runtime::recover`](crate::Runtime::recover) can rebuild the call from
//! the store and finish the effect with nobody calling: verify it, re-run it
//! if that is safe, or escalate it.
//!
//! ```
//! use agent_effects::handler::{EffectHandler, Handler, VerifiableEffect};
//! use agent_effects::{EffectContext, EffectFailure, EffectKind, EffectOutcome, Runtime, Verification};
//! use agent_effects_memory::MemoryStore;
//!
//! struct SendInvoice;
//!
//! impl EffectHandler for SendInvoice {
//!     const NAME: &'static str = "invoice.send";
//!     type Input = String; // the customer's email
//!     type Output = String; // the provider's message id
//!     type Error = EffectFailure;
//!
//!     fn kind(&self) -> EffectKind {
//!         EffectKind::IrreversibleWrite
//!     }
//!
//!     async fn execute(&self, ctx: &EffectContext, to: &String) -> Result<String, EffectFailure> {
//!         // Send through the provider, forwarding ctx.idempotency_key().
//!         Ok(format!("msg-for-{to}"))
//!     }
//! }
//!
//! impl VerifiableEffect for SendInvoice {
//!     async fn verify(&self, _: &EffectContext, to: &String) -> Result<Verification<String>, EffectFailure> {
//!         // Look the message up in the provider's outbox.
//!         Ok(Verification::Confirmed(format!("msg-for-{to}")))
//!     }
//! }
//!
//! # #[tokio::main(flavor = "current_thread")]
//! # async fn main() -> Result<(), agent_effects::RuntimeError> {
//! let runtime = Runtime::builder(MemoryStore::new())
//!     .register(Handler::new(SendInvoice).verifiable())
//!     .build();
//!
//! let outcome = runtime
//!     .submit::<SendInvoice>("invoice-1001", "ada@example.com".to_string())
//!     .actor("agent:billing")
//!     .await?;
//! assert_eq!(outcome, EffectOutcome::Committed("msg-for-ada@example.com".to_string()));
//! # Ok(())
//! # }
//! ```
//!
//! Optional capabilities are separate traits ([`VerifiableEffect`]; the
//! compensation trait follows), so a handler only implements what it can
//! actually do. The [`Handler`] builder only offers `.verifiable()` for
//! handlers that implement it.

use std::any::Any;
use std::collections::HashMap;
use std::fmt::Display;
use std::future::{Future, IntoFuture};
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;

use serde::Serialize;
use serde::de::DeserializeOwned;
use serde_json::Value;

use crate::compensation::{
    CompensationContext, CompensationOutcome, CompensationSpec, Compensator,
};
use crate::effect::{EffectContext, EffectFailure, EffectOutcome, EffectSpec, Precondition};
use crate::error::RuntimeError;
use crate::id::{EffectKey, EffectName, LogicalKey};
use crate::kind::EffectKind;
use crate::policy::{Capabilities, RiskLevel};
use crate::retry::RetryPolicy;
use crate::runtime::Runtime;
use crate::state::EffectStatus;
use crate::store::{EffectRecord, EffectStore, StoreError};
use crate::verification::{NoVerification, Verification, VerificationMode, VerifyWith};

/// An effect the runtime can execute, and re-execute after a crash, from
/// its stored input.
///
/// Implement `execute`; override the other methods to describe the effect.
/// The defaults are the most cautious choices.
pub trait EffectHandler: Send + Sync + 'static {
    /// The effect's name, e.g. `payment.charge`. One handler per name.
    const NAME: &'static str;

    /// What the effect acts on. Stored with the record so recovery can
    /// rebuild the call: keep credentials in the handler, not here.
    type Input: Serialize + DeserializeOwned + Send + Sync + 'static;

    /// What the effect produces. Stored and replayed to later callers.
    type Output: Serialize + DeserializeOwned + Send + 'static;

    /// The handler's error, classified for the runtime.
    type Error: Into<EffectFailure> + Send + 'static;

    /// What re-executing does to the outside world. Defaults to
    /// [`EffectKind::IrreversibleWrite`].
    fn kind(&self) -> EffectKind {
        EffectKind::IrreversibleWrite
    }

    /// Whether the remote system deduplicates on
    /// [`EffectContext::idempotency_key`], which `execute` must send.
    fn remote_idempotency(&self) -> bool {
        false
    }

    /// The retry policy; `None` uses the runtime's.
    fn retry_policy(&self) -> Option<RetryPolicy> {
        None
    }

    /// How long to wait for one attempt; `None` waits indefinitely.
    fn attempt_timeout(&self) -> Option<Duration> {
        None
    }

    /// How much damage the effect could do; see
    /// [`RiskPolicy`](crate::RiskPolicy). Defaults to [`RiskLevel::Low`].
    fn risk(&self) -> RiskLevel {
        RiskLevel::Low
    }

    /// Whether a human must approve the effect before its first attempt;
    /// see [`approval`](crate::approval).
    fn requires_approval(&self) -> bool {
        false
    }

    /// Checked before the first attempt; see
    /// [`EffectBuilder::precondition`](crate::EffectBuilder::precondition).
    fn precondition(
        &self,
        ctx: &EffectContext,
        input: &Self::Input,
    ) -> impl Future<Output = Precondition> + Send {
        let _ = (ctx, input);
        async { Precondition::Satisfied }
    }

    /// Performs the effect once.
    fn execute(
        &self,
        ctx: &EffectContext,
        input: &Self::Input,
    ) -> impl Future<Output = Result<Self::Output, Self::Error>> + Send;
}

/// A handler whose effect can be looked up in the remote system.
///
/// Register it with [`Handler::verifiable`] to verify after every success and
/// to resolve unknown outcomes; see
/// [`EffectBuilder::verify`](crate::EffectBuilder::verify).
pub trait VerifiableEffect: EffectHandler {
    /// How far a lookup can be trusted. Defaults to
    /// [`VerificationMode::Authoritative`]; use
    /// [`VerificationMode::EventuallyConsistent`] for lookups that lag.
    fn verification_mode(&self) -> VerificationMode {
        VerificationMode::Authoritative
    }

    /// Asks the remote system whether the effect applied.
    fn verify(
        &self,
        ctx: &EffectContext,
        input: &Self::Input,
    ) -> impl Future<Output = Result<Verification<Self::Output>, Self::Error>> + Send;
}

/// A handler whose effect can be undone.
///
/// Register it with [`Handler::compensable`]; undo an effect with
/// [`Runtime::compensate`](crate::Runtime::compensate). The compensation must
/// be idempotent: a crashed or ambiguous attempt is run again. Send
/// [`CompensationContext::idempotency_key`] to the remote system.
pub trait CompensableEffect: EffectHandler {
    /// Undoes the effect. `output` is what `execute` returned, if it was
    /// stored.
    fn compensate(
        &self,
        ctx: &CompensationContext,
        input: &Self::Input,
        output: Option<&Self::Output>,
    ) -> impl Future<Output = Result<(), Self::Error>> + Send;
}

type BoxFuture<T> = Pin<Box<dyn Future<Output = T> + Send>>;

type VerifyFn<H> = Arc<
    dyn Fn(
            Arc<H>,
            EffectContext,
            Arc<<H as EffectHandler>::Input>,
        ) -> BoxFuture<Result<Verification<<H as EffectHandler>::Output>, EffectFailure>>
        + Send
        + Sync,
>;

/// A handler and the capabilities it is registered with. Pass it to
/// [`RuntimeBuilder::register`](crate::RuntimeBuilder::register).
pub struct Handler<H: EffectHandler> {
    effect: Arc<H>,
    verify: Option<(VerificationMode, VerifyFn<H>)>,
    compensate: Option<Compensator>,
}

impl<H: EffectHandler> Clone for Handler<H> {
    fn clone(&self) -> Self {
        Self {
            effect: Arc::clone(&self.effect),
            verify: self.verify.clone(),
            compensate: self.compensate.clone(),
        }
    }
}

impl<H: EffectHandler> Handler<H> {
    /// Registers `handler` with no optional capabilities.
    pub fn new(handler: H) -> Self {
        Self {
            effect: Arc::new(handler),
            verify: None,
            compensate: None,
        }
    }
}

impl<H: CompensableEffect> Handler<H> {
    /// Lets the effect be undone with
    /// [`Runtime::compensate`](crate::Runtime::compensate), using
    /// [`CompensableEffect::compensate`], and lets recovery finish an
    /// interrupted compensation.
    #[must_use]
    pub fn compensable(mut self) -> Self {
        let handler = Arc::clone(&self.effect);
        self.compensate = Some(Arc::new(move |ctx, input, output| {
            let handler = Arc::clone(&handler);
            Box::pin(async move {
                let input: H::Input = serde_json::from_value(input.unwrap_or(Value::Null))
                    .map_err(|e| {
                        EffectFailure::permanent(format!("stored input does not match: {e}"))
                    })?;
                let output: Option<H::Output> = output
                    .map(serde_json::from_value)
                    .transpose()
                    .map_err(|e| {
                        EffectFailure::permanent(format!("stored output does not match: {e}"))
                    })?;
                handler
                    .compensate(&ctx, &input, output.as_ref())
                    .await
                    .map_err(Into::into)
            })
        }));
        self
    }
}

impl<H: VerifiableEffect> Handler<H> {
    /// Verifies the effect after every success and to resolve unknown
    /// outcomes, using [`VerifiableEffect::verify`].
    #[must_use]
    pub fn verifiable(mut self) -> Self {
        let mode = self.effect.verification_mode();
        let verify: VerifyFn<H> = Arc::new(|handler, ctx, input| {
            Box::pin(async move { handler.verify(&ctx, &input).await.map_err(Into::into) })
        });
        self.verify = Some((mode, verify));
        self
    }
}

/// A registered handler, erased so the runtime can hold handlers of many
/// types and resume their effects by name.
pub(crate) struct Registered<S> {
    /// The `Handler<H>`, for typed submission.
    typed: Arc<dyn Any + Send + Sync>,
    /// Finishes an effect of this handler from its stored record.
    pub(crate) resume: Resume<S>,
}

pub(crate) type Resume<S> = Arc<
    dyn Fn(Runtime<S>, EffectRecord) -> BoxFuture<Result<EffectStatus, RuntimeError>> + Send + Sync,
>;

pub(crate) type Registry<S> = HashMap<&'static str, Registered<S>>;

impl<S: EffectStore> Registered<S> {
    pub(crate) fn new<H: EffectHandler>(handler: Handler<H>) -> Self {
        let typed: Arc<dyn Any + Send + Sync> = Arc::new(handler.clone());
        let resume: Resume<S> = Arc::new(move |runtime: Runtime<S>, record: EffectRecord| {
            let handler = handler.clone();
            Box::pin(async move { resume(&runtime, &handler, record).await })
        });
        Self { typed, resume }
    }

    pub(crate) fn typed<H: EffectHandler>(&self) -> Option<Handler<H>> {
        self.typed.downcast_ref::<Handler<H>>().cloned()
    }
}

/// A submission of input to a registered handler. Await it to run the
/// effect. Created by [`Runtime::submit`].
#[must_use = "a submission does nothing until it is awaited"]
pub struct Submission<'a, S, H: EffectHandler> {
    runtime: &'a Runtime<S>,
    key: String,
    input: H::Input,
    actor: Option<String>,
}

impl<'a, S, H: EffectHandler> Submission<'a, S, H> {
    pub(crate) fn new(runtime: &'a Runtime<S>, key: impl Display, input: H::Input) -> Self {
        Self {
            runtime,
            key: key.to_string(),
            input,
            actor: None,
        }
    }

    /// Who is asking for the effect, e.g. `agent:billing`.
    pub fn actor(mut self, actor: impl Into<String>) -> Self {
        self.actor = Some(actor.into());
        self
    }
}

impl<'a, S: EffectStore, H: EffectHandler> IntoFuture for Submission<'a, S, H> {
    type Output = Result<EffectOutcome<H::Output>, RuntimeError>;
    type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>;

    fn into_future(self) -> Self::IntoFuture {
        Box::pin(async move {
            let handler = self
                .runtime
                .handler::<H>()
                .ok_or(RuntimeError::NotRegistered { name: H::NAME })?;
            let key = EffectKey::new(EffectName::new(H::NAME)?, LogicalKey::new(self.key)?);
            let json = serde_json::to_value(&self.input).map_err(RuntimeError::Input)?;
            let stored = Stored {
                // Computed by the runtime, after redaction.
                fingerprint: None,
                json,
                actor: self.actor,
                from_record: false,
            };
            run(self.runtime, &handler, key, Arc::new(self.input), stored).await
        })
    }
}

/// Finishes an effect of `handler` from its record, with no caller: its
/// compensation if one is under way, else the effect itself.
async fn resume<S: EffectStore, H: EffectHandler>(
    runtime: &Runtime<S>,
    handler: &Handler<H>,
    record: EffectRecord,
) -> Result<EffectStatus, RuntimeError> {
    let id = record.id;
    if record.status == EffectStatus::Compensating {
        let compensate = handler
            .compensate
            .clone()
            .ok_or(RuntimeError::NotCompensable { name: H::NAME })?;
        let spec = CompensationSpec {
            key: record.key.clone(),
            reason: None,
            actor: Some(format!("recovery:{}", runtime.worker_id())),
            retry: handler.effect.retry_policy(),
            attempt_timeout: handler.effect.attempt_timeout(),
        };
        runtime.compensate_effect(spec, compensate).await?;
        let settled = runtime
            .store()
            .get(id)
            .await?
            .ok_or(StoreError::NotFound(id))?;
        return Ok(settled.status);
    }
    let json = record.input.clone().unwrap_or(Value::Null);
    let input: H::Input = serde_json::from_value(json.clone())
        .map_err(|source| RuntimeError::StoredInput { id, source })?;
    // Rebuilt from the record itself, so it is the same effect by
    // definition: reuse its fingerprint rather than recompute it, which
    // would refuse every older record if canonicalization ever changed.
    let stored = Stored {
        json,
        fingerprint: record.input_fingerprint.clone(),
        actor: record.created_by.clone(),
        from_record: true,
    };
    run(
        runtime,
        handler,
        record.key.clone(),
        Arc::new(input),
        stored,
    )
    .await?;
    let settled = runtime
        .store()
        .get(id)
        .await?
        .ok_or(StoreError::NotFound(id))?;
    Ok(settled.status)
}

/// What is recorded about an effect besides its key.
struct Stored {
    json: Value,
    fingerprint: Option<String>,
    actor: Option<String>,
    /// From an existing record: already redacted, fingerprint kept as is.
    from_record: bool,
}

/// Runs `handler`'s effect through the same machinery as a closure effect.
async fn run<S: EffectStore, H: EffectHandler>(
    runtime: &Runtime<S>,
    handler: &Handler<H>,
    key: EffectKey,
    input: Arc<H::Input>,
    stored: Stored,
) -> Result<EffectOutcome<H::Output>, RuntimeError> {
    let effect = Arc::clone(&handler.effect);
    let precondition = {
        let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
        Arc::new(move |ctx: EffectContext| -> BoxFuture<Precondition> {
            let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
            Box::pin(async move { effect.precondition(&ctx, &input).await })
        })
    };
    let spec = EffectSpec {
        fingerprint: stored.fingerprint,
        input: Some(stored.json),
        key,
        capabilities: Capabilities {
            kind: effect.kind(),
            remote_idempotency: effect.remote_idempotency(),
            verification: handler
                .verify
                .as_ref()
                .map_or(VerificationMode::None, |(mode, _)| *mode),
        },
        actor: stored.actor,
        retry: effect
            .retry_policy()
            .unwrap_or_else(|| runtime.default_retry()),
        attempt_timeout: effect.attempt_timeout(),
        precondition: Some(precondition),
        require_approval: effect.requires_approval(),
        risk: effect.risk(),
        automatic_retry: true,
        input_stored: stored.from_record,
    };
    let action = {
        let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
        move |ctx: EffectContext| {
            let (effect, input) = (Arc::clone(&effect), Arc::clone(&input));
            async move { effect.execute(&ctx, &input).await.map_err(Into::into) }
        }
    };
    match &handler.verify {
        Some((_, verify)) => {
            let verify = Arc::clone(verify);
            let checker =
                VerifyWith(move |ctx| verify(Arc::clone(&effect), ctx, Arc::clone(&input)));
            runtime.execute(spec, action, checker).await
        }
        None => runtime.execute(spec, action, NoVerification).await,
    }
}

/// A request to undo a registered handler's effect. Await it. Created by
/// [`Runtime::compensate`](crate::Runtime::compensate).
#[must_use = "a compensation does nothing until it is awaited"]
pub struct CompensationSubmission<'a, S, H: EffectHandler> {
    runtime: &'a Runtime<S>,
    key: String,
    reason: Option<String>,
    actor: Option<String>,
    _handler: std::marker::PhantomData<fn() -> H>,
}

impl<'a, S, H: EffectHandler> CompensationSubmission<'a, S, H> {
    pub(crate) fn new(runtime: &'a Runtime<S>, key: impl Display) -> Self {
        Self {
            runtime,
            key: key.to_string(),
            reason: None,
            actor: None,
            _handler: std::marker::PhantomData,
        }
    }

    /// Why the effect is being undone; recorded in the audit trail.
    pub fn reason(mut self, reason: impl Into<String>) -> Self {
        self.reason = Some(reason.into());
        self
    }

    /// Who is undoing it.
    pub fn actor(mut self, actor: impl Into<String>) -> Self {
        self.actor = Some(actor.into());
        self
    }
}

impl<'a, S: EffectStore, H: EffectHandler> IntoFuture for CompensationSubmission<'a, S, H> {
    type Output = Result<CompensationOutcome, RuntimeError>;
    type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>;

    fn into_future(self) -> Self::IntoFuture {
        Box::pin(async move {
            let handler = self
                .runtime
                .handler::<H>()
                .ok_or(RuntimeError::NotRegistered { name: H::NAME })?;
            let compensate = handler
                .compensate
                .clone()
                .ok_or(RuntimeError::NotCompensable { name: H::NAME })?;
            let spec = CompensationSpec {
                key: EffectKey::new(EffectName::new(H::NAME)?, LogicalKey::new(self.key)?),
                reason: self.reason,
                actor: self.actor,
                retry: handler.effect.retry_policy(),
                attempt_timeout: handler.effect.attempt_timeout(),
            };
            self.runtime.compensate_effect(spec, compensate).await
        })
    }
}