Skip to main content

agent_effects_store/
state.rs

1//! The effect state machine.
2//!
3//! [`EffectStatus::apply`] is the single source of truth for which
4//! transitions are legal. Stores and the runtime never change a status
5//! without going through it. The table is documented in
6//! `docs/design.md#state-machine`.
7
8use std::fmt;
9
10use serde::{Deserialize, Serialize};
11
12/// Where an effect is in its lifecycle.
13#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
14#[serde(rename_all = "snake_case")]
15#[non_exhaustive]
16pub enum EffectStatus {
17    /// Recorded and waiting to run, either for the first time or for a
18    /// scheduled retry.
19    Pending,
20    /// Waiting for a human or policy decision before it may run.
21    AwaitingApproval,
22    /// An attempt is in flight under a lease.
23    Executing,
24    /// An attempt may or may not have applied.
25    Unknown,
26    /// Checking the remote system for the effect's outcome.
27    Verifying,
28    /// The effect applied. Final, unless it is later compensated.
29    Committed,
30    /// The effect definitely did not apply and will not be retried. Terminal.
31    Failed,
32    /// A precondition or approval decision refused the effect before it ran.
33    /// Terminal.
34    Rejected,
35    /// The runtime cannot resolve the outcome on its own; an operator must.
36    NeedsIntervention,
37    /// A committed effect is being undone. Attempts are retried; the
38    /// compensation must be idempotent.
39    Compensating,
40    /// The effect was undone. Terminal.
41    Compensated,
42    /// Undoing the effect failed for good; an operator must finish it or
43    /// order a retry.
44    CompensationFailed,
45}
46
47/// Every status, for exhaustive tests and store migrations.
48pub const ALL_STATUSES: [EffectStatus; 12] = [
49    EffectStatus::Pending,
50    EffectStatus::AwaitingApproval,
51    EffectStatus::Executing,
52    EffectStatus::Unknown,
53    EffectStatus::Verifying,
54    EffectStatus::Committed,
55    EffectStatus::Failed,
56    EffectStatus::Rejected,
57    EffectStatus::NeedsIntervention,
58    EffectStatus::Compensating,
59    EffectStatus::Compensated,
60    EffectStatus::CompensationFailed,
61];
62
63impl EffectStatus {
64    /// Whether no further transition is possible.
65    ///
66    /// `Committed` is not terminal: compensation can start from it.
67    pub const fn is_terminal(self) -> bool {
68        matches!(self, Self::Failed | Self::Rejected | Self::Compensated)
69    }
70
71    /// Whether nothing is left to do without a new request: the effect
72    /// committed, failed, was rejected or was undone. Only settled records
73    /// may be pruned ([`PruneQuery`](crate::PruneQuery)).
74    ///
75    /// Unlike [`Self::is_terminal`] this includes `Committed`, whose
76    /// compensation would be a new request. `CompensationFailed` is not
77    /// settled: it waits for an operator.
78    pub const fn is_settled(self) -> bool {
79        matches!(
80            self,
81            Self::Committed | Self::Failed | Self::Rejected | Self::Compensated
82        )
83    }
84
85    /// Whether the effect may have changed the outside world without the
86    /// runtime knowing the result, so it must not simply be started again.
87    pub const fn is_in_doubt(self) -> bool {
88        matches!(
89            self,
90            Self::Executing | Self::Unknown | Self::Verifying | Self::NeedsIntervention
91        )
92    }
93
94    /// The stable storage representation.
95    pub const fn as_str(self) -> &'static str {
96        match self {
97            Self::Pending => "pending",
98            Self::AwaitingApproval => "awaiting_approval",
99            Self::Executing => "executing",
100            Self::Unknown => "unknown",
101            Self::Verifying => "verifying",
102            Self::Committed => "committed",
103            Self::Failed => "failed",
104            Self::Rejected => "rejected",
105            Self::NeedsIntervention => "needs_intervention",
106            Self::Compensating => "compensating",
107            Self::Compensated => "compensated",
108            Self::CompensationFailed => "compensation_failed",
109        }
110    }
111
112    /// Parses the storage representation produced by [`Self::as_str`].
113    pub fn parse(s: &str) -> Option<Self> {
114        ALL_STATUSES.into_iter().find(|status| status.as_str() == s)
115    }
116
117    /// The status after `event`, or an error if the event is not legal here.
118    ///
119    /// # Errors
120    ///
121    /// Returns [`InvalidTransition`] for any pair not in the transition table.
122    pub const fn apply(self, event: Transition) -> Result<Self, InvalidTransition> {
123        use EffectStatus as S;
124        use Transition as T;
125
126        let next = match (self, event) {
127            (S::Pending, T::RequestApproval) => S::AwaitingApproval,
128            (S::AwaitingApproval, T::Approve) => S::Pending,
129            (S::AwaitingApproval, T::Deny) | (S::Pending, T::PreconditionRejected) => S::Rejected,
130
131            (S::Pending, T::StartAttempt) => S::Executing,
132
133            (S::Executing, T::Succeeded)
134            | (S::Verifying, T::VerificationConfirmed)
135            | (S::Unknown | S::NeedsIntervention, T::ResolvedApplied) => S::Committed,
136
137            (S::Executing, T::FailedDefinitively)
138            | (S::Verifying, T::VerificationNotApplied)
139            | (S::Unknown | S::NeedsIntervention, T::ResolvedNotApplied) => S::Failed,
140
141            (S::Executing | S::Verifying | S::Unknown, T::ScheduleRetry)
142            | (S::Unknown | S::NeedsIntervention, T::ResolvedRetry) => S::Pending,
143
144            (S::Executing | S::Verifying, T::OutcomeUnknown | T::LeaseExpired) => S::Unknown,
145
146            (S::Executing | S::Unknown, T::StartVerification) => S::Verifying,
147
148            (S::Unknown, T::Escalate) | (S::Verifying, T::VerificationConflict) => {
149                S::NeedsIntervention
150            }
151
152            (S::Committed, T::StartCompensation)
153            | (S::Compensating, T::ScheduleCompensationRetry | T::StartCompensationRetry)
154            | (S::CompensationFailed, T::ResolvedRetry) => S::Compensating,
155            (S::Compensating, T::CompensationSucceeded)
156            | (S::CompensationFailed, T::ResolvedCompensated) => S::Compensated,
157            (S::Compensating, T::CompensationFailed) => S::CompensationFailed,
158
159            _ => {
160                return Err(InvalidTransition { from: self, event });
161            }
162        };
163        Ok(next)
164    }
165}
166
167impl fmt::Display for EffectStatus {
168    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
169        f.write_str(self.as_str())
170    }
171}
172
173/// Something that happens to an effect and may change its status.
174///
175/// Each variant is also an audit event; see [`Transition::as_str`].
176#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
177#[serde(rename_all = "snake_case")]
178#[non_exhaustive]
179pub enum Transition {
180    /// Policy requires approval before the first attempt.
181    RequestApproval,
182    /// Approval was granted.
183    Approve,
184    /// Approval was refused.
185    Deny,
186    /// A precondition permanently refused the effect.
187    PreconditionRejected,
188    /// An attempt begins. Persisted before the action is invoked.
189    StartAttempt,
190    /// The action succeeded and no verification is configured.
191    Succeeded,
192    /// The action failed and the effect definitely did not apply, with no
193    /// retry left or allowed.
194    FailedDefinitively,
195    /// Another attempt is scheduled, because the effect definitely did not
196    /// apply or re-executing it is known to be safe.
197    ScheduleRetry,
198    /// The attempt ended without a definitive answer.
199    OutcomeUnknown,
200    /// The lease holder disappeared mid-attempt (recovery).
201    LeaseExpired,
202    /// Verification begins, after a success or to reconcile an unknown outcome.
203    StartVerification,
204    /// Verification found the effect applied.
205    VerificationConfirmed,
206    /// Verification found, authoritatively, that the effect did not apply,
207    /// and it will not be retried.
208    VerificationNotApplied,
209    /// Verification found remote state that contradicts the effect.
210    VerificationConflict,
211    /// The runtime gave up resolving an unknown outcome on its own.
212    Escalate,
213    /// An operator confirmed the effect applied.
214    ResolvedApplied,
215    /// An operator confirmed the effect did not apply.
216    ResolvedNotApplied,
217    /// An operator asserted it is safe to run the effect again, or to try
218    /// its failed compensation again.
219    ResolvedRetry,
220    /// Undoing a committed effect begins. Persisted before the first
221    /// compensation attempt.
222    StartCompensation,
223    /// Another compensation attempt is scheduled, after a failed attempt.
224    ScheduleCompensationRetry,
225    /// Another compensation attempt begins: a scheduled retry, or a rerun
226    /// of an attempt a crash interrupted. Persisted before it runs.
227    StartCompensationRetry,
228    /// The compensation succeeded.
229    CompensationSucceeded,
230    /// The compensation failed for good.
231    CompensationFailed,
232    /// An operator undid the effect by hand.
233    ResolvedCompensated,
234}
235
236/// Every transition, for exhaustive tests.
237pub const ALL_TRANSITIONS: [Transition; 24] = [
238    Transition::RequestApproval,
239    Transition::Approve,
240    Transition::Deny,
241    Transition::PreconditionRejected,
242    Transition::StartAttempt,
243    Transition::Succeeded,
244    Transition::FailedDefinitively,
245    Transition::ScheduleRetry,
246    Transition::OutcomeUnknown,
247    Transition::LeaseExpired,
248    Transition::StartVerification,
249    Transition::VerificationConfirmed,
250    Transition::VerificationNotApplied,
251    Transition::VerificationConflict,
252    Transition::Escalate,
253    Transition::ResolvedApplied,
254    Transition::ResolvedNotApplied,
255    Transition::ResolvedRetry,
256    Transition::StartCompensation,
257    Transition::ScheduleCompensationRetry,
258    Transition::StartCompensationRetry,
259    Transition::CompensationSucceeded,
260    Transition::CompensationFailed,
261    Transition::ResolvedCompensated,
262];
263
264impl Transition {
265    /// Parses the audit-event name produced by [`Self::as_str`].
266    pub fn parse(s: &str) -> Option<Self> {
267        ALL_TRANSITIONS.into_iter().find(|t| t.as_str() == s)
268    }
269
270    /// The stable audit-event name, e.g. `effect.attempt_started`.
271    pub const fn as_str(self) -> &'static str {
272        match self {
273            Self::RequestApproval => "effect.approval_requested",
274            Self::Approve => "effect.approved",
275            Self::Deny => "effect.denied",
276            Self::PreconditionRejected => "effect.precondition_rejected",
277            Self::StartAttempt => "effect.attempt_started",
278            Self::Succeeded => "effect.succeeded",
279            Self::FailedDefinitively => "effect.failed",
280            Self::ScheduleRetry => "effect.retry_scheduled",
281            Self::OutcomeUnknown => "effect.outcome_unknown",
282            Self::LeaseExpired => "effect.lease_expired",
283            Self::StartVerification => "effect.verification_started",
284            Self::VerificationConfirmed => "effect.verification_confirmed",
285            Self::VerificationNotApplied => "effect.verification_not_applied",
286            Self::VerificationConflict => "effect.verification_conflict",
287            Self::Escalate => "effect.escalated",
288            Self::ResolvedApplied => "effect.resolved_applied",
289            Self::ResolvedNotApplied => "effect.resolved_not_applied",
290            Self::ResolvedRetry => "effect.resolved_retry",
291            Self::StartCompensation => "effect.compensation_started",
292            Self::ScheduleCompensationRetry => "effect.compensation_retry_scheduled",
293            Self::StartCompensationRetry => "effect.compensation_retry_started",
294            Self::CompensationSucceeded => "effect.compensated",
295            Self::CompensationFailed => "effect.compensation_failed",
296            Self::ResolvedCompensated => "effect.resolved_compensated",
297        }
298    }
299}
300
301impl fmt::Display for Transition {
302    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
303        f.write_str(self.as_str())
304    }
305}
306
307/// A transition that is not legal from the current status.
308#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
309#[error("transition {event} is not allowed from status {from}")]
310pub struct InvalidTransition {
311    /// The status the effect was in.
312    pub from: EffectStatus,
313    /// The rejected transition.
314    pub event: Transition,
315}
316
317#[cfg(test)]
318mod tests {
319    use std::collections::{HashSet, VecDeque};
320
321    use proptest::prelude::*;
322
323    use super::*;
324
325    fn successors(status: EffectStatus) -> impl Iterator<Item = EffectStatus> {
326        ALL_TRANSITIONS
327            .into_iter()
328            .filter_map(move |t| status.apply(t).ok())
329    }
330
331    fn reachable_from(start: EffectStatus) -> HashSet<EffectStatus> {
332        let mut seen = HashSet::from([start]);
333        let mut queue = VecDeque::from([start]);
334        while let Some(status) = queue.pop_front() {
335            for next in successors(status) {
336                if seen.insert(next) {
337                    queue.push_back(next);
338                }
339            }
340        }
341        seen
342    }
343
344    fn sources_of(target: EffectStatus) -> HashSet<(EffectStatus, Transition)> {
345        ALL_STATUSES
346            .into_iter()
347            .flat_map(|s| ALL_TRANSITIONS.into_iter().map(move |t| (s, t)))
348            .filter(|(s, t)| s.apply(*t) == Ok(target))
349            .collect()
350    }
351
352    #[test]
353    fn compensation_starts_only_from_committed() {
354        let sources = sources_of(EffectStatus::Compensating);
355        let starts: HashSet<_> = sources
356            .iter()
357            .filter(|(from, _)| *from != EffectStatus::Compensating)
358            .copied()
359            .collect();
360        assert_eq!(
361            starts,
362            HashSet::from([
363                (EffectStatus::Committed, Transition::StartCompensation),
364                (EffectStatus::CompensationFailed, Transition::ResolvedRetry),
365            ])
366        );
367        let exits_of_committed: HashSet<_> = ALL_TRANSITIONS
368            .into_iter()
369            .filter(|t| EffectStatus::Committed.apply(*t).is_ok())
370            .collect();
371        assert_eq!(
372            exits_of_committed,
373            HashSet::from([Transition::StartCompensation])
374        );
375    }
376
377    #[test]
378    fn compensated_requires_evidence() {
379        let events: HashSet<Transition> = sources_of(EffectStatus::Compensated)
380            .into_iter()
381            .map(|(_, t)| t)
382            .collect();
383        assert_eq!(
384            events,
385            HashSet::from([
386                Transition::CompensationSucceeded,
387                Transition::ResolvedCompensated
388            ])
389        );
390    }
391
392    #[test]
393    fn terminal_states_have_no_exits() {
394        for status in ALL_STATUSES.into_iter().filter(|s| s.is_terminal()) {
395            assert_eq!(successors(status).count(), 0, "{status} has an exit");
396        }
397    }
398
399    #[test]
400    fn every_state_is_reachable_from_pending() {
401        let reachable = reachable_from(EffectStatus::Pending);
402        for status in ALL_STATUSES {
403            assert!(reachable.contains(&status), "{status} is unreachable");
404        }
405    }
406
407    #[test]
408    fn every_state_can_settle() {
409        for status in ALL_STATUSES {
410            assert!(
411                reachable_from(status).iter().any(|s| s.is_terminal()),
412                "{status} can never reach a terminal state"
413            );
414        }
415    }
416
417    #[test]
418    fn only_pending_starts_an_attempt() {
419        let sources = sources_of(EffectStatus::Executing);
420        assert_eq!(
421            sources,
422            HashSet::from([(EffectStatus::Pending, Transition::StartAttempt)])
423        );
424    }
425
426    #[test]
427    fn in_doubt_states_never_restart_without_a_retry_decision() {
428        // Leaving doubt for Pending must always be an explicit decision that
429        // re-executing is safe, never a side effect of some other event.
430        for status in ALL_STATUSES.into_iter().filter(|s| s.is_in_doubt()) {
431            for t in ALL_TRANSITIONS {
432                if status.apply(t) == Ok(EffectStatus::Pending) {
433                    assert!(
434                        matches!(t, Transition::ScheduleRetry | Transition::ResolvedRetry),
435                        "{status} --{t}--> pending"
436                    );
437                }
438            }
439        }
440    }
441
442    #[test]
443    fn committed_requires_evidence() {
444        let events: HashSet<Transition> = sources_of(EffectStatus::Committed)
445            .into_iter()
446            .map(|(_, t)| t)
447            .collect();
448        assert_eq!(
449            events,
450            HashSet::from([
451                Transition::Succeeded,
452                Transition::VerificationConfirmed,
453                Transition::ResolvedApplied,
454            ])
455        );
456    }
457
458    #[test]
459    fn a_lost_lease_never_counts_as_failure() {
460        for status in ALL_STATUSES {
461            if let Ok(next) = status.apply(Transition::LeaseExpired) {
462                assert_eq!(next, EffectStatus::Unknown);
463            }
464        }
465    }
466
467    #[test]
468    fn storage_representation_round_trips_and_matches_serde() {
469        for status in ALL_STATUSES {
470            assert_eq!(EffectStatus::parse(status.as_str()), Some(status));
471            assert_eq!(
472                serde_json::to_string(&status).unwrap(),
473                format!("\"{}\"", status.as_str())
474            );
475        }
476        let names: HashSet<_> = ALL_TRANSITIONS.iter().map(|t| t.as_str()).collect();
477        assert_eq!(names.len(), ALL_TRANSITIONS.len());
478        for transition in ALL_TRANSITIONS {
479            assert_eq!(Transition::parse(transition.as_str()), Some(transition));
480        }
481    }
482
483    fn transition() -> impl Strategy<Value = Transition> {
484        proptest::sample::select(ALL_TRANSITIONS.to_vec())
485    }
486
487    proptest! {
488        #[test]
489        fn random_histories_respect_invariants(events in prop::collection::vec(transition(), 0..64)) {
490            let mut status = EffectStatus::Pending;
491            for event in events {
492                match status.apply(event) {
493                    Ok(next) => {
494                        prop_assert!(!status.is_terminal());
495                        if status == EffectStatus::Committed {
496                            prop_assert_ne!(next, EffectStatus::Executing);
497                        }
498                        status = next;
499                    }
500                    Err(err) => {
501                        prop_assert_eq!(err.from, status);
502                    }
503                }
504            }
505        }
506    }
507}