Skip to main content

agent_effects/
retention.rs

1//! Retention: pruning settled records.
2//!
3//! Records are kept forever unless the runtime has a [`RetentionPolicy`].
4//! With one, [`Runtime::prune`] (and every round of
5//! [`Runtime::run_recovery`]) deletes records that settled long enough ago,
6//! with their audit trail:
7//!
8//! ```
9//! use std::time::Duration;
10//! use agent_effects::{RetentionPolicy, Runtime};
11//! use agent_effects_memory::MemoryStore;
12//!
13//! const DAY: Duration = Duration::from_secs(86_400);
14//! let runtime = Runtime::builder(MemoryStore::new())
15//!     .retention(RetentionPolicy::settled(30 * DAY).failed(7 * DAY))
16//!     .build();
17//! ```
18//!
19//! Only settled records are pruned: `Committed`, `Failed`, `Rejected` and
20//! `Compensated`. An effect that is unknown, waits for an operator or
21//! approval, or is still running is kept however old it is, and so is one
22//! whose lease is live (a compensation about to start).
23//!
24//! **Pruning forgets.** A pruned key is free: the next call with it starts a
25//! new effect and runs it again, and a pruned committed effect can no longer
26//! be compensated. Keep committed records at least as long as any caller
27//! might retry with the same key, which is the same rule as for a remote
28//! system's idempotency keys.
29
30use std::time::Duration;
31
32use tracing::debug;
33
34use crate::error::RuntimeError;
35use crate::runtime::Runtime;
36use crate::state::EffectStatus;
37use crate::store::{EffectStore, PruneQuery};
38
39/// How long settled records are kept, per status. `None` keeps them
40/// forever, which is the default.
41#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
42pub struct RetentionPolicy {
43    committed: Option<Duration>,
44    failed: Option<Duration>,
45    rejected: Option<Duration>,
46    compensated: Option<Duration>,
47}
48
49impl RetentionPolicy {
50    /// Keeps every record forever.
51    pub const KEEP_ALL: Self = Self {
52        committed: None,
53        failed: None,
54        rejected: None,
55        compensated: None,
56    };
57
58    /// Prunes every settled record `older_than` after it settled.
59    pub const fn settled(older_than: Duration) -> Self {
60        Self {
61            committed: Some(older_than),
62            failed: Some(older_than),
63            rejected: Some(older_than),
64            compensated: Some(older_than),
65        }
66    }
67
68    /// Prunes committed effects `older_than` after they committed. A
69    /// later call with a pruned key runs the effect again.
70    #[must_use]
71    pub const fn committed(mut self, older_than: Duration) -> Self {
72        self.committed = Some(older_than);
73        self
74    }
75
76    /// Prunes failed effects `older_than` after they failed.
77    #[must_use]
78    pub const fn failed(mut self, older_than: Duration) -> Self {
79        self.failed = Some(older_than);
80        self
81    }
82
83    /// Prunes rejected effects (precondition or approver said no)
84    /// `older_than` after they were rejected.
85    #[must_use]
86    pub const fn rejected(mut self, older_than: Duration) -> Self {
87        self.rejected = Some(older_than);
88        self
89    }
90
91    /// Prunes compensated effects `older_than` after they were undone.
92    #[must_use]
93    pub const fn compensated(mut self, older_than: Duration) -> Self {
94        self.compensated = Some(older_than);
95        self
96    }
97
98    /// Whether this policy never prunes anything.
99    pub const fn keeps_all(&self) -> bool {
100        self.committed.is_none()
101            && self.failed.is_none()
102            && self.rejected.is_none()
103            && self.compensated.is_none()
104    }
105
106    /// Each status this policy prunes, with its retention.
107    pub fn rules(&self) -> impl Iterator<Item = (EffectStatus, Duration)> {
108        [
109            (EffectStatus::Committed, self.committed),
110            (EffectStatus::Failed, self.failed),
111            (EffectStatus::Rejected, self.rejected),
112            (EffectStatus::Compensated, self.compensated),
113        ]
114        .into_iter()
115        .filter_map(|(status, age)| Some((status, age?)))
116    }
117}
118
119/// What one [`Runtime::prune`] pass deleted, per status.
120#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
121#[non_exhaustive]
122pub struct PruneReport {
123    /// Committed effects deleted.
124    pub committed: u64,
125    /// Failed effects deleted.
126    pub failed: u64,
127    /// Rejected effects deleted.
128    pub rejected: u64,
129    /// Compensated effects deleted.
130    pub compensated: u64,
131}
132
133impl PruneReport {
134    /// All records deleted.
135    pub const fn total(&self) -> u64 {
136        self.committed + self.failed + self.rejected + self.compensated
137    }
138
139    fn add(&mut self, status: EffectStatus, deleted: u64) {
140        match status {
141            EffectStatus::Committed => self.committed += deleted,
142            EffectStatus::Failed => self.failed += deleted,
143            EffectStatus::Rejected => self.rejected += deleted,
144            EffectStatus::Compensated => self.compensated += deleted,
145            _ => {}
146        }
147    }
148}
149
150impl<S: EffectStore> Runtime<S> {
151    /// Deletes the settled records the [retention policy](RetentionPolicy)
152    /// says are old enough, in batches, with their audit trails. Does
153    /// nothing without a policy. [`Self::run_recovery`] calls it every
154    /// round; call it yourself to prune on another schedule.
155    ///
156    /// # Errors
157    ///
158    /// [`RuntimeError::Store`] if the store fails. Batches deleted before
159    /// the failure stay deleted.
160    pub async fn prune(&self) -> Result<PruneReport, RuntimeError> {
161        let mut report = PruneReport::default();
162        let now = self.now();
163        for (status, older_than) in self.retention().rules() {
164            let query = PruneQuery::new(status, older_than, now);
165            loop {
166                let deleted = self.store().prune(query.clone()).await?;
167                report.add(status, deleted);
168                if deleted < query.limit as u64 {
169                    break;
170                }
171            }
172        }
173        if report.total() > 0 {
174            debug!(?report, "pruned settled effects");
175        }
176        Ok(report)
177    }
178}