use tracing::Instrument;
use zeph_llm::provider::{Message, Role};
use zeph_tools::executor::ToolCall;
use super::retry_backoff_ms;
use crate::agent::Agent;
use crate::channel::Channel;
const REFORMAT_DEFAULT_TIMEOUT_SECS: u64 = 30;
impl<C: Channel> Agent<C> {
#[tracing::instrument(name = "core.tool.handle_retry_phase", skip_all, level = "debug", err)]
pub(super) async fn handle_retry_phase(
&mut self,
tool_calls: &[zeph_llm::provider::ToolUseRequest],
calls: &[ToolCall],
tool_results: &mut [Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>],
max_retries: usize,
cancel: &tokio_util::sync::CancellationToken,
) -> Result<bool, crate::agent::error::AgentError> {
if max_retries == 0 {
return Ok(false);
}
let max_retry_duration_secs = self.tool_orchestrator.max_retry_duration_secs;
let retry_base_ms = self.tool_orchestrator.retry_base_ms;
let retry_max_ms = self.tool_orchestrator.retry_max_ms;
for idx in 0..tool_results.len() {
if cancel.is_cancelled() {
self.cancel_tool_batch(tool_calls, "tool execution cancelled by user")
.await?;
return Ok(true);
}
let is_transient = matches!(
tool_results[idx],
Err(ref e) if e.kind() == zeph_tools::ErrorKind::Transient
);
if !is_transient {
continue;
}
let tc = &tool_calls[idx];
if !self
.tool_executor
.is_tool_retryable_erased(tc.name.as_str())
{
continue;
}
let call = &calls[idx];
let mut attempt = 0_usize;
let retry_start = std::time::Instant::now();
let result = loop {
let exec_result = tokio::select! {
r = self.tool_executor.execute_tool_call_erased(call).instrument(
tracing::info_span!("tool_exec_retry", tool_name = %tc.name, idx = %tc.id)
) => r,
() = cancel.cancelled() => {
self.cancel_tool_batch(tool_calls, "tool retry cancelled by user")
.await?;
return Ok(true);
}
};
match exec_result {
Err(ref e)
if e.kind() == zeph_tools::ErrorKind::Transient
&& attempt < max_retries =>
{
let elapsed_secs = retry_start.elapsed().as_secs();
if max_retry_duration_secs > 0 && elapsed_secs >= max_retry_duration_secs {
tracing::warn!(
tool = %tc.name, elapsed_secs, max_retry_duration_secs,
"tool retry budget exceeded, aborting retries"
);
break exec_result;
}
attempt += 1;
let delay_ms = retry_backoff_ms(attempt - 1, retry_base_ms, retry_max_ms);
tracing::warn!(
tool = %tc.name, attempt, delay_ms, error = %e,
"transient tool error, retrying with backoff"
);
self.channel
.send_status_best_effort(&format!("Retrying {}...", tc.name))
.await;
tokio::select! {
() = tokio::time::sleep(std::time::Duration::from_millis(delay_ms)) => {}
() = cancel.cancelled() => {
self.cancel_tool_batch(
tool_calls,
"retry backoff interrupted by cancellation",
)
.await?;
return Ok(true);
}
}
self.channel.send_status_best_effort("").await;
}
result => break result,
}
};
tool_results[idx] = result;
}
Ok(false)
}
#[tracing::instrument(
name = "core.tool.handle_reformat_phase",
skip_all,
level = "debug",
err
)]
pub(super) async fn handle_reformat_phase(
&mut self,
tool_calls: &[zeph_llm::provider::ToolUseRequest],
calls: &[ToolCall],
tool_results: &mut [Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>],
cancel: &tokio_util::sync::CancellationToken,
) -> Result<bool, crate::agent::error::AgentError> {
if self
.tool_orchestrator
.parameter_reformat_provider
.is_empty()
{
return Ok(false);
}
let budget_secs = self.tool_orchestrator.max_retry_duration_secs;
let reformat_start = std::time::Instant::now();
for idx in 0..tool_results.len() {
if cancel.is_cancelled() {
self.cancel_tool_batch(tool_calls, "parameter reformat phase cancelled by user")
.await?;
return Ok(true);
}
let needs_reformat = matches!(
tool_results[idx],
Err(ref e) if e.category().needs_parameter_reformat()
);
if !needs_reformat {
continue;
}
let tc = &tool_calls[idx];
if budget_secs > 0 && reformat_start.elapsed().as_secs() >= budget_secs {
tracing::warn!(tool = %tc.name, "parameter reformat budget exhausted, skipping");
continue;
}
let error_message = tool_results[idx]
.as_ref()
.err()
.map(std::string::ToString::to_string)
.unwrap_or_default();
self.channel
.send_status_best_effort(&format!("Reformatting parameters for {}...", tc.name))
.await;
let new_result = self
.reformat_tool_call(&calls[idx], tc, &error_message, cancel)
.await;
self.channel.send_status_best_effort("").await;
if cancel.is_cancelled() {
self.cancel_tool_batch(tool_calls, "parameter reformat phase cancelled by user")
.await?;
return Ok(true);
}
if let Some(result) = new_result {
if let Err(ref e) = result
&& let Some(ref d) = self.runtime.debug.debug_dumper
{
d.dump_tool_error(tc.name.as_str(), e);
}
tool_results[idx] = result;
}
}
Ok(false)
}
async fn reformat_tool_call(
&mut self,
call: &ToolCall,
tc: &zeph_llm::provider::ToolUseRequest,
error_message: &str,
cancel: &tokio_util::sync::CancellationToken,
) -> Option<Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>> {
#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
struct ReformattedArguments {
arguments: serde_json::Value,
}
let Some(schema) = self
.tool_executor
.tool_definitions_erased()
.into_iter()
.find(|d| d.id.as_ref() == tc.name.as_str())
.map(|d| d.schema)
else {
tracing::warn!(tool = %tc.name, "parameter reformat: tool schema not found, skipping");
return None;
};
let provider_name = self.tool_orchestrator.parameter_reformat_provider.clone();
let provider = match self.resolve_pool_entry_provider(&provider_name) {
super::super::learning::PoolProviderResolution::Resolved(p) => *p,
super::super::learning::PoolProviderResolution::RegistryNotWired => {
self.resolve_background_provider(&provider_name)
}
super::super::learning::PoolProviderResolution::Unresolvable => {
tracing::warn!(
tool = %tc.name,
provider = %provider_name,
"parameter reformat: configured provider unresolvable, keeping original error"
);
return None;
}
};
let original_args = serde_json::Value::Object(call.params.clone());
let prompt = format!(
"A tool call failed parameter validation. Propose corrected arguments as a JSON \
object under the `arguments` key.\n\n\
Tool: {}\nJSON schema:\n{}\n\nOriginal arguments:\n{}\n\nError: {error_message}",
tc.name,
serde_json::to_string_pretty(&schema).unwrap_or_default(),
serde_json::to_string_pretty(&original_args).unwrap_or_default(),
);
let messages = [Message::from_legacy(Role::User, prompt)];
let timeout_secs = match self.tool_orchestrator.max_retry_duration_secs {
0 => REFORMAT_DEFAULT_TIMEOUT_SECS,
secs => secs,
};
let reformat = tokio::select! {
r = tokio::time::timeout(
std::time::Duration::from_secs(timeout_secs),
provider.chat_typed_erased::<ReformattedArguments>(&messages),
) => match r {
Ok(Ok(reformat)) => reformat,
Ok(Err(e)) => {
tracing::warn!(tool = %tc.name, error = %e, "parameter reformat: provider call failed");
return None;
}
Err(_) => {
tracing::warn!(tool = %tc.name, timeout_secs, "parameter reformat: provider call timed out");
return None;
}
},
() = cancel.cancelled() => return None,
};
let serde_json::Value::Object(corrected_params) = reformat.arguments else {
tracing::warn!(
tool = %tc.name,
"parameter reformat: corrected arguments were not a JSON object, skipping retry"
);
return None;
};
let mut retry_call = call.clone();
retry_call.params = corrected_params;
tokio::select! {
r = self.tool_executor.execute_tool_call_erased(&retry_call).instrument(
tracing::info_span!("tool_exec_reformat", tool_name = %tc.name)
) => Some(r),
() = cancel.cancelled() => None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
mod reformat_phase_tests {
use zeph_tools::registry::{InvocationHint, ToolDef};
use super::*;
use crate::agent::agent_tests::*;
fn test_tool_def() -> ToolDef {
ToolDef {
id: "test_tool".into(),
description: "a test tool".into(),
schema: schemars::Schema::default(),
invocation: InvocationHint::ToolCall,
output_schema: None,
server_id: None,
}
}
fn invalid_params_error() -> zeph_tools::ToolError {
zeph_tools::ToolError::InvalidParams {
message: "path must be a string".to_owned(),
}
}
fn bad_tool_call() -> ToolCall {
let mut params = serde_json::Map::new();
params.insert("path".to_owned(), serde_json::json!(123));
ToolCall {
tool_id: zeph_common::ToolName::new("test_tool"),
params,
..Default::default()
}
}
fn tool_use_request() -> zeph_llm::provider::ToolUseRequest {
zeph_llm::provider::ToolUseRequest {
id: "call-1".to_owned(),
name: "test_tool".to_owned().into(),
input: serde_json::json!({"path": 123}),
}
}
#[tokio::test]
async fn retries_with_corrected_arguments_on_success() {
let provider = mock_provider(vec![r#"{"arguments":{"path":"/corrected"}}"#.into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::new(vec![Ok(Some(zeph_tools::ToolOutput {
tool_name: "test_tool".to_owned().into(),
summary: "done".to_owned(),
blocks_executed: 1,
filter_stats: None,
diff: None,
streamed: false,
terminal_id: None,
locations: None,
raw_response: None,
claim_source: None,
..Default::default()
}))])
.with_definitions(vec![test_tool_def()]);
let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
agent.tool_orchestrator.parameter_reformat_provider = "fast".to_owned();
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
let output = tool_results
.remove(0)
.expect("reformat retry should succeed")
.expect("tool output should be present");
assert_eq!(output.summary, "done");
}
#[tokio::test]
async fn keeps_original_error_when_provider_returns_malformed_json() {
let provider = mock_provider(vec!["not json at all".into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor =
MockToolExecutor::new(vec![Ok(None)]).with_definitions(vec![test_tool_def()]);
let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
agent.tool_orchestrator.parameter_reformat_provider = "fast".to_owned();
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
assert!(
matches!(
tool_results[0],
Err(zeph_tools::ToolError::InvalidParams { .. })
),
"original error must be preserved on parse failure"
);
}
#[tokio::test]
async fn keeps_original_error_when_tool_schema_unknown() {
let provider = mock_provider(vec![r#"{"arguments":{"path":"/corrected"}}"#.into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::new(vec![Ok(None)]);
let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
agent.tool_orchestrator.parameter_reformat_provider = "fast".to_owned();
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
assert!(
matches!(
tool_results[0],
Err(zeph_tools::ToolError::InvalidParams { .. })
),
"original error must be preserved when schema is unknown"
);
}
#[tokio::test]
async fn is_a_noop_when_provider_not_configured() {
let provider = mock_provider(vec![r#"{"arguments":{"path":"/corrected"}}"#.into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor =
MockToolExecutor::new(vec![Ok(None)]).with_definitions(vec![test_tool_def()]);
let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
assert!(
matches!(
tool_results[0],
Err(zeph_tools::ToolError::InvalidParams { .. })
),
"reformat must not run when parameter_reformat_provider is empty"
);
}
#[tokio::test]
async fn falls_back_to_primary_when_registry_never_wired() {
let provider = mock_provider(vec![r#"{"arguments":{"path":"/corrected"}}"#.into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::new(vec![Ok(Some(zeph_tools::ToolOutput {
tool_name: "test_tool".to_owned().into(),
summary: "done".to_owned(),
blocks_executed: 1,
filter_stats: None,
diff: None,
streamed: false,
terminal_id: None,
locations: None,
raw_response: None,
claim_source: None,
..Default::default()
}))])
.with_definitions(vec![test_tool_def()]);
let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
agent.tool_orchestrator.parameter_reformat_provider = "unregistered".to_owned();
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
let output = tool_results
.remove(0)
.expect("reformat should run using the primary provider fallback")
.expect("tool output should be present");
assert_eq!(output.summary, "done");
}
fn empty_snapshot() -> crate::agent::state::ProviderConfigSnapshot {
crate::agent::state::ProviderConfigSnapshot {
claude_api_key: None,
openai_api_key: None,
gemini_api_key: None,
compatible_api_keys: std::collections::HashMap::new(),
llm_request_timeout_secs: 30,
embedding_model: String::new(),
gonka_private_key: None,
gonka_address: None,
cocoon_access_hash: None,
}
}
#[tokio::test]
async fn is_a_noop_when_registry_wired_but_name_absent_from_pool() {
let provider = mock_provider(vec![r#"{"arguments":{"path":"/corrected"}}"#.into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor =
MockToolExecutor::new(vec![Ok(None)]).with_definitions(vec![test_tool_def()]);
let other_entry = crate::config::ProviderEntry {
provider_type: crate::config::ProviderKind::Ollama,
name: Some("other".into()),
..Default::default()
};
let mut agent = Agent::new(provider, channel, registry, None, 5, executor)
.with_provider_pool(vec![other_entry], empty_snapshot());
agent.tool_orchestrator.parameter_reformat_provider = "unregistered".to_owned();
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
assert!(
matches!(
tool_results[0],
Err(zeph_tools::ToolError::InvalidParams { .. })
),
"reformat must no-op when the registry is wired but the configured name is \
absent from the pool, not silently fall back to the primary provider"
);
}
#[tokio::test]
async fn is_a_noop_when_registry_wired_but_provider_build_fails() {
let provider = mock_provider(vec![r#"{"arguments":{"path":"/corrected"}}"#.into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor =
MockToolExecutor::new(vec![Ok(None)]).with_definitions(vec![test_tool_def()]);
let broken_entry = crate::config::ProviderEntry {
provider_type: crate::config::ProviderKind::Claude,
name: Some("broken".into()),
..Default::default()
};
let mut agent = Agent::new(provider, channel, registry, None, 5, executor)
.with_provider_pool(vec![broken_entry], empty_snapshot());
agent.tool_orchestrator.parameter_reformat_provider = "broken".to_owned();
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
assert!(
matches!(
tool_results[0],
Err(zeph_tools::ToolError::InvalidParams { .. })
),
"reformat must no-op when the in-pool provider fails to build, not silently \
fall back to the primary provider"
);
}
#[tokio::test]
async fn replaces_original_error_when_retry_still_fails() {
let provider = mock_provider(vec![r#"{"arguments":{"path":"/still-bad"}}"#.into()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::new(vec![Err(zeph_tools::ToolError::InvalidParams {
message: "still not a valid path".to_owned(),
})])
.with_definitions(vec![test_tool_def()]);
let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
agent.tool_orchestrator.parameter_reformat_provider = "fast".to_owned();
let tool_calls = vec![tool_use_request()];
let calls = vec![bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
match &tool_results[0] {
Err(zeph_tools::ToolError::InvalidParams { message }) => {
assert_eq!(
message, "still not a valid path",
"the retried failure must replace the original error, not leave the \
pre-reformat error in place"
);
}
other => {
panic!("expected the retried failure to replace the original, got {other:?}")
}
}
}
#[tokio::test]
async fn budget_exhausted_skips_remaining_calls_in_same_phase() {
let provider = mock_provider(vec![
r#"{"arguments":{"path":"/corrected"}}"#.into(),
r#"{"arguments":{"path":"/corrected"}}"#.into(),
]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::new(vec![Ok(Some(zeph_tools::ToolOutput {
tool_name: "test_tool".to_owned().into(),
summary: "done".to_owned(),
blocks_executed: 1,
filter_stats: None,
diff: None,
streamed: false,
terminal_id: None,
locations: None,
raw_response: None,
claim_source: None,
..Default::default()
}))])
.with_definitions(vec![test_tool_def()])
.with_delay(1_100);
let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
agent.tool_orchestrator.parameter_reformat_provider = "fast".to_owned();
agent.tool_orchestrator.max_retry_duration_secs = 1;
let tool_calls = vec![tool_use_request(), tool_use_request()];
let calls = vec![bad_tool_call(), bad_tool_call()];
let mut tool_results: Vec<
Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError>,
> = vec![Err(invalid_params_error()), Err(invalid_params_error())];
let cancel = tokio_util::sync::CancellationToken::new();
agent
.handle_reformat_phase(&tool_calls, &calls, &mut tool_results, &cancel)
.await
.unwrap();
assert!(
tool_results[0].is_ok(),
"first call is within budget and should be reformatted successfully"
);
assert!(
matches!(
tool_results[1],
Err(zeph_tools::ToolError::InvalidParams { .. })
),
"second call must be skipped once the whole-phase budget is exhausted by the \
first call's real elapsed time"
);
}
}
}