Skip to main content

agent_effects/
recovery.rs

1//! Recovery and the operator API.
2//!
3//! - [`Runtime::recover`] marks attempts whose worker died as unknown, then
4//!   finishes every unsettled effect that has a registered
5//!   [handler](crate::handler) from its stored input. Closure effects are
6//!   not durable, so they wait for a caller to re-run them;
7//! - [`Runtime::pending`] lists the effects waiting for a caller or an
8//!   operator;
9//! - [`Runtime::resolve`] records an operator's decision.
10
11use std::time::Duration;
12
13use serde::Serialize;
14use serde_json::{Value, json};
15use tokio::time::MissedTickBehavior;
16use tracing::{info, warn};
17
18use crate::error::RuntimeError;
19use crate::id::EffectId;
20use crate::runtime::Runtime;
21use crate::state::{EffectStatus, Transition};
22use crate::store::{
23    EffectRecord, EffectStore, ErrorRecord, ListQuery, StoreError, TransitionRequest,
24};
25
26/// Page size for recovery scans.
27const PAGE: usize = 100;
28
29/// What one [`Runtime::recover`] pass did.
30#[derive(Clone, Debug, Default, PartialEq, Eq)]
31#[non_exhaustive]
32pub struct RecoveryReport {
33    /// Effects whose worker's lease expired mid-attempt or mid-verification,
34    /// now marked unknown.
35    pub marked_unknown: Vec<EffectId>,
36    /// Effects that matched the scan but were taken by another caller or
37    /// worker before this pass reached them.
38    pub skipped: Vec<EffectId>,
39    /// Effects with a registered handler that this pass moved forward, and
40    /// the status each ended in. `Committed`, `Failed` and `Rejected` are
41    /// settled; `NeedsIntervention` waits for an operator; anything else is
42    /// still in progress or still unknown.
43    pub resumed: Vec<(EffectId, EffectStatus)>,
44    /// Unsettled effects with no registered handler (closure effects). They
45    /// need a caller to run them again with the same key.
46    pub unhandled: Vec<EffectId>,
47    /// Effects whose handler could not be run, with the reason, e.g. a
48    /// stored input that no longer deserializes.
49    pub resume_errors: Vec<(EffectId, String)>,
50}
51
52/// An operator's decision about an effect the runtime could not resolve.
53#[derive(Clone, Debug, PartialEq)]
54pub enum Resolution {
55    /// The effect applied. `output` becomes its result for every later call,
56    /// so it should deserialize into the type those calls expect.
57    Applied {
58        /// The effect's output, as JSON.
59        output: Option<Value>,
60    },
61    /// The effect did not apply. It ends `Failed`, with the operator's note
62    /// as its error.
63    NotApplied,
64    /// Running the effect again is safe. The next call runs it, even if its
65    /// retry budget is spent, since this is an explicit decision. On a
66    /// failed compensation: try the compensation again.
67    Retry,
68    /// For a failed compensation: the operator undid the effect by hand. It
69    /// ends `Compensated`.
70    Compensated,
71}
72
73impl Resolution {
74    /// [`Resolution::Applied`] with `output` serialized.
75    ///
76    /// # Errors
77    ///
78    /// If `output` cannot be serialized to JSON.
79    pub fn applied<T: Serialize + ?Sized>(output: &T) -> Result<Self, serde_json::Error> {
80        Ok(Self::Applied {
81            output: Some(serde_json::to_value(output)?),
82        })
83    }
84}
85
86impl<S: EffectStore> Runtime<S> {
87    /// One recovery pass, in two steps.
88    ///
89    /// 1. Every effect whose worker's lease expired mid-attempt or
90    ///    mid-verification becomes `Unknown`, never `Failed`: it may have
91    ///    changed the outside world.
92    /// 2. Every unsettled effect nobody holds (see [`Self::pending`]) whose
93    ///    name has a registered [handler](crate::handler) is finished from its
94    ///    stored input, exactly as if its caller had called again: verified,
95    ///    re-run if that is safe, or escalated. This includes effects left
96    ///    `Pending` by a crash. Effects waiting for an operator, and retries
97    ///    scheduled for later, are left alone. Unsettled closure effects are
98    ///    reported as `unhandled`: only a caller can re-run them.
99    ///
100    /// Effects are resumed one at a time, so a pass lasts as long as their
101    /// retries and verifications take. Safe to run from several workers at
102    /// once: every change happens under a lease.
103    ///
104    /// # Errors
105    ///
106    /// [`RuntimeError::Store`] if the store fails. Progress made before the
107    /// failure is kept. A handler that cannot run is reported in
108    /// `resume_errors` instead.
109    pub async fn recover(&self) -> Result<RecoveryReport, RuntimeError> {
110        let mut report = RecoveryReport::default();
111        self.mark_abandoned(&mut report).await?;
112        self.resume_pending(&mut report).await?;
113        Ok(report)
114    }
115
116    async fn mark_abandoned(&self, report: &mut RecoveryReport) -> Result<(), RuntimeError> {
117        let mut after = None;
118        loop {
119            let mut query = ListQuery::expired_leases(self.now()).limit(PAGE);
120            if let Some(id) = after {
121                query = query.after(id);
122            }
123            let page = self.store().list(query).await?;
124            let full = page.len() == PAGE;
125            after = page.last().map(|record| record.id);
126            for record in page {
127                if self.mark_unknown(record.id).await? {
128                    report.marked_unknown.push(record.id);
129                } else {
130                    report.skipped.push(record.id);
131                }
132            }
133            if !full {
134                return Ok(());
135            }
136        }
137    }
138
139    async fn resume_pending(&self, report: &mut RecoveryReport) -> Result<(), RuntimeError> {
140        let mut after = None;
141        loop {
142            let page = self.pending(after, PAGE).await?;
143            let full = page.len() == PAGE;
144            after = page.last().map(|record| record.id);
145            for record in page {
146                let id = record.id;
147                let not_due = record.next_attempt_at.is_some_and(|at| at > self.now());
148                let operator = matches!(
149                    record.status,
150                    EffectStatus::NeedsIntervention
151                        | EffectStatus::CompensationFailed
152                        | EffectStatus::AwaitingApproval
153                );
154                if operator || not_due {
155                    continue;
156                }
157                let Some(resume) = self.resumer(record.key.name.as_str()) else {
158                    report.unhandled.push(id);
159                    continue;
160                };
161                match resume(self.clone(), record).await {
162                    Ok(status) => {
163                        info!(effect.id = %id, %status, "recovery resumed effect");
164                        report.resumed.push((id, status));
165                    }
166                    Err(e) => {
167                        warn!(effect.id = %id, error = %e, "recovery could not resume effect");
168                        report.resume_errors.push((id, e.to_string()));
169                    }
170                }
171            }
172            if !full {
173                return Ok(());
174            }
175        }
176    }
177
178    /// Runs [`Self::recover`] every `interval`, forever, followed by
179    /// [`Self::prune`] when the runtime has a
180    /// [retention policy](crate::RetentionPolicy). Spawn it:
181    ///
182    /// ```ignore
183    /// tokio::spawn({
184    ///     let runtime = runtime.clone();
185    ///     async move { runtime.run_recovery(Duration::from_secs(30)).await }
186    /// });
187    /// ```
188    ///
189    /// A failed pass is logged and retried at the next tick.
190    pub async fn run_recovery(&self, interval: Duration) {
191        let mut ticker = tokio::time::interval(interval);
192        ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);
193        loop {
194            ticker.tick().await;
195            match self.recover().await {
196                Ok(report) if !report.marked_unknown.is_empty() || !report.resumed.is_empty() => {
197                    info!(
198                        marked_unknown = report.marked_unknown.len(),
199                        resumed = report.resumed.len(),
200                        unhandled = report.unhandled.len(),
201                        "recovery pass"
202                    );
203                }
204                Ok(_) => {}
205                Err(e) => warn!(error = %e, "recovery pass failed"),
206            }
207            match self.prune().await {
208                Ok(report) if report.total() > 0 => {
209                    info!(pruned = report.total(), "pruned settled effects");
210                }
211                Ok(_) => {}
212                Err(e) => warn!(error = %e, "pruning failed"),
213            }
214        }
215    }
216
217    /// Effects waiting for a caller to re-run them or an operator to decide,
218    /// ordered by creation, `limit` at a time after `after`.
219    ///
220    /// These are the effects that are unsettled with nobody working on
221    /// them: `Pending` (for example a worker died during a backoff wait),
222    /// `Executing` or `Verifying` with an expired lease (not yet recovered),
223    /// `Unknown`, `NeedsIntervention`, `AwaitingApproval`, an interrupted
224    /// `Compensating`, and `CompensationFailed`.
225    ///
226    /// # Errors
227    ///
228    /// [`RuntimeError::Store`] if the store fails.
229    pub async fn pending(
230        &self,
231        after: Option<EffectId>,
232        limit: usize,
233    ) -> Result<Vec<EffectRecord>, RuntimeError> {
234        let mut query = ListQuery::statuses([
235            EffectStatus::Pending,
236            EffectStatus::AwaitingApproval,
237            EffectStatus::Executing,
238            EffectStatus::Verifying,
239            EffectStatus::Unknown,
240            EffectStatus::NeedsIntervention,
241            EffectStatus::Compensating,
242            EffectStatus::CompensationFailed,
243        ])
244        .limit(limit);
245        query.lease_expired_at = Some(self.now());
246        query.after = after;
247        Ok(self.store().list(query).await?)
248    }
249
250    /// Records an operator's decision about an effect that is `Unknown` or
251    /// `NeedsIntervention`. `actor` identifies the operator (e.g.
252    /// `operator:alice`); `note` says why, for the audit trail.
253    ///
254    /// # Errors
255    ///
256    /// [`RuntimeError::Store`] wrapping:
257    ///
258    /// - [`StoreError::NotFound`] for an unknown id;
259    /// - [`StoreError::InvalidTransition`] if the effect is not unresolved,
260    ///   for example already committed;
261    /// - [`StoreError::LeaseHeld`] while a caller or worker is working on it;
262    /// - [`StoreError::VersionConflict`] if it changed during the call.
263    pub async fn resolve(
264        &self,
265        id: EffectId,
266        resolution: Resolution,
267        actor: impl Into<String>,
268        note: impl Into<String>,
269    ) -> Result<EffectRecord, RuntimeError> {
270        let record = self
271            .store()
272            .get(id)
273            .await?
274            .ok_or(StoreError::NotFound(id))?;
275        let note = note.into();
276        let transition = match resolution {
277            Resolution::Applied { .. } => Transition::ResolvedApplied,
278            Resolution::NotApplied => Transition::ResolvedNotApplied,
279            Resolution::Retry => Transition::ResolvedRetry,
280            Resolution::Compensated => Transition::ResolvedCompensated,
281        };
282        let mut request = TransitionRequest::new(&record, None, transition, self.now());
283        request.actor = Some(actor.into());
284        request.payload = Some(json!({ "note": note }));
285        match resolution {
286            Resolution::Applied { output } => request.output = output,
287            Resolution::NotApplied => {
288                request.error = Some(ErrorRecord {
289                    class: None,
290                    message: note,
291                });
292            }
293            Resolution::Retry | Resolution::Compensated => {}
294        }
295        let record = self.commit_transition(&record, request).await?;
296        info!(effect.id = %id, %transition, "effect resolved by an operator");
297        Ok(record)
298    }
299
300    /// Approves an effect waiting in `AwaitingApproval`. It becomes
301    /// `Pending`: a caller's next call runs it, and so does recovery for a
302    /// registered handler. `actor` is recorded as the approver.
303    ///
304    /// # Errors
305    ///
306    /// As for [`Self::resolve`]; `InvalidTransition` if it is not awaiting
307    /// approval.
308    pub async fn approve(
309        &self,
310        id: EffectId,
311        actor: impl Into<String>,
312        note: impl Into<String>,
313    ) -> Result<EffectRecord, RuntimeError> {
314        self.decide(id, Transition::Approve, actor.into(), note.into())
315            .await
316    }
317
318    /// Denies an effect waiting in `AwaitingApproval`. It ends `Rejected`,
319    /// with `reason` as its error.
320    ///
321    /// # Errors
322    ///
323    /// As for [`Self::approve`].
324    pub async fn deny(
325        &self,
326        id: EffectId,
327        actor: impl Into<String>,
328        reason: impl Into<String>,
329    ) -> Result<EffectRecord, RuntimeError> {
330        self.decide(id, Transition::Deny, actor.into(), reason.into())
331            .await
332    }
333
334    async fn decide(
335        &self,
336        id: EffectId,
337        transition: Transition,
338        actor: String,
339        note: String,
340    ) -> Result<EffectRecord, RuntimeError> {
341        let record = self
342            .store()
343            .get(id)
344            .await?
345            .ok_or(StoreError::NotFound(id))?;
346        let mut request = TransitionRequest::new(&record, None, transition, self.now());
347        request.actor = Some(actor);
348        if transition == Transition::Deny {
349            request.error = Some(ErrorRecord {
350                class: None,
351                message: note.clone(),
352            });
353        }
354        request.payload = Some(json!({ "note": note }));
355        let record = self.commit_transition(&record, request).await?;
356        info!(effect.id = %id, %transition, "approval decided by an operator");
357        Ok(record)
358    }
359
360    /// Moves one effect with an expired lease to `Unknown`. Returns `false`
361    /// if someone else got to it first.
362    async fn mark_unknown(&self, id: EffectId) -> Result<bool, RuntimeError> {
363        let store = self.store();
364        let lease = match store
365            .acquire_lease(id, self.worker_id(), self.now(), self.lease_ttl())
366            .await
367        {
368            Ok(lease) => lease,
369            Err(StoreError::LeaseHeld { .. }) => return Ok(false),
370            Err(e) => return Err(e.into()),
371        };
372        let marked = async {
373            let record = store.get(id).await?.ok_or(StoreError::NotFound(id))?;
374            if !matches!(
375                record.status,
376                EffectStatus::Executing | EffectStatus::Verifying
377            ) {
378                return Ok(false);
379            }
380            let mut request =
381                TransitionRequest::new(&record, Some(&lease), Transition::LeaseExpired, self.now());
382            request.actor = Some(format!("recovery:{}", self.worker_id()));
383            self.commit_transition(&record, request).await?;
384            info!(
385                effect.id = %id,
386                effect.name = %record.key.name,
387                "worker lost mid-effect; outcome is now unknown"
388            );
389            Ok::<_, RuntimeError>(true)
390        }
391        .await;
392        if let Err(e) = store.release_lease(&lease).await {
393            warn!(error = %e, "could not release recovery lease; it will expire");
394        }
395        marked
396    }
397}