Skip to main content

agentplane/core/
retry.rs

1//! Retry policy — how many times, how long apart, and when not at all.
2//!
3//! # Retrying is a safety decision before it is a reliability one
4//!
5//! Most runtimes treat retry as a reliability knob: the call failed, try again,
6//! and push idempotency onto whoever wrote the tool. That is a defensible
7//! trade when the worst case is a duplicate log line. It is the wrong trade
8//! when the worst case is a duplicate payment, and it is wrong in a way that
9//! only shows up in production, on the one call that timed out.
10//!
11//! So a policy here never decides alone. Three things do, in order:
12//!
13//! 1. The failure's [`Disposition`](crate::core::Disposition) — did the call
14//!    reach the outside world? A refused connection and a timed-out request are
15//!    both transient, and only one of them is safe to repeat.
16//! 2. The effect's [`Recovery`](crate::core::Recovery) — for an in-doubt
17//!    failure, is guessing permitted at all?
18//! 3. This policy — and only then, how many times and how far apart.
19//!
20//! A policy cannot authorise a repeat that the first two refuse. Raising
21//! `max_attempts` never makes a mutating in-doubt call retryable.
22//!
23//! # Backoff waits in-process, and that is deliberate
24//!
25//! A retry's backoff is a `tokio` sleep, so it holds the worker for its
26//! duration. `max_backoff` is therefore not just a schedule ceiling — it is the
27//! bound on how long one effect can occupy a frame, and it is why the default
28//! is seconds rather than minutes.
29//!
30//! Suspending instead would be strictly worse here. Waking a suspended run
31//! replays it from the beginning, so a run fifty steps deep pays fifty steps of
32//! replay to avoid a five-second sleep. Durable suspension wins for waits
33//! measured in minutes and hours; in-process wins for the seconds-scale jitter
34//! that retry actually needs.
35//!
36//! So the boundary is drawn by purpose, not by duration:
37//!
38//! * **Retrying a flaky call** — this module. Bounded by `max_backoff`.
39//! * **Waiting for the world** — a settlement date, five Werktage — is not a
40//!   retry. Use [`StepCtx::sleep`](crate::runtime::StepCtx::sleep), which
41//!   suspends the run and costs a row rather than a thread.
42//!
43//! Setting `max_backoff` to an hour is legal and will hold a worker for an
44//! hour. That is stated rather than prevented, because a deployment that knows
45//! its own concurrency may want exactly that.
46//!
47//! # Backoff is computed, not drawn — unless the peer names it
48//!
49//! The runtime forbids ambient randomness, so jitter cannot come from an RNG.
50//! It is derived instead from the hash of the run, the effect key, and the
51//! attempt number — which decorrelates concurrent runs the way jitter is
52//! supposed to, while staying a pure function of things already in the journal.
53//!
54//! A computed schedule is a **guess about when a service recovers**, and one
55//! failure does not need guessing at: a peer that answers *rate limited* and
56//! names its own window. So advice wins where a peer gives it, bounded by
57//! [`max_advice`](RetryPolicy::max_advice). The parsing rule is
58//! [`retry_after_seconds`], shared by every wire here so no surface acts on
59//! advice another would refuse.
60
61use std::time::Duration;
62
63use serde::{Deserialize, Serialize};
64
65use crate::core::{Digest, EffectKey, RunId};
66
67/// How many times to repeat a failed effect, and how far apart.
68///
69/// The defaults are safe rather than timid: three attempts with exponential
70/// backoff. They are safe *because* [`Disposition`](crate::core::Disposition)
71/// gates them — an effect whose failure is in-doubt is not repeated under this
72/// policy unless its [`Recovery`](crate::core::Recovery) permits guessing, no
73/// matter what `max_attempts` says.
74#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
75pub struct RetryPolicy {
76    /// Total attempts including the first. `1` means "never repeat".
77    pub max_attempts: u32,
78    /// Delay before the second attempt.
79    pub initial_backoff: Duration,
80    /// Ceiling on the delay, however many attempts have passed.
81    ///
82    /// Also the bound on how long one effect holds a worker — see the module
83    /// docs. For waits longer than a few seconds, the run should suspend
84    /// instead: that is [`StepCtx::sleep`](crate::runtime::StepCtx::sleep), not
85    /// a retry policy.
86    pub max_backoff: Duration,
87    /// Growth factor per attempt. An integer, so the schedule is exactly
88    /// reproducible on any platform without depending on float rounding.
89    pub multiplier: u32,
90    /// Whether to spread the delay across runs. See the module docs — this is
91    /// derived from a hash, not drawn from an RNG.
92    pub jitter: bool,
93    /// The longest this run will wait because a *peer* named a time.
94    ///
95    /// A separate ceiling from [`max_backoff`](Self::max_backoff), because the
96    /// two bound different risks. `max_backoff` bounds a **guess**: nobody knows
97    /// when the service recovers, so waiting long is waste. This bounds
98    /// **somebody else's word** — a `Retry-After` from the one party with an
99    /// interest in never being called again — and the number that makes the
100    /// wait useful is theirs, not ours. Sharing one ceiling would mean either
101    /// guessing for as long as a peer may demand, or discarding advice that is
102    /// merely longer than a guess would have been.
103    ///
104    /// A minute covers every published provider rate-limit window and is short
105    /// enough that a hostile or broken value costs one wasted minute of one
106    /// worker. Advice past it is **clamped, not discarded**: waiting part of a
107    /// window is closer to right than ignoring it, and if the window really was
108    /// longer the next refusal names what is left.
109    ///
110    /// Raising it holds a worker for that long. A rate limit measured in more
111    /// than minutes is not a retry at all — that is
112    /// [`StepCtx::sleep`](crate::runtime::StepCtx::sleep), which costs a row.
113    pub max_advice: Duration,
114}
115
116impl Default for RetryPolicy {
117    fn default() -> Self {
118        Self {
119            max_attempts: 3,
120            initial_backoff: Duration::from_millis(100),
121            max_backoff: Duration::from_secs(10),
122            multiplier: 2,
123            jitter: true,
124            max_advice: Duration::from_secs(60),
125        }
126    }
127}
128
129impl RetryPolicy {
130    /// Never repeat. The first failure is final.
131    ///
132    /// The right choice for an effect that is expensive, externally rate
133    /// limited, or whose driver already retries internally — a policy stacked
134    /// on a driver that retries is a multiplication, not an addition.
135    #[must_use]
136    pub const fn never() -> Self {
137        Self {
138            max_attempts: 1,
139            initial_backoff: Duration::ZERO,
140            max_backoff: Duration::ZERO,
141            multiplier: 1,
142            jitter: false,
143            // Nothing to advise: there is no second attempt to schedule.
144            max_advice: Duration::ZERO,
145        }
146    }
147
148    /// `n` total attempts with the default backoff schedule.
149    #[must_use]
150    pub fn attempts(n: u32) -> Self {
151        Self {
152            max_attempts: n.max(1),
153            ..Self::default()
154        }
155    }
156
157    /// Replace the backoff schedule, keeping the attempt count.
158    #[must_use]
159    pub fn with_backoff(mut self, initial: Duration, max: Duration) -> Self {
160        self.initial_backoff = initial;
161        self.max_backoff = max;
162        self
163    }
164
165    /// Turn jitter off, making the schedule identical across runs.
166    #[must_use]
167    pub fn without_jitter(mut self) -> Self {
168        self.jitter = false;
169        self
170    }
171
172    /// Change how long a peer's own `Retry-After` may hold a worker.
173    ///
174    /// See [`max_advice`](Self::max_advice) for why this is not `max_backoff`.
175    #[must_use]
176    pub const fn wait_at_most(mut self, advice: Duration) -> Self {
177        self.max_advice = advice;
178        self
179    }
180
181    /// How long to wait before `attempt`, given whatever the peer said.
182    ///
183    /// `advice` wins when present, clamped to [`max_advice`](Self::max_advice),
184    /// **including when it is shorter** than the computed backoff: the peer is
185    /// describing its own recovery, and waiting longer than it asked buys
186    /// nothing.
187    #[must_use]
188    pub fn wait_before(
189        &self,
190        run: RunId,
191        key: EffectKey,
192        attempt: u32,
193        advice: Option<Duration>,
194    ) -> Duration {
195        // Attempt 1 is not a retry, so nothing schedules it — not even a peer.
196        // Advice can only arrive attached to a failure, so this guard is
197        // belt-and-braces rather than reachable, and it keeps the one property
198        // every caller relies on: the first attempt is immediate.
199        if attempt <= 1 {
200            return Duration::ZERO;
201        }
202        match advice {
203            Some(named) if !named.is_zero() => named.min(self.max_advice),
204            _ => self.backoff(run, key, attempt),
205        }
206    }
207
208    /// Whether another attempt is left after `attempt` has failed.
209    ///
210    /// `attempt` is 1-based, matching what the journal records.
211    #[must_use]
212    pub fn permits(&self, attempt: u32) -> bool {
213        attempt < self.max_attempts
214    }
215
216    /// How long to wait before `attempt` (1-based; attempt 1 never waits).
217    ///
218    /// Exponential with an integer multiplier and a ceiling, then optionally
219    /// spread by a hash-derived factor in `[0.5, 1.0]` of the computed delay.
220    /// Halving rather than scaling from zero keeps a floor under the schedule:
221    /// full jitter can pick a near-zero delay and hammer a service that is
222    /// already struggling.
223    #[must_use]
224    pub fn backoff(&self, run: RunId, key: EffectKey, attempt: u32) -> Duration {
225        if attempt <= 1 {
226            return Duration::ZERO;
227        }
228
229        // Saturating throughout: a large multiplier and a large attempt count
230        // must produce `max_backoff`, not an overflow panic.
231        // A multiplier of one or less never grows the delay, so there is
232        // nothing to iterate.
233        let steps = if self.multiplier > 1 { attempt - 2 } else { 0 };
234        let mut delay = self.initial_backoff;
235        for _ in 0..steps {
236            delay = delay.saturating_mul(self.multiplier).min(self.max_backoff);
237            if delay >= self.max_backoff {
238                break;
239            }
240        }
241        let delay = delay.min(self.max_backoff);
242
243        if !self.jitter || delay.is_zero() {
244            return delay;
245        }
246
247        // Deterministic stand-in for an RNG draw. Including the run id is what
248        // decorrelates two runs of the same plan retrying the same effect —
249        // without it, identical keys would produce identical schedules and
250        // reconverge into exactly the thundering herd jitter exists to prevent.
251        let mut seed = Vec::with_capacity(64);
252        seed.extend_from_slice(run.to_string().as_bytes());
253        seed.extend_from_slice(&key.to_hex().into_bytes());
254        seed.extend_from_slice(&attempt.to_be_bytes());
255        let digest = Digest::of(&seed);
256        let spread = u64::from_be_bytes(digest.as_bytes()[..8].try_into().expect("8 bytes"));
257
258        // Map into [0.5, 1.0] of the computed delay, in integer arithmetic.
259        let half = delay.as_nanos() / 2;
260        let extra = (half.saturating_mul(u128::from(spread))) / u128::from(u64::MAX);
261        Duration::from_nanos(u64::try_from(half + extra).unwrap_or(u64::MAX))
262    }
263}
264
265/// Parse a `Retry-After` value into seconds.
266///
267/// Only the delta-seconds form is read. The HTTP-date form is equally legal and
268/// is deliberately ignored: acting on it means trusting the sender's clock
269/// against ours, and a peer whose clock is a day fast would park a run for a
270/// day. A value this cannot read is *no advice*, which is the same answer as an
271/// absent header — the caller's own schedule applies.
272///
273/// Here rather than beside either caller because it is the same question on
274/// every wire this crate speaks, and two spellings of a bound are two bounds:
275/// the one that drifts is whichever surface nobody probed.
276#[must_use]
277pub fn retry_after_seconds(value: &str) -> Option<u64> {
278    let seconds = value.trim().parse::<u64>().ok()?;
279    // Zero is a legal encoding of "come back immediately", and taking it
280    // literally would replace the caller's backoff with no wait at all — which
281    // is the one schedule a rate limit must not produce.
282    (seconds > 0).then_some(seconds)
283}
284
285#[cfg(test)]
286mod tests {
287    use super::*;
288    use crate::core::StepId;
289
290    fn key() -> EffectKey {
291        EffectKey::derive(
292            StepId(1),
293            crate::core::Phase::Forward,
294            0,
295            1,
296            "test.effect",
297            b"{}",
298        )
299    }
300
301    #[test]
302    fn the_first_attempt_never_waits() {
303        assert_eq!(
304            RetryPolicy::default().backoff(RunId::generate(), key(), 1),
305            Duration::ZERO
306        );
307    }
308
309    /// `without_jitter` turns jitter off and changes nothing else.
310    ///
311    /// Every existing schedule test builds `RetryPolicy` by struct literal with
312    /// `jitter: false`, so the builder a caller actually reaches for had no test
313    /// at all — and a `without_jitter` that set the wrong field would leave the
314    /// schedule non-deterministic while reading as though it had fixed it.
315    #[test]
316    fn without_jitter_makes_the_schedule_reproducible() {
317        let jittered = RetryPolicy::default();
318        let fixed = RetryPolicy::default().without_jitter();
319
320        assert!(jittered.jitter, "the default is jittered");
321        assert!(!fixed.jitter);
322        assert_eq!(
323            fixed.max_attempts, jittered.max_attempts,
324            "turning jitter off must not move the attempt ceiling"
325        );
326        assert_eq!(fixed.initial_backoff, jittered.initial_backoff);
327        assert_eq!(fixed.max_backoff, jittered.max_backoff);
328
329        // The property jitter removal is *for*: two runs of the same effect now
330        // wait the same amount. With jitter on, the run id decorrelates them.
331        let key = key();
332        let (a, b) = (RunId::generate(), RunId::generate());
333        assert_eq!(fixed.backoff(a, key, 3), fixed.backoff(b, key, 3));
334    }
335
336    #[test]
337    fn backoff_grows_and_then_stops_at_the_ceiling() {
338        let p = RetryPolicy {
339            max_attempts: 10,
340            initial_backoff: Duration::from_millis(100),
341            max_backoff: Duration::from_millis(800),
342            multiplier: 2,
343            jitter: false,
344            ..RetryPolicy::default()
345        };
346        let run = RunId::generate();
347        let at = |n| p.backoff(run, key(), n);
348        assert_eq!(at(2), Duration::from_millis(100));
349        assert_eq!(at(3), Duration::from_millis(200));
350        assert_eq!(at(4), Duration::from_millis(400));
351        assert_eq!(at(5), Duration::from_millis(800));
352        assert_eq!(at(6), Duration::from_millis(800), "ceiling holds");
353        assert_eq!(at(50), Duration::from_millis(800), "and keeps holding");
354    }
355
356    /// A flat multiplier holds the initial delay, and answers without walking
357    /// every attempt to get there.
358    #[test]
359    fn a_flat_multiplier_holds_the_initial_delay() {
360        let p = RetryPolicy {
361            initial_backoff: Duration::from_millis(100),
362            max_backoff: Duration::from_secs(10),
363            multiplier: 1,
364            jitter: false,
365            ..RetryPolicy::default()
366        };
367        assert_eq!(
368            p.backoff(RunId::generate(), key(), u32::MAX),
369            Duration::from_millis(100)
370        );
371    }
372
373    #[test]
374    fn an_absurd_schedule_saturates_instead_of_panicking() {
375        let p = RetryPolicy {
376            max_attempts: u32::MAX,
377            initial_backoff: Duration::from_secs(1),
378            max_backoff: Duration::from_mins(1),
379            multiplier: u32::MAX,
380            jitter: true,
381            ..RetryPolicy::default()
382        };
383        assert!(p.backoff(RunId::generate(), key(), u32::MAX) <= Duration::from_mins(1));
384    }
385
386    #[test]
387    fn jitter_stays_within_half_the_delay() {
388        let p = RetryPolicy {
389            jitter: true,
390            ..RetryPolicy::default()
391        };
392        let run = RunId::generate();
393        for attempt in 2..8 {
394            let d = p.backoff(run, key(), attempt);
395            let plain = RetryPolicy { jitter: false, ..p }.backoff(run, key(), attempt);
396            assert!(
397                d >= plain / 2 && d <= plain,
398                "attempt {attempt}: {d:?} outside [{:?}, {plain:?}]",
399                plain / 2
400            );
401        }
402    }
403
404    #[test]
405    fn jitter_decorrelates_runs_but_repeats_for_one_run() {
406        let (a, b) = (RunId::generate(), RunId::generate());
407        let p = RetryPolicy::default();
408        assert_ne!(
409            p.backoff(a, key(), 3),
410            p.backoff(b, key(), 3),
411            "two runs retrying the same effect must not reconverge"
412        );
413        assert_eq!(
414            p.backoff(a, key(), 3),
415            p.backoff(a, key(), 3),
416            "and the schedule must be a pure function, not a draw"
417        );
418    }
419
420    #[test]
421    fn never_permits_no_second_attempt() {
422        assert!(!RetryPolicy::never().permits(1));
423        assert!(RetryPolicy::attempts(3).permits(1));
424        assert!(RetryPolicy::attempts(3).permits(2));
425        assert!(!RetryPolicy::attempts(3).permits(3));
426    }
427}