diff --git a/crates/tinyagents-harness/src/agent_loop/lifecycle_tests.rs b/crates/tinyagents-harness/src/agent_loop/lifecycle_tests.rs index e51f187d0..7149755ad 100644 --- a/crates/tinyagents-harness/src/agent_loop/lifecycle_tests.rs +++ b/crates/tinyagents-harness/src/agent_loop/lifecycle_tests.rs @@ -337,6 +337,7 @@ async fn a_recovery_pop_is_retracted_and_the_mirror_stays_exact() { let model_script = vec![ truncated_empty(2048), truncated_empty(4096), + truncated_empty(8192), tool_turn(&["c1"]), response(vec![], "done"), ]; @@ -358,7 +359,7 @@ async fn a_recovery_pop_is_retracted_and_the_mirror_stays_exact() { .iter() .filter(|e| matches!(e, AgentEvent::MessageRetracted { index: 1 })) .count(); - assert_eq!(retractions, 2, "both blank replies were dropped"); + assert_eq!(retractions, 3, "all blank replies were dropped"); // The recovery nudge is announced as an ordinary append. assert!( events.iter().any(|e| matches!( diff --git a/crates/tinyagents-harness/src/agent_loop/mod_tests.rs b/crates/tinyagents-harness/src/agent_loop/mod_tests.rs index fc9be4bbe..5f6ea26fa 100644 --- a/crates/tinyagents-harness/src/agent_loop/mod_tests.rs +++ b/crates/tinyagents-harness/src/agent_loop/mod_tests.rs @@ -1137,12 +1137,14 @@ async fn truncated_empty_response_retries_then_succeeds() { "the retry should be observable; got kinds {:?}", recorder.kinds() ); - // The first attempt sent the configured 2048 cap; the retry doubled it. + // The first retry re-sends at the same cap: most dead calls are a long + // think the model does not repeat, and a cap raised for them makes every + // later dead call cost two to four times as long. let sent: Vec> = model.requests().iter().map(|r| r.max_tokens).collect(); assert_eq!( sent, - vec![Some(2048), Some(4096)], - "the retried request should carry double the original token budget" + vec![Some(2048), Some(2048)], + "the first retry keeps the original token budget" ); } @@ -1152,7 +1154,10 @@ async fn truncated_empty_boost_does_not_leak_into_later_turns() { // recovery state, but they used to live for the whole run — so every turn // after a recovered one was dispatched at the boosted cap (overriding the // caller's `max_turn_output_tokens`) and a later truncation got no retry. + // Two deaths at the cap raise it for the second retry; the next turn is + // back at the configured per-turn cap. let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), truncated_empty_response(2048), tool_call_response("c1", "fake", json!({})), text_response("done", 4, 3), @@ -1175,12 +1180,288 @@ async fn truncated_empty_boost_does_not_leak_into_later_turns() { let sent: Vec> = model.requests().iter().map(|r| r.max_tokens).collect(); assert_eq!( sent, - vec![Some(2048), Some(4096), Some(2048)], - "only the retry of the truncated turn carries the boost; the next turn \ + vec![Some(2048), Some(2048), Some(4096), Some(2048)], + "only the second retry of the truncated turn carries the boost; the next turn \ is back at the configured per-turn cap" ); } +#[tokio::test] +async fn truncated_empty_retry_stops_at_the_cap_ceiling_and_nudges() { + // Once the cap sits at its 4x ceiling, a retry re-sends the transcript + // that just died at that very cap, so it fails the same way at the same + // cost (two 65k dead calls of three minutes each ended one run with + // nothing written). With retries still in hand, the loop must skip the + // retry, say why, and nudge the model to take a smaller step instead. + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), + truncated_empty_response(2048), + truncated_empty_response(4096), + truncated_empty_response(8192), + text_response("recovered", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + harness.with_policy(RunPolicy { + truncated_empty_retries: 4, + ..RunPolicy::default() + }); + + use crate::testkit::EventRecorder; + let recorder = EventRecorder::new(); + let ctx = RunContext::new( + RunConfig::new("truncated-ceiling").with_max_turn_output_tokens(2048), + (), + ) + .with_events(recorder.sink()); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the nudged call recovers"); + + assert_eq!(run.text(), Some("recovered".to_string())); + assert_eq!( + run.model_calls, 5, + "a same-cap retry, two climbs to the ceiling, then the dead call at the ceiling is nudged, not retried" + ); + let sent: Vec> = model.requests().iter().map(|r| r.max_tokens).collect(); + assert_eq!( + sent, + vec![Some(2048), Some(2048), Some(4096), Some(8192), Some(4096)] + ); + let last = model + .requests() + .last() + .expect("five requests") + .messages + .clone(); + assert_eq!( + last.last().map(|m| m.text()), + Some("Your reply was cut off again before any tool call or answer. The output limit for the next call is 4096 tokens: reason in a few sentences at most, then act: write a short answer from what you already have.".to_string()), + "the fifth call carries the nudge at a reduced cap rather than a bare re-send" + ); + let skipped: Vec = recorder + .events() + .iter() + .filter_map(|e| match e { + AgentEvent::ControlApplied { control, detail } + if control == "truncated_empty_retry_skipped" => + { + Some(detail.clone()) + } + _ => None, + }) + .collect(); + assert_eq!( + skipped.len(), + 1, + "exactly one skip, for the dead call at the ceiling" + ); + assert!( + skipped[0].contains("already at its ceiling (8192)"), + "the skip names the cap: {}", + skipped[0] + ); +} + +/// A scripted model whose first call takes `delay`: a dead call that cost +/// real wall-clock time, which is what the retry plan measures. +struct SlowFirstCall { + inner: Arc, + delay: std::time::Duration, + calls: Mutex, +} + +#[async_trait] +impl ChatModel for SlowFirstCall { + async fn invoke( + &self, + state: &State, + request: ModelRequest, + ) -> tinyinference_llm::Result { + let first = { + let mut calls = self.calls.lock().expect("calls lock"); + *calls += 1; + *calls == 1 + }; + if first { + tokio::time::sleep(self.delay).await; + } + self.inner.invoke(state, request).await + } +} + +#[tokio::test] +async fn truncated_empty_retry_yields_to_the_clock() { + // A dead call of 400 ms for 4096 tokens says a retry at the same cap + // would take about 400 ms. With 600 ms of a 1 s run left, that is more + // than half the clock, so the loop skips the retry. A 2048-token nudge + // takes about 200 ms and fits the remaining clock. + let scripted = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(4096), + text_response("recovered", 4, 3), + ])); + let model = Arc::new(SlowFirstCall { + inner: Arc::clone(&scripted), + delay: std::time::Duration::from_millis(400), + calls: Mutex::new(0), + }); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", model as _); + harness.with_policy(RunPolicy { + limits: RunLimits::default().with_max_wall_clock_ms(Some(1_000)), + truncated_empty_nudges: 1, + ..RunPolicy::default() + }); + + use crate::testkit::EventRecorder; + let recorder = EventRecorder::new(); + let ctx = RunContext::new( + RunConfig::new("truncated-clock").with_max_turn_output_tokens(4096), + (), + ) + .with_events(recorder.sink()); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the nudged call recovers inside the clock"); + + assert_eq!(run.text(), Some("recovered".to_string())); + assert_eq!(run.model_calls, 2, "no retry: straight to the nudge"); + let sent: Vec> = scripted.requests().iter().map(|r| r.max_tokens).collect(); + assert_eq!( + sent, + vec![Some(4096), Some(2048)], + "the nudged call runs at the cap the clock affords" + ); + let skipped: Vec = recorder + .events() + .iter() + .filter_map(|e| match e { + AgentEvent::ControlApplied { control, detail } + if control == "truncated_empty_retry_skipped" => + { + Some(detail.clone()) + } + _ => None, + }) + .collect(); + assert_eq!(skipped.len(), 1); + assert!( + skipped[0].contains("would run about") && skipped[0].contains("left"), + "the skip explains the clock arithmetic: {}", + skipped[0] + ); +} + +#[tokio::test] +async fn truncated_empty_nudges_repeat_while_the_clock_allows_with_a_halving_cap() { + // A nudged call that dies too used to end the turn (one run closed with + // 21 of 30 minutes unused). With a clock that still has room, the loop + // nudges again and halves the cap each time, down to the floor. + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), + truncated_empty_response(2048), + truncated_empty_response(4096), + truncated_empty_response(4096), + truncated_empty_response(2048), + text_response("recovered", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + harness.with_policy(RunPolicy { + limits: RunLimits::default().with_max_wall_clock_ms(Some(60_000)), + ..RunPolicy::default() + }); + let ctx = RunContext::new( + RunConfig::new("truncated-clock-nudges").with_max_turn_output_tokens(2048), + (), + ); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the third nudge recovers"); + + assert_eq!(run.text(), Some("recovered".to_string())); + assert_eq!(run.model_calls, 6); + let sent: Vec> = model.requests().iter().map(|r| r.max_tokens).collect(); + assert_eq!( + sent, + vec![ + Some(2048), + Some(2048), + Some(4096), + Some(4096), + Some(2048), + Some(2048) + ], + "a same-cap retry, one climb, then nudges at the held cap, half of it, and the floor" + ); + let last = model + .requests() + .last() + .expect("six requests") + .messages + .clone(); + let text = last.last().map(|m| m.text()).unwrap_or_default(); + assert!( + text.contains("cut off again") && text.contains("2048 tokens"), + "a repeated nudge names the new limit: {text}" + ); +} + +#[tokio::test] +async fn truncated_empty_nudges_stop_at_their_limit_even_with_clock_left() { + // A model that dies at every cap must not be nudged for the rest of the + // run: one first attempt, two retries and six nudges, then the blank is + // surfaced and the run closes. + let responses: Vec = (0..12).map(|_| truncated_empty_response(2048)).collect(); + let model = Arc::new(crate::testkit::ScriptedModel::new(responses)); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + harness.with_policy(RunPolicy { + limits: RunLimits::default().with_max_wall_clock_ms(Some(60_000)), + ..RunPolicy::default() + }); + let ctx = RunContext::new( + RunConfig::new("truncated-clock-limit").with_max_turn_output_tokens(2048), + (), + ); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("a blank is surfaced, not an error, by default"); + assert_eq!( + run.model_calls, + 1 + 2 + 6, + "first attempt, two retries, six nudges" + ); + assert_eq!(run.text().unwrap_or_default().trim(), ""); +} + +#[tokio::test] +async fn truncated_empty_nudges_do_not_repeat_without_a_clock() { + // No wall clock: nothing to spend on a longer run, so the policy's one + // nudge still applies and a second dead nudge surfaces the blank. + let responses: Vec = (0..6).map(|_| truncated_empty_response(2048)).collect(); + let model = Arc::new(crate::testkit::ScriptedModel::new(responses)); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + let ctx = RunContext::new( + RunConfig::new("truncated-no-clock").with_max_turn_output_tokens(2048), + (), + ); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("a blank is surfaced by default"); + assert_eq!( + run.model_calls, + 1 + 2 + 1, + "first attempt, two retries, one nudge" + ); +} + #[tokio::test] async fn truncated_empty_retry_budget_is_restored_for_a_later_turn() { // The retry budget is per turn too: a second truncated-empty completion in a @@ -1236,6 +1517,7 @@ async fn truncated_empty_retries_exhausted_returns_blank_by_default() { // `truncated_empty_nudges` the run gave up after two calls; the third call // is the nudge. let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), truncated_empty_response(2048), truncated_empty_response(4096), truncated_empty_response(4096), @@ -1254,8 +1536,8 @@ async fn truncated_empty_retries_exhausted_returns_blank_by_default() { assert_eq!(run.text(), Some(String::new())); assert_eq!( - run.model_calls, 3, - "one original attempt, one boosted retry, one nudge" + run.model_calls, 4, + "one original attempt, a same-cap retry, a boosted retry, one nudge" ); } @@ -1267,6 +1549,7 @@ async fn truncated_empty_nudge_after_spent_retries_reaches_the_tool_call() { // is spent, the loop must tell the model plainly what happened and keep // going, so the tool call it was deliberating about still happens. let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), truncated_empty_response(2048), truncated_empty_response(4096), tool_call_response("c1", "fake", json!({})), @@ -1287,20 +1570,20 @@ async fn truncated_empty_nudge_after_spent_retries_reaches_the_tool_call() { .expect("the nudged run reaches the tool call and finishes"); assert_eq!(run.text(), Some("done".to_string())); - assert_eq!(run.model_calls, 4); + assert_eq!(run.model_calls, 5); let requests = model.requests(); - let nudged = &requests[2].messages; + let nudged = &requests[3].messages; let nudge = nudged.last().map(Message::text).unwrap_or_default(); assert!( nudge.contains("ran out of output tokens") && nudge.contains("tool call"), - "the third request ends with the truncation nudge; got {nudge:?}" + "the fourth request ends with the truncation nudge; got {nudge:?}" ); assert!( nudged.iter().all(|m| !matches!(m, Message::Assistant(_))), "the blank assistant rows are not replayed" ); assert!( - requests[3] + requests[4] .messages .iter() .any(|m| m.text().contains("tool output")), @@ -1313,6 +1596,7 @@ async fn truncated_empty_nudge_asks_for_an_answer_when_no_tool_is_callable() { // A turn with nothing to call (here: no tools registered) must not be told // to "make the next tool call" — that only invites a call that cannot run. let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), truncated_empty_response(2048), truncated_empty_response(4096), text_response("short answer", 4, 3), @@ -1327,7 +1611,7 @@ async fn truncated_empty_nudge_asks_for_an_answer_when_no_tool_is_callable() { assert_eq!(run.text(), Some("short answer".to_string())); let requests = model.requests(); - let nudge = requests[2] + let nudge = requests[3] .messages .last() .map(Message::text) @@ -1390,6 +1674,7 @@ async fn truncated_empty_nudge_disabled_by_zero_policy() { // truncated_empty_nudges=0 keeps the pre-nudge behavior: after the boosted // retry the blank reply is the final answer. let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), truncated_empty_response(2048), truncated_empty_response(4096), ])); @@ -1404,7 +1689,7 @@ async fn truncated_empty_nudge_disabled_by_zero_policy() { .invoke_default(&(), vec![Message::user("hi")]) .await .expect("with nudges disabled the blank final is returned"); - assert_eq!(run.model_calls, 2); + assert_eq!(run.model_calls, 3); assert_eq!(run.text(), Some(String::new())); } @@ -1413,6 +1698,7 @@ async fn truncated_empty_retries_exhausted_errors_when_guard_enabled() { // With the empty-response guard on, exhausting the retries and the nudge // surfaces the typed EmptyResponse error rather than a blank success. let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), truncated_empty_response(2048), truncated_empty_response(4096), truncated_empty_response(4096), @@ -1826,7 +2112,11 @@ async fn an_empty_max_tokens_reply_takes_the_truncated_empty_retry_path() { .expect("retried, not surfaced"); assert_eq!(run.text(), Some("recovered".to_string())); - assert_eq!(model.requests()[1].max_tokens, Some(4096)); + assert_eq!( + model.requests()[1].max_tokens, + Some(2048), + "the first retry re-sends at the same cap" + ); } #[tokio::test] diff --git a/crates/tinyagents-harness/src/agent_loop/response_recovery.rs b/crates/tinyagents-harness/src/agent_loop/response_recovery.rs index 73b1ff8b1..80061049d 100644 --- a/crates/tinyagents-harness/src/agent_loop/response_recovery.rs +++ b/crates/tinyagents-harness/src/agent_loop/response_recovery.rs @@ -15,6 +15,7 @@ use super::run_loop::{ DROPPED_TOOL_CALL_NUDGE, TRUNCATED_EMPTY_ANSWER_NUDGE, TRUNCATED_EMPTY_TOOL_NUDGE, UNDECODABLE_TOOL_CALL_NUDGE, WITHHELD_TOOL_CALL_NUDGE, truncated_call_positions, }; +use super::turn_recovery::TRUNCATED_CLOCK_NUDGE_LIMIT; use super::*; impl AgentHarness { @@ -166,6 +167,7 @@ impl AgentHarness { response, tool_calls, attempt_max_tokens, + started_at_ms, recovery, tools_available: tools_available_this_turn, text_dialect_calls_recoverable, @@ -227,20 +229,38 @@ impl AgentHarness { let truncated_empty = tool_calls.is_empty() && crate::finish_reason::is_length_stop(response.finish_reason.as_deref()) && response.text().trim().is_empty(); + // What a retry would cost and whether it can change anything. The + // dead call's own duration and output count give the rate this model + // emits at here; the next cap divided by that rate is how long the + // retry would run if it dies the same way. + let truncated_retry = truncated_empty.then(|| { + let dead_ms = crate::ids::now_ms().saturating_sub(started_at_ms); + let dead_tokens = response + .usage + .as_ref() + .map(|u| u.output_tokens) + .unwrap_or(0); + let remaining = crate::middleware::library::TurnClock::of( + ctx, + self.policy + .limits + .max_wall_clock_ms + .map(std::time::Duration::from_millis), + ) + .map(|clock| clock.remaining()); + turn_recovery.truncated_retry_plan(attempt_max_tokens, dead_tokens, dead_ms, remaining) + }); if truncated_empty && turn_recovery.truncated_empty_retries_used < self.policy.truncated_empty_retries && ctx.limits.remaining_model_calls() > 0 + && truncated_retry.as_ref().is_some_and(|plan| plan.worth_it()) { // Drop the useless empty assistant row appended above so the // retry re-sends the identical transcript. messages.pop(); ctx.retract_transcript(messages.len()); - turn_recovery.truncated_empty_retries_used += 1; - // Grow the token budget when the request set one: double it, - // clamped at 4x the original cap. An unset budget stays unset - // (a plain retry is still worthwhile — the failure is - // stochastic). - turn_recovery.boost_max_tokens(attempt_max_tokens); + let plan = truncated_retry.expect("a truncated-empty reply has a plan"); + turn_recovery.take_truncated_retry(attempt_max_tokens, &plan); let record = ctx.emit(AgentEvent::RetryScheduled { call_id: call_id.clone(), attempt: turn_recovery.truncated_empty_retries_used as usize, @@ -248,25 +268,101 @@ impl AgentHarness { status.set_last_event(record.id); return true; } + // A retry that cannot grow the cap re-sends the transcript that just + // died at this cap, and one that would run past half the remaining + // clock leaves no time to use its result either way (one 13-minute + // run spent ten of them in four such calls, the last two at the + // ceiling, and ended with nothing written). Say why, and let the + // nudge below ask for a smaller step instead. The cap it runs under is + // pinned to what the clock affords. + if let Some(plan) = truncated_retry + .as_ref() + .filter(|plan| !plan.worth_it()) + .filter(|_| { + turn_recovery.truncated_empty_retries_used < self.policy.truncated_empty_retries + }) + { + let reason = plan.skip_reason(); + tracing::info!( + target: "tinyagents::agent_loop", + run_id = %ctx.run_id(), + call_id = %call_id, + dead_ms = plan.dead_ms, + dead_tokens = plan.dead_tokens, + current_cap = ?plan.current, + next_cap = ?plan.next, + remaining_ms = ?plan.remaining.map(|d| d.as_millis()), + "[agent_loop] truncated-empty retry skipped: {reason}" + ); + ctx.emit(AgentEvent::ControlApplied { + control: "truncated_empty_retry_skipped".to_string(), + detail: format!("model call `{call_id}`: {reason}"), + }); + if let Some(cap) = plan.affordable_cap() { + turn_recovery.boosted_max_tokens = Some(cap); + } + } - // The boosted retry is spent and the model still deliberated - // past its output budget. Re-sending the same transcript keeps - // failing the same way (a high-effort reasoning model thinks - // as long as it is allowed to), and finishing here hands the - // host a blank reply it can only close as if the work were - // done. Say plainly what happened and ask for the next step, - // then carry on with the loop. The boosted cap stays in force. + // The retries are spent and the model still deliberated past its + // output budget. Re-sending the same transcript keeps failing the + // same way (a high-effort reasoning model thinks as long as it is + // allowed to), and finishing here hands the host a blank reply it + // can only close as if the work were done. Say plainly what happened + // and ask for the next step, then carry on with the loop. + // + // A nudged call that dies too used to end the turn: the loop + // surfaced the blank and the host closed the run, in one case with + // 21 of 30 minutes unused and the deliverable unwritten. While the + // run has a clock with room for another bounded call, keep nudging, + // and halve the output cap each time: the cap is the one limit on + // deliberation these providers honour, and a smaller one makes the + // model act sooner. + let clock_allows_another_nudge = truncated_retry + .as_ref() + .is_some_and(|plan| plan.another_nudge_fits()) + && turn_recovery.truncated_empty_nudges_used < TRUNCATED_CLOCK_NUDGE_LIMIT; + let clock_only_nudge = + turn_recovery.truncated_empty_nudges_used >= self.policy.truncated_empty_nudges; + let policy_nudge_fits = truncated_retry + .as_ref() + .is_some_and(|plan| plan.remaining.is_none() || plan.another_nudge_fits()); if truncated_empty - && turn_recovery.truncated_empty_nudges_used < self.policy.truncated_empty_nudges + && (turn_recovery.truncated_empty_nudges_used < self.policy.truncated_empty_nudges + || clock_allows_another_nudge) + && policy_nudge_fits && ctx.limits.remaining_model_calls() > 0 { messages.pop(); ctx.retract_transcript(messages.len()); turn_recovery.truncated_empty_nudges_used += 1; - let nudge = if tools_available_this_turn { - TRUNCATED_EMPTY_TOOL_NUDGE - } else { - TRUNCATED_EMPTY_ANSWER_NUDGE + let retry_was_skipped_for_clock = truncated_retry + .as_ref() + .is_some_and(|plan| !plan.fits_clock()); + let retry_was_skipped_at_ceiling = truncated_retry + .as_ref() + .is_some_and(|plan| !plan.first_retry && !plan.cap_grows()); + let repeat_cap = (turn_recovery.truncated_empty_nudges_used > 1 + || clock_only_nudge + || retry_was_skipped_for_clock + || retry_was_skipped_at_ceiling) + .then(|| truncated_retry.as_ref().and_then(|plan| plan.nudge_cap())) + .flatten(); + if let Some(cap) = repeat_cap { + turn_recovery.boosted_max_tokens = Some(cap); + } + let nudge: String = match repeat_cap { + Some(cap) => format!( + "Your reply was cut off again before any tool call or answer. The output \ + limit for the next call is {cap} tokens: reason in a few sentences at most, \ + then act: {}.", + if tools_available_this_turn { + "make one tool call" + } else { + "write a short answer from what you already have" + } + ), + None if tools_available_this_turn => TRUNCATED_EMPTY_TOOL_NUDGE.to_string(), + None => TRUNCATED_EMPTY_ANSWER_NUDGE.to_string(), }; tracing::info!( target: "tinyagents::agent_loop", diff --git a/crates/tinyagents-harness/src/agent_loop/run_loop.rs b/crates/tinyagents-harness/src/agent_loop/run_loop.rs index b1bbd6dca..4643eda4c 100644 --- a/crates/tinyagents-harness/src/agent_loop/run_loop.rs +++ b/crates/tinyagents-harness/src/agent_loop/run_loop.rs @@ -941,6 +941,7 @@ impl AgentHarness { response: &response, tool_calls: &tool_calls, attempt_max_tokens, + started_at_ms: model_started_at_ms, recovery: &recovery, tools_available: tools_available_this_turn, text_dialect_calls_recoverable: forced_text_dialect diff --git a/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs b/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs index 7d5d9ec7c..b03a2defb 100644 --- a/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs +++ b/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs @@ -14,7 +14,218 @@ use super::types::TurnRecovery; +/// Share of the remaining wall clock a truncated-empty retry may take. A +/// retry that would run longer than this leaves no time to act on whatever +/// it returns, so the loop nudges instead. +pub(super) const TRUNCATED_RETRY_CLOCK_SHARE: f64 = 0.5; + +/// Lowest output cap a repeated nudge drives the call down to. +pub(super) const TRUNCATED_NUDGE_CAP_FLOOR: u32 = 2048; + +/// Most nudges one turn may spend on a model that keeps dying at the cap, +/// however much clock is left: beyond this the run is better closed than +/// prolonged. +pub(super) const TRUNCATED_CLOCK_NUDGE_LIMIT: u32 = 6; + +/// Whether retrying a call that died at its output cap with nothing to show +/// can help, judged from that call's own rate and the run's clock. +/// +/// Two patterns of dead call exist and both have to work. Most are a long +/// think the model does not repeat: in one run five of seven dead calls were +/// followed by a live call of about 3k tokens, so the doubled cap was never +/// needed, and a cap raised for them made every later dead call cost two to +/// four times as long. Some are a step that genuinely no longer fits the cap +/// (23k tokens after a 16k death). So the first retry re-sends at the same +/// cap, growth waits for a second death at that cap, and a retry that would +/// run past half the remaining clock, or that cannot grow the cap at all, is +/// not made: the loop nudges instead. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(super) struct TruncatedRetryPlan { + /// The cap the dead call ran under. + pub(super) current: Option, + /// The cap a retry would run under. + pub(super) next: Option, + /// The cap the first call of the turn ran under: the floor for any cap + /// the clock pins. + pub(super) base: Option, + /// Output tokens the dead call emitted (its hidden reasoning included). + pub(super) dead_tokens: u64, + /// How long the dead call took. + pub(super) dead_ms: u64, + /// Wall clock left for the run, when it has one. + pub(super) remaining: Option, + /// The first retry of this turn re-sends at the same cap on purpose. + pub(super) first_retry: bool, +} + +impl TruncatedRetryPlan { + /// Milliseconds the dead call spent per output token, when both are known. + fn ms_per_token(&self) -> Option { + (self.dead_tokens > 0 && self.dead_ms > 0) + .then(|| self.dead_ms as f64 / self.dead_tokens as f64) + } + + /// The retry can grow the cap, or the call never had one (a plain retry of + /// an uncapped call is still worth one attempt: the failure is stochastic). + pub(super) fn cap_grows(&self) -> bool { + match (self.current, self.next) { + (Some(current), Some(next)) => next > current, + _ => true, + } + } + + /// How long the retry would run if it dies at its cap the same way. + fn expected_ms(&self) -> u64 { + match (self.ms_per_token(), self.next) { + (Some(rate), Some(next)) => (rate * next as f64) as u64, + (None, Some(next)) if self.current.is_some_and(|current| current > 0) => { + let current = self.current.unwrap_or_default(); + (self.dead_ms as f64 * next as f64 / current as f64) as u64 + } + _ => self.dead_ms, + } + } + + /// The retry finishes inside its share of the remaining clock, or there is + /// no clock to keep. + pub(super) fn fits_clock(&self) -> bool { + match self.remaining { + Some(remaining) => { + self.expected_ms() as f64 + <= remaining.as_millis() as f64 * TRUNCATED_RETRY_CLOCK_SHARE + } + None => true, + } + } + + pub(super) fn worth_it(&self) -> bool { + (self.first_retry || self.cap_grows()) && self.fits_clock() + } + + pub(super) fn skip_reason(&self) -> String { + if !self.first_retry && !self.cap_grows() { + format!( + "the output cap is already at its ceiling ({}); re-sending the same transcript at the same cap fails the same way", + self.current.map_or("unset".to_string(), |c| c.to_string()) + ) + } else { + format!( + "a retry at {} tokens would run about {}s at the rate of the call that just died, more than half the {}s left", + self.next + .map_or("the same cap".to_string(), |c| c.to_string()), + self.expected_ms() / 1000, + self.remaining.map_or(0, |d| d.as_secs()) + ) + } + } + + /// The largest cap, between the first cap and the current one, whose call + /// fits the clock share at the dead call's rate: the cap the nudged call + /// should run under so it cannot spend the rest of the run the same way. + /// `None` when there is nothing to learn from (no cap, no rate, no clock). + pub(super) fn affordable_cap(&self) -> Option { + let current = self.current?; + let rate = self.ms_per_token()?; + let remaining = self.remaining?; + let affordable = (remaining.as_millis() as f64 * TRUNCATED_RETRY_CLOCK_SHARE / rate) as u64; + let floor = self.base.unwrap_or(current).min(current).max(1); + (affordable >= u64::from(floor)).then_some(affordable.min(u64::from(current)) as u32) + } + + /// The cap the next nudged call runs under: half the current one, never + /// below [`TRUNCATED_NUDGE_CAP_FLOOR`]. `None` when the call had no cap. + pub(super) fn halved_cap(&self) -> Option { + let current = self.current?; + Some((current / 2).max(TRUNCATED_NUDGE_CAP_FLOOR).min(current)) + } + + /// The nudge cap is the smaller of the halved cap and the cap the clock + /// can afford. If the known rate says even the original cap cannot fit, + /// there is no affordable nudge cap. + pub(super) fn nudge_cap(&self) -> Option { + let halved = self.halved_cap()?; + match (self.ms_per_token(), self.remaining) { + (Some(rate), Some(remaining)) => { + let affordable = + (remaining.as_millis() as f64 * TRUNCATED_RETRY_CLOCK_SHARE / rate) as u64; + let minimum = halved.min(TRUNCATED_NUDGE_CAP_FLOOR); + (affordable >= u64::from(minimum)) + .then_some(halved.min(affordable.min(u64::from(u32::MAX)) as u32)) + } + _ => Some(halved), + } + } + + /// A nudged call at the halved cap, at the dead call's rate, fits inside + /// its share of the remaining clock. Without a clock there is nothing to + /// spend, so the answer is no: the policy's own nudge count applies. + pub(super) fn another_nudge_fits(&self) -> bool { + let Some(remaining) = self.remaining else { + return false; + }; + let expected_ms = match (self.current, self.ms_per_token(), self.nudge_cap()) { + (None, _, _) => { + // An uncapped call cannot be shortened by changing its output + // limit. Still allow a clock-bounded nudge when the previous + // call's duration leaves enough room for another attempt. + self.dead_ms + } + (_, Some(rate), Some(cap)) => (rate * cap as f64) as u64, + // A capped call with no affordable nudge cap cannot safely repeat. + (_, _, None) => return false, + (_, None, Some(_)) => self.dead_ms, + }; + expected_ms as f64 <= remaining.as_millis() as f64 * TRUNCATED_RETRY_CLOCK_SHARE + } +} + impl TurnRecovery { + /// What a retry of the dead call that just returned would cost and + /// whether it can change anything (see [`TruncatedRetryPlan`]). + pub(super) fn truncated_retry_plan( + &self, + attempt_max_tokens: Option, + dead_tokens: u64, + dead_ms: u64, + remaining: Option, + ) -> TruncatedRetryPlan { + let first_retry = self.truncated_empty_retries_used == 0; + let current = self.boosted_max_tokens.or(attempt_max_tokens); + let base = self.truncation_base.or(attempt_max_tokens); + let next = attempt_max_tokens.map(|sent| { + let base = self.truncation_base.unwrap_or(sent); + let current = self.boosted_max_tokens.unwrap_or(sent); + if first_retry { + current + } else { + current.saturating_mul(2).min(base.saturating_mul(4)) + } + }); + TruncatedRetryPlan { + current, + next, + base, + dead_tokens, + dead_ms, + remaining, + first_retry, + } + } + + /// Take the retry a plan describes: remember the original cap and set the + /// cap the retry runs under. An unset cap stays unset. + pub(super) fn take_truncated_retry( + &mut self, + attempt_max_tokens: Option, + plan: &TruncatedRetryPlan, + ) { + self.truncated_empty_retries_used += 1; + if let Some(sent) = attempt_max_tokens { + self.truncation_base.get_or_insert(sent); + self.boosted_max_tokens = plan.next; + } + } + /// Grows the next request's output cap after a length-truncated reply: /// double the cap last sent, clamped at 4x the original. An unset cap /// stays unset (a plain retry is still worthwhile — the failure is diff --git a/crates/tinyagents-harness/src/agent_loop/turn_recovery_tests.rs b/crates/tinyagents-harness/src/agent_loop/turn_recovery_tests.rs index dd254d12a..ae6532946 100644 --- a/crates/tinyagents-harness/src/agent_loop/turn_recovery_tests.rs +++ b/crates/tinyagents-harness/src/agent_loop/turn_recovery_tests.rs @@ -37,6 +37,75 @@ fn boost_leaves_an_unset_cap_unset() { assert_eq!(recovery, TurnRecovery::default()); } +#[test] +fn nudge_cap_is_clamped_to_the_clock_affordable_cap() { + let plan = TruncatedRetryPlan { + current: Some(16_000), + next: Some(16_000), + base: Some(2_048), + dead_tokens: 16_000, + dead_ms: 16_000, + remaining: Some(std::time::Duration::from_millis(6_000)), + first_retry: false, + }; + + assert_eq!(plan.affordable_cap(), Some(3_000)); + assert_eq!(plan.nudge_cap(), Some(3_000)); + assert!(plan.another_nudge_fits()); +} + +#[test] +fn nudge_is_rejected_when_even_the_minimum_cap_exceeds_the_clock_budget() { + let plan = TruncatedRetryPlan { + current: Some(4_096), + next: Some(4_096), + base: Some(2_048), + dead_tokens: 4_096, + dead_ms: 4_096, + remaining: Some(std::time::Duration::from_millis(2_000)), + first_retry: false, + }; + + assert_eq!(plan.affordable_cap(), None); + assert_eq!(plan.nudge_cap(), None); + assert!(!plan.another_nudge_fits()); +} + +#[test] +fn uncapped_nudge_uses_the_dead_call_duration_as_its_clock_estimate() { + let enough_time = TruncatedRetryPlan { + current: None, + next: None, + base: None, + dead_tokens: 0, + dead_ms: 1_000, + remaining: Some(std::time::Duration::from_millis(3_000)), + first_retry: false, + }; + let too_little_time = TruncatedRetryPlan { + remaining: Some(std::time::Duration::from_millis(1_000)), + ..enough_time + }; + + assert!(enough_time.another_nudge_fits()); + assert!(!too_little_time.another_nudge_fits()); +} + +#[test] +fn retry_clock_estimate_scales_with_the_candidate_cap_without_usage() { + let plan = TruncatedRetryPlan { + current: Some(4_000), + next: Some(8_000), + base: Some(4_000), + dead_tokens: 0, + dead_ms: 1_000, + remaining: Some(std::time::Duration::from_millis(10_000)), + first_retry: false, + }; + + assert_eq!(plan.expected_ms(), 2_000); +} + #[test] fn reset_truncated_empty_clears_only_the_truncated_empty_state() { let mut recovery = spent(); diff --git a/crates/tinyagents-harness/src/agent_loop/types.rs b/crates/tinyagents-harness/src/agent_loop/types.rs index ac9751d37..fdb1f9bbd 100644 --- a/crates/tinyagents-harness/src/agent_loop/types.rs +++ b/crates/tinyagents-harness/src/agent_loop/types.rs @@ -163,6 +163,9 @@ pub(super) struct ResponseTurn<'a> { pub(super) tool_calls: &'a [ToolCall], /// The output cap actually sent with the request that produced `response`. pub(super) attempt_max_tokens: Option, + /// When the model call that produced `response` started (`ids::now_ms`), + /// so a recovery can weigh a retry against how long the call just took. + pub(super) started_at_ms: u64, /// What the dialect layer recovered or withheld from the response text. pub(super) recovery: &'a super::dialect::TextRecovery, /// Whether the request offered a callable tool this turn. diff --git a/crates/tinyagents-harness/src/runtime/types.rs b/crates/tinyagents-harness/src/runtime/types.rs index 53583fa09..3ef994187 100644 --- a/crates/tinyagents-harness/src/runtime/types.rs +++ b/crates/tinyagents-harness/src/runtime/types.rs @@ -326,8 +326,14 @@ pub struct RunPolicy { /// The retry runs *before* [`Self::error_on_empty_response`]; only once /// these retries are exhausted does that guard (if enabled) apply. /// - /// Defaults to `1` (one retry, two attempts total). Set to `0` to disable - /// for exact-replay callers that must not re-issue a call. + /// The first retry re-sends at the same cap (most dead calls are a long + /// think the model does not repeat); the second doubles it, clamped at 4x, + /// for the step that genuinely no longer fits. A retry that would run past + /// half the run's remaining wall clock is skipped in favour of the nudge. + /// + /// Defaults to `2` (two retries, three attempts total): one at the same + /// cap and one with room. Set to `0` to disable for exact-replay callers + /// that must not re-issue a call. pub truncated_empty_retries: u32, /// Re-prompts after [`Self::truncated_empty_retries`] are spent and the /// model *still* returned a truncated-empty completion. @@ -652,7 +658,7 @@ impl Default for RunPolicy { // On by default: a truncated-empty completion is useless to every // caller, so one stochastic-failure retry is strictly better than a // blank final. - truncated_empty_retries: 1, + truncated_empty_retries: 2, // A truncation that survives the boosted retry is a model that // keeps deliberating; one plain "stop and act" re-prompt recovers // the step instead of ending the run on a blank reply.