Skip to main content

agent_effects_store/
lib.rs

1//! Storage contract for [`agent-effects`](https://crates.io/crates/agent-effects).
2//!
3//! This crate holds what a store persists and enforces: effect identity
4//! ([`id`]), classification ([`kind`], [`failure`]), the state machine
5//! ([`state`]), records, leases and audit events, and the [`EffectStore`]
6//! trait. Store backends depend on this crate only, never on the runtime.
7//!
8//! The rules (lease fencing, version checks, the transition table, attempt
9//! bookkeeping) live in [`EffectRecord`]'s pure methods. A store's job is
10//! only to run "load, apply, save the record and its audit event"
11//! atomically, so every backend behaves the same. The `testkit` feature's
12//! `testkit::conformance` suite checks that.
13
14pub mod failure;
15pub mod id;
16pub mod kind;
17pub mod state;
18#[cfg(feature = "testkit")]
19pub mod testkit;
20
21mod error;
22mod record;
23
24use std::future::Future;
25use std::time::{Duration, SystemTime};
26
27use serde::{Deserialize, Serialize};
28use serde_json::Value;
29
30pub use error::StoreError;
31pub use failure::{Disposition, FailureClass};
32pub use id::{
33    EffectId, EffectKey, EffectName, IdempotencyKey, IdentityError, LogicalKey, WorkerId,
34};
35pub use kind::EffectKind;
36pub use record::EffectRecord;
37pub use state::{EffectStatus, InvalidTransition, Transition};
38
39/// Persistence for effect records, leases and audit events.
40///
41/// Times are passed in by the caller rather than read from a clock, so the
42/// runtime's clock governs leases. Stores must keep at least millisecond
43/// precision.
44///
45/// Every method that changes a record must apply the change through the
46/// matching [`EffectRecord`] method inside one atomic unit (a transaction or
47/// a lock), so concurrent callers can never interleave between the check and
48/// the write.
49pub trait EffectStore: Send + Sync + 'static {
50    /// Inserts a new effect, or returns the existing record with the same
51    /// [`EffectKey`] untouched.
52    ///
53    /// Atomic on the key: of any number of concurrent calls for one key,
54    /// exactly one reports [`InsertOutcome::inserted`].
55    fn insert_or_get(
56        &self,
57        new: NewEffect,
58    ) -> impl Future<Output = Result<InsertOutcome, StoreError>> + Send;
59
60    /// Loads a record by id.
61    fn get(
62        &self,
63        id: EffectId,
64    ) -> impl Future<Output = Result<Option<EffectRecord>, StoreError>> + Send;
65
66    /// Loads a record by its logical identity.
67    fn get_by_key(
68        &self,
69        key: &EffectKey,
70    ) -> impl Future<Output = Result<Option<EffectRecord>, StoreError>> + Send;
71
72    /// Takes the execution lease, via [`EffectRecord::acquire_lease`].
73    fn acquire_lease(
74        &self,
75        id: EffectId,
76        owner: &WorkerId,
77        now: SystemTime,
78        ttl: Duration,
79    ) -> impl Future<Output = Result<Lease, StoreError>> + Send;
80
81    /// Extends a held lease, via [`EffectRecord::renew_lease`].
82    fn renew_lease(
83        &self,
84        lease: &Lease,
85        now: SystemTime,
86        ttl: Duration,
87    ) -> impl Future<Output = Result<Lease, StoreError>> + Send;
88
89    /// Gives up a lease, via [`EffectRecord::release_lease`]. Releasing a
90    /// lease that was already lost is not an error.
91    fn release_lease(&self, lease: &Lease) -> impl Future<Output = Result<(), StoreError>> + Send;
92
93    /// Applies a status transition, via [`EffectRecord::apply`], and appends
94    /// its audit event in the same atomic unit. Returns the updated record.
95    fn transition(
96        &self,
97        request: TransitionRequest,
98    ) -> impl Future<Output = Result<EffectRecord, StoreError>> + Send;
99
100    /// Lists records matching `query`, ordered by id (creation time).
101    fn list(
102        &self,
103        query: ListQuery,
104    ) -> impl Future<Output = Result<Vec<EffectRecord>, StoreError>> + Send;
105
106    /// The audit events of one effect, ordered by sequence.
107    fn events(
108        &self,
109        id: EffectId,
110    ) -> impl Future<Output = Result<Vec<EffectEvent>, StoreError>> + Send;
111
112    /// Deletes up to `query.limit` records matching `query`, lowest id
113    /// first, together with their audit events, and returns how many it
114    /// deleted. Each record is checked and deleted atomically, so a record
115    /// that a worker leases or changes concurrently is either deleted before
116    /// that change or not at all.
117    ///
118    /// A deleted key is free: the next `insert_or_get` for it inserts a new
119    /// record. A query for a status that is not
120    /// [settled](EffectStatus::is_settled) deletes nothing.
121    fn prune(&self, query: PruneQuery) -> impl Future<Output = Result<u64, StoreError>> + Send;
122}
123
124/// A new effect to record.
125#[derive(Clone, Debug, PartialEq)]
126pub struct NewEffect {
127    /// The id the record gets if it is inserted.
128    pub id: EffectId,
129    /// The logical identity; unique per store.
130    pub key: EffectKey,
131    /// The effect's kind.
132    pub kind: EffectKind,
133    /// The input, already redacted for storage.
134    pub input: Option<Value>,
135    /// A stable hash of the unredacted input, used to reject a reused key
136    /// with a different input.
137    pub input_fingerprint: Option<String>,
138    /// Who asked for the effect, e.g. `agent:refund-agent`.
139    pub created_by: Option<String>,
140    /// The creation time.
141    pub now: SystemTime,
142}
143
144impl NewEffect {
145    /// A new effect with a fresh id and no input or actor.
146    pub fn new(key: EffectKey, kind: EffectKind, now: SystemTime) -> Self {
147        Self {
148            id: EffectId::new(),
149            key,
150            kind,
151            input: None,
152            input_fingerprint: None,
153            created_by: None,
154            now,
155        }
156    }
157}
158
159/// The result of [`EffectStore::insert_or_get`].
160#[derive(Clone, Debug, PartialEq)]
161pub struct InsertOutcome {
162    /// The new or existing record.
163    pub record: EffectRecord,
164    /// Whether this call created it.
165    pub inserted: bool,
166}
167
168/// Proof of holding an effect's execution lease.
169///
170/// Every acquisition increments the record's lease epoch. A store accepts a
171/// lease only while the record still carries its owner and epoch and it has
172/// not expired. That fences off a worker whose lease was taken over: all its
173/// later writes fail with [`StoreError::LeaseLost`].
174#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
175pub struct Lease {
176    /// The effect the lease is for.
177    pub effect_id: EffectId,
178    /// The holder.
179    pub owner: WorkerId,
180    /// The fencing token.
181    pub epoch: u64,
182    /// When the lease lapses unless renewed.
183    pub expires_at: SystemTime,
184}
185
186/// A request to move an effect through one [`Transition`].
187#[derive(Clone, Debug, PartialEq)]
188pub struct TransitionRequest {
189    /// The effect to change.
190    pub id: EffectId,
191    /// The record version the caller last saw; the change fails with
192    /// [`StoreError::VersionConflict`] if it moved on.
193    pub expected_version: u64,
194    /// The caller's lease. `None` is only accepted while nobody holds a live
195    /// lease, e.g. for an operator resolving an effect.
196    pub lease: Option<Lease>,
197    /// The transition to apply.
198    pub transition: Transition,
199    /// The time of the change.
200    pub now: SystemTime,
201    /// The action's result, stored when present.
202    pub output: Option<Value>,
203    /// A failure to record as the effect's last error.
204    pub error: Option<ErrorRecord>,
205    /// For [`Transition::ScheduleRetry`] and [`Transition::ResolvedRetry`]:
206    /// when the next attempt may start. Defaults to `now`.
207    pub next_attempt_at: Option<SystemTime>,
208    /// Who made the change, for the audit event.
209    pub actor: Option<String>,
210    /// Extra audit detail, already redacted.
211    pub payload: Option<Value>,
212}
213
214impl TransitionRequest {
215    /// A transition with no output, error, actor or payload.
216    pub fn new(
217        record: &EffectRecord,
218        lease: Option<&Lease>,
219        transition: Transition,
220        now: SystemTime,
221    ) -> Self {
222        Self {
223            id: record.id,
224            expected_version: record.version,
225            lease: lease.cloned(),
226            transition,
227            now,
228            output: None,
229            error: None,
230            next_attempt_at: None,
231            actor: None,
232            payload: None,
233        }
234    }
235}
236
237/// A recorded failure.
238#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
239pub struct ErrorRecord {
240    /// How the failure was classified, if it was.
241    pub class: Option<FailureClass>,
242    /// A description, already redacted.
243    pub message: String,
244}
245
246/// One entry of an effect's audit trail.
247#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
248pub struct EffectEvent {
249    /// The effect.
250    pub effect_id: EffectId,
251    /// Position in the trail, starting at 1. Equal to the record's version
252    /// after the transition.
253    pub sequence: u64,
254    /// What happened.
255    pub transition: Transition,
256    /// The status before.
257    pub from: EffectStatus,
258    /// The status after.
259    pub to: EffectStatus,
260    /// The record's attempt count after the transition.
261    pub attempt: u32,
262    /// Who made the change.
263    pub actor: Option<String>,
264    /// Extra detail, already redacted.
265    pub payload: Option<Value>,
266    /// When it happened.
267    pub at: SystemTime,
268}
269
270/// Which records [`EffectStore::prune`] deletes: those in a settled
271/// `status` that last changed at least `older_than` before `now` and hold no
272/// live lease at `now`.
273///
274/// The age is a duration, not a cutoff time, so a store that reads its own
275/// clock (like PostgreSQL's) can apply it to that clock.
276#[derive(Clone, Debug, PartialEq, Eq)]
277pub struct PruneQuery {
278    /// The status to prune. Must be [settled](EffectStatus::is_settled).
279    pub status: EffectStatus,
280    /// The minimum time since the record's last transition (`updated_at`).
281    pub older_than: Duration,
282    /// The current time.
283    pub now: SystemTime,
284    /// At most this many records.
285    pub limit: usize,
286}
287
288impl PruneQuery {
289    /// The default batch size.
290    pub const DEFAULT_LIMIT: usize = 500;
291
292    /// Records in `status` settled at least `older_than` before `now`.
293    pub fn new(status: EffectStatus, older_than: Duration, now: SystemTime) -> Self {
294        Self {
295            status,
296            older_than,
297            now,
298            limit: Self::DEFAULT_LIMIT,
299        }
300    }
301
302    /// Caps the batch size.
303    #[must_use]
304    pub fn limit(mut self, limit: usize) -> Self {
305        self.limit = limit;
306        self
307    }
308
309    /// The latest `updated_at` a record may have to be pruned, or `None` if
310    /// that time cannot be represented (nothing is that old).
311    pub fn cutoff(&self) -> Option<SystemTime> {
312        self.now.checked_sub(self.older_than)
313    }
314
315    /// Whether `record` should be pruned (ignoring `limit`).
316    pub fn matches(&self, record: &EffectRecord) -> bool {
317        self.status.is_settled()
318            && record.status == self.status
319            && self
320                .cutoff()
321                .is_some_and(|cutoff| record.updated_at <= cutoff)
322            && record.live_lease_owner(self.now).is_none()
323    }
324}
325
326/// A filter for [`EffectStore::list`].
327#[derive(Clone, Debug, PartialEq, Eq)]
328pub struct ListQuery {
329    /// Only these statuses. Empty means any status.
330    pub statuses: Vec<EffectStatus>,
331    /// Only records without a live lease at this time.
332    pub lease_expired_at: Option<SystemTime>,
333    /// Only records with a greater id, for paging.
334    pub after: Option<EffectId>,
335    /// At most this many records.
336    pub limit: usize,
337}
338
339impl ListQuery {
340    /// The default page size.
341    pub const DEFAULT_LIMIT: usize = 100;
342
343    /// Records in any of `statuses`.
344    pub fn statuses(statuses: impl IntoIterator<Item = EffectStatus>) -> Self {
345        Self {
346            statuses: statuses.into_iter().collect(),
347            lease_expired_at: None,
348            after: None,
349            limit: Self::DEFAULT_LIMIT,
350        }
351    }
352
353    /// Records stuck mid-attempt: executing or verifying with no live lease
354    /// at `now`. Recovery moves these to Unknown.
355    pub fn expired_leases(now: SystemTime) -> Self {
356        Self {
357            lease_expired_at: Some(now),
358            ..Self::statuses([EffectStatus::Executing, EffectStatus::Verifying])
359        }
360    }
361
362    /// Continues after `id`.
363    #[must_use]
364    pub fn after(mut self, id: EffectId) -> Self {
365        self.after = Some(id);
366        self
367    }
368
369    /// Caps the page size.
370    #[must_use]
371    pub fn limit(mut self, limit: usize) -> Self {
372        self.limit = limit;
373        self
374    }
375
376    /// Whether `record` passes the filter (ignoring `limit`).
377    pub fn matches(&self, record: &EffectRecord) -> bool {
378        (self.statuses.is_empty() || self.statuses.contains(&record.status))
379            && self
380                .lease_expired_at
381                .is_none_or(|now| record.live_lease_owner(now).is_none())
382            && self.after.is_none_or(|after| record.id > after)
383    }
384}