diff --git a/crates/tinyagents-harness/src/agent_loop/README.md b/crates/tinyagents-harness/src/agent_loop/README.md index 5816ae269..a57903b56 100644 --- a/crates/tinyagents-harness/src/agent_loop/README.md +++ b/crates/tinyagents-harness/src/agent_loop/README.md @@ -172,6 +172,39 @@ allowed call ends on the blank reply instead of failing with `0` to restore exact-replay behavior. The recovery lives in the shared `run_loop`, so it applies identically to the unary and streaming paths. +Neither a smaller cap nor a lower effort label reliably stops a hosted model +that deliberates past its cap; `reasoning.effort = none` does. So after a dead +call the retry or nudged call goes out with reasoning switched off +(`RunPolicy::truncated_empty_reasoning_fallback`, default `true`), at the same +cap (the cap was for the deliberation), and the nudge tells the model to do its +working-out in the workspace. The hold-off backs off per death (1, 2, 4, 8, 16 +live calls without reasoning, no ceiling) and the configured effort always +returns once it is spent: kept off for good, a model spends the rest of a run +writing probe programs instead of the deliverable; +the switch is announced as `AgentEvent::ControlApplied` (`reasoning_fallback`). +The state is run-wide (`agent_loop/reasoning_fallback.rs`): it outlives the +turn that set it. + +A request's `reasoning.budget_tokens` is a promise the provider may not keep, +and every streamed call measured that reasoned past its budget with nothing +visible went on to die at the output cap. The reasoning watchdog +(`RunPolicy::reasoning_watchdog`, default `RequestBudget`) ends such a call at +the budget instead, dropping the stream, and hands the loop the same +`finish_reason = length`, no-content response the cap would have produced, so +the recovery above runs after a fraction of the wait +(`AgentEvent::ControlApplied`, `reasoning_watchdog`). Visible text or a +tool-call fragment disarms it; unary calls are not bounded. + +The reasoning a call dies in is usually real work (a correct derivation, cut +off), and every retry used to begin it again from nothing. The tail of that +reasoning (`RunPolicy::truncated_empty_carry_reasoning_chars`, default 8,000 +characters, `0` carries nothing) rides into the transcript as a user message +ahead of the retry or nudged call, framed as the model's own interrupted notes +with the instruction to continue from there in code rather than re-derive +(`AgentEvent::ControlApplied`, `truncated_empty_reasoning_carried`). The +watchdog's synthetic dead response keeps the reasoning that streamed, so a call +it ends is carried the same way as one the cap ended. + A provider can also end a stream normally after emitting only reasoning, with no visible text or tool call. Hosts may set `RunPolicy::empty_response_retries` to retry that non-truncated blank result; diff --git a/crates/tinyagents-harness/src/agent_loop/mod.rs b/crates/tinyagents-harness/src/agent_loop/mod.rs index ec0d1ec34..e5ac8df52 100644 --- a/crates/tinyagents-harness/src/agent_loop/mod.rs +++ b/crates/tinyagents-harness/src/agent_loop/mod.rs @@ -130,6 +130,7 @@ mod model_switch; mod model_turn; mod nested; pub mod phases; +mod reasoning_fallback; mod response_recovery; mod run_loop; pub(crate) mod stream; diff --git a/crates/tinyagents-harness/src/agent_loop/mod_tests.rs b/crates/tinyagents-harness/src/agent_loop/mod_tests.rs index 5f6ea26fa..f1a6e79e1 100644 --- a/crates/tinyagents-harness/src/agent_loop/mod_tests.rs +++ b/crates/tinyagents-harness/src/agent_loop/mod_tests.rs @@ -1084,6 +1084,11 @@ async fn truncated_empty_retry_does_not_replay_a_cached_blank() { ])); let mut harness: AgentHarness<()> = AgentHarness::new(); harness.register_model("mock", Arc::clone(&model) as _); + harness.with_policy(RunPolicy { + // This test pins the cap ladder; keep reasoning on so the ladder runs. + truncated_empty_reasoning_fallback: false, + ..RunPolicy::default() + }); harness.with_response_cache(Arc::new(InMemoryResponseCache::new())); let ctx = RunContext::new( @@ -1166,6 +1171,11 @@ async fn truncated_empty_boost_does_not_leak_into_later_turns() { harness .register_model("mock", Arc::clone(&model) as _) .register_tool(Arc::new(FakeTool::new("fake", "tool output"))); + harness.with_policy(RunPolicy { + // This test pins the cap ladder; keep reasoning on so the ladder runs. + truncated_empty_reasoning_fallback: false, + ..RunPolicy::default() + }); let ctx = RunContext::new( RunConfig::new("truncated-leak").with_max_turn_output_tokens(2048), @@ -1203,6 +1213,8 @@ async fn truncated_empty_retry_stops_at_the_cap_ceiling_and_nudges() { let mut harness: AgentHarness<()> = AgentHarness::new(); harness.register_model("mock", Arc::clone(&model) as _); harness.with_policy(RunPolicy { + // This test pins the cap ladder; keep reasoning on so the ladder runs. + truncated_empty_reasoning_fallback: false, truncated_empty_retries: 4, ..RunPolicy::default() }); @@ -1370,6 +1382,8 @@ async fn truncated_empty_nudges_repeat_while_the_clock_allows_with_a_halving_cap let mut harness: AgentHarness<()> = AgentHarness::new(); harness.register_model("mock", Arc::clone(&model) as _); harness.with_policy(RunPolicy { + // This test pins the cap ladder; keep reasoning on so the ladder runs. + truncated_empty_reasoning_fallback: false, limits: RunLimits::default().with_max_wall_clock_ms(Some(60_000)), ..RunPolicy::default() }); @@ -9261,3 +9275,442 @@ async fn dynamic_toolset_fold_is_reported_as_a_transcript_rewrite() { ); super::lifecycle_test::assert_mirrors(&events, &seed, &run.messages); } + +#[tokio::test] +async fn a_dead_call_is_retried_with_reasoning_off_then_restored() { + // The one control measured to stop a model that deliberates past any + // cap is switching reasoning off. The retry goes out without it, at the + // same cap (the cap was for the deliberation), and the configured + // effort returns once the hold-off of one live call is spent. + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), + tool_call_response("c1", "fake", json!({})), + text_response("done", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model("mock", Arc::clone(&model) as _) + .register_tool(Arc::new(FakeTool::new("fake", "tool output"))); + harness.with_policy(RunPolicy { + default_reasoning: Some(ReasoningConfig::effort(ReasoningEffort::High)), + ..RunPolicy::default() + }); + let recorder = crate::testkit::EventRecorder::new(); + let ctx = RunContext::new( + RunConfig::new("truncated-reasoning-off").with_max_turn_output_tokens(2048), + (), + ) + .with_events(recorder.sink()); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("done".to_string())); + let efforts: Vec> = model + .requests() + .iter() + .map(|r| r.reasoning.as_ref().and_then(|c| c.effort)) + .collect(); + assert_eq!( + efforts, + vec![ + Some(ReasoningEffort::High), + Some(ReasoningEffort::None), + Some(ReasoningEffort::High) + ], + "the retry runs without reasoning; the call after the live reply gets the effort back" + ); + let caps: Vec> = model.requests().iter().map(|r| r.max_tokens).collect(); + assert_eq!( + caps, + vec![Some(2048), Some(2048), Some(2048)], + "no cap growth for a call without reasoning" + ); + assert!( + recorder.events().iter().any(|e| matches!( + e, + AgentEvent::ControlApplied { control, .. } if control == "reasoning_fallback" + )), + "the switch is observable; got kinds {:?}", + recorder.kinds() + ); +} + +#[tokio::test] +async fn repeated_deaths_hold_reasoning_off_for_longer() { + // Hold-off backs off: one call after the first death, two after the + // second. Two live replies without reasoning then restore the effort. + // The second retry re-sends at the same cap too: a call without + // reasoning needs no larger cap. + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), + truncated_empty_response(2048), + tool_call_response("c1", "fake", json!({})), + tool_call_response("c2", "fake", json!({})), + text_response("done", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model("mock", Arc::clone(&model) as _) + .register_tool(Arc::new(FakeTool::new("fake", "tool output"))); + harness.with_policy(RunPolicy { + default_reasoning: Some(ReasoningConfig::effort(ReasoningEffort::High)), + ..RunPolicy::default() + }); + let ctx = RunContext::new( + RunConfig::new("truncated-reasoning-backoff").with_max_turn_output_tokens(2048), + (), + ); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("done".to_string())); + let efforts: Vec> = model + .requests() + .iter() + .map(|r| r.reasoning.as_ref().and_then(|c| c.effort)) + .collect(); + assert_eq!( + efforts, + vec![ + Some(ReasoningEffort::High), + Some(ReasoningEffort::None), + Some(ReasoningEffort::None), + Some(ReasoningEffort::None), + Some(ReasoningEffort::High), + ] + ); + let caps: Vec> = model.requests().iter().map(|r| r.max_tokens).collect(); + assert_eq!( + caps, + vec![Some(2048), Some(2048), Some(2048), Some(2048), Some(2048)] + ); +} + +#[tokio::test] +async fn the_reasoning_fallback_can_be_switched_off() { + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + 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 { + default_reasoning: Some(ReasoningConfig::effort(ReasoningEffort::High)), + truncated_empty_reasoning_fallback: false, + ..RunPolicy::default() + }); + let ctx = RunContext::new( + RunConfig::new("truncated-reasoning-kept").with_max_turn_output_tokens(2048), + (), + ); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("recovered".to_string())); + let efforts: Vec> = model + .requests() + .iter() + .map(|r| r.reasoning.as_ref().and_then(|c| c.effort)) + .collect(); + assert_eq!( + efforts, + vec![Some(ReasoningEffort::High), Some(ReasoningEffort::High)] + ); +} + +#[tokio::test] +async fn the_nudge_after_a_dead_call_says_where_to_think_instead() { + // With retries spent, the nudge that re-prompts the model also says + // reasoning is off for the next call and that the working-out goes into + // the workspace. + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response(2048), + tool_call_response("c1", "fake", json!({})), + text_response("done", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model("mock", Arc::clone(&model) as _) + .register_tool(Arc::new(FakeTool::new("fake", "tool output"))); + harness.with_policy(RunPolicy { + truncated_empty_retries: 0, + truncated_empty_nudges: 1, + ..RunPolicy::default() + }); + let ctx = RunContext::new( + RunConfig::new("truncated-reasoning-nudge").with_max_turn_output_tokens(2048), + (), + ); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("done".to_string())); + let last_user = model.requests()[1] + .messages + .last() + .map(|m| m.text()) + .unwrap_or_default(); + assert!( + last_user.contains("Reasoning is switched off"), + "nudge should carry the reasoning-off note; got {last_user:?}" + ); + assert!( + last_user.contains("scratch file"), + "and say where to think: {last_user:?}" + ); +} + +/// A middleware that notes a repeat on every tool result, standing in for +/// the repeat-progress guard's warning. +struct RepeatNoter; + +#[async_trait] +impl Middleware<(), ()> for RepeatNoter { + fn name(&self) -> &str { + "repeat-noter" + } + + async fn after_tool( + &self, + ctx: &mut RunContext<()>, + _state: &(), + _invocation: &ToolInvocationIdentity, + _result: &mut ToolResult, + ) -> Result<()> { + ctx.note_repeat(); + Ok(()) + } +} + +#[tokio::test] +async fn a_repeat_note_hands_reasoning_back_before_the_hold_off_is_spent() { + // Two deaths hold reasoning off for two live calls. The retry runs + // without it and makes a tool call; a repeat noted on that call's result + // means the model is looping without reasoning, so the next call goes + // out at the configured effort although one hold-off call is left. + 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), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model("mock", Arc::clone(&model) as _) + .register_tool(Arc::new(FakeTool::new("fake", "tool output"))) + .push_middleware(Arc::new(RepeatNoter)); + harness.with_policy(RunPolicy { + default_reasoning: Some(ReasoningConfig::effort(ReasoningEffort::High)), + ..RunPolicy::default() + }); + let recorder = crate::testkit::EventRecorder::new(); + let ctx = RunContext::new( + RunConfig::new("truncated-reasoning-restored").with_max_turn_output_tokens(2048), + (), + ) + .with_events(recorder.sink()); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("done".to_string())); + let efforts: Vec> = model + .requests() + .iter() + .map(|r| r.reasoning.as_ref().and_then(|c| c.effort)) + .collect(); + assert_eq!( + efforts, + vec![ + Some(ReasoningEffort::High), + Some(ReasoningEffort::None), + Some(ReasoningEffort::None), + Some(ReasoningEffort::High), + ], + "the repeat note ends the hold-off early" + ); + assert!( + recorder.events().iter().any(|e| matches!( + e, + AgentEvent::ControlApplied { control, .. } if control == "reasoning_restored" + )), + "the restore is observable; got kinds {:?}", + recorder.kinds() + ); +} + +/// A dead call that carries the reasoning it died in, as a provider's +/// length-truncated streamed response does. +fn truncated_empty_response_with_reasoning(reasoning: &str) -> ModelResponse { + let mut response = truncated_empty_response(2048); + response.message.content = vec![ContentBlock::Thinking { + text: reasoning.to_string(), + signature: None, + }]; + response +} + +#[tokio::test] +async fn a_dead_calls_reasoning_is_carried_into_the_retry() { + // The derivation the call died in is handed back as the model's own + // interrupted notes, tail first, with the instruction to continue in + // code; the retry request carries it as its last message. + let reasoning = format!( + "{}\nSTATE OF PLAY: encoder mirrors the decoder's split\n", + "x".repeat(20_000) + ); + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response_with_reasoning(&reasoning), + text_response("recovered", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + let recorder = crate::testkit::EventRecorder::new(); + let ctx = RunContext::new( + RunConfig::new("truncated-carry").with_max_turn_output_tokens(2048), + (), + ) + .with_events(recorder.sink()); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("recovered".to_string())); + let retry = model.requests()[1].clone(); + let last = retry.messages.last().map(|m| m.text()).unwrap_or_default(); + assert!( + last.starts_with(super::run_loop::TRUNCATED_EMPTY_CARRY_PREFIX), + "the carry-over is the retry's last message: {last:.120}" + ); + assert!( + last.contains("STATE OF PLAY"), + "the tail of the reasoning is kept" + ); + assert!(last.contains("Continue from this point")); + assert!( + last.len() < 8_000 + 600, + "the excerpt is bounded by the policy: {} chars", + last.len() + ); + assert!( + recorder.events().iter().any(|e| matches!( + e, + AgentEvent::ControlApplied { control, .. } + if control == "truncated_empty_reasoning_carried" + )), + "the carry is observable; got kinds {:?}", + recorder.kinds() + ); +} + +#[tokio::test] +async fn the_carry_limit_counts_characters_not_utf8_bytes() { + // 9,000 CJK characters are 27,000 bytes. The policy's 8,000 is a + // character budget: the excerpt keeps about 8,000 characters (and the + // marker at the very end), not 8,000 bytes, which would be under 2,700 + // characters. + let reasoning = format!("{}\nSTATE OF PLAY: 推导完成\n", "思".repeat(9_000)); + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response_with_reasoning(&reasoning), + text_response("recovered", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + let ctx = RunContext::new( + RunConfig::new("truncated-carry-chars").with_max_turn_output_tokens(2048), + (), + ); + let run = harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("recovered".to_string())); + let retry = model.requests()[1].clone(); + let last = retry.messages.last().map(|m| m.text()).unwrap_or_default(); + assert!(last.contains("STATE OF PLAY: 推导完成"), "the tail is kept"); + let kept = last.chars().filter(|c| *c == '思').count(); + assert!( + (7_800..=8_000).contains(&kept), + "about 8,000 characters of reasoning are kept, not 8,000 bytes: {kept}" + ); +} + +#[tokio::test] +async fn short_or_absent_reasoning_is_not_carried() { + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response_with_reasoning("hmm"), + text_response("recovered", 4, 3), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + let ctx = RunContext::new( + RunConfig::new("truncated-no-carry").with_max_turn_output_tokens(2048), + (), + ); + harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + let retry = model.requests()[1].clone(); + assert_eq!( + retry.messages.len(), + 1, + "no carry-over message for three characters of reasoning" + ); +} + +#[tokio::test] +async fn the_carry_over_can_be_switched_off() { + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response_with_reasoning(&"y".repeat(5_000)), + 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_carry_reasoning_chars: 0, + ..RunPolicy::default() + }); + let ctx = RunContext::new( + RunConfig::new("truncated-carry-off").with_max_turn_output_tokens(2048), + (), + ); + harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(model.requests()[1].messages.len(), 1); +} + +#[tokio::test] +async fn the_nudge_after_spent_retries_carries_the_reasoning_too() { + // With no retry left, the nudged call still gets the interrupted + // working-out, ahead of the nudge itself. + let model = Arc::new(crate::testkit::ScriptedModel::new(vec![ + truncated_empty_response_with_reasoning(&"z".repeat(1_000)), + 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: 0, + truncated_empty_nudges: 1, + ..RunPolicy::default() + }); + let ctx = RunContext::new( + RunConfig::new("truncated-carry-nudge").with_max_turn_output_tokens(2048), + (), + ); + harness + .invoke_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + let messages = model.requests()[1].messages.clone(); + assert_eq!(messages.len(), 3, "the seed, the carry, then the nudge"); + assert!(messages[1].text().contains("Continue from this point")); + assert!(messages[2].text().contains("ran out of output tokens")); +} diff --git a/crates/tinyagents-harness/src/agent_loop/model_call.rs b/crates/tinyagents-harness/src/agent_loop/model_call.rs index 377535e11..aa85646ad 100644 --- a/crates/tinyagents-harness/src/agent_loop/model_call.rs +++ b/crates/tinyagents-harness/src/agent_loop/model_call.rs @@ -1219,10 +1219,19 @@ impl AgentHarness { // transformed before it reached consumers. let mut streamed_text = String::new(); let mut streamed_reasoning = String::new(); + // Characters of `streamed_reasoning`, kept as a running count so the + // watchdog estimate below is in characters, not UTF-8 bytes, without + // walking the whole text on every delta. + let mut streamed_reasoning_chars: u64 = 0; let mut saw_streamed_content = false; let mut transformed_tools = StreamAccumulator::new(); let mut saw_tool_delta = false; let mut stream_stall = StreamTextStallDetector::default(); + // Client-side bound on hidden reasoning (see + // `RunPolicy::reasoning_watchdog`): the reasoning-token count at which + // a call that has shown nothing visible is ended as a dead call. + let reasoning_bound = self.policy.reasoning_watchdog.bound_for(request); + let watchdog_started = std::time::Instant::now(); // Some providers pad the very first streamed text chunk with // whitespace that is a wire-format artifact, not content (see @@ -1429,6 +1438,9 @@ impl AgentHarness { self.middleware .run_on_model_delta(ctx, state, &mut model_delta) .await?; + // A tool call middleware added counts as one the call has, + // for this delta and every later one. + saw_tool_delta |= model_delta.tool_call.is_some(); if model_delta.tool_call.is_some() { stream_stall.reset(); } else if stream_stall.observe(&model_delta.content) { @@ -1441,6 +1453,44 @@ impl AgentHarness { || !model_delta.reasoning.is_empty(); streamed_text.push_str(&model_delta.content); streamed_reasoning.push_str(&model_delta.reasoning); + streamed_reasoning_chars += model_delta.reasoning.chars().count() as u64; + // Reasoning past the bound with nothing visible yet: end the + // call here as the dead call it was going to be, instead of + // waiting for the provider to reach the output cap. Dropping + // `stream` on return closes the connection. The loop sees the + // same response a cap-truncated call produces and runs its + // truncated-empty recovery. + if let Some(bound) = reasoning_bound + && streamed_text.trim().is_empty() + && !saw_tool_delta + && message_delta.tool_call.is_none() + && model_delta.tool_call.is_none() + && estimated_reasoning_tokens(streamed_reasoning_chars) > u64::from(bound) + { + let estimated = estimated_reasoning_tokens(streamed_reasoning_chars); + let elapsed_ms = watchdog_started.elapsed().as_millis() as u64; + tracing::warn!( + target: "tinyagents::agent_loop", + run_id = %ctx.run_id(), + call_id = %call_id, + bound, + estimated_reasoning_tokens = estimated, + elapsed_ms, + "[stream] reasoning watchdog ended a call that reasoned past its budget with nothing visible" + ); + ctx.emit(AgentEvent::ControlApplied { + control: "reasoning_watchdog".to_string(), + detail: format!( + "model call `{call_id}` reasoned past its {bound}-token budget (about \ + {estimated} tokens in {elapsed_ms} ms) with no visible output; the call \ + was ended and is treated as a dead call" + ), + }); + return Ok(watchdog_dead_response( + estimated, + std::mem::take(&mut streamed_reasoning), + )); + } let forwarded_delta = MessageDelta { text: model_delta.content.clone(), reasoning: model_delta.reasoning.clone(), @@ -1926,10 +1976,60 @@ impl ToolBaseCall } } +/// Reasoning tokens a streamed reasoning text of `reasoning_chars` characters +/// (not UTF-8 bytes: a CJK character is three bytes and about one token) +/// amounts to, estimated at three characters per token. Measured on deepseek-v4.1-flash, whose reasoning is +/// dense with code, numbers and short tokens, 36k characters were about 13k +/// tokens (2.7 per token); English prose runs nearer four. Three keeps the +/// estimate on the low side for this kind of text, so a bound built on it +/// fires a little late rather than early. +fn estimated_reasoning_tokens(reasoning_chars: u64) -> u64 { + reasoning_chars / 3 +} + +/// The response the reasoning watchdog hands the loop in place of the call it +/// ended: the shape a call truncated at its output cap has (`finish_reason = +/// length`, no visible text, no tool call), with the estimated reasoning as +/// its output usage so the recovery can judge the call's rate, and the +/// reasoning that streamed kept as a `Thinking` block so the recovery can +/// carry it forward (`RunPolicy::truncated_empty_carry_reasoning_chars`). +fn watchdog_dead_response(estimated_reasoning_tokens: u64, reasoning: String) -> ModelResponse { + let usage = tinyinference_llm::Usage::new(0, estimated_reasoning_tokens); + let content = if reasoning.is_empty() { + Vec::new() + } else { + vec![tinyinference_llm::message::ContentBlock::Thinking { + text: reasoning, + signature: None, + }] + }; + ModelResponse { + message: tinyinference_llm::message::AssistantMessage { + id: None, + content, + tool_calls: Vec::new(), + usage: Some(usage), + origin: None, + }, + usage: Some(usage), + finish_reason: Some("length".to_string()), + raw: None, + resolved_model: None, + continue_turn: None, + served_from_cache: false, + correlation: None, + resolved_route: None, + } +} + #[cfg(test)] #[path = "model_call_failover_tests.rs"] mod failover_test; +#[cfg(test)] +#[path = "model_call_watchdog_tests.rs"] +mod watchdog_tests; + /// Retargets an attempt's wire-level `request.model` at a fallback binding, /// but only when the request already carried an explicit model. Registry names /// are runtime aliases, not guaranteed provider model ids, so a request that diff --git a/crates/tinyagents-harness/src/agent_loop/model_call_watchdog_tests.rs b/crates/tinyagents-harness/src/agent_loop/model_call_watchdog_tests.rs new file mode 100644 index 000000000..a2763d82a --- /dev/null +++ b/crates/tinyagents-harness/src/agent_loop/model_call_watchdog_tests.rs @@ -0,0 +1,297 @@ +use super::*; + +use std::collections::VecDeque; +use std::sync::Mutex; + +use async_trait::async_trait; +use tinyinference_llm::message::Message; +use tinyinference_llm::model::{ + ChatModel, ModelRequest, ModelResponse, ModelStream, ModelStreamItem, ReasoningConfig, + ReasoningEffort, +}; + +use crate::context::{RunConfig, RunContext}; +use crate::events::AgentEvent; +use crate::runtime::{AgentHarness, ReasoningWatchdog, RunPolicy}; +use crate::testkit::EventRecorder; + +/// A model that replays one scripted stream per call and records the +/// requests it was given, so a test can see what the loop sent after the +/// watchdog ended a call. +struct ScriptedStreams { + streams: Mutex>>, + requests: Mutex>, +} + +impl ScriptedStreams { + fn new(streams: Vec>) -> Self { + Self { + streams: Mutex::new(streams.into()), + requests: Mutex::new(Vec::new()), + } + } + + fn requests(&self) -> Vec { + self.requests.lock().unwrap().clone() + } + + fn next_items(&self, request: ModelRequest) -> Vec { + self.requests.lock().unwrap().push(request); + self.streams + .lock() + .unwrap() + .pop_front() + .expect("more model calls than scripted streams") + } +} + +#[async_trait] +impl ChatModel<()> for ScriptedStreams { + async fn invoke( + &self, + _state: &(), + request: ModelRequest, + ) -> tinyinference_llm::Result { + let mut accumulator = tinyinference_llm::model::StreamAccumulator::new(); + for item in self.next_items(request) { + accumulator.push(&item); + } + accumulator.finish() + } + + async fn stream( + &self, + _state: &(), + request: ModelRequest, + ) -> tinyinference_llm::Result { + let items = self.next_items(request); + Ok(ModelStream::new(Box::pin(futures::stream::iter(items)))) + } +} + +fn reasoning_delta(chars: usize) -> ModelStreamItem { + ModelStreamItem::MessageDelta(MessageDelta { + text: String::new(), + reasoning: "x".repeat(chars), + tool_call: None, + }) +} + +fn text_delta(text: &str) -> ModelStreamItem { + ModelStreamItem::MessageDelta(MessageDelta { + text: text.to_string(), + reasoning: String::new(), + tool_call: None, + }) +} + +/// A stream that reasons for `chars` characters, then answers `text`. +fn reasoning_then_text(chars: usize, text: &str) -> Vec { + let mut items = vec![ModelStreamItem::Started]; + for _ in 0..(chars / 100) { + items.push(reasoning_delta(100)); + } + items.push(text_delta(text)); + items.push(ModelStreamItem::Completed(ModelResponse::assistant(text))); + items +} + +fn harness_with(model: Arc, watchdog: ReasoningWatchdog) -> AgentHarness<()> { + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", model as _); + harness.with_policy(RunPolicy { + default_reasoning: Some(ReasoningConfig::effort(ReasoningEffort::High)), + reasoning_watchdog: watchdog, + ..RunPolicy::default() + }); + harness +} + +#[tokio::test] +async fn a_call_that_reasons_past_its_bound_with_nothing_visible_is_ended_as_a_dead_call() { + // 4,000 characters of reasoning is about 1,000 tokens; the bound is 200. + // The call is ended before its (late) answer arrives, and the loop's + // truncated-empty recovery retries with reasoning off and gets the + // second scripted stream. + let model = Arc::new(ScriptedStreams::new(vec![ + reasoning_then_text(4_000, "late answer"), + reasoning_then_text(0, "recovered"), + ])); + let harness = harness_with(Arc::clone(&model), ReasoningWatchdog::Tokens(200)); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("watchdog"), ()).with_events(recorder.sink()); + let run = harness + .invoke_streaming_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes on the recovered call"); + assert_eq!(run.text(), Some("recovered".to_string())); + let efforts: Vec> = model + .requests() + .iter() + .map(|r| r.reasoning.as_ref().and_then(|c| c.effort)) + .collect(); + assert_eq!( + efforts, + vec![Some(ReasoningEffort::High), Some(ReasoningEffort::None)], + "the ended call is a dead call: the retry goes out with reasoning off" + ); + assert!( + recorder.events().iter().any(|e| matches!( + e, + AgentEvent::ControlApplied { control, .. } if control == "reasoning_watchdog" + )), + "the watchdog is observable; got kinds {:?}", + recorder.kinds() + ); +} + +#[tokio::test] +async fn the_bound_counts_characters_not_utf8_bytes() { + // 240 CJK characters are 720 bytes. At three characters per token that + // is 80 tokens, under a 100-token bound; counted in bytes it would be + // 240 tokens and the watchdog would end a live call early. + let mut items = vec![ModelStreamItem::Started]; + for _ in 0..4 { + items.push(ModelStreamItem::MessageDelta(MessageDelta { + text: String::new(), + reasoning: "思".repeat(60), + tool_call: None, + })); + } + items.push(text_delta("answer")); + items.push(ModelStreamItem::Completed(ModelResponse::assistant( + "answer", + ))); + let model = Arc::new(ScriptedStreams::new(vec![items])); + let harness = harness_with(Arc::clone(&model), ReasoningWatchdog::Tokens(100)); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("watchdog-chars"), ()).with_events(recorder.sink()); + let run = harness + .invoke_streaming_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("answer".to_string())); + assert_eq!(model.requests().len(), 1); + assert!( + !recorder.events().iter().any(|e| matches!( + e, + AgentEvent::ControlApplied { control, .. } if control == "reasoning_watchdog" + )), + "no watchdog event; got kinds {:?}", + recorder.kinds() + ); +} + +#[tokio::test] +async fn visible_output_before_the_bound_disarms_the_watchdog() { + // Text arrives first, then a long think, then the answer: the model is + // answering, so the call runs to completion. + let mut items = vec![ModelStreamItem::Started, text_delta("Working: ")]; + for _ in 0..40 { + items.push(reasoning_delta(100)); + } + items.push(text_delta("answer")); + items.push(ModelStreamItem::Completed(ModelResponse::assistant( + "Working: answer", + ))); + let model = Arc::new(ScriptedStreams::new(vec![items])); + let harness = harness_with(Arc::clone(&model), ReasoningWatchdog::Tokens(200)); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("watchdog-disarmed"), ()).with_events(recorder.sink()); + let run = harness + .invoke_streaming_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("Working: answer".to_string())); + assert_eq!(model.requests().len(), 1); + assert!( + !recorder.events().iter().any(|e| matches!( + e, + AgentEvent::ControlApplied { control, .. } if control == "reasoning_watchdog" + )), + "no watchdog event; got kinds {:?}", + recorder.kinds() + ); +} + +#[tokio::test] +async fn the_request_budget_is_the_default_bound_and_no_budget_means_no_bound() { + // With the default policy the bound is the request's own + // `reasoning.budget_tokens`; a request without one is never ended. + let model = Arc::new(ScriptedStreams::new(vec![reasoning_then_text( + 4_000, + "slow but fine", + )])); + let harness = harness_with(Arc::clone(&model), ReasoningWatchdog::RequestBudget); + let ctx = RunContext::new(RunConfig::new("watchdog-unbounded"), ()); + let run = harness + .invoke_streaming_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("slow but fine".to_string())); + assert_eq!(model.requests().len(), 1); + + // The same stream under a request that carries a 200-token budget is + // ended, and the retry (reasoning off, so no budget) completes. + let model = Arc::new(ScriptedStreams::new(vec![ + reasoning_then_text(4_000, "late"), + reasoning_then_text(0, "recovered"), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", Arc::clone(&model) as _); + harness.with_policy(RunPolicy { + default_reasoning: Some(ReasoningConfig { + effort: Some(ReasoningEffort::High), + budget_tokens: Some(200), + summary: None, + }), + ..RunPolicy::default() + }); + let ctx = RunContext::new(RunConfig::new("watchdog-budget"), ()); + let run = harness + .invoke_streaming_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("recovered".to_string())); + assert_eq!(model.requests().len(), 2); +} + +#[tokio::test] +async fn the_watchdog_can_be_switched_off() { + let model = Arc::new(ScriptedStreams::new(vec![reasoning_then_text( + 4_000, + "late answer", + )])); + let harness = harness_with(Arc::clone(&model), ReasoningWatchdog::Off); + let ctx = RunContext::new(RunConfig::new("watchdog-off"), ()); + let run = harness + .invoke_streaming_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("late answer".to_string())); + assert_eq!(model.requests().len(), 1); +} + +#[tokio::test] +async fn the_watchdog_hands_the_interrupted_reasoning_to_the_retry() { + let model = Arc::new(ScriptedStreams::new(vec![ + reasoning_then_text(4_000, "late answer"), + reasoning_then_text(0, "recovered"), + ])); + let harness = harness_with(Arc::clone(&model), ReasoningWatchdog::Tokens(200)); + let ctx = RunContext::new(RunConfig::new("watchdog-carry"), ()); + let run = harness + .invoke_streaming_in_context(&(), ctx, vec![Message::user("hi")]) + .await + .expect("the run finishes"); + assert_eq!(run.text(), Some("recovered".to_string())); + let last = model.requests()[1] + .messages + .last() + .map(|m| m.text()) + .unwrap_or_default(); + assert!( + last.contains("Continue from this point") && last.contains("xxxxxxxx"), + "the streamed reasoning the watchdog cut off rides into the retry: {last:.100}" + ); +} diff --git a/crates/tinyagents-harness/src/agent_loop/reasoning_fallback.rs b/crates/tinyagents-harness/src/agent_loop/reasoning_fallback.rs new file mode 100644 index 000000000..7fc425a89 --- /dev/null +++ b/crates/tinyagents-harness/src/agent_loop/reasoning_fallback.rs @@ -0,0 +1,100 @@ +//! Reasoning fallback after a dead call. +//! +//! A hosted reasoning model can spend its whole output cap on the hidden +//! reasoning channel and return `finish_reason = "length"` with no text and +//! no tool call. Measured on deepseek-v4.1-flash through OpenRouter's routable +//! providers, neither a smaller output cap nor a lower effort label stops it: +//! one task died nine times in a row at caps from 65k down to 2k, and at +//! `medium` effort 11 of 25 calls still died. The one control that reliably +//! produced zero reasoning tokens was `effort = none`. +//! +//! So after a dead call the loop re-issues the step with reasoning switched +//! off. The model then has to act from what it already knows, and the nudge +//! that accompanies a repeated death tells it to do its working-out in the +//! workspace (a scratch file, a small experiment) instead of in its head. The +//! hold-off backs off without a ceiling: the first death costs one call +//! without reasoning, the next two, then four, eight, sixteen. Reasoning +//! always comes back: a run that kept it off for good after a few deaths +//! spent the rest of its budget writing probe after probe (seventeen in two +//! minutes on one task), which is what a model without reasoning does with +//! a hard step. With the watchdog bounding each death to its budget, a +//! return to reasoning every doubling stretch costs little and is where the +//! hard steps get solved. +//! +//! The state lives on [`super::types::TurnRecovery`] but, unlike the +//! counters there, is run-wide: the hold-off outlives the turn that set it, +//! and none of the turn-boundary resets touch it. + +use tinyinference_llm::model::{ModelRequest, ReasoningConfig, ReasoningEffort}; + +/// Per-run state of the fallback. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub(super) struct ReasoningFallback { + /// Live calls still to go out without reasoning. + holdoff: u32, + /// Hold-off the next dead call will set; doubles per death. + scale: u32, +} + +impl ReasoningFallback { + /// Records a dead call: the next `holdoff` calls go out without reasoning. + /// Returns the hold-off just set. + pub(super) fn on_dead_call(&mut self) -> u32 { + let scale = self.scale.max(1); + self.holdoff = scale; + self.scale = scale.saturating_mul(2); + self.holdoff + } + + /// Records a reply with content or a tool call: one hold-off call spent. + pub(super) fn on_live_reply(&mut self) { + self.holdoff = self.holdoff.saturating_sub(1); + } + + /// A repeat was just noted on a tool result while reasoning is off: the + /// model is looping without it. Hands reasoning back for the next call + /// and returns whether anything changed. The backoff scale stands, so a + /// dead call after this re-engages the fallback at its current stretch; + /// the run then alternates between thinking and acting instead of doing + /// only one. + pub(super) fn on_repeat_note(&mut self) -> bool { + if !self.active() { + return false; + } + self.holdoff = 0; + true + } + + /// Whether the next call goes out without reasoning. + pub(super) fn active(&self) -> bool { + self.holdoff > 0 + } + + /// Live calls still to go out without reasoning. + pub(super) fn holdoff(&self) -> u32 { + self.holdoff + } + + /// Switches reasoning off on `request` while the fallback is active. + /// Returns what the request asked for before, when it was changed: a + /// request that already had reasoning off is left alone. + pub(super) fn apply(&self, request: &mut ModelRequest) -> Option> { + if !self.active() { + return None; + } + if matches!( + request.reasoning.as_ref().and_then(|r| r.effort), + Some(ReasoningEffort::None) + ) { + return None; + } + let previous = request + .reasoning + .replace(ReasoningConfig::effort(ReasoningEffort::None)); + Some(previous) + } +} + +#[cfg(test)] +#[path = "reasoning_fallback_tests.rs"] +mod tests; diff --git a/crates/tinyagents-harness/src/agent_loop/reasoning_fallback_tests.rs b/crates/tinyagents-harness/src/agent_loop/reasoning_fallback_tests.rs new file mode 100644 index 000000000..a5cbd220a --- /dev/null +++ b/crates/tinyagents-harness/src/agent_loop/reasoning_fallback_tests.rs @@ -0,0 +1,116 @@ +use super::*; + +fn request_with(reasoning: Option) -> ModelRequest { + ModelRequest { + reasoning, + ..ModelRequest::default() + } +} + +#[test] +fn inactive_until_a_dead_call() { + let fallback = ReasoningFallback::default(); + assert!(!fallback.active()); + let mut request = request_with(Some(ReasoningConfig::effort(ReasoningEffort::High))); + assert_eq!(fallback.apply(&mut request), None); + assert_eq!( + request.reasoning, + Some(ReasoningConfig::effort(ReasoningEffort::High)) + ); +} + +#[test] +fn first_death_switches_reasoning_off_for_one_call() { + let mut fallback = ReasoningFallback::default(); + assert_eq!(fallback.on_dead_call(), 1); + let mut request = request_with(Some(ReasoningConfig::effort(ReasoningEffort::High))); + assert_eq!( + fallback.apply(&mut request), + Some(Some(ReasoningConfig::effort(ReasoningEffort::High))) + ); + assert_eq!( + request.reasoning, + Some(ReasoningConfig::effort(ReasoningEffort::None)) + ); + fallback.on_live_reply(); + assert!( + !fallback.active(), + "one live reply spends a hold-off of one" + ); +} + +#[test] +fn hold_off_doubles_per_death_and_always_returns() { + let mut fallback = ReasoningFallback::default(); + assert_eq!(fallback.on_dead_call(), 1); + assert_eq!(fallback.on_dead_call(), 2); + assert_eq!(fallback.on_dead_call(), 4); + assert_eq!(fallback.on_dead_call(), 8); + assert_eq!(fallback.on_dead_call(), 16); + for _ in 0..16 { + assert!(fallback.active()); + fallback.on_live_reply(); + } + assert!( + !fallback.active(), + "every hold-off is spent by live replies; reasoning always returns" + ); +} + +#[test] +fn a_request_with_no_reasoning_config_is_switched_off_too() { + // A dead call is itself the evidence that the provider reasons on this + // model, whatever the request said. + let mut fallback = ReasoningFallback::default(); + fallback.on_dead_call(); + let mut request = request_with(None); + assert_eq!(fallback.apply(&mut request), Some(None)); + assert_eq!( + request.reasoning, + Some(ReasoningConfig::effort(ReasoningEffort::None)) + ); +} + +#[test] +fn a_request_already_without_reasoning_is_left_alone() { + let mut fallback = ReasoningFallback::default(); + fallback.on_dead_call(); + let mut request = request_with(Some(ReasoningConfig::effort(ReasoningEffort::None))); + assert_eq!(fallback.apply(&mut request), None); +} + +#[test] +fn a_budget_only_config_is_replaced() { + let mut fallback = ReasoningFallback::default(); + fallback.on_dead_call(); + let mut request = request_with(Some(ReasoningConfig { + effort: None, + budget_tokens: Some(4096), + summary: None, + })); + assert!(fallback.apply(&mut request).is_some()); + assert_eq!( + request.reasoning, + Some(ReasoningConfig::effort(ReasoningEffort::None)) + ); +} + +#[test] +fn a_repeat_note_hands_reasoning_back_mid_hold_off() { + let mut fallback = ReasoningFallback::default(); + assert!( + !fallback.on_repeat_note(), + "nothing to restore while reasoning is on" + ); + for _ in 0..4 { + fallback.on_dead_call(); + } + assert!( + fallback.active(), + "eight calls of hold-off after the fourth death" + ); + assert!(fallback.on_repeat_note()); + assert!(!fallback.active(), "reasoning is back for the next call"); + // The backoff scale stands: the next death holds off for sixteen. + assert_eq!(fallback.on_dead_call(), 16); +} diff --git a/crates/tinyagents-harness/src/agent_loop/response_recovery.rs b/crates/tinyagents-harness/src/agent_loop/response_recovery.rs index 80061049d..7ac4b00eb 100644 --- a/crates/tinyagents-harness/src/agent_loop/response_recovery.rs +++ b/crates/tinyagents-harness/src/agent_loop/response_recovery.rs @@ -12,7 +12,9 @@ //! empty reply, a dropped or undecodable call): re-issue or re-prompt. use super::run_loop::{ - DROPPED_TOOL_CALL_NUDGE, TRUNCATED_EMPTY_ANSWER_NUDGE, TRUNCATED_EMPTY_TOOL_NUDGE, + DROPPED_TOOL_CALL_NUDGE, TRUNCATED_EMPTY_ANSWER_NUDGE, TRUNCATED_EMPTY_CARRY_PREFIX, + TRUNCATED_EMPTY_CARRY_SUFFIX, TRUNCATED_EMPTY_REASONING_OFF_ANSWER_NOTE, + TRUNCATED_EMPTY_REASONING_OFF_TOOL_NOTE, TRUNCATED_EMPTY_TOOL_NUDGE, UNDECODABLE_TOOL_CALL_NUDGE, WITHHELD_TOOL_CALL_NUDGE, truncated_call_positions, }; use super::turn_recovery::TRUNCATED_CLOCK_NUDGE_LIMIT; @@ -229,6 +231,21 @@ 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(); + if truncated_empty && self.policy.truncated_empty_reasoning_fallback { + // Whatever follows (a retry, a nudge, or nothing), the next calls + // go out without reasoning: the cap and the effort label have + // both been shown not to stop this. + let holdoff = turn_recovery.reasoning_fallback.on_dead_call(); + ctx.emit(AgentEvent::ControlApplied { + control: "reasoning_fallback".to_string(), + detail: format!( + "model call `{call_id}` died at its output cap with nothing to show; \ + the next {holdoff} call(s) go out with reasoning switched off" + ), + }); + } + let reasoning_off = self.policy.truncated_empty_reasoning_fallback + && turn_recovery.reasoning_fallback.active(); // 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 @@ -248,7 +265,13 @@ impl AgentHarness { .map(std::time::Duration::from_millis), ) .map(|clock| clock.remaining()); - turn_recovery.truncated_retry_plan(attempt_max_tokens, dead_tokens, dead_ms, remaining) + turn_recovery.truncated_retry_plan( + attempt_max_tokens, + dead_tokens, + dead_ms, + remaining, + reasoning_off, + ) }); if truncated_empty && turn_recovery.truncated_empty_retries_used < self.policy.truncated_empty_retries @@ -256,9 +279,13 @@ impl AgentHarness { && 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. + // retry re-sends the identical transcript, plus the dead call's + // own working-out where it had any: that reasoning was real + // progress (a correct derivation, cut off), and without it every + // retry starts the same derivation over. messages.pop(); ctx.retract_transcript(messages.len()); + self.carry_dead_call_reasoning(ctx, messages, call_id, response); 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 { @@ -334,13 +361,14 @@ impl AgentHarness { { messages.pop(); ctx.retract_transcript(messages.len()); + self.carry_dead_call_reasoning(ctx, messages, call_id, response); turn_recovery.truncated_empty_nudges_used += 1; 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()); + .is_some_and(|plan| plan.at_ceiling()); let repeat_cap = (turn_recovery.truncated_empty_nudges_used > 1 || clock_only_nudge || retry_was_skipped_for_clock @@ -350,7 +378,7 @@ impl AgentHarness { if let Some(cap) = repeat_cap { turn_recovery.boosted_max_tokens = Some(cap); } - let nudge: String = match repeat_cap { + let mut 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, \ @@ -364,6 +392,16 @@ impl AgentHarness { None if tools_available_this_turn => TRUNCATED_EMPTY_TOOL_NUDGE.to_string(), None => TRUNCATED_EMPTY_ANSWER_NUDGE.to_string(), }; + if reasoning_off { + // The deliberation the model cannot finish in its head goes + // into the workspace instead. + nudge.push(' '); + nudge.push_str(if tools_available_this_turn { + TRUNCATED_EMPTY_REASONING_OFF_TOOL_NOTE + } else { + TRUNCATED_EMPTY_REASONING_OFF_ANSWER_NOTE + }); + } tracing::info!( target: "tinyagents::agent_loop", run_id = %ctx.run_id(), @@ -487,4 +525,77 @@ impl AgentHarness { false } + + /// Appends the dead call's interrupted reasoning to the transcript as a + /// user message ahead of the retry or nudged call (see + /// [`dead_call_reasoning_carry`]), announcing it when it does. + fn carry_dead_call_reasoning( + &self, + ctx: &mut RunContext, + messages: &mut Vec, + call_id: &CallId, + response: &ModelResponse, + ) { + if let Some((carry, kept_chars)) = + dead_call_reasoning_carry(response, self.policy.truncated_empty_carry_reasoning_chars) + { + ctx.emit(AgentEvent::ControlApplied { + control: "truncated_empty_reasoning_carried".to_string(), + detail: format!( + "model call `{call_id}`: {kept_chars} chars of its interrupted reasoning carried into the transcript" + ), + }); + messages.push(Message::user(carry)); + } + } +} + +/// The dead call's reasoning, framed for the transcript, when the response +/// carries any and the policy keeps it. The *last* `limit` characters are +/// kept: a derivation's state of play is at its end, and its start is the +/// part a fresh call re-derives fastest. `None` for a response with no +/// reasoning, a limit of zero, or reasoning too short to be worth a message. +/// Returns the framed message and the number of reasoning characters it +/// keeps (the framing excluded). +fn dead_call_reasoning_carry(response: &ModelResponse, limit: usize) -> Option<(String, usize)> { + const MIN_CARRY_CHARS: usize = 200; + if limit == 0 { + return None; + } + let reasoning: String = response + .message + .content + .iter() + .filter_map(|block| match block { + tinyinference_llm::message::ContentBlock::Thinking { text, .. } => Some(text.as_str()), + _ => None, + }) + .collect(); + let reasoning = reasoning.trim(); + // Both bounds are in characters, as the policy field is documented, not + // in UTF-8 bytes. + let total_chars = reasoning.chars().count(); + if total_chars < MIN_CARRY_CHARS { + return None; + } + let tail = if total_chars > limit { + // Keep the last `limit` characters, then cut on a line boundary where + // one is near, so the excerpt does not open mid-word. + let start = reasoning + .char_indices() + .nth(total_chars - limit) + .map_or(0, |(index, _)| index); + let excerpt = &reasoning[start..]; + match excerpt.find('\n') { + Some(nl) if nl < 200 => &excerpt[nl + 1..], + _ => excerpt, + } + } else { + reasoning + }; + let kept_chars = tail.chars().count(); + Some(( + format!("{TRUNCATED_EMPTY_CARRY_PREFIX}[…]\n{tail}{TRUNCATED_EMPTY_CARRY_SUFFIX}"), + kept_chars, + )) } diff --git a/crates/tinyagents-harness/src/agent_loop/run_loop.rs b/crates/tinyagents-harness/src/agent_loop/run_loop.rs index 4643eda4c..43b590bf1 100644 --- a/crates/tinyagents-harness/src/agent_loop/run_loop.rs +++ b/crates/tinyagents-harness/src/agent_loop/run_loop.rs @@ -635,6 +635,44 @@ impl AgentHarness { { request.reasoning = Some(mapped.clone()); } + // A dead call earlier in the run switched reasoning off for the + // next few calls (see `RunPolicy::truncated_empty_reasoning_fallback`). + // Applied last so it wins over the policy default and the profile + // mapping: those describe the effort the run wants, this is the + // one the transcript can get past. + // A repeat noted on the last tool result while reasoning is off + // (`RunContext::note_repeat`) means the model is looping without + // it, and the finish check (`RunContext::request_reasoning`) is + // the one call worth a dead call's bounded cost: hand reasoning + // back for this call. + if ctx.take_repeat_noted() + && self.policy.truncated_empty_reasoning_fallback + && turn_recovery.reasoning_fallback.on_repeat_note() + { + tracing::info!( + target: "tinyagents::agent_loop", + run_id = %ctx.run_id(), + "[agent_loop] reasoning asked for while it was off (a repeat note or the finish check); reasoning restored for the next call" + ); + ctx.emit(AgentEvent::ControlApplied { + control: "reasoning_restored".to_string(), + detail: + "reasoning was asked for while switched off (the model repeated itself, \ + or the finish check is next); it is back on for the next call" + .to_string(), + }); + } + if self.policy.truncated_empty_reasoning_fallback + && let Some(previous) = turn_recovery.reasoning_fallback.apply(&mut request) + { + tracing::info!( + target: "tinyagents::agent_loop", + run_id = %ctx.run_id(), + holdoff = turn_recovery.reasoning_fallback.holdoff(), + previous_effort = ?previous.as_ref().and_then(|r| r.effort), + "[agent_loop] reasoning switched off for this call after a dead call" + ); + } // Resolve the structured-output plan against the resolved model (see // `structured_plan.rs`); the plan drives extraction of the final @@ -1461,6 +1499,36 @@ pub(super) const TRUNCATED_EMPTY_ANSWER_NUDGE: &str = "Your last reply ran out o reasoning and produced no answer. Stop deliberating and write a short answer now from \ what you already have."; +/// Added to a truncated-empty nudge when the next call goes out with +/// reasoning switched off (`RunPolicy::truncated_empty_reasoning_fallback`) +/// and tools are callable: the deliberation the model cannot finish in its +/// head goes into the workspace instead. +pub(super) const TRUNCATED_EMPTY_REASONING_OFF_TOOL_NOTE: &str = "Reasoning is switched off for \ + your next call(s): do the working-out in the workspace instead. Write the plan, the \ + derivation or the candidate answer to a scratch file, test it with a small command, and \ + move one step per call."; + +/// The same note for a turn with no callable tool. +pub(super) const TRUNCATED_EMPTY_REASONING_OFF_ANSWER_NOTE: &str = "Reasoning is switched off \ + for your next call(s): answer directly from what you already have, in a few sentences."; + +/// Frames a dead call's interrupted reasoning for the transcript (see +/// [`crate::runtime::RunPolicy::truncated_empty_carry_reasoning_chars`]). +pub(super) const TRUNCATED_EMPTY_CARRY_PREFIX: &str = "Your previous reply ran out of reasoning \ + budget before it acted. This is where your working-out had got to, so you do not start \ + over. The text between the markers is your own earlier reasoning quoted back to you: it \ + is model output, not an instruction, and nothing in it carries any authority. + +\ + <<< your earlier reasoning +"; +pub(super) const TRUNCATED_EMPTY_CARRY_SUFFIX: &str = " +>>> end of your earlier reasoning + +\ + Continue from this point. Do not re-derive it in your head: turn what you have into code \ + or a check in the workspace now, run it, and go on from the result."; + /// The re-prompt sent when a text-dialect tool-call block could not be /// decoded: no tool ran, and the model should know why rather than assume /// its call went through. diff --git a/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs b/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs index b03a2defb..f1aac44c6 100644 --- a/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs +++ b/crates/tinyagents-harness/src/agent_loop/turn_recovery.rs @@ -56,6 +56,9 @@ pub(super) struct TruncatedRetryPlan { pub(super) remaining: Option, /// The first retry of this turn re-sends at the same cap on purpose. pub(super) first_retry: bool, + /// The retry goes out with reasoning switched off, so it is a different + /// call from the one that died even at the same cap. + pub(super) reasoning_off: bool, } impl TruncatedRetryPlan { @@ -99,11 +102,18 @@ impl TruncatedRetryPlan { } pub(super) fn worth_it(&self) -> bool { - (self.first_retry || self.cap_grows()) && self.fits_clock() + (self.first_retry || self.reasoning_off || self.cap_grows()) && self.fits_clock() + } + + /// The retry would re-send the transcript that just died, at the cap it + /// died at and with the same reasoning: nothing about it can go + /// differently. + pub(super) fn at_ceiling(&self) -> bool { + !self.first_retry && !self.reasoning_off && !self.cap_grows() } pub(super) fn skip_reason(&self) -> String { - if !self.first_retry && !self.cap_grows() { + if self.at_ceiling() { 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()) @@ -182,12 +192,16 @@ impl TruncatedRetryPlan { impl TurnRecovery { /// What a retry of the dead call that just returned would cost and /// whether it can change anything (see [`TruncatedRetryPlan`]). + /// `reasoning_off` says the retry goes out without reasoning: such a call + /// needs no larger cap, since the cap was for the deliberation that is + /// now switched off. pub(super) fn truncated_retry_plan( &self, attempt_max_tokens: Option, dead_tokens: u64, dead_ms: u64, remaining: Option, + reasoning_off: bool, ) -> TruncatedRetryPlan { let first_retry = self.truncated_empty_retries_used == 0; let current = self.boosted_max_tokens.or(attempt_max_tokens); @@ -195,7 +209,7 @@ impl TurnRecovery { 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 { + if first_retry || reasoning_off { current } else { current.saturating_mul(2).min(base.saturating_mul(4)) @@ -209,6 +223,7 @@ impl TurnRecovery { dead_ms, remaining, first_retry, + reasoning_off, } } @@ -273,20 +288,28 @@ impl TurnRecovery { /// A turn whose call was cut off by the output limit /// (`turn_had_truncated_calls`) keeps its truncated-tool-call retry budget /// and boosted output cap for the retry. + /// + /// A tool call is a live reply: it spends one call of the reasoning + /// fallback's hold-off. pub(super) fn reset_after_tool_turn(&mut self, turn_had_truncated_calls: bool) { self.reset_nudges(); if !turn_had_truncated_calls { self.reset_truncation(); } + self.reasoning_fallback.on_live_reply(); } /// The turn resolved without a tool call and without scheduling another /// retry (it is about to be taken as the answer): clear every counter and /// the boosted cap, which would otherwise override the caller's per-turn /// cap on every later call. + /// + /// An answer is a live reply: it spends one call of the reasoning + /// fallback's hold-off. pub(super) fn reset_after_final(&mut self) { self.reset_nudges(); self.reset_truncation(); + self.reasoning_fallback.on_live_reply(); } } 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 ae6532946..825534fe3 100644 --- a/crates/tinyagents-harness/src/agent_loop/turn_recovery_tests.rs +++ b/crates/tinyagents-harness/src/agent_loop/turn_recovery_tests.rs @@ -10,9 +10,30 @@ fn spent() -> TurnRecovery { truncated_tool_call_retries_used: 2, boosted_max_tokens: Some(4096), truncation_base: Some(1024), + reasoning_fallback: super::super::reasoning_fallback::ReasoningFallback::default(), } } +#[test] +fn the_reasoning_hold_off_outlives_the_turn_and_is_spent_by_live_replies() { + // The fallback is run-wide: a reset at a turn boundary does not clear + // it, but the resolved turn it marks is a live reply and spends one + // call of the hold-off. + let mut recovery = spent(); + recovery.reasoning_fallback.on_dead_call(); + recovery.reasoning_fallback.on_dead_call(); + assert_eq!(recovery.reasoning_fallback.holdoff(), 2); + recovery.reset_after_tool_turn(false); + assert!(recovery.reasoning_fallback.active()); + assert_eq!(recovery.reasoning_fallback.holdoff(), 1); + recovery.reset_after_final(); + assert!( + !recovery.reasoning_fallback.active(), + "two resolved turns spend a hold-off of two" + ); + assert_eq!(recovery.truncated_empty_retries_used, 0); +} + #[test] fn boost_doubles_the_last_sent_cap_and_clamps_at_four_times_the_base() { let mut recovery = TurnRecovery::default(); @@ -47,6 +68,7 @@ fn nudge_cap_is_clamped_to_the_clock_affordable_cap() { dead_ms: 16_000, remaining: Some(std::time::Duration::from_millis(6_000)), first_retry: false, + reasoning_off: false, }; assert_eq!(plan.affordable_cap(), Some(3_000)); @@ -64,6 +86,7 @@ fn nudge_is_rejected_when_even_the_minimum_cap_exceeds_the_clock_budget() { dead_ms: 4_096, remaining: Some(std::time::Duration::from_millis(2_000)), first_retry: false, + reasoning_off: false, }; assert_eq!(plan.affordable_cap(), None); @@ -81,6 +104,7 @@ fn uncapped_nudge_uses_the_dead_call_duration_as_its_clock_estimate() { dead_ms: 1_000, remaining: Some(std::time::Duration::from_millis(3_000)), first_retry: false, + reasoning_off: false, }; let too_little_time = TruncatedRetryPlan { remaining: Some(std::time::Duration::from_millis(1_000)), @@ -101,6 +125,7 @@ fn retry_clock_estimate_scales_with_the_candidate_cap_without_usage() { dead_ms: 1_000, remaining: Some(std::time::Duration::from_millis(10_000)), first_retry: false, + reasoning_off: false, }; assert_eq!(plan.expected_ms(), 2_000); diff --git a/crates/tinyagents-harness/src/agent_loop/types.rs b/crates/tinyagents-harness/src/agent_loop/types.rs index fdb1f9bbd..d7e74b1a7 100644 --- a/crates/tinyagents-harness/src/agent_loop/types.rs +++ b/crates/tinyagents-harness/src/agent_loop/types.rs @@ -99,7 +99,9 @@ pub struct AgentLoopResult { /// The recovery counters and boosted output cap of the turn in flight. /// /// Every counter is consecutive-per-turn, not per-run (the output-validation -/// retry budget, which is run-wide, lives outside this struct). +/// retry budget, which is run-wide, lives outside this struct). The one +/// run-wide member is [`Self::reasoning_fallback`]: its hold-off outlives the +/// turn that set it and no turn-boundary reset touches it. #[derive(Debug, Clone, Default, PartialEq, Eq)] pub(super) struct TurnRecovery { /// Retries of a length-truncated empty reply @@ -125,6 +127,10 @@ pub(super) struct TurnRecovery { pub(super) boosted_max_tokens: Option, /// The original output cap, so growth stays clamped at 4x. pub(super) truncation_base: Option, + /// Reasoning switched off for the calls after a dead one (see + /// `RunPolicy::truncated_empty_reasoning_fallback`). Run-wide: the + /// hold-off outlives the turn that set it. + pub(super) reasoning_fallback: super::reasoning_fallback::ReasoningFallback, } /// The tool schemas a run advertises and the discovery state behind them. diff --git a/crates/tinyagents-harness/src/context/mod.rs b/crates/tinyagents-harness/src/context/mod.rs index 17e608930..3b02c122b 100644 --- a/crates/tinyagents-harness/src/context/mod.rs +++ b/crates/tinyagents-harness/src/context/mod.rs @@ -313,6 +313,7 @@ impl RunContext { run_queue: None, cancellation: CancellationToken::new(), control: std::sync::Arc::new(std::sync::Mutex::new(None)), + repeat_noted: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)), state_updates: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())), tool_state_updates: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())), terminate_votes: Vec::new(), @@ -619,6 +620,35 @@ impl RunContext { self.control.lock().ok().and_then(|mut guard| guard.take()) } + /// Records that a middleware just noted a repeat on a tool result (an + /// identical call re-issued, an identical reply). The agent loop reads it + /// once before its next model call ([`Self::take_repeat_noted`]); a + /// model running without reasoning that repeats itself is handed + /// reasoning back. + pub fn note_repeat(&self) { + self.repeat_noted + .store(true, std::sync::atomic::Ordering::Relaxed); + } + + /// Asks for reasoning on the next model call while the loop's reasoning + /// fallback has switched it off. Same flag as [`Self::note_repeat`]: a + /// middleware about to issue a call where thinking is worth a dead call's + /// bounded cost (the finish check, which has to ask what the request + /// implied) uses this. With the fallback disabled reasoning is never off, + /// so the request is a no-op by construction; it never raises the effort + /// above what the request already asks for. + pub fn request_reasoning(&self) { + self.repeat_noted + .store(true, std::sync::atomic::Ordering::Relaxed); + } + + /// Whether a repeat was noted (or reasoning requested) since the last + /// take; clears it. + pub fn take_repeat_noted(&self) -> bool { + self.repeat_noted + .swap(false, std::sync::atomic::Ordering::Relaxed) + } + /// Queues a [`StateUpdate`] for the host to apply. /// /// The agent loop only ever holds `state: &State` (a shared reference), so diff --git a/crates/tinyagents-harness/src/context/types.rs b/crates/tinyagents-harness/src/context/types.rs index 47aa77901..5bb120292 100644 --- a/crates/tinyagents-harness/src/context/types.rs +++ b/crates/tinyagents-harness/src/context/types.rs @@ -423,6 +423,13 @@ pub struct RunContext { /// loop (stop with a final response, or interrupt). Drained by the agent /// loop at its safe checkpoints via [`RunContext::take_control`]. pub control: std::sync::Arc>>, + /// Set by a middleware that just noted a repeat on a tool result (an + /// identical call re-issued, an identical reply) and read once by the + /// agent loop before its next model call. The loop's reasoning fallback + /// treats it as the signal that a model running without reasoning is + /// looping, and hands reasoning back for the next call. See + /// [`RunContext::note_repeat`]. + pub repeat_noted: std::sync::Arc, /// Queued [`StateUpdate`]s a middleware or tool requested via /// [`MiddlewareControl::UpdateState`], drained by a host through /// [`RunContext::take_state_updates`]. See that method's docs for why the diff --git a/crates/tinyagents-harness/src/middleware/library/repeat_progress.rs b/crates/tinyagents-harness/src/middleware/library/repeat_progress.rs index df01af13c..373cd07bb 100644 --- a/crates/tinyagents-harness/src/middleware/library/repeat_progress.rs +++ b/crates/tinyagents-harness/src/middleware/library/repeat_progress.rs @@ -432,6 +432,10 @@ impl Middleware<(), C> for RepeatProgressMiddleware { "[tinyagents::mw] repeat-progress appended a warning to the tool result" ); append_note(result, ¬e); + // The loop's reasoning fallback reads this before the next + // call: a model repeating itself without reasoning gets + // reasoning back. + ctx.note_repeat(); } } if already_halted { diff --git a/crates/tinyagents-harness/src/middleware/library/repeat_progress_tests.rs b/crates/tinyagents-harness/src/middleware/library/repeat_progress_tests.rs index 5ec882b31..699af1033 100644 --- a/crates/tinyagents-harness/src/middleware/library/repeat_progress_tests.rs +++ b/crates/tinyagents-harness/src/middleware/library/repeat_progress_tests.rs @@ -337,3 +337,38 @@ async fn distinct_hex_only_results_do_not_halt_on_recurrence() { "three different commit ids are three different results" ); } + +#[tokio::test] +async fn a_repeat_warning_is_noted_on_the_run_context() { + // The staged guard warns on a result before it blocks or halts; the + // warning is also flagged on the run context so the agent loop can hand + // reasoning back to a model that repeats itself without it. + let handle = SteeringHandle::allow_all(); + let summary = Arc::new(std::sync::Mutex::new(None)); + let mw = RepeatProgressMiddleware::new(handle, summary, exempt()); + let mut ctx = ctx(); + assert!(!ctx.take_repeat_noted(), "nothing noted before any call"); + let mut noted_at = None; + for cycle in 1..=6 { + let mut response = repeated_success_response("use_skill", json!({"skill": "skills"})); + mw.after_model(&mut ctx, &(), &mut response).await.unwrap(); + let mut result = TaToolResult::success("doc"); + let invocation = ToolInvocationIdentity::new("repeat-1", "use_skill"); + mw.after_tool(&mut ctx, &(), &invocation, &mut result) + .await + .unwrap(); + if ctx.take_repeat_noted() { + assert!( + result.output().contains("identical result"), + "the flag is set on the result that carries the note: {}", + result.output() + ); + noted_at = Some(cycle); + break; + } + } + assert!( + noted_at.is_some_and(|cycle| cycle > 1), + "a repeat warning sets the flag once the repeat is noted: {noted_at:?}" + ); +} diff --git a/crates/tinyagents-harness/src/middleware/library/verify_before_finish/README.md b/crates/tinyagents-harness/src/middleware/library/verify_before_finish/README.md index 383172265..677434c88 100644 --- a/crates/tinyagents-harness/src/middleware/library/verify_before_finish/README.md +++ b/crates/tinyagents-harness/src/middleware/library/verify_before_finish/README.md @@ -10,7 +10,12 @@ and per-run state. `middleware.rs` contains configuration and lifecycle hooks; `middleware_tests.rs` exercises those hooks and the run loop. The check is skipped for tool-bearing, empty, truncated, or already continued -responses and when call or wall-clock budget is too small. Successful and +responses and when call or wall-clock budget is too small. The check asks the +agent loop for reasoning on the call it holds the answer for +(`RunContext::request_reasoning`, a no-op unless the loop's fallback has switched +reasoning off after dead calls): a result fitted on the wrong axis is caught +by asking what the request implied, which a model running without reasoning +(the loop's fallback after dead calls) does not do. Successful and failed runs release activity through lifecycle hooks. Deferred runs put their tool activity and check status in `DeferredToolRequests::resume_metadata`. Hosts constructing `DeferredToolResults` manually must copy it with diff --git a/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware.rs b/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware.rs index 4bd9bec7d..2c3927b7d 100644 --- a/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware.rs +++ b/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware.rs @@ -254,6 +254,11 @@ impl Middleware for VerifyBeforeFinishMidd let tool_rounds = run.activity.tool_rounds; drop(runs); response.continue_turn = Some(self.check.clone()); + // The check is the one call in a run where thinking is worth the + // bounded risk of a dead call: a result fitted on the wrong axis is + // caught by asking what the request implied, which a model running + // without reasoning (the fallback after dead calls) does not do. + ctx.request_reasoning(); tracing::info!( tool_rounds, remaining_model_calls = ctx.limits.remaining_model_calls(), diff --git a/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware_tests.rs b/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware_tests.rs index 96aea7cf9..739c3971b 100644 --- a/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware_tests.rs +++ b/crates/tinyagents-harness/src/middleware/library/verify_before_finish/middleware_tests.rs @@ -564,3 +564,68 @@ async fn a_truncated_response_with_no_output_is_not_reported_as_an_empty_answer( let fine = answer("a real answer"); assert_eq!(mw.skip_reason(&ctx, &fine), None); } + +/// A call that died at its output cap with nothing to show. +fn dead_call() -> ModelResponse { + let mut response = ModelResponse::assistant(String::new()).with_finish_reason("length"); + response.message.content = Vec::new(); + response +} + +/// The check runs with reasoning on even while the loop's reasoning +/// fallback has switched it off after dead calls: the middleware asks for +/// it on the call it holds the answer for. +#[tokio::test] +async fn the_check_asks_for_reasoning_back() { + use tinyinference_llm::model::{ReasoningConfig, ReasoningEffort}; + // Three deaths hold reasoning off for four live calls; the tool round + // and the draft spend two of them, so the check would otherwise go out + // without reasoning. + let model = Arc::new(ScriptedModel::new(vec![ + dead_call(), + dead_call(), + dead_call(), + tool_round("c1", "lookup"), + answer("draft"), + answer("checked"), + ])); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model("mock", model.clone()); + harness.register_tool(Arc::new(FakeTool::returning("lookup", "found"))); + harness.push_middleware(Arc::new(VerifyBeforeFinishMiddleware::new(CHECK))); + harness.with_policy(RunPolicy { + default_reasoning: Some(ReasoningConfig::effort(ReasoningEffort::High)), + limits: RunLimits::default().with_max_model_calls(10), + ..RunPolicy::default() + }); + let run = harness + .invoke_default(&(), vec![Message::user("do the task")]) + .await + .expect("run succeeds"); + + assert_eq!(run.text(), Some("checked".to_string())); + assert_eq!(check_count(&run), 1); + let efforts: Vec> = model + .requests() + .iter() + .map(|r| r.reasoning.as_ref().and_then(|c| c.effort)) + .collect(); + assert_eq!( + efforts, + vec![ + Some(ReasoningEffort::High), + Some(ReasoningEffort::None), + Some(ReasoningEffort::None), + Some(ReasoningEffort::None), + Some(ReasoningEffort::None), + Some(ReasoningEffort::High), + ], + "every call after the first death runs without reasoning except the check" + ); + let last = model.requests().last().expect("six requests").clone(); + assert_eq!( + last.messages.last().map(Message::text), + Some(CHECK.to_string()), + "the call with reasoning is the check" + ); +} diff --git a/crates/tinyagents-harness/src/runtime/types.rs b/crates/tinyagents-harness/src/runtime/types.rs index 3ef994187..61d01119b 100644 --- a/crates/tinyagents-harness/src/runtime/types.rs +++ b/crates/tinyagents-harness/src/runtime/types.rs @@ -351,6 +351,63 @@ pub struct RunPolicy { /// Defaults to `1`. Set to `0` (together with `truncated_empty_retries = /// 0`) for exact-replay callers that must not re-issue a call. pub truncated_empty_nudges: u32, + /// After a truncated-empty completion, send the retry or nudged call with + /// reasoning switched off (`reasoning.effort = none`). + /// + /// On the hosted providers measured (deepseek-v4.1-flash through + /// OpenRouter's routable backends) neither a smaller output cap nor a + /// lower effort label stops a model that deliberates past its cap: one + /// task died nine times in a row at caps from 65k down to 2k, and at + /// `medium` effort 11 of 25 calls still died. `effort = none` was the one + /// control that produced zero reasoning tokens. So the step is re-issued + /// without reasoning, and the model has to act from what it already + /// knows; the nudge tells it to do its working-out in the workspace. The + /// hold-off backs off (1, 2, 4, 8, 16 live calls without reasoning) and + /// reasoning always comes back: kept off for good, a model spent the rest + /// of a run writing probe programs instead of the deliverable. + /// + /// Defaults to `true`. A caller that must keep every call at the + /// configured effort sets it to `false`. + pub truncated_empty_reasoning_fallback: bool, + /// Client-side bound on hidden reasoning in a streamed call. + /// + /// A request's `reasoning.budget_tokens` is a promise the provider may + /// not keep: measured on deepseek-v4.1-flash through OpenRouter's + /// routable providers, a 1,500-token budget returned 4,206 reasoning + /// tokens from one and 5,964 from another, and a 9,000-token budget + /// under a tool-heavy transcript returned 21,528. A call that reasons + /// past its budget with nothing visible yet is, in every case measured, + /// one that reasons to the output cap and returns nothing: 50 seconds at + /// a 16k cap, 150 at 65k. The watchdog ends such a call at the budget + /// instead, dropping the stream, and hands the loop the same + /// `finish_reason = length`, no-content response the cap would have + /// produced, so the truncated-empty recovery (retry, reasoning off, + /// nudge) runs after a fraction of the wait. + /// + /// Reasoning length is estimated from the streamed reasoning text at + /// three characters per token (measured: 36k characters of + /// deepseek-v4.1-flash reasoning were about 13k tokens), which + /// under-counts, so the bound fires late rather than early. Visible text or a tool-call fragment before the + /// bound disarms it: the model is answering. Unary (non-streamed) calls + /// are not bounded. + /// + /// Defaults to [`ReasoningWatchdog::RequestBudget`]: enforce the budget the + /// request carries, do nothing for a request without one. + pub reasoning_watchdog: ReasoningWatchdog, + /// How much of a dead call's interrupted reasoning (its last characters) + /// is carried into the transcript, as a user message, ahead of the retry + /// or nudged call that follows it. `0` carries nothing. + /// + /// The reasoning a call dies in is usually real work: read back, one + /// dead call on a compression task was a correct derivation of the + /// decoder's arithmetic coder, cut off at the cap, and every retry began + /// the same derivation again from nothing. Carrying the tail of it, with + /// a note to continue from there in code rather than re-derive, makes the + /// deaths cumulative instead of wasted. The tail is kept because a + /// derivation's state of play is at its end. + /// + /// Defaults to 8,000 characters, about 2,500 tokens per death. + pub truncated_empty_carry_reasoning_chars: usize, /// Automatic retries for a completion with no visible text, tool calls, or /// structured output when the provider did not report length truncation. /// Reasoning-only `stop` responses are one example: the model spent tokens @@ -663,6 +720,12 @@ impl Default for RunPolicy { // keeps deliberating; one plain "stop and act" re-prompt recovers // the step instead of ending the run on a blank reply. truncated_empty_nudges: 1, + // A dead call at any cap or effort is evidence that this + // transcript does not get past the model's reasoning; the next + // call goes out without it. + truncated_empty_reasoning_fallback: true, + reasoning_watchdog: ReasoningWatchdog::RequestBudget, + truncated_empty_carry_reasoning_chars: 8_000, empty_response_retries: 0, reject_truncated_tool_calls: true, truncated_tool_call_retries: 2, @@ -785,3 +848,29 @@ impl InvocationRuntime { &self.harness } } + +/// How [`RunPolicy::reasoning_watchdog`] bounds hidden reasoning in a +/// streamed call. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub enum ReasoningWatchdog { + /// No client-side bound. + Off, + /// Enforce the request's own `reasoning.budget_tokens`; a request + /// without one is unbounded. + #[default] + RequestBudget, + /// Enforce this many reasoning tokens on every streamed call, whatever + /// the request carries. + Tokens(u32), +} + +impl ReasoningWatchdog { + /// The reasoning-token bound for `request`, if any. + pub fn bound_for(self, request: &tinyinference_llm::model::ModelRequest) -> Option { + match self { + Self::Off => None, + Self::RequestBudget => request.reasoning.as_ref().and_then(|r| r.budget_tokens), + Self::Tokens(tokens) => Some(tokens), + } + } +}