| 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 |