返回 CodeWhale
turn_budget.rs
根目录 / crates / tui / src / core / engine / turn_budget.rs
1 //! Turn budgets shared by interactive hosts and headless execution.
2 //!
3 //! Model steps and the cumulative per-turn wall clock are uncapped by
4 //! default. Explicit positive limits still apply; neither a step counter nor
5 //! elapsed time is a measure of useful progress, and a long autonomous turn
6 //! should not stop at an arbitrary hour. Per-step stream budgets remain
7 //! finite: they bound one stuck request, not the whole turn.
8 //!
9 //! `EngineConfig` retains its integer representation for embedders:
10 //! `u32::MAX` represents no model-step limit. `TurnContext::step_limit`
11 //! resolves that representation to `None` before checking a ceiling or
12 //! emitting diagnostics. It is not a very large finite fallback.
13 //!
14 //! ## Honesty at the limit
15 //!
16 //! Hitting a budget is never a clean success. The step ceiling already ends
17 //! the turn as `TurnOutcomeStatus::Failed` with the limit named (or, when
18 //! the model still owes work, grants exactly one bounded final-report turn
19 //! first). [`TurnWallClock`] follows the same contract in `run_turn`.
20
21 use std::time::{Duration, Instant};
22
23 /// No model-step limit unless the caller configures one. This is the
24 /// compatibility representation, not a ceiling checked at `u32::MAX`.
25 pub const DEFAULT_MAX_MODEL_STEPS: u32 = u32::MAX;
26 /// Smallest accepted model-step ceiling. One step still lets the model
27 /// answer once.
28 pub const MIN_MAX_MODEL_STEPS: u32 = 1;
29 /// Largest explicitly configured model-step ceiling.
30 pub const MAX_MAX_MODEL_STEPS: u32 = 100_000;
31
32 /// No per-turn wall-clock limit unless the caller configures one.
33 ///
34 /// `Duration::MAX` is the representation, mirroring
35 /// [`DEFAULT_MAX_MODEL_STEPS`]: [`TurnWallClock::exhausted`] can never reach
36 /// it, and callers turning it into an absolute deadline must use
37 /// `checked_add` (see `exec_agent`). When configured, the budget is measured
38 /// across every model step of one turn, not per request, and time blocked on
39 /// a human approval decision is excluded (see
40 /// [`TurnWallClock::begin_human_wait`]).
41 pub const DEFAULT_TURN_WALL_CLOCK: Duration = Duration::MAX;
42 /// Smallest accepted per-turn wall-clock budget. Below this a single slow
43 /// reasoning request would trip the budget before it could finish.
44 pub const MIN_TURN_WALL_CLOCK_SECS: u64 = 30;
45 /// Largest accepted per-turn wall-clock budget (24 hours).
46 pub const MAX_TURN_WALL_CLOCK_SECS: u64 = 86_400;
47
48 /// Default per-step cap on accumulated streamed content, in bytes.
49 /// Preserves the pre-R1 hard-coded value; R1 only makes it overridable.
50 pub const DEFAULT_STREAM_MAX_CONTENT_BYTES: usize = super::streaming::STREAM_MAX_CONTENT_BYTES;
51 /// Smallest accepted per-step stream content cap (64 KiB).
52 pub const MIN_STREAM_MAX_CONTENT_BYTES: usize = 64 * 1024;
53 /// Largest accepted per-step stream content cap (512 MiB).
54 pub const MAX_STREAM_MAX_CONTENT_BYTES: usize = 512 * 1024 * 1024;
55
56 /// Default per-step cap on a single stream's wall-clock duration, in
57 /// seconds. Preserves the pre-R1 hard-coded value.
58 pub const DEFAULT_STREAM_MAX_DURATION_SECS: u64 = super::streaming::STREAM_MAX_DURATION_SECS;
59 /// Smallest accepted per-step stream duration cap.
60 pub const MIN_STREAM_MAX_DURATION_SECS: u64 = 10;
61 /// Largest accepted per-step stream duration cap (24 hours).
62 pub const MAX_STREAM_MAX_DURATION_SECS: u64 = 86_400;
63
64 /// Default whole-request resume budget after a failed stream (open failure,
65 /// dead stream, sleep, or mid-stream network drop). Preserves the
66 /// pre-#6700 hard-coded `MAX_STREAM_RETRIES`.
67 pub const DEFAULT_STREAM_MAX_RESUMES: u32 = super::streaming::MAX_STREAM_RETRIES;
68 /// Largest accepted resume budget. `0` is accepted and disables resumes.
69 pub const MAX_STREAM_MAX_RESUMES: u32 = 10;
70 /// Default in-stream transparent retry budget (nothing streamed yet).
71 /// Preserves the pre-#6700 hard-coded `MAX_TRANSPARENT_STREAM_RETRIES`.
72 pub const DEFAULT_STREAM_MAX_TRANSPARENT_RETRIES: u32 =
73 super::streaming::MAX_TRANSPARENT_STREAM_RETRIES;
74 /// Largest accepted transparent retry budget. `0` disables them.
75 pub const MAX_STREAM_MAX_TRANSPARENT_RETRIES: u32 = 10;
76 /// Default streak of recoverable stream errors tolerated in one stream.
77 /// Preserves the pre-#6700 hard-coded `MAX_STREAM_ERRORS_BEFORE_FAIL`.
78 pub const DEFAULT_STREAM_MAX_ERRORS: u32 = super::streaming::MAX_STREAM_ERRORS_BEFORE_FAIL;
79 /// Smallest accepted error streak: the first error ends the stream. `0` is
80 /// not a value here; like the other finite stream budgets it selects
81 /// [`DEFAULT_STREAM_MAX_ERRORS`].
82 pub const MIN_STREAM_MAX_ERRORS: u32 = 1;
83 /// Largest accepted error streak.
84 pub const MAX_STREAM_MAX_ERRORS: u32 = 50;
85
86 /// Stream-level retry budgets for one engine (#6700). Every field stays
87 /// finite; the defaults are the historical compiled-in values.
88 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
89 pub struct StreamRetryLimits {
90 /// Whole-request re-issues after a failed stream, shared by every
91 /// resume path and by stream-open failures (#6699).
92 pub max_resumes: u32,
93 /// In-stream re-requests while nothing has streamed yet (#103).
94 pub max_transparent_retries: u32,
95 /// Recoverable errors tolerated in one stream before it ends.
96 pub max_errors: u32,
97 }
98
99 impl Default for StreamRetryLimits {
100 fn default() -> Self {
101 resolve_stream_retry_limits(None, None, None)
102 }
103 }
104
105 /// Resolve configured stream retry budgets.
106 ///
107 /// `None` selects each default. The two retry counts accept `0` (no
108 /// retries) and clamp to their maximum. The error streak is a finite stream
109 /// budget like `stream_max_content_mb`: `0` selects its default, and other
110 /// values clamp to `MIN_STREAM_MAX_ERRORS..=MAX_STREAM_MAX_ERRORS`.
111 #[must_use]
112 pub fn resolve_stream_retry_limits(
113 max_resumes: Option<u32>,
114 max_transparent_retries: Option<u32>,
115 max_errors: Option<u32>,
116 ) -> StreamRetryLimits {
117 StreamRetryLimits {
118 max_resumes: max_resumes
119 .unwrap_or(DEFAULT_STREAM_MAX_RESUMES)
120 .min(MAX_STREAM_MAX_RESUMES),
121 max_transparent_retries: max_transparent_retries
122 .unwrap_or(DEFAULT_STREAM_MAX_TRANSPARENT_RETRIES)
123 .min(MAX_STREAM_MAX_TRANSPARENT_RETRIES),
124 max_errors: match max_errors {
125 None | Some(0) => DEFAULT_STREAM_MAX_ERRORS,
126 Some(value) => value.clamp(MIN_STREAM_MAX_ERRORS, MAX_STREAM_MAX_ERRORS),
127 },
128 }
129 }
130
131 /// Resolve a configured model-step ceiling.
132 ///
133 /// `None` and `0` select the uncapped default. Explicit positive values
134 /// clamp into `MIN_MAX_MODEL_STEPS..=MAX_MAX_MODEL_STEPS`. Resolve raw input
135 /// once: passing an already resolved default back as an explicit value
136 /// would incorrectly install the maximum configurable ceiling.
137 #[must_use]
138 pub fn resolve_max_model_steps(raw: Option<u32>) -> u32 {
139 match raw {
140 None | Some(0) => DEFAULT_MAX_MODEL_STEPS,
141 Some(value) => value.clamp(MIN_MAX_MODEL_STEPS, MAX_MAX_MODEL_STEPS),
142 }
143 }
144
145 /// Resolve a configured per-turn wall-clock budget.
146 ///
147 /// `None` and `0` select the uncapped default ([`DEFAULT_TURN_WALL_CLOCK`]).
148 /// Explicit positive values, in seconds, clamp into
149 /// `MIN_TURN_WALL_CLOCK_SECS..=MAX_TURN_WALL_CLOCK_SECS`.
150 #[must_use]
151 pub fn resolve_turn_wall_clock(raw: Option<u64>) -> Duration {
152 match raw {
153 None | Some(0) => DEFAULT_TURN_WALL_CLOCK,
154 Some(value) => {
155 Duration::from_secs(value.clamp(MIN_TURN_WALL_CLOCK_SECS, MAX_TURN_WALL_CLOCK_SECS))
156 }
157 }
158 }
159
160 /// Resolve a configured per-step stream content cap, given megabytes.
161 ///
162 /// `None` and `0` both resolve to [`DEFAULT_STREAM_MAX_CONTENT_BYTES`].
163 #[must_use]
164 pub fn resolve_stream_max_content_bytes(raw_mb: Option<u64>) -> usize {
165 match raw_mb {
166 None | Some(0) => DEFAULT_STREAM_MAX_CONTENT_BYTES,
167 Some(mb) => usize::try_from(mb.saturating_mul(1024 * 1024))
168 .unwrap_or(MAX_STREAM_MAX_CONTENT_BYTES)
169 .clamp(MIN_STREAM_MAX_CONTENT_BYTES, MAX_STREAM_MAX_CONTENT_BYTES),
170 }
171 }
172
173 /// Resolve a configured per-step stream duration cap, in seconds.
174 ///
175 /// `None` and `0` both resolve to [`DEFAULT_STREAM_MAX_DURATION_SECS`].
176 #[must_use]
177 pub fn resolve_stream_max_duration_secs(raw: Option<u64>) -> u64 {
178 match raw {
179 None | Some(0) => DEFAULT_STREAM_MAX_DURATION_SECS,
180 Some(value) => value.clamp(MIN_STREAM_MAX_DURATION_SECS, MAX_STREAM_MAX_DURATION_SECS),
181 }
182 }
183
184 /// Cumulative wall-clock budget for one turn.
185 ///
186 /// Started once at the top of `Engine::run_turn` and checked at the
187 /// provider-request boundary, so a turn that trips the budget stops before
188 /// authorizing another billable request and keeps every tool result already
189 /// in the transcript.
190 ///
191 /// Time blocked on a human approval decision is excluded: the budget bounds
192 /// what the agent spends on its own, not how long a person takes to answer.
193 /// Without that exclusion an approval prompt left open overnight would fail
194 /// the turn — and discard the work the user just approved — the moment they
195 /// came back.
196 #[derive(Debug)]
197 pub(crate) struct TurnWallClock {
198 budget: Duration,
199 started_at: Instant,
200 /// Total time already excluded because the turn was blocked on a human.
201 excluded: Duration,
202 /// Set while currently blocked on a human decision.
203 blocked_since: Option<Instant>,
204 }
205
206 impl TurnWallClock {
207 /// Start a fresh budget. A zero budget is legal here (and only here):
208 /// it is how tests assert the stop path without sleeping. Configuration
209 /// never produces one — [`resolve_turn_wall_clock`] reads `0` as no limit.
210 pub(crate) fn start(budget: Duration) -> Self {
211 Self {
212 budget,
213 started_at: Instant::now(),
214 excluded: Duration::ZERO,
215 blocked_since: None,
216 }
217 }
218
219 /// The budget this clock was started with.
220 pub(crate) fn budget(&self) -> Duration {
221 self.budget
222 }
223
224 /// Wall-clock time this turn has spent on its own work, excluding time
225 /// blocked on a human decision.
226 pub(crate) fn spent(&self) -> Duration {
227 let blocked_now = self
228 .blocked_since
229 .map_or(Duration::ZERO, |since| since.elapsed());
230 self.started_at
231 .elapsed()
232 .saturating_sub(self.excluded)
233 .saturating_sub(blocked_now)
234 }
235
236 /// Whether the cumulative budget is spent.
237 pub(crate) fn exhausted(&self) -> bool {
238 self.spent() >= self.budget
239 }
240
241 /// Stop counting: the turn is now waiting on a human decision.
242 /// Idempotent — a second call while already blocked does nothing, so a
243 /// nested or re-entered approval cannot double-exclude.
244 pub(crate) fn begin_human_wait(&mut self) {
245 if self.blocked_since.is_none() {
246 self.blocked_since = Some(Instant::now());
247 }
248 }
249
250 /// Resume counting after a human decision, banking the blocked time.
251 pub(crate) fn end_human_wait(&mut self) {
252 if let Some(since) = self.blocked_since.take() {
253 self.excluded = self.excluded.saturating_add(since.elapsed());
254 }
255 }
256
257 /// Test-only: pretend `elapsed` more wall-clock time has passed, so the
258 /// exhaustion path can be exercised without sleeping. An in-progress
259 /// human wait is rewound too — otherwise the simulated time would count
260 /// as agent-owned work that never actually happened.
261 #[cfg(test)]
262 pub(crate) fn rewind_for_test(&mut self, elapsed: Duration) {
263 self.started_at = self
264 .started_at
265 .checked_sub(elapsed)
266 .unwrap_or(self.started_at);
267 if let Some(since) = self.blocked_since {
268 self.blocked_since = Some(since.checked_sub(elapsed).unwrap_or(since));
269 }
270 }
271 }
272
273 #[cfg(test)]
274 mod tests {
275 use super::*;
276
277 #[test]
278 fn model_step_defaults_do_not_install_a_ceiling() {
279 assert_eq!(resolve_max_model_steps(None), DEFAULT_MAX_MODEL_STEPS);
280 assert_eq!(resolve_max_model_steps(Some(0)), DEFAULT_MAX_MODEL_STEPS);
281 assert_eq!(DEFAULT_MAX_MODEL_STEPS, u32::MAX);
282 const { assert!(MAX_MAX_MODEL_STEPS < u32::MAX) };
283 }
284
285 #[test]
286 fn model_step_ceiling_is_overridable_and_clamped() {
287 assert_eq!(resolve_max_model_steps(Some(7)), 7);
288 assert_eq!(resolve_max_model_steps(Some(1)), 1);
289 assert_eq!(
290 resolve_max_model_steps(Some(u32::MAX)),
291 MAX_MAX_MODEL_STEPS,
292 "explicit positive overrides retain the configured ceiling"
293 );
294 }
295
296 #[test]
297 fn turn_wall_clock_defaults_to_no_limit() {
298 assert_eq!(resolve_turn_wall_clock(None), DEFAULT_TURN_WALL_CLOCK);
299 assert_eq!(resolve_turn_wall_clock(Some(0)), DEFAULT_TURN_WALL_CLOCK);
300 let mut clock = TurnWallClock::start(resolve_turn_wall_clock(None));
301 clock.rewind_for_test(Duration::from_secs(10 * MAX_TURN_WALL_CLOCK_SECS));
302 assert!(!clock.exhausted(), "the default never stops a turn");
303 }
304
305 /// Code mode and headless exec sleep toward what is left of the default
306 /// budget; a near-`Duration::MAX` sleep must park, not panic.
307 #[tokio::test]
308 async fn an_unbounded_remaining_budget_is_a_safe_sleep() {
309 let remaining = DEFAULT_TURN_WALL_CLOCK.saturating_sub(Duration::from_secs(1));
310 let woke =
311 tokio::time::timeout(Duration::from_millis(20), tokio::time::sleep(remaining)).await;
312 assert!(woke.is_err(), "the sleep parks until cancelled");
313 let started = std::time::Instant::now();
314 assert!(started.checked_add(remaining).is_none());
315 }
316
317 #[test]
318 fn turn_wall_clock_is_overridable_and_clamped() {
319 assert_eq!(resolve_turn_wall_clock(Some(120)), Duration::from_secs(120));
320 assert_eq!(
321 resolve_turn_wall_clock(Some(1)),
322 Duration::from_secs(MIN_TURN_WALL_CLOCK_SECS)
323 );
324 assert_eq!(
325 resolve_turn_wall_clock(Some(u64::MAX)),
326 Duration::from_secs(MAX_TURN_WALL_CLOCK_SECS)
327 );
328 }
329
330 #[test]
331 fn stream_caps_default_reject_zero_and_are_overridable() {
332 assert_eq!(
333 resolve_stream_max_content_bytes(None),
334 DEFAULT_STREAM_MAX_CONTENT_BYTES
335 );
336 assert_eq!(
337 resolve_stream_max_content_bytes(Some(0)),
338 DEFAULT_STREAM_MAX_CONTENT_BYTES
339 );
340 assert_eq!(resolve_stream_max_content_bytes(Some(1)), 1024 * 1024);
341 assert_eq!(
342 resolve_stream_max_content_bytes(Some(u64::MAX)),
343 MAX_STREAM_MAX_CONTENT_BYTES
344 );
345
346 assert_eq!(
347 resolve_stream_max_duration_secs(None),
348 DEFAULT_STREAM_MAX_DURATION_SECS
349 );
350 assert_eq!(
351 resolve_stream_max_duration_secs(Some(0)),
352 DEFAULT_STREAM_MAX_DURATION_SECS
353 );
354 assert_eq!(resolve_stream_max_duration_secs(Some(60)), 60);
355 assert_eq!(
356 resolve_stream_max_duration_secs(Some(u64::MAX)),
357 MAX_STREAM_MAX_DURATION_SECS
358 );
359 }
360
361 #[test]
362 fn stream_retry_limits_default_to_historical_values_and_clamp() {
363 let defaults = resolve_stream_retry_limits(None, None, None);
364 assert_eq!(defaults, StreamRetryLimits::default());
365 assert_eq!(defaults.max_resumes, 3);
366 assert_eq!(defaults.max_transparent_retries, 2);
367 assert_eq!(defaults.max_errors, 5);
368
369 let disabled = resolve_stream_retry_limits(Some(0), Some(0), Some(0));
370 assert_eq!(disabled.max_resumes, 0, "0 disables resumes");
371 assert_eq!(disabled.max_transparent_retries, 0);
372 assert_eq!(
373 disabled.max_errors, DEFAULT_STREAM_MAX_ERRORS,
374 "0 selects the error-streak default, like the other finite stream budgets"
375 );
376 assert_eq!(
377 resolve_stream_retry_limits(None, None, Some(1)).max_errors,
378 MIN_STREAM_MAX_ERRORS
379 );
380
381 let huge = resolve_stream_retry_limits(Some(u32::MAX), Some(u32::MAX), Some(u32::MAX));
382 assert_eq!(huge.max_resumes, MAX_STREAM_MAX_RESUMES);
383 assert_eq!(
384 huge.max_transparent_retries,
385 MAX_STREAM_MAX_TRANSPARENT_RETRIES
386 );
387 assert_eq!(huge.max_errors, MAX_STREAM_MAX_ERRORS);
388 }
389
390 #[test]
391 fn wall_clock_exhausts_once_the_budget_is_spent() {
392 let mut clock = TurnWallClock::start(Duration::from_secs(60));
393 assert!(!clock.exhausted());
394 clock.rewind_for_test(Duration::from_secs(61));
395 assert!(clock.exhausted());
396 assert!(clock.spent() >= clock.budget());
397 }
398
399 #[test]
400 fn wall_clock_excludes_time_blocked_on_a_human_decision() {
401 let mut clock = TurnWallClock::start(Duration::from_secs(60));
402 clock.begin_human_wait();
403 // The human took two minutes; the agent spent none of its budget.
404 clock.rewind_for_test(Duration::from_secs(120));
405 assert!(
406 !clock.exhausted(),
407 "an unanswered approval prompt must not burn the turn budget"
408 );
409 clock.end_human_wait();
410 assert!(!clock.exhausted());
411 // Agent-owned time after the decision still counts.
412 clock.rewind_for_test(Duration::from_secs(61));
413 assert!(clock.exhausted());
414 }
415
416 #[test]
417 fn nested_human_waits_cannot_double_exclude() {
418 let mut clock = TurnWallClock::start(Duration::from_secs(60));
419 clock.begin_human_wait();
420 clock.begin_human_wait();
421 clock.rewind_for_test(Duration::from_secs(30));
422 clock.end_human_wait();
423 // A second unmatched end is a no-op, not another exclusion.
424 clock.end_human_wait();
425 clock.rewind_for_test(Duration::from_secs(61));
426 assert!(clock.exhausted());
427 }
428 }
429
429 lines RUST