Skip to main content

agent_effects_store/
record.rs

1use std::time::{Duration, SystemTime};
2
3use serde::{Deserialize, Serialize};
4use serde_json::Value;
5
6use crate::id::{EffectId, EffectKey, WorkerId};
7use crate::kind::EffectKind;
8use crate::state::{EffectStatus, InvalidTransition, Transition};
9use crate::{EffectEvent, ErrorRecord, Lease, NewEffect, StoreError, TransitionRequest};
10
11/// The durable state of one effect.
12///
13/// The methods are pure: they check a change against the record and apply
14/// it in memory. Stores call them inside a transaction and persist the
15/// result, so the rules are the same for every backend.
16#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
17pub struct EffectRecord {
18    /// Record id.
19    pub id: EffectId,
20    /// Logical identity, unique per store.
21    pub key: EffectKey,
22    /// The effect's kind.
23    pub kind: EffectKind,
24    /// Current status.
25    pub status: EffectStatus,
26    /// The input, redacted for storage.
27    pub input: Option<Value>,
28    /// Hash of the unredacted input.
29    pub input_fingerprint: Option<String>,
30    /// The latest action or verification result.
31    pub output: Option<Value>,
32    /// The latest recorded failure.
33    pub last_error: Option<ErrorRecord>,
34    /// Who asked for the effect.
35    pub created_by: Option<String>,
36    /// Attempts started so far.
37    pub attempt_count: u32,
38    /// An attempt may have applied the effect, and no evidence has shown
39    /// otherwise since. Set when an outcome becomes unknown; cleared only by
40    /// a trusted verification or an operator. While it is set, the effect
41    /// cannot become `Failed` through a failed attempt (see [`Self::apply`]).
42    pub may_have_applied: bool,
43    /// Compensation attempts started so far.
44    pub compensation_attempts: u32,
45    /// The effect was approved, so it is not asked again, even after a
46    /// restart.
47    pub approved: bool,
48    /// When a scheduled retry may start.
49    pub next_attempt_at: Option<SystemTime>,
50    /// When the latest attempt started.
51    pub attempt_started_at: Option<SystemTime>,
52    /// When the latest attempt was last known to be in flight: the first
53    /// transition out of `Executing` after it started. Settle delays count
54    /// from here.
55    pub attempt_ended_at: Option<SystemTime>,
56    /// Current lease holder, if any. The lease may have expired.
57    pub lease_owner: Option<WorkerId>,
58    /// Fencing token; incremented by every lease acquisition.
59    pub lease_epoch: u64,
60    /// When the current lease lapses.
61    pub lease_expires_at: Option<SystemTime>,
62    /// Incremented by every transition. Lease operations do not change it.
63    pub version: u64,
64    /// Creation time.
65    pub created_at: SystemTime,
66    /// Time of the latest change.
67    pub updated_at: SystemTime,
68    /// When the effect committed.
69    pub committed_at: Option<SystemTime>,
70}
71
72impl EffectRecord {
73    /// A fresh `Pending` record.
74    pub fn new(new: NewEffect) -> Self {
75        Self {
76            id: new.id,
77            key: new.key,
78            kind: new.kind,
79            status: EffectStatus::Pending,
80            input: new.input,
81            input_fingerprint: new.input_fingerprint,
82            output: None,
83            last_error: None,
84            created_by: new.created_by,
85            attempt_count: 0,
86            may_have_applied: false,
87            compensation_attempts: 0,
88            approved: false,
89            next_attempt_at: None,
90            attempt_started_at: None,
91            attempt_ended_at: None,
92            lease_owner: None,
93            lease_epoch: 0,
94            lease_expires_at: None,
95            version: 0,
96            created_at: new.now,
97            updated_at: new.now,
98            committed_at: None,
99        }
100    }
101
102    /// The holder of a lease that is still live at `now`.
103    pub fn live_lease_owner(&self, now: SystemTime) -> Option<&WorkerId> {
104        match (&self.lease_owner, self.lease_expires_at) {
105            (Some(owner), Some(expires_at)) if expires_at > now => Some(owner),
106            _ => None,
107        }
108    }
109
110    /// Takes the lease for `ttl`, provided nobody holds a live one.
111    ///
112    /// # Errors
113    ///
114    /// [`StoreError::LeaseHeld`] if a live lease exists, even one held by
115    /// `owner`: two tasks of one worker must not run the same effect either.
116    pub fn acquire_lease(
117        &mut self,
118        owner: &WorkerId,
119        now: SystemTime,
120        ttl: Duration,
121    ) -> Result<Lease, StoreError> {
122        if let Some(holder) = self.live_lease_owner(now) {
123            return Err(StoreError::LeaseHeld {
124                owner: holder.clone(),
125                expires_at: self.lease_expires_at.unwrap_or(now),
126            });
127        }
128        let expires_at = expiry(now, ttl);
129        self.lease_epoch += 1;
130        self.lease_owner = Some(owner.clone());
131        self.lease_expires_at = Some(expires_at);
132        Ok(Lease {
133            effect_id: self.id,
134            owner: owner.clone(),
135            epoch: self.lease_epoch,
136            expires_at,
137        })
138    }
139
140    /// Extends `lease` to `now + ttl`.
141    ///
142    /// # Errors
143    ///
144    /// [`StoreError::LeaseLost`] if the lease expired or was taken over.
145    /// Renewal is strict: an expired lease cannot be revived, even if nobody
146    /// else took it.
147    pub fn renew_lease(
148        &mut self,
149        lease: &Lease,
150        now: SystemTime,
151        ttl: Duration,
152    ) -> Result<Lease, StoreError> {
153        self.check_lease(lease, now)?;
154        let expires_at = expiry(now, ttl);
155        self.lease_expires_at = Some(expires_at);
156        Ok(Lease {
157            expires_at,
158            ..lease.clone()
159        })
160    }
161
162    /// Clears `lease` if it is still the current one, expired or not.
163    /// Returns whether anything changed.
164    pub fn release_lease(&mut self, lease: &Lease) -> bool {
165        let current = lease.effect_id == self.id
166            && self.lease_epoch == lease.epoch
167            && self.lease_owner.as_ref() == Some(&lease.owner);
168        if current {
169            self.lease_owner = None;
170            self.lease_expires_at = None;
171        }
172        current
173    }
174
175    /// Applies `request` and returns the audit event to persist with it.
176    /// On error the record is unchanged.
177    ///
178    /// Checks, in order: the lease (or the absence of a live one), the
179    /// version, and the transition table. Then it updates the bookkeeping:
180    ///
181    /// - [`Transition::StartAttempt`] increments the attempt count, stamps
182    ///   `attempt_started_at`, and clears `attempt_ended_at` and
183    ///   `next_attempt_at`; any transition out of `Executing` stamps
184    ///   `attempt_ended_at`;
185    /// - [`Transition::StartCompensation`] sets `compensation_attempts` to 1
186    ///   and [`Transition::StartCompensationRetry`] increments it; both
187    ///   clear `next_attempt_at`;
188    /// - a retry transition, including
189    ///   [`Transition::ScheduleCompensationRetry`], sets `next_attempt_at`
190    ///   (default `now`);
191    /// - reaching `Committed` stamps `committed_at`;
192    /// - reaching `Unknown` sets `may_have_applied`; evidence that the effect
193    ///   did not apply clears it (`VerificationNotApplied`,
194    ///   `ResolvedNotApplied`, or a `ScheduleRetry` out of `Verifying`, which
195    ///   the runtime issues only after a trusted "not applied");
196    /// - `output` and `error`, when given, replace the stored ones.
197    ///
198    /// `Failed` must mean the effect did not apply. A failed attempt proves
199    /// that only for itself, so `FailedDefinitively` is refused while
200    /// `may_have_applied` is set (except for `Read` effects, which apply
201    /// nothing).
202    ///
203    /// # Errors
204    ///
205    /// [`StoreError::LeaseLost`], [`StoreError::LeaseHeld`],
206    /// [`StoreError::VersionConflict`] or [`StoreError::InvalidTransition`].
207    pub fn apply(&mut self, request: TransitionRequest) -> Result<EffectEvent, StoreError> {
208        let now = request.now;
209        match &request.lease {
210            Some(lease) => self.check_lease(lease, now)?,
211            None => {
212                if let Some(owner) = self.live_lease_owner(now) {
213                    return Err(StoreError::LeaseHeld {
214                        owner: owner.clone(),
215                        expires_at: self.lease_expires_at.unwrap_or(now),
216                    });
217                }
218            }
219        }
220        if request.expected_version != self.version {
221            return Err(StoreError::VersionConflict {
222                expected: request.expected_version,
223                actual: self.version,
224            });
225        }
226        let from = self.status;
227        let to = from.apply(request.transition)?;
228        if request.transition == Transition::FailedDefinitively
229            && self.may_have_applied
230            && self.kind != EffectKind::Read
231        {
232            return Err(InvalidTransition {
233                from,
234                event: request.transition,
235            }
236            .into());
237        }
238
239        match request.transition {
240            Transition::StartAttempt => {
241                self.attempt_count = self.attempt_count.saturating_add(1);
242                self.attempt_started_at = Some(now);
243                self.attempt_ended_at = None;
244                self.next_attempt_at = None;
245            }
246            Transition::ScheduleRetry
247            | Transition::ResolvedRetry
248            | Transition::ScheduleCompensationRetry => {
249                self.next_attempt_at = Some(request.next_attempt_at.unwrap_or(now));
250            }
251            Transition::Approve => self.approved = true,
252            Transition::StartCompensation => {
253                self.compensation_attempts = 1;
254                self.next_attempt_at = None;
255            }
256            Transition::StartCompensationRetry => {
257                self.compensation_attempts = self.compensation_attempts.saturating_add(1);
258                self.next_attempt_at = None;
259            }
260            _ => {}
261        }
262        if from == EffectStatus::Executing {
263            self.attempt_ended_at = Some(now);
264        }
265        if to == EffectStatus::Committed {
266            self.committed_at = Some(now);
267        }
268        if to == EffectStatus::Unknown {
269            self.may_have_applied = true;
270        }
271        let shown_not_applied = matches!(
272            request.transition,
273            Transition::VerificationNotApplied | Transition::ResolvedNotApplied
274        ) || (from == EffectStatus::Verifying
275            && request.transition == Transition::ScheduleRetry);
276        if shown_not_applied {
277            self.may_have_applied = false;
278        }
279        if let Some(output) = request.output {
280            self.output = Some(output);
281        }
282        if let Some(error) = request.error {
283            self.last_error = Some(error);
284        }
285        self.status = to;
286        self.version += 1;
287        self.updated_at = now;
288
289        Ok(EffectEvent {
290            effect_id: self.id,
291            sequence: self.version,
292            transition: request.transition,
293            from,
294            to,
295            attempt: self.attempt_count,
296            actor: request.actor,
297            payload: request.payload,
298            at: now,
299        })
300    }
301
302    fn check_lease(&self, lease: &Lease, now: SystemTime) -> Result<(), StoreError> {
303        let valid = lease.effect_id == self.id
304            && self.lease_epoch == lease.epoch
305            && self.lease_owner.as_ref() == Some(&lease.owner)
306            && self
307                .lease_expires_at
308                .is_some_and(|expires_at| expires_at > now);
309        if valid {
310            Ok(())
311        } else {
312            Err(StoreError::LeaseLost)
313        }
314    }
315}
316
317/// `now + ttl`. An overflowing TTL yields an already-expired lease, which
318/// fails safe: the holder is fenced off immediately.
319fn expiry(now: SystemTime, ttl: Duration) -> SystemTime {
320    now.checked_add(ttl).unwrap_or(now)
321}
322
323#[cfg(test)]
324mod tests {
325    use super::*;
326    use crate::id::{EffectName, LogicalKey};
327
328    const TTL: Duration = Duration::from_secs(30);
329
330    fn t(secs: u64) -> SystemTime {
331        SystemTime::UNIX_EPOCH + Duration::from_secs(1_700_000_000 + secs)
332    }
333
334    fn record() -> EffectRecord {
335        let key = EffectKey::new(
336            EffectName::new("payment.charge").unwrap(),
337            LogicalKey::new("order_1").unwrap(),
338        );
339        EffectRecord::new(NewEffect::new(key, EffectKind::IrreversibleWrite, t(0)))
340    }
341
342    fn worker(name: &str) -> WorkerId {
343        WorkerId::new(name)
344    }
345
346    #[test]
347    fn failed_apply_leaves_the_record_unchanged() {
348        let mut rec = record();
349        let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
350        let before = rec.clone();
351
352        let mut bad = TransitionRequest::new(&rec, Some(&lease), Transition::Succeeded, t(1));
353        bad.output = Some(Value::from(1));
354        assert!(matches!(
355            rec.apply(bad),
356            Err(StoreError::InvalidTransition(_))
357        ));
358        assert_eq!(rec, before);
359    }
360
361    #[test]
362    fn stale_epoch_is_fenced_off_after_takeover() {
363        let mut rec = record();
364        let old = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
365        let new = rec.acquire_lease(&worker("b"), t(30), TTL).unwrap();
366        assert_eq!(new.epoch, old.epoch + 1);
367
368        let request = TransitionRequest::new(&rec, Some(&old), Transition::StartAttempt, t(31));
369        assert!(matches!(rec.apply(request), Err(StoreError::LeaseLost)));
370        assert!(
371            !rec.release_lease(&old),
372            "stale release must not clear b's lease"
373        );
374        assert_eq!(rec.live_lease_owner(t(31)), Some(&worker("b")));
375    }
376
377    #[test]
378    fn same_owner_cannot_double_acquire() {
379        let mut rec = record();
380        rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
381        assert!(matches!(
382            rec.acquire_lease(&worker("a"), t(1), TTL),
383            Err(StoreError::LeaseHeld { .. })
384        ));
385    }
386
387    #[test]
388    fn overflowing_ttl_fails_safe() {
389        let mut rec = record();
390        let lease = rec
391            .acquire_lease(&worker("a"), t(0), Duration::MAX)
392            .unwrap();
393        assert_eq!(lease.expires_at, t(0));
394        assert_eq!(rec.live_lease_owner(t(0)), None);
395    }
396
397    #[test]
398    fn a_failed_attempt_cannot_fail_an_effect_that_may_have_applied() {
399        let mut rec = record();
400        let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
401        for transition in [
402            Transition::StartAttempt,
403            Transition::OutcomeUnknown,
404            Transition::ScheduleRetry,
405            Transition::StartAttempt,
406        ] {
407            rec.apply(TransitionRequest::new(&rec, Some(&lease), transition, t(1)))
408                .unwrap();
409        }
410        assert!(
411            rec.may_have_applied,
412            "the unknown outcome is remembered across retries"
413        );
414        let before = rec.clone();
415        let fail = TransitionRequest::new(&rec, Some(&lease), Transition::FailedDefinitively, t(2));
416        assert!(matches!(
417            rec.apply(fail),
418            Err(StoreError::InvalidTransition(_))
419        ));
420        assert_eq!(rec, before);
421
422        // A trusted verification clears it; then a failure is a failure.
423        for transition in [
424            Transition::StartVerification,
425            Transition::ScheduleRetry,
426            Transition::StartAttempt,
427        ] {
428            rec.apply(TransitionRequest::new(&rec, Some(&lease), transition, t(3)))
429                .unwrap();
430        }
431        assert!(!rec.may_have_applied);
432        rec.apply(TransitionRequest::new(
433            &rec,
434            Some(&lease),
435            Transition::FailedDefinitively,
436            t(4),
437        ))
438        .unwrap();
439        assert_eq!(rec.status, EffectStatus::Failed);
440    }
441
442    #[test]
443    fn bookkeeping_follows_transitions() {
444        let mut rec = record();
445        let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
446
447        let event = rec
448            .apply(TransitionRequest::new(
449                &rec,
450                Some(&lease),
451                Transition::StartAttempt,
452                t(1),
453            ))
454            .unwrap();
455        assert_eq!((event.sequence, event.attempt), (1, 1));
456        assert_eq!(rec.attempt_started_at, Some(t(1)));
457
458        let mut retry = TransitionRequest::new(&rec, Some(&lease), Transition::ScheduleRetry, t(2));
459        retry.next_attempt_at = Some(t(10));
460        retry.error = Some(ErrorRecord {
461            class: Some(crate::FailureClass::Transient),
462            message: "connection refused".into(),
463        });
464        rec.apply(retry).unwrap();
465        assert_eq!(rec.next_attempt_at, Some(t(10)));
466        assert_eq!(rec.attempt_ended_at, Some(t(2)));
467        assert!(rec.last_error.is_some());
468
469        rec.apply(TransitionRequest::new(
470            &rec,
471            Some(&lease),
472            Transition::StartAttempt,
473            t(10),
474        ))
475        .unwrap();
476        assert_eq!(rec.attempt_count, 2);
477        assert_eq!(rec.next_attempt_at, None);
478        assert_eq!(rec.attempt_ended_at, None);
479
480        let mut done = TransitionRequest::new(&rec, Some(&lease), Transition::Succeeded, t(11));
481        done.output = Some(serde_json::json!({ "payment": "pi_1" }));
482        let event = rec.apply(done).unwrap();
483        assert_eq!(event.to, EffectStatus::Committed);
484        assert_eq!(rec.committed_at, Some(t(11)));
485        assert_eq!(rec.version, 4);
486    }
487}