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}