use super::completion_runtime::CompletionFlow;
use super::execution_state::ExecutionLoopState;
use super::llm_turn::LlmTurnRequest;
use super::queue_forwarder::QueueEventForwarder;
use super::{AgentEvent, AgentLoop, AgentResult};
use crate::llm::{ContentBlock, Message};
use crate::prompts::AgentStyle;
use anyhow::Result;
use tokio::sync::mpsc;
const TOOL_BUDGET_FINALIZATION: &str = "Tool-use budget reached. Stop gathering evidence and return the best complete final answer now using only the tool results already present. Do not call any tool.";
impl AgentLoop {
#[allow(clippy::too_many_arguments)]
pub(super) async fn execute_loop(
&self,
history: &[Message],
prompt: &str,
effective_style: AgentStyle,
session_id: Option<&str>,
event_tx: Option<mpsc::Sender<AgentEvent>>,
cancel_token: &tokio_util::sync::CancellationToken,
emit_end: bool,
) -> Result<AgentResult> {
self.execute_loop_inner(
history,
prompt,
prompt,
Some(effective_style),
session_id,
event_tx,
cancel_token,
emit_end,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn execute_loop_inner(
&self,
history: &[Message],
msg_prompt: &str,
effective_prompt: &str,
effective_style: Option<AgentStyle>,
session_id: Option<&str>,
event_tx: Option<mpsc::Sender<AgentEvent>>,
cancel_token: &tokio_util::sync::CancellationToken,
emit_end: bool,
seed: Option<super::execution_state::ExecutionSeed>,
) -> Result<AgentResult> {
let mut state = ExecutionLoopState::new_seeded(history, seed);
let style_prompt = if effective_prompt.is_empty() {
msg_prompt
} else {
effective_prompt
};
let prompt_mode = self
.resolve_prompt_mode(effective_style, style_prompt, &event_tx)
.await;
let effective_system_prompt = prompt_mode.system_prompt;
if let Some(tx) = &event_tx {
tx.send(AgentEvent::Start {
prompt: effective_prompt.to_string(),
})
.await
.ok();
}
let _queue_forwarder = QueueEventForwarder::start(
self.command_queue.as_ref(),
event_tx.as_ref(),
cancel_token,
);
let prompt_before_hooks = effective_prompt;
let turn_context = self
.prepare_turn_context(
&effective_system_prompt,
effective_prompt,
state.messages.len(),
session_id,
&event_tx,
)
.await?;
let effective_prompt = turn_context.effective_prompt.as_str();
let augmented_system = turn_context.augmented_system;
self.config.rl_trajectory_recorder.record_execution_start(
crate::rl_trajectory::ExecutionStartRecord {
session_id: session_id.unwrap_or(""),
workspace: &self.tool_context.workspace,
prompt: effective_prompt,
history,
system_prompt: augmented_system.as_deref(),
max_tool_rounds: self.config.max_tool_rounds,
planning_mode: &format!("{:?}", self.config.planning_mode),
},
);
if !msg_prompt.is_empty() {
state.messages.push(Message::user(effective_prompt));
} else if effective_prompt != prompt_before_hooks {
rewrite_latest_user_prompt(&mut state.messages, prompt_before_hooks, effective_prompt);
}
loop {
let force_finalization = state.current_turn() >= self.config.max_tool_rounds;
if force_finalization {
state.messages.push(Message::user(TOOL_BUDGET_FINALIZATION));
}
let turn = state.next_turn();
let capability_turn = self
.capability_runtime
.as_ref()
.map(|runtime| runtime.begin_turn(turn))
.transpose()?;
let scoped_cancellation = capability_turn
.as_ref()
.map(crate::capability::AgentCapabilityTurn::cancellation)
.unwrap_or_else(|| cancel_token.clone());
let scoped_tool_context = capability_turn.as_ref().map_or_else(
|| self.tool_context.clone(),
|scope| {
self.tool_context
.clone()
.with_capability_context(scope.tool_context())
},
);
let llm_turn = match self
.execute_llm_turn(
&mut state,
LlmTurnRequest {
turn,
augmented_system: &augmented_system,
effective_prompt,
session_id,
event_tx: &event_tx,
cancel_token: &scoped_cancellation,
force_no_tools: force_finalization,
},
)
.await
{
Ok(turn) => turn,
Err(_) if scoped_cancellation.is_cancelled() => {
close_capability_turn(capability_turn.as_ref()).await?;
return Ok(state.finish_interrupted());
}
Err(error) => {
if let Err(close_error) = close_capability_turn(capability_turn.as_ref()).await
{
tracing::warn!(
error = %close_error,
"Capability Turn close also failed after provider failure"
);
}
return Err(error);
}
};
debug_assert_eq!(llm_turn.turn, turn);
let response = llm_turn.response;
let tool_calls = llm_turn.tool_calls;
if force_finalization && !tool_calls.is_empty() {
let error = format!(
"Max tool rounds ({}) exceeded; the reserved finalization turn attempted another tool call",
self.config.max_tool_rounds
);
self.emit_error(&event_tx, error.clone()).await;
close_capability_turn(capability_turn.as_ref()).await?;
anyhow::bail!(error);
}
if tool_calls.is_empty() {
match self
.complete_no_tool_response(
&mut state,
turn,
&response,
effective_prompt,
session_id,
&event_tx,
emit_end,
&scoped_cancellation,
&scoped_tool_context,
force_finalization,
)
.await
{
CompletionFlow::Continue => {
close_capability_turn(capability_turn.as_ref()).await?;
continue;
}
CompletionFlow::Finished(final_text) => {
close_capability_turn(capability_turn.as_ref()).await?;
return Ok(state.finish(final_text));
}
}
}
if let Err(e) = self
.execute_tool_turn(
tool_calls,
&mut state,
&event_tx,
session_id,
&scoped_cancellation,
&scoped_tool_context,
)
.await
{
if scoped_cancellation.is_cancelled() {
close_capability_turn(capability_turn.as_ref()).await?;
return Ok(state.finish_interrupted());
}
if let Err(close_error) = close_capability_turn(capability_turn.as_ref()).await {
tracing::warn!(
error = %close_error,
"Capability Turn close also failed after Tool failure"
);
}
return Err(e);
}
close_capability_turn(capability_turn.as_ref()).await?;
self.persist_loop_checkpoint(turn, &state, session_id).await;
}
}
async fn persist_loop_checkpoint(
&self,
turn: usize,
state: &super::execution_state::ExecutionLoopState,
session_id: Option<&str>,
) {
let Some(sink) = self.checkpoint_sink.as_ref() else {
return;
};
let Some(run_id) = self.checkpoint_run_id.as_ref() else {
return;
};
let checkpoint = crate::loop_checkpoint::LoopCheckpoint {
schema_version: crate::loop_checkpoint::LOOP_CHECKPOINT_SCHEMA_VERSION,
run_id: run_id.clone(),
session_id: session_id.unwrap_or("").to_string(),
capability_binding: self.checkpoint_capability_binding.clone(),
turn,
messages: state.messages.clone(),
total_usage: state.total_usage.clone(),
tool_calls_count: state.tool_calls_count,
verification_reports: state.verification_reports.clone(),
convergence: state.convergence_checkpoint(),
checkpoint_ms: self.config.host_env.now_ms(),
};
sink.save_checkpoint(&checkpoint).await;
}
}
async fn close_capability_turn(
turn: Option<&crate::capability::AgentCapabilityTurn>,
) -> anyhow::Result<()> {
let Some(turn) = turn else {
return Ok(());
};
let report = turn.close().await?;
if !report.is_clean() {
anyhow::bail!(
"Capability Turn close was incomplete (tasks failed: {}, tasks timed out: {}, child scopes failed: {}, child scopes timed out: {}, effects failed: {}, effects timed out: {})",
report.tasks_failed,
report.tasks_timed_out,
report.child_scopes_failed,
report.child_scopes_timed_out,
report.effects_failed,
report.effects_timed_out,
);
}
Ok(())
}
fn rewrite_latest_user_prompt(messages: &mut [Message], original: &str, replacement: &str) {
let candidate = messages
.iter()
.rposition(|message| message.role == "user" && message.text() == original)
.or_else(|| {
messages
.iter()
.rposition(|message| message.role == "user" && !message.text().is_empty())
});
let Some(index) = candidate else {
tracing::warn!("PrePrompt modified input but no user message could be rewritten");
return;
};
let message = &mut messages[index];
let mut wrote_text = false;
message.content.retain_mut(|block| match block {
ContentBlock::Text { text } if !wrote_text => {
*text = replacement.to_string();
wrote_text = true;
true
}
ContentBlock::Text { .. } => false,
_ => true,
});
if !wrote_text {
message.content.push(ContentBlock::Text {
text: replacement.to_string(),
});
}
}