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}