返回 CodeWhale
retry_status.rs
根目录 / crates / runtime / src / retry_status.rs
1 //! Process-wide retry-state surface (#499).
2 //!
3 //! Read-side caveat (0.9.4): the renderer this module was written for was
4 //! the legacy footer's retry banner, which went with `FooterWidget`. The
5 //! *producer* — `client::send_with_retry` — is still live and still records
6 //! every retry, and `client`'s own tests read it back through `snapshot`.
7 //! The read surface (`snapshot`, the countdown, the banner fields) is
8 //! therefore test-gated (`cfg(any(test, feature = "test-support"))` where
9 //! possible, so the TUI's tests reach it through the `test-support`
10 //! feature; `cfg_attr(not(test))` allows where prod constructs); a renderer
11 //! restores it by dropping the gates. Give the banner a renderer, or
12 //! delete the producer too — but not half of it.
13 //!
14 //! The HTTP retry path in `client::send_with_retry` already times its
15 //! waits and knows the error category. This module gives the TUI a way
16 //! to observe that state — `start`, `succeeded`, and `failed` flip a
17 //! global `RetryState` that the footer / status panel reads each frame.
18 //!
19 //! Why a process-wide global: the user-facing TUI runs as one engine
20 //! per process, and the only retry state we want to surface is the one
21 //! the user is staring at. Sub-agent retries in background tasks
22 //! deliberately do **not** light up the foreground banner — they're
23 //! supposed to be invisible. If a future feature ever needs per-engine
24 //! retry surfaces, swap this for an `Arc<RwLock<...>>` carried on the
25 //! `EngineHandle`; the public API stays the same.
26 //!
27 //! The `Retry-After` pause is process-wide too, but keyed by provider scope:
28 //! every request to the rate-limited route waits it out, and requests to any
29 //! other route (another provider, a local runtime) do not.
30
31 use std::collections::HashMap;
32 use std::sync::{Mutex, OnceLock};
33 use std::time::{Duration, Instant};
34
35 /// Rate-limit pause deadlines, one per provider scope (see
36 /// [`note_rate_limit`]). A 429 from one provider must not stall requests to
37 /// another provider, a local runtime, or a sub-agent on a different route.
38 type RateLimitPauses = HashMap<String, Instant>;
39
40 /// One in-flight retry attempt. `deadline` is the wall-clock time the
41 /// next request will fire — the UI subtracts `Instant::now()` from it
42 /// to render a live countdown.
43 #[derive(Debug, Clone)]
44 #[cfg_attr(not(test), allow(dead_code))]
45 pub struct RetryBanner {
46 /// 1-indexed retry attempt number (the first retry is attempt 1).
47 pub attempt: u32,
48 /// Time at which the next request will be sent.
49 pub deadline: Instant,
50 /// Short human-readable reason ("rate limited", "server error", …).
51 pub reason: String,
52 }
53
54 /// Snapshot of the retry surface for the UI to render.
55 #[derive(Debug, Clone, Default)]
56 pub enum RetryState {
57 /// No retry in flight. Banner hidden.
58 #[default]
59 Idle,
60 /// A request is sleeping before retrying. Show countdown banner.
61 Active(#[cfg_attr(not(test), allow(dead_code))] RetryBanner),
62 /// All retries exhausted; show failure row until the next turn
63 /// starts. `since` records when the row was set so a future polish
64 /// pass can age it out automatically; today the engine clears it on
65 /// `TurnStarted`.
66 Failed { reason: String, since: Instant },
67 }
68
69 impl RetryState {
70 /// Wall-clock seconds remaining on the active banner, or `None` if
71 /// not active. Saturates at zero — the renderer should treat any
72 /// negative remaining as "firing now".
73 #[cfg(any(test, feature = "test-support"))]
74 #[must_use]
75 pub fn seconds_remaining(&self) -> Option<u64> {
76 match self {
77 Self::Active(banner) => Some(
78 banner
79 .deadline
80 .saturating_duration_since(Instant::now())
81 .as_secs(),
82 ),
83 _ => None,
84 }
85 }
86
87 /// Whether the failure row should still be shown. Mirrors the
88 /// "until next turn" rule in the issue spec; the engine clears it
89 /// explicitly via [`clear`] on `TurnStarted`.
90 #[cfg(any(test, feature = "test-support"))]
91 #[must_use]
92 pub fn is_failed(&self) -> bool {
93 matches!(self, Self::Failed { .. })
94 }
95 }
96
97 /// Lazy-init the cell on first read so callers don't have to initialize
98 /// process-wide state at boot.
99 #[cfg(not(any(test, all(feature = "test-thread-scoped-state", debug_assertions))))]
100 fn with_state<R>(f: impl FnOnce(&mut RetryState) -> R) -> R {
101 static STATE: OnceLock<Mutex<RetryState>> = OnceLock::new();
102 let mut state = STATE
103 .get_or_init(|| Mutex::new(RetryState::Idle))
104 .lock()
105 .unwrap_or_else(|error| error.into_inner());
106 f(&mut state)
107 }
108
109 #[cfg(not(any(test, all(feature = "test-thread-scoped-state", debug_assertions))))]
110 fn with_rate_limit<R>(f: impl FnOnce(&mut RateLimitPauses) -> R) -> R {
111 static STATE: OnceLock<Mutex<RateLimitPauses>> = OnceLock::new();
112 let mut state = STATE
113 .get_or_init(|| Mutex::new(HashMap::new()))
114 .lock()
115 .unwrap_or_else(|error| error.into_inner());
116 f(&mut state)
117 }
118
119 /// Under test, this state is per-thread.
120 ///
121 /// Production has exactly one foreground engine per process, so a global is the
122 /// right model there. The test harness does not: retry state is written as a
123 /// side effect of *any* client request, by production code that has no test
124 /// guard to take, so a real request in one test could publish a banner or a
125 /// provider pause into another test's assertions. Scoping by thread removes the
126 /// race at its source rather than asking every future test that happens to
127 /// perform HTTP to remember a lock.
128 ///
129 /// This changes behavior (a `Retry-After` pause no longer holds across tokio
130 /// worker threads), so it is *not* part of `test-support`, which only adds
131 /// helpers. The TUI's tests opt in with `test-thread-scoped-state`, and even
132 /// then only a build with debug assertions gets it: an optimized build, e.g.
133 /// `cargo build --release --workspace --all-features`, keeps the process-wide
134 /// pause. Cargo unifies features across one build, so under
135 /// `cargo test --workspace` (or a debug `--all-features` build) the debug
136 /// `codewhale` binary that integration tests spawn gets the per-thread pause
137 /// too; release and optimized binaries never do.
138 #[cfg(any(test, all(feature = "test-thread-scoped-state", debug_assertions)))]
139 fn with_state<R>(f: impl FnOnce(&mut RetryState) -> R) -> R {
140 #[allow(clippy::type_complexity)]
141 static STATE: OnceLock<Mutex<std::collections::HashMap<std::thread::ThreadId, RetryState>>> =
142 OnceLock::new();
143 let mut by_thread = STATE
144 .get_or_init(|| Mutex::new(std::collections::HashMap::new()))
145 .lock()
146 .unwrap_or_else(|error| error.into_inner());
147 f(by_thread
148 .entry(std::thread::current().id())
149 .or_insert(RetryState::Idle))
150 }
151
152 #[cfg(any(test, all(feature = "test-thread-scoped-state", debug_assertions)))]
153 fn with_rate_limit<R>(f: impl FnOnce(&mut RateLimitPauses) -> R) -> R {
154 static STATE: OnceLock<Mutex<HashMap<std::thread::ThreadId, RateLimitPauses>>> =
155 OnceLock::new();
156 let mut by_thread = STATE
157 .get_or_init(|| Mutex::new(HashMap::new()))
158 .lock()
159 .unwrap_or_else(|error| error.into_inner());
160 f(by_thread.entry(std::thread::current().id()).or_default())
161 }
162
163 /// Read snapshot for renderers. No production renderer exists since the
164 /// legacy footer went away; `client` retry tests are the only readers.
165 #[cfg(any(test, feature = "test-support"))]
166 #[must_use]
167 pub fn snapshot() -> RetryState {
168 with_state(|state| state.clone())
169 }
170
171 /// Extend the rate-limit pause window for one provider `scope` (the caller's
172 /// route identity and host). This is separate from the footer banner so one
173 /// successful concurrent request cannot clear another request's active
174 /// `Retry-After` window, and it is keyed so a 429 from one provider never
175 /// pauses requests to another. Expired scopes are dropped here, so the map
176 /// holds at most the providers currently rate limited.
177 pub fn note_rate_limit(scope: &str, delay: Duration) {
178 let now = Instant::now();
179 let deadline = now + delay;
180 with_rate_limit(|pauses| {
181 pauses.retain(|_, existing| *existing > now);
182 let current = pauses.entry(scope.to_string()).or_insert(deadline);
183 if *current < deadline {
184 *current = deadline;
185 }
186 });
187 }
188
189 /// Remaining rate-limit pause for `scope`, if any.
190 #[must_use]
191 pub fn rate_limit_remaining(scope: &str) -> Option<Duration> {
192 let now = Instant::now();
193 with_rate_limit(|pauses| match pauses.get(scope).copied() {
194 Some(deadline) if deadline > now => Some(deadline.duration_since(now)),
195 Some(_) => {
196 pauses.remove(scope);
197 None
198 }
199 None => None,
200 })
201 }
202
203 /// Mark an in-flight retry. `attempt` is the number of the *upcoming*
204 /// retry (1 for the first); `delay` is how long the client will sleep
205 /// before firing.
206 pub fn start(attempt: u32, delay: Duration, reason: impl Into<String>) {
207 let banner = RetryBanner {
208 attempt,
209 deadline: Instant::now() + delay,
210 reason: reason.into(),
211 };
212 with_state(|state| *state = RetryState::Active(banner));
213 }
214
215 /// Mark the retry chain as having succeeded. Hides the banner.
216 pub fn succeeded() {
217 with_state(|state| *state = RetryState::Idle);
218 }
219
220 /// Mark the retry chain as having exhausted retries. The renderer keeps
221 /// the failure row until [`clear`] (typically called on `TurnStarted`).
222 pub fn failed(reason: impl Into<String>) {
223 with_state(|state| {
224 *state = RetryState::Failed {
225 reason: reason.into(),
226 since: Instant::now(),
227 };
228 });
229 }
230
231 /// Reset to idle. Called on `TurnStarted` so the previous turn's
232 /// failure row doesn't bleed into the next turn.
233 pub fn clear() {
234 with_state(|state| *state = RetryState::Idle);
235 }
236
237 /// Drop every scope's rate-limit pause.
238 #[cfg(any(test, feature = "test-support"))]
239 pub fn clear_rate_limit() {
240 with_rate_limit(HashMap::clear);
241 }
242
243 /// Test helper: serialize tests that touch the global state so cargo's
244 /// parallel runner can't observe a torn read. The guard is exported so
245 /// tests in *other* modules (e.g. footer rendering tests) can hold the
246 /// same lock as the ones in `retry_status::tests`.
247 #[cfg(any(test, feature = "test-support"))]
248 pub fn test_guard() -> std::sync::MutexGuard<'static, ()> {
249 static GUARD: Mutex<()> = Mutex::new(());
250 GUARD.lock().unwrap_or_else(|e| e.into_inner())
251 }
252
253 #[cfg(test)]
254 mod tests {
255 use super::*;
256
257 /// Acquire the cross-module test guard from [`super::test_guard`] and
258 /// reset state to `Idle` before yielding to the test body.
259 fn setup() -> std::sync::MutexGuard<'static, ()> {
260 let g = test_guard();
261 clear();
262 clear_rate_limit();
263 g
264 }
265
266 #[test]
267 fn idle_by_default_after_clear() {
268 let _g = setup();
269 assert!(matches!(snapshot(), RetryState::Idle));
270 assert_eq!(snapshot().seconds_remaining(), None);
271 }
272
273 #[test]
274 fn start_then_succeeded_returns_to_idle() {
275 let _g = setup();
276 start(1, Duration::from_secs(5), "rate limited");
277 let s = snapshot();
278 assert!(matches!(s, RetryState::Active(_)));
279 let remaining = s.seconds_remaining().unwrap();
280 assert!(remaining <= 5, "{remaining}");
281 succeeded();
282 assert!(matches!(snapshot(), RetryState::Idle));
283 }
284
285 #[test]
286 fn failed_persists_until_clear() {
287 let _g = setup();
288 failed("upstream 500");
289 let s = snapshot();
290 assert!(s.is_failed());
291 if let RetryState::Failed { reason, .. } = s {
292 assert_eq!(reason, "upstream 500");
293 } else {
294 panic!("expected Failed");
295 }
296 clear();
297 assert!(matches!(snapshot(), RetryState::Idle));
298 }
299
300 #[test]
301 fn deadline_in_past_yields_zero_remaining() {
302 let _g = setup();
303 // Bypass `start` so we can plant a deadline already in the past.
304 with_state(|state| {
305 *state = RetryState::Active(RetryBanner {
306 attempt: 2,
307 deadline: Instant::now() - Duration::from_secs(1),
308 reason: "test".into(),
309 });
310 });
311 assert_eq!(snapshot().seconds_remaining(), Some(0));
312 clear();
313 }
314
315 #[test]
316 fn rate_limit_deadline_survives_banner_clear() {
317 let _g = setup();
318 note_rate_limit("provider-a", Duration::from_secs(5));
319 start(1, Duration::from_secs(5), "rate limited");
320 succeeded();
321 assert!(
322 rate_limit_remaining("provider-a").is_some(),
323 "provider rate limit pause must not be cleared by an unrelated success"
324 );
325 clear_rate_limit();
326 }
327
328 #[test]
329 fn one_providers_rate_limit_does_not_pause_another() {
330 let _g = setup();
331 note_rate_limit("provider-a@api.a.example", Duration::from_secs(30));
332 assert!(rate_limit_remaining("provider-a@api.a.example").is_some());
333 assert_eq!(rate_limit_remaining("ollama@127.0.0.1:11434"), None);
334
335 // A shorter Retry-After never shortens an existing window.
336 note_rate_limit("provider-a@api.a.example", Duration::from_secs(1));
337 assert!(
338 rate_limit_remaining("provider-a@api.a.example").unwrap() > Duration::from_secs(20)
339 );
340 clear_rate_limit();
341 assert_eq!(rate_limit_remaining("provider-a@api.a.example"), None);
342 }
343 }
344
344 lines RUST