返回 CodeWhale
governor.rs
根目录 / crates / tui / src / tools / subagent / governor.rs
1 //! Rate-limit aware adaptive scheduling for sub-agent fan-out ("swarm mode").
2 //!
3 //! A swarm can launch an unbounded number of sub-agents against one shared
4 //! LLM provider, so parallel 429s are the steady state rather than an edge
5 //! case. This module gives the sub-agent module two cooperating pieces:
6 //!
7 //! 1. [`DynamicGate`] — a launch gate with a *dynamically adjustable
8 //! capacity*. The previous gate was a `tokio::sync::Semaphore`, whose
9 //! capacity is fixed at construction; the only way to "shrink" it was to
10 //! replace the `Arc`, which silently fails while any child still holds a
11 //! permit (that is exactly why `update_runtime_limits` only applied
12 //! launch-concurrency changes when no sub-agent was running). A
13 //! custom gate can drop its capacity below the number of active holders:
14 //! existing children keep running to completion, while new admissions
15 //! block until `active < capacity`.
16 //!
17 //! 2. [`RateLimitGovernor`] — a sliding-window observer fed by the sub-agent
18 //! LLM call path. Every rate-limited attempt and every successful attempt
19 //! is reported; when the recent failure rate crosses a threshold the
20 //! governor shrinks the gate (multiplicative decrease), and under a
21 //! sustained burst it pauses new admissions entirely. Sustained success
22 //! recovers capacity additively (AIMD), which converges without the
23 //! oscillation a symmetric controller would show.
24 //!
25 //! Retries themselves stay in the LLM call path (see
26 //! the canonical Engine model-step dispatch): the governor never
27 //! delays an in-flight call, it only decides whether *new* launches may be
28 //! admitted. `QuotaExhausted` is deliberately not reported — quota is a
29 //! billing condition, not a transient throttle, and must keep following the
30 //! existing fatal/checkpoint path.
31
32 use std::collections::VecDeque;
33 use std::sync::Mutex;
34 use std::time::{Duration, Instant};
35
36 use tokio::sync::oneshot;
37
38 /// Observation window for rate-limit events. Events older than this are
39 /// pruned on every governor interaction.
40 const RATE_LIMIT_WINDOW: Duration = Duration::from_secs(60);
41
42 /// Rate-limit events inside [`RATE_LIMIT_WINDOW`] at which the governor
43 /// starts shrinking launch concurrency (AIMD multiplicative decrease).
44 const THROTTLE_EVENT_THRESHOLD: usize = 2;
45
46 /// Recent rate-limit *ratio* (limited attempts / attempts) at which the
47 /// governor also shrinks launch concurrency, even below the absolute count
48 /// threshold. With very few in-flight calls, two 429s may be 100% of traffic.
49 const THROTTLE_RATIO_THRESHOLD: f64 = 0.3;
50
51 /// Rate-limit events inside the window at which the governor pauses new
52 /// admissions entirely (gate capacity 0). Held permits are unaffected.
53 const PAUSE_EVENT_THRESHOLD: usize = 4;
54
55 /// Successful attempts required to add one unit of launch capacity back
56 /// (AIMD additive increase). Successes are counted per gate-holder, so a
57 /// shrunken fleet still recovers at a controlled pace.
58 const SUCCESS_PER_INCREASE_STEP: u32 = 3;
59
60 /// Full-jitter exponential backoff for a rate-limited sub-agent API attempt
61 /// (`retry_number` is 1-based): the raw backoff is
62 /// `initial * 2^(n-1)` capped at [`RATE_LIMIT_MAX_BACKOFF`], and the actual
63 /// delay is drawn uniformly from `[0, backoff)` (AWS "full jitter"). Full
64 /// jitter de-synchronizes a fan-out of children that were all 429'd by the
65 /// same provider response; the cap keeps a retrying child inside its
66 /// wall-time budget instead of giving up.
67 const RATE_LIMIT_MAX_BACKOFF: Duration = Duration::from_secs(120);
68 const RATE_LIMIT_BACKOFF_JITTER_FACTOR: f64 = 1.0; // full jitter
69
70 /// Uniformly random factor in `[0, 1)` derived from UUID v4 entropy, the
71 /// same idiom as `llm_client::RetryConfig::delay_for_attempt`.
72 fn random_unit_factor() -> f64 {
73 let bytes = *uuid::Uuid::new_v4().as_bytes();
74 let sample = u16::from_le_bytes([bytes[0], bytes[1]]);
75 f64::from(sample) / f64::from(u16::MAX)
76 }
77
78 /// Raw (pre-jitter) exponential backoff for a rate-limited attempt.
79 fn rate_limit_backoff_base(retry_number: u32) -> Duration {
80 let multiplier = 1u32
81 .checked_shl(retry_number.saturating_sub(1))
82 .unwrap_or(u32::MAX);
83 Duration::from_millis(250)
84 .saturating_mul(multiplier)
85 .min(RATE_LIMIT_MAX_BACKOFF)
86 }
87
88 /// Full-jitter retry delay for a rate-limited attempt.
89 pub(crate) fn rate_limit_retry_delay(retry_number: u32) -> Duration {
90 let base = rate_limit_backoff_base(retry_number).as_secs_f64();
91 // Full jitter: uniform in [0, base). Reaching exactly `base` is fine and
92 // only sharpens de-synchronization; the draw can never exceed it.
93 Duration::from_secs_f64(base * (1.0 - RATE_LIMIT_BACKOFF_JITTER_FACTOR * random_unit_factor()))
94 }
95
96 // === DynamicGate ===
97
98 #[derive(Debug)]
99 struct GateWaiter {
100 sender: oneshot::Sender<DynamicGatePermit>,
101 }
102
103 #[derive(Debug)]
104 struct GateInner {
105 capacity: usize,
106 active: usize,
107 waiters: VecDeque<GateWaiter>,
108 }
109
110 /// A launch gate with runtime-adjustable capacity (see module docs).
111 ///
112 /// `acquire` returns a [`DynamicGatePermit`] whose `Drop` releases the slot
113 /// and wakes one waiter. Reducing capacity below `active` is allowed: the
114 /// surplus holders finish naturally and no new permit is granted until the
115 /// active count drops under the new capacity.
116 ///
117 /// Waiters receive an *already granted* permit through a oneshot channel, so
118 /// a waiter future that is cancelled after the grant is dispatched simply
119 /// drops the permit, whose `Drop` hands the slot to the next waiter. (A
120 /// wake-and-recheck design would lose that wakeup — the cancelled waiter
121 /// never re-checks, and with no remaining holders there is no later release
122 /// to re-dispatch it.)
123 #[derive(Debug)]
124 pub(crate) struct DynamicGate {
125 inner: Mutex<GateInner>,
126 }
127
128 impl DynamicGate {
129 pub(crate) fn new(capacity: usize) -> Self {
130 Self {
131 inner: Mutex::new(GateInner {
132 capacity: capacity.max(1),
133 active: 0,
134 waiters: VecDeque::new(),
135 }),
136 }
137 }
138
139 pub(crate) fn capacity(&self) -> usize {
140 self.inner.lock().expect("launch gate poisoned").capacity
141 }
142
143 /// Free admission slots right now (`capacity - active`). Diagnostics and
144 /// tests only; racy by design.
145 #[cfg(test)]
146 pub(crate) fn available_permits(&self) -> usize {
147 let inner = self.inner.lock().expect("launch gate poisoned");
148 inner.capacity.saturating_sub(inner.active)
149 }
150
151 /// Adjust the gate capacity. Raising it grants queued waiters the new
152 /// headroom immediately; lowering it simply stops new admissions until
153 /// the active count drains below the new capacity.
154 pub(crate) fn set_capacity(self: &std::sync::Arc<Self>, capacity: usize) {
155 let mut inner = self.inner.lock().expect("launch gate poisoned");
156 inner.capacity = capacity;
157 Self::wake_locked(self, &mut inner);
158 }
159
160 /// Grant queued waiters while there is headroom. Called with the lock
161 /// held; each waiter receives an already-counted permit, so a cancelled
162 /// receiver's permit is disarmed (never `Drop`ped under the lock) and the
163 /// slot flows to the next waiter.
164 fn wake_locked(gate: &std::sync::Arc<Self>, inner: &mut GateInner) {
165 while inner.active < inner.capacity {
166 let Some(waiter) = inner.waiters.pop_front() else {
167 break;
168 };
169 let permit = DynamicGatePermit {
170 gate: Some(std::sync::Arc::clone(gate)),
171 };
172 match waiter.sender.send(permit) {
173 Ok(()) => inner.active += 1,
174 Err(mut returned) => {
175 // The waiter future was cancelled before receiving the
176 // grant. Disarm instead of dropping: `Drop` would call
177 // `release()` and re-enter the lock we are holding.
178 let _ = returned.disarm();
179 }
180 }
181 }
182 }
183
184 fn release(self: &std::sync::Arc<Self>) {
185 let mut inner = self.inner.lock().expect("launch gate poisoned");
186 inner.active = inner.active.saturating_sub(1);
187 Self::wake_locked(self, &mut inner);
188 }
189
190 /// Try to acquire a permit without waiting.
191 pub(crate) fn try_acquire(self: &std::sync::Arc<Self>) -> Option<DynamicGatePermit> {
192 let mut inner = self.inner.lock().expect("launch gate poisoned");
193 (inner.active < inner.capacity).then(|| {
194 inner.active += 1;
195 DynamicGatePermit {
196 gate: Some(std::sync::Arc::clone(self)),
197 }
198 })
199 }
200
201 /// Acquire a permit, waiting until capacity is available. Cancellation
202 /// safe: a dropped future either leaves a stale queue entry (skipped and
203 /// disarmed by the granter) or drops an already-dispatched permit (whose
204 /// `Drop` re-releases the slot).
205 pub(crate) async fn acquire(self: &std::sync::Arc<Self>) -> DynamicGatePermit {
206 loop {
207 let rx = {
208 let mut inner = self.inner.lock().expect("launch gate poisoned");
209 if inner.active < inner.capacity {
210 inner.active += 1;
211 return DynamicGatePermit {
212 gate: Some(std::sync::Arc::clone(self)),
213 };
214 }
215 let (tx, rx) = oneshot::channel();
216 inner.waiters.push_back(GateWaiter { sender: tx });
217 rx
218 };
219 // Defensive: a failed receive requires the queued sender to be
220 // dropped without sending — which requires the gate itself to be
221 // dropped, impossible while this future holds an `Arc` to it.
222 // Loop anyway so a future refactor that breaks that invariant
223 // degrades to re-queueing instead of unwrapping.
224 if let Ok(permit) = rx.await {
225 return permit;
226 }
227 }
228 }
229 }
230
231 /// One held launch slot. Released on drop.
232 ///
233 /// The gate is an `Option` so the wake path can disarm a permit whose
234 /// receiver vanished without running `Drop` (which would re-enter the locked
235 /// `release()`).
236 pub(crate) struct DynamicGatePermit {
237 gate: Option<std::sync::Arc<DynamicGate>>,
238 }
239
240 impl DynamicGatePermit {
241 fn disarm(&mut self) -> Option<std::sync::Arc<DynamicGate>> {
242 self.gate.take()
243 }
244 }
245
246 impl std::fmt::Debug for DynamicGatePermit {
247 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
248 f.debug_struct("DynamicGatePermit").finish()
249 }
250 }
251
252 impl Drop for DynamicGatePermit {
253 fn drop(&mut self) {
254 if let Some(gate) = self.gate.take() {
255 gate.release();
256 }
257 }
258 }
259
260 // === RateLimitGovernor ===
261
262 #[derive(Debug)]
263 struct GovernorState {
264 /// Ceiling additive increase may climb to (configured launch
265 /// concurrency).
266 max_capacity: usize,
267 /// Timestamps of rate-limited attempts inside the window.
268 limited: VecDeque<Instant>,
269 /// Timestamps of all reported attempts inside the window (successes and
270 /// rate limits) — the denominator of the recent rate-limit ratio.
271 attempts: VecDeque<Instant>,
272 consecutive_successes: u32,
273 paused: bool,
274 }
275
276 /// Rate-limit aware scheduler over a [`DynamicGate`] (see module docs).
277 #[derive(Debug)]
278 pub(crate) struct RateLimitGovernor {
279 gate: std::sync::Arc<DynamicGate>,
280 state: Mutex<GovernorState>,
281 }
282
283 impl RateLimitGovernor {
284 pub(crate) fn new(max_capacity: usize) -> (std::sync::Arc<Self>, std::sync::Arc<DynamicGate>) {
285 let gate = std::sync::Arc::new(DynamicGate::new(max_capacity.max(1)));
286 let governor = std::sync::Arc::new(Self {
287 gate: std::sync::Arc::clone(&gate),
288 state: Mutex::new(GovernorState {
289 max_capacity: max_capacity.max(1),
290 limited: VecDeque::new(),
291 attempts: VecDeque::new(),
292 consecutive_successes: 0,
293 paused: false,
294 }),
295 });
296 (governor, gate)
297 }
298
299 /// The governor's launch gate. `SubAgentManager` hands this to spawned
300 /// tasks in place of the old fixed `Semaphore`. (Directly exercised by
301 /// governor unit tests.)
302 #[cfg(test)]
303 pub(crate) fn gate(&self) -> std::sync::Arc<DynamicGate> {
304 std::sync::Arc::clone(&self.gate)
305 }
306
307 /// Apply a new configured launch capacity: the AIMD ceiling and the gate
308 /// capacity while not throttled. Applies to the live gate immediately
309 /// (raising and lowering alike) unless the governor is paused — a pause
310 /// keeps capacity 0 until recovery, so an external limit change cannot
311 /// silently lift a rate-limit pause.
312 pub(crate) fn set_max_capacity(&self, max_capacity: usize) {
313 let mut state = self.state.lock().expect("rate limit governor poisoned");
314 state.max_capacity = max_capacity.max(1);
315 if !state.paused {
316 self.gate.set_capacity(state.max_capacity);
317 }
318 }
319
320 fn prune(state: &mut GovernorState, now: Instant) {
321 while state
322 .limited
323 .front()
324 .is_some_and(|at| now.duration_since(*at) > RATE_LIMIT_WINDOW)
325 {
326 state.limited.pop_front();
327 }
328 while state
329 .attempts
330 .front()
331 .is_some_and(|at| now.duration_since(*at) > RATE_LIMIT_WINDOW)
332 {
333 state.attempts.pop_front();
334 }
335 }
336
337 /// Report that a sub-agent LLM attempt is starting. Contributes to the
338 /// recent-attempt denominator for the ratio heuristic.
339 pub(crate) fn record_attempt(&self, now: Instant) {
340 let mut state = self.state.lock().expect("rate limit governor poisoned");
341 Self::prune(&mut state, now);
342 state.attempts.push_back(now);
343 }
344
345 /// Lift a pause whose rate-limit events have all aged out of the window,
346 /// resuming at a conservative quarter of the configured capacity so
347 /// additive increase climbs the rest of the way. Callers must hold the
348 /// state lock; `prune` first.
349 fn unpause_if_window_drained(&self, state: &mut GovernorState) {
350 if !state.paused || !state.limited.is_empty() {
351 return;
352 }
353 state.paused = false;
354 let capacity = (state.max_capacity / 4).max(1);
355 self.gate.set_capacity(capacity);
356 tracing::info!(
357 target: "subagent",
358 launch_capacity = capacity,
359 max_capacity = state.max_capacity,
360 "rate-limit governor resumed launches after window drained"
361 );
362 }
363
364 /// Time-driven recovery probe for queued launches. A pause is normally
365 /// lifted by a successful LLM attempt from an in-flight child, but if the
366 /// entire in-flight fleet finishes while 429 events are still inside the
367 /// window, no success ever arrives — without this probe the queue would
368 /// freeze until each queued child hits its wall-time deadline. Once every
369 /// limit event has aged out, the next probe resumes launches.
370 pub(crate) fn recover_if_window_drained(&self, now: Instant) {
371 let mut state = self.state.lock().expect("rate limit governor poisoned");
372 Self::prune(&mut state, now);
373 self.unpause_if_window_drained(&mut state);
374 }
375
376 /// Report a successful sub-agent LLM attempt. Drives AIMD additive
377 /// increase and clears the pause once the window has drained.
378 pub(crate) fn record_success(&self, now: Instant) {
379 let mut state = self.state.lock().expect("rate limit governor poisoned");
380 Self::prune(&mut state, now);
381 state.consecutive_successes = state.consecutive_successes.saturating_add(1);
382
383 self.unpause_if_window_drained(&mut state);
384
385 if !state.paused
386 && state.consecutive_successes >= SUCCESS_PER_INCREASE_STEP
387 && self.gate.capacity() < state.max_capacity
388 {
389 state.consecutive_successes = 0;
390 let capacity = (self.gate.capacity() + 1).min(state.max_capacity);
391 self.gate.set_capacity(capacity);
392 tracing::debug!(
393 target: "subagent",
394 launch_capacity = capacity,
395 "rate-limit governor additively increased launch capacity"
396 );
397 }
398 }
399
400 /// Report a rate-limited (429) sub-agent LLM attempt. May shrink or pause
401 /// the launch gate; never touches in-flight calls or retries.
402 pub(crate) fn record_rate_limited(&self, now: Instant) {
403 let mut state = self.state.lock().expect("rate limit governor poisoned");
404 Self::prune(&mut state, now);
405 state.limited.push_back(now);
406 // The denominator (`attempts`) already contains this attempt — the
407 // call path reports `record_attempt` before every LLM call, retries
408 // included. Pushing again would double-count failures and skew the
409 // ratio.
410 state.consecutive_successes = 0;
411
412 if state.paused {
413 return;
414 }
415
416 let events = state.limited.len();
417 let attempts = state.attempts.len().max(1);
418 let ratio = f64::from(events as u32) / f64::from(attempts as u32);
419
420 if events >= PAUSE_EVENT_THRESHOLD {
421 state.paused = true;
422 // Capacity 0 blocks all *new* admissions; children already holding
423 // permits keep running to completion.
424 self.gate.set_capacity(0);
425 tracing::warn!(
426 target: "subagent",
427 window_events = events,
428 window_attempts = attempts,
429 "rate-limit governor paused new sub-agent launches (sustained provider 429s); \
430 queued children wait for the window to drain"
431 );
432 return;
433 }
434
435 // The ratio heuristic only fires once the window has real volume
436 // (>= 2 observed attempts): with a single attempt every 429 is 100%
437 // and would shrink the gate on the first blip, fighting the absolute
438 // count threshold that is meant to own small-fleet behavior.
439 if events >= THROTTLE_EVENT_THRESHOLD
440 || (state.attempts.len() >= 2 && ratio > THROTTLE_RATIO_THRESHOLD)
441 {
442 let current = self.gate.capacity();
443 if current > 1 {
444 let capacity = (current / 2).max(1);
445 self.gate.set_capacity(capacity);
446 tracing::warn!(
447 target: "subagent",
448 window_events = events,
449 window_ratio = format!("{ratio:.2}"),
450 previous_capacity = current,
451 launch_capacity = capacity,
452 "rate-limit governor multiplicatively decreased launch capacity"
453 );
454 }
455 }
456 }
457
458 /// Whether new launches are currently paused because of sustained 429s.
459 pub(crate) fn is_paused(&self, now: Instant) -> bool {
460 let mut state = self.state.lock().expect("rate limit governor poisoned");
461 Self::prune(&mut state, now);
462 state.paused
463 }
464
465 /// Observability snapshot: gate capacity, window limit events, paused.
466 /// Feeds the fleet throttling line (addendum F5) and tests.
467 pub(crate) fn snapshot(&self, now: Instant) -> GovernorSnapshot {
468 let mut state = self.state.lock().expect("rate limit governor poisoned");
469 Self::prune(&mut state, now);
470 GovernorSnapshot {
471 launch_capacity: self.gate.capacity(),
472 max_capacity: state.max_capacity,
473 window_limited: state.limited.len(),
474 window_attempts: state.attempts.len(),
475 paused: state.paused,
476 }
477 }
478 }
479
480 /// Point-in-time view of the governor for status surfaces and tests.
481 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
482 pub(crate) struct GovernorSnapshot {
483 pub(crate) launch_capacity: usize,
484 pub(crate) max_capacity: usize,
485 pub(crate) window_limited: usize,
486 pub(crate) window_attempts: usize,
487 pub(crate) paused: bool,
488 }
489
490 impl GovernorSnapshot {
491 /// One line for the fleet header and queued rows (addendum F5), or `None`
492 /// while launches run at the full configured concurrency.
493 #[must_use]
494 pub(crate) fn status_line(&self) -> Option<String> {
495 let window = RATE_LIMIT_WINDOW.as_secs();
496 if self.paused {
497 return Some(format!(
498 "launches paused after {} provider rate limit(s) in the last {window}s",
499 self.window_limited
500 ));
501 }
502 (self.launch_capacity < self.max_capacity).then(|| {
503 format!(
504 "launch slots throttled to {}/{} after {} provider rate limit(s) in the last {window}s",
505 self.launch_capacity, self.max_capacity, self.window_limited
506 )
507 })
508 }
509 }
510
511 #[cfg(test)]
512 mod tests {
513 use super::*;
514
515 fn ms(n: u64) -> Duration {
516 Duration::from_millis(n)
517 }
518
519 #[test]
520 fn window_counts_and_prunes_events() {
521 let (governor, _gate) = RateLimitGovernor::new(4);
522 let t0 = Instant::now();
523 for i in 0..5 {
524 governor.record_attempt(t0 + ms(i * 10));
525 governor.record_rate_limited(t0 + ms(i * 10));
526 }
527 let snap = governor.snapshot(t0 + ms(60));
528 assert_eq!(snap.window_limited, 5);
529 assert_eq!(snap.window_attempts, 5);
530
531 // Events older than the 60s window drop out (strictly past the
532 // window edge: the newest event is at t0+40ms).
533 let snap = governor.snapshot(t0 + RATE_LIMIT_WINDOW + ms(50));
534 assert_eq!(snap.window_limited, 0);
535 assert_eq!(snap.window_attempts, 0);
536 }
537
538 #[test]
539 fn multiplicative_decrease_halves_capacity_on_threshold() {
540 let (governor, _gate) = RateLimitGovernor::new(8);
541 let t0 = Instant::now();
542 // First event: below both thresholds, no change.
543 governor.record_attempt(t0);
544 governor.record_rate_limited(t0);
545 assert_eq!(governor.snapshot(t0).launch_capacity, 8);
546 // Second event: hits the count threshold, halve.
547 governor.record_attempt(t0 + ms(1));
548 governor.record_rate_limited(t0 + ms(1));
549 assert_eq!(governor.snapshot(t0).launch_capacity, 4);
550 // Third: halve again.
551 governor.record_attempt(t0 + ms(2));
552 governor.record_rate_limited(t0 + ms(2));
553 assert_eq!(governor.snapshot(t0).launch_capacity, 2);
554 // Fourth: hits the pause threshold.
555 governor.record_attempt(t0 + ms(3));
556 governor.record_rate_limited(t0 + ms(3));
557 let snap = governor.snapshot(t0);
558 assert!(snap.paused);
559 }
560
561 #[test]
562 fn status_line_names_throttle_and_pause_only_when_active() {
563 let (governor, _gate) = RateLimitGovernor::new(8);
564 let t0 = Instant::now();
565 assert_eq!(governor.snapshot(t0).status_line(), None);
566 for i in 0..2 {
567 governor.record_attempt(t0 + ms(i));
568 governor.record_rate_limited(t0 + ms(i));
569 }
570 let throttled = governor
571 .snapshot(t0 + ms(2))
572 .status_line()
573 .expect("throttled");
574 assert!(throttled.contains("throttled to 4/8"), "{throttled}");
575 for i in 2..4 {
576 governor.record_attempt(t0 + ms(i));
577 governor.record_rate_limited(t0 + ms(i));
578 }
579 let paused = governor.snapshot(t0 + ms(4)).status_line().expect("paused");
580 assert!(paused.starts_with("launches paused"), "{paused}");
581 }
582
583 #[test]
584 fn ratio_threshold_triggers_decrease_even_with_few_events() {
585 let (governor, _gate) = RateLimitGovernor::new(8);
586 let t0 = Instant::now();
587 // One success then one 429: the absolute event count is below the
588 // threshold, but the 50% limit ratio must still shrink the gate.
589 governor.record_attempt(t0);
590 governor.record_success(t0);
591 governor.record_attempt(t0 + ms(1));
592 governor.record_rate_limited(t0 + ms(1));
593 assert!(
594 governor.snapshot(t0 + ms(2)).launch_capacity < 8,
595 "50% limit ratio should trigger a decrease"
596 );
597 }
598
599 #[test]
600 fn additive_increase_recovers_capacity_gradually() {
601 let (governor, _gate) = RateLimitGovernor::new(8);
602 let t0 = Instant::now();
603 // Drive capacity down to 4 via two events.
604 governor.record_attempt(t0);
605 governor.record_rate_limited(t0);
606 governor.record_attempt(t0 + ms(1));
607 governor.record_rate_limited(t0 + ms(1));
608 assert_eq!(governor.snapshot(t0).launch_capacity, 4);
609
610 // Three consecutive successes add exactly one unit of capacity.
611 for i in 0..3u32 {
612 governor.record_attempt(t0 + ms(10 + u64::from(i)));
613 governor.record_success(t0 + ms(10 + u64::from(i)));
614 }
615 assert_eq!(governor.snapshot(t0 + ms(20)).launch_capacity, 5);
616 for i in 0..3u32 {
617 governor.record_attempt(t0 + ms(30 + u64::from(i)));
618 governor.record_success(t0 + ms(30 + u64::from(i)));
619 }
620 assert_eq!(governor.snapshot(t0 + ms(40)).launch_capacity, 6);
621
622 // A rate limit resets the success streak.
623 governor.record_attempt(t0 + ms(50));
624 governor.record_rate_limited(t0 + ms(50));
625 for i in 0..2u32 {
626 governor.record_attempt(t0 + ms(60 + u64::from(i)));
627 governor.record_success(t0 + ms(60 + u64::from(i)));
628 }
629 governor.record_attempt(t0 + ms(80));
630 governor.record_success(t0 + ms(80));
631 // 2 successes before the limit + 1 after = 3 successes, but the limit
632 // reset the streak, and the third event in the window halved again
633 // (6 -> 3) before successes could climb.
634 assert!(governor.snapshot(t0 + ms(90)).launch_capacity <= 6);
635 }
636
637 #[test]
638 fn pause_releases_only_after_window_drains() {
639 let (governor, gate) = RateLimitGovernor::new(8);
640 let t0 = Instant::now();
641 for i in 0..4 {
642 governor.record_attempt(t0 + ms(i));
643 governor.record_rate_limited(t0 + ms(i));
644 }
645 assert!(governor.is_paused(t0 + ms(10)));
646 assert_eq!(governor.snapshot(t0 + ms(10)).launch_capacity, 0);
647
648 // Successes before the window drains do NOT unpause.
649 governor.record_success(t0 + ms(20));
650 assert!(governor.is_paused(t0 + ms(30)));
651
652 // Once every limit event ages out, the next success resumes at a
653 // quarter of capacity.
654 let late = t0 + RATE_LIMIT_WINDOW + ms(10);
655 governor.record_success(late);
656 assert!(!governor.is_paused(late));
657 assert_eq!(governor.snapshot(late).launch_capacity, 2);
658 assert_eq!(gate.capacity(), 2);
659 }
660
661 #[test]
662 fn capacity_increase_is_capped_at_max() {
663 let (governor, _gate) = RateLimitGovernor::new(2);
664 let t0 = Instant::now();
665 for i in 0..12u32 {
666 governor.record_attempt(t0 + ms(u64::from(i)));
667 governor.record_success(t0 + ms(u64::from(i)));
668 }
669 assert_eq!(governor.snapshot(t0).launch_capacity, 2);
670 }
671
672 #[test]
673 fn gate_blocks_when_full_and_releases_on_drop() {
674 let (governor, gate) = RateLimitGovernor::new(1);
675 let rt = tokio::runtime::Builder::new_current_thread()
676 .enable_time()
677 .build()
678 .expect("test runtime");
679 rt.block_on(async move {
680 let first = governor.gate().try_acquire().expect("first permit");
681 assert!(gate.try_acquire().is_none(), "capacity 1 must be full");
682
683 let g2 = std::sync::Arc::clone(&gate);
684 let waiter = tokio::spawn(async move { g2.acquire().await });
685
686 // Waiter stays blocked while the first permit is held.
687 tokio::time::sleep(ms(20)).await;
688 assert!(!waiter.is_finished());
689
690 drop(first);
691 let _second = waiter.await.expect("waiter task");
692 });
693 }
694
695 #[test]
696 fn gate_set_capacity_shrinks_below_active_and_re_admits_later() {
697 let (governor, gate) = RateLimitGovernor::new(4);
698 let rt = tokio::runtime::Builder::new_current_thread()
699 .enable_time()
700 .build()
701 .expect("test runtime");
702 rt.block_on(async move {
703 let mut held: Vec<_> = (0..4)
704 .map(|_| gate.try_acquire().expect("permit within capacity"))
705 .collect();
706 assert_eq!(gate.capacity(), 4);
707
708 // Shrink below the active count: no new permit is granted.
709 governor.gate().set_capacity(1);
710 assert_eq!(gate.capacity(), 1);
711 assert!(gate.try_acquire().is_none());
712
713 let g2 = std::sync::Arc::clone(&gate);
714 let waiter = tokio::spawn(async move { g2.acquire().await });
715 tokio::time::sleep(ms(20)).await;
716 assert!(!waiter.is_finished(), "must wait while active >= capacity");
717
718 // Releasing holders drains `active` toward the new capacity; the
719 // waiter is admitted only once every held permit is released
720 // (active 4 -> 0 < capacity 1).
721 drop(held.swap_remove(0));
722 drop(held.swap_remove(0));
723 drop(held.swap_remove(0));
724 drop(held);
725 let _permit = waiter.await.expect("waiter admitted after drain");
726 assert!(gate.try_acquire().is_none(), "capacity 1 is now full");
727 drop(_permit);
728 });
729 }
730
731 #[test]
732 fn rate_limit_retry_delay_is_full_jitter_within_base() {
733 for retry in 1..=12u32 {
734 let base = rate_limit_backoff_base(retry);
735 for _ in 0..64 {
736 let delay = rate_limit_retry_delay(retry);
737 assert!(delay <= base, "full jitter must not exceed the base");
738 }
739 }
740 // The cap holds for absurd retry numbers.
741 assert_eq!(rate_limit_backoff_base(40), RATE_LIMIT_MAX_BACKOFF);
742 }
743
744 /// A pause must lift via the time-driven probe even when no in-flight
745 /// child ever reports another success (the in-flight fleet drained before
746 /// the window did): otherwise queued children freeze until their
747 /// wall-time deadline.
748 #[test]
749 fn forkguard_rate_limit_governor_pauses_and_time_recovers_after_window_drains() {
750 let (governor, _gate) = RateLimitGovernor::new(8);
751 let t0 = Instant::now();
752 for i in 0..4 {
753 governor.record_attempt(t0 + ms(i));
754 governor.record_rate_limited(t0 + ms(i));
755 }
756 assert!(governor.is_paused(t0 + ms(10)));
757
758 // Probe while 429 events are still inside the window: stays paused.
759 governor.recover_if_window_drained(t0 + ms(20));
760 assert!(governor.is_paused(t0 + ms(30)));
761
762 // Once every limit event has aged out, the probe resumes launches at
763 // a quarter of the configured capacity — no success event required.
764 let late = t0 + RATE_LIMIT_WINDOW + ms(10);
765 governor.recover_if_window_drained(late);
766 assert!(!governor.is_paused(late));
767 assert_eq!(governor.snapshot(late).launch_capacity, 2);
768 }
769
770 /// A runtime launch-concurrency change must not silently lift a pause:
771 /// the gate stays at capacity 0 until the window drains, then resumes at
772 /// a quarter of the *new* configured capacity.
773 #[test]
774 fn forkguard_rate_limit_governor_limit_change_keeps_pause_capacity_zero() {
775 let (governor, gate) = RateLimitGovernor::new(8);
776 let t0 = Instant::now();
777 for i in 0..4 {
778 governor.record_attempt(t0 + ms(i));
779 governor.record_rate_limited(t0 + ms(i));
780 }
781 assert!(governor.is_paused(t0 + ms(1)));
782
783 governor.set_max_capacity(4);
784 assert_eq!(gate.capacity(), 0, "pause must keep capacity 0");
785
786 let late = t0 + RATE_LIMIT_WINDOW + ms(10);
787 governor.recover_if_window_drained(late);
788 assert_eq!(
789 gate.capacity(),
790 1,
791 "resume at a quarter of the new capacity"
792 );
793 }
794
795 /// A waiter cancelled *after* its grant was dispatched must not swallow
796 /// the slot: the permit is dropped with the cancelled future and its
797 /// `Drop` re-releases it for the next waiter.
798 #[test]
799 fn forkguard_dynamic_gate_redispatches_grant_of_cancelled_waiter() {
800 let (_governor, gate) = RateLimitGovernor::new(1);
801 let rt = tokio::runtime::Builder::new_current_thread()
802 .enable_time()
803 .build()
804 .expect("test runtime");
805 rt.block_on(async move {
806 let holder = gate.try_acquire().expect("holder");
807 let g2 = std::sync::Arc::clone(&gate);
808 let waiter = tokio::spawn(async move { g2.acquire().await });
809 tokio::time::sleep(ms(20)).await;
810 assert!(!waiter.is_finished(), "waiter must be queued");
811
812 // Releasing the holder dispatches the grant into the waiter's
813 // channel; on a current-thread runtime the waiter has not polled
814 // yet when we abort it, so the permit is dropped mid-flight.
815 drop(holder);
816 waiter.abort();
817 tokio::time::sleep(ms(20)).await;
818
819 assert!(
820 gate.try_acquire().is_some(),
821 "grant of cancelled waiter must be re-released, not leaked"
822 );
823 });
824 }
825
826 /// The mirror case of the redispatch test: a waiter cancelled *before*
827 /// its grant was dispatched leaves a stale queue entry with a dead
828 /// receiver. The granter must skip that entry — disarming the already
829 /// built permit instead of dropping it, which would re-enter the gate
830 /// lock held by `wake_locked` — and the slot must stay usable.
831 #[test]
832 fn forkguard_dynamic_gate_skips_stale_queued_waiter_without_leaking_slot() {
833 let (_governor, gate) = RateLimitGovernor::new(1);
834 let rt = tokio::runtime::Builder::new_current_thread()
835 .enable_time()
836 .build()
837 .expect("test runtime");
838 rt.block_on(async move {
839 let holder = gate.try_acquire().expect("holder");
840 let g2 = std::sync::Arc::clone(&gate);
841 let waiter = tokio::spawn(async move { g2.acquire().await });
842 tokio::time::sleep(ms(20)).await;
843 assert!(
844 !waiter.is_finished(),
845 "waiter must be queued behind the holder"
846 );
847
848 // Cancel while the gate is full: no grant was ever dispatched,
849 // so the stale entry stays queued with a dead receiver.
850 waiter.abort();
851 tokio::time::sleep(ms(20)).await;
852
853 // Releasing the holder runs the granter over the stale entry.
854 drop(holder);
855 assert_eq!(
856 gate.available_permits(),
857 1,
858 "cancelled queued waiter must neither swallow nor leak the slot"
859 );
860 let permit = gate
861 .try_acquire()
862 .expect("slot usable after the stale entry is skipped");
863 drop(permit);
864 });
865 }
866
867 /// Stress: concurrent acquire/release with aborts and capacity
868 /// oscillation through 0 (a pause). Whatever the interleaving, every
869 /// slot must come home — a lost wakeup or a leaked (never released)
870 /// permit leaves the gate short of full capacity at the end of a round,
871 /// and an over-granted permit keeps a slot alive after all owners are
872 /// gone. Both fail the drain assertion.
873 #[test]
874 fn forkguard_dynamic_gate_stress_drains_to_full_capacity_despite_aborts() {
875 let (_governor, gate) = RateLimitGovernor::new(4);
876 let rt = tokio::runtime::Builder::new_multi_thread()
877 .worker_threads(2)
878 .enable_time()
879 .build()
880 .expect("test runtime");
881 rt.block_on(async move {
882 for round in 0..24usize {
883 // Start every round with live headroom, then briefly drop to
884 // 0 mid-round on every third round: the pause case keeps a
885 // full queue parked while nothing holds a permit. Capacity 0
886 // is never left in place while joining — with nobody holding
887 // a permit a permanent 0 would deadlock the round by design,
888 // so the restore below is part of the scenario.
889 gate.set_capacity(1 + (round % 2));
890 let mut handles = Vec::new();
891 for i in 0..16u32 {
892 let g = std::sync::Arc::clone(&gate);
893 handles.push(tokio::spawn(async move {
894 let _permit = g.acquire().await;
895 tokio::time::sleep(ms(u64::from(i % 4))).await;
896 }));
897 }
898 // Abort every third task: some while still queued (stale
899 // queue entries), some already holding a permit (the
900 // drop-releases-and-rewakes path).
901 for handle in handles.iter().step_by(3) {
902 handle.abort();
903 }
904 if round % 3 == 0 {
905 gate.set_capacity(0);
906 tokio::time::sleep(ms(2)).await;
907 }
908 gate.set_capacity(4);
909 for handle in handles {
910 // A task that cannot finish inside the budget means a
911 // lost wakeup, a leaked permit, or a slot swallowed by a
912 // stale entry — fail the round instead of hanging.
913 tokio::time::timeout(ms(2000), handle)
914 .await
915 .expect("task must finish: stuck rounds mean lost wakeups or leaked slots")
916 .ok();
917 }
918 // Let straggler permit drops (cancelled waiter re-release)
919 // run before asserting the drain.
920 tokio::time::sleep(ms(5)).await;
921 assert_eq!(
922 gate.available_permits(),
923 4,
924 "round {round}: gate must drain to full capacity despite aborts and pauses"
925 );
926 }
927 });
928 }
929 }
930
930 lines RUST