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}