use super::*;
impl AgentSession {
pub(crate) fn run_turn_untracked<P: Provider + ?Sized>(
&self,
provider: &P,
request: AgentRunRequest<'_, '_>,
steering: Option<AgentSteering>,
defer_final_steering: bool,
request_sequence: Arc<AtomicU64>,
) -> anyhow::Result<AgentRunOutput> {
self.run_turn(
provider,
request,
steering,
defer_final_steering,
request_sequence,
None,
)
}
pub(super) fn run_turn<P: Provider + ?Sized>(
&self,
provider: &P,
run: AgentRunRequest<'_, '_>,
steering: Option<AgentSteering>,
defer_final_steering: bool,
request_sequence: Arc<AtomicU64>,
preflight_projection: Option<PreflightRequestProjection>,
) -> anyhow::Result<AgentRunOutput> {
let cancellation = run.cancellation.clone();
cancellation.check()?;
let scoped_tools = self.tools_for_task_scope(run.session, run.tools)?;
let mut run = AgentRunRequest {
tools: scoped_tools.as_ref().or(run.tools),
..run
};
let request_tools = tool_request_configuration(run.tools);
let mut prepared = self.prepare_initial_run(&run, request_tools, preflight_projection)?;
self.record_initial_user_input(&mut prepared, &mut run)?;
let mut tool_execution_expectation = ToolExecutionExpectation::from_request(
run.prompt,
run.tools.is_some(),
&prepared.tool_request_configuration.advertised_tool_names,
);
if matches!(run.prompt_origin, crate::output::UserPromptOrigin::Steering) {
tool_execution_expectation.invalidate();
if let Some(steering) = steering.as_ref()
&& let Some(prompts) = steering.acknowledge_reserved_prompt(run.prompt)
&& let Some(sink) = run.output_sink.as_deref_mut()
{
sink.steering_prompts_acknowledged(&prompts);
}
}
let title_guard = self.start_title_generation_guard(&mut run, cancellation.clone());
let replay_diagnostics = std::mem::take(&mut prepared.session_read_diagnostics);
let replay_warnings = std::mem::take(&mut prepared.replay_warnings);
let provider_prompt = std::mem::take(&mut prepared.provider_prompt);
let mut state = self.start_print_turn_state(prepared, &run, title_guard, request_sequence);
let mut auto_continue_used = false;
let mut staged_dangling_tool_recovery = false;
let steering_tools = run.tools;
let steering_cwd = run.cwd;
let steering_mode = run.invocation_mode;
let expand_steering_prompt = |text: &str| {
Self::expand_effective_prompt(
text,
steering_tools,
steering_mode,
steering_cwd,
&cancellation,
)
.unwrap_or_else(|| text.to_string())
};
let turn_steering = TurnSteering {
steering: steering.as_ref(),
defer_final: defer_final_steering,
expand: &expand_steering_prompt,
};
macro_rules! terminal_try {
($expression:expr) => {
match $expression {
Ok(value) => value,
Err(error) => {
let _ = self.record_terminal_failure(&mut state, &mut run, &error);
return Err(error);
}
}
};
}
terminal_try!(self.update_tool_checkpoint_context(&run));
terminal_try!(self.emit_initial_user_prompt_and_replay_diagnostics(
&mut run,
provider_prompt,
&replay_diagnostics,
&replay_warnings,
));
terminal_try!(self.apply_skill_suggestion(&mut state, &mut run, &cancellation));
loop {
terminal_try!(cancellation.check());
if staged_dangling_tool_recovery {
let boundary_steering =
steering.as_ref().and_then(AgentSteering::observe_collapsed);
if let Some(batch) = boundary_steering {
staged_dangling_tool_recovery = false;
if defer_final_steering {
if let Err(error) = state.assistant_chunk_batch.flush() {
state
.session_persistence
.warn_once(&mut run.output_sink, &error)?;
}
break;
}
terminal_try!(self.complete_assistant_turn_without_tools(
&mut state,
&mut run,
&cancellation,
));
terminal_try!(cancellation.check());
state.turn_state.finish_text_action_for_continuation();
terminal_try!(inject_steering_update_at_continuation_boundary(
steering.as_ref(),
Some(batch),
&mut state.turn_state,
&state.session_persistence,
&mut run.output_sink,
expand_steering_prompt,
));
tool_execution_expectation.invalidate();
continue;
}
staged_dangling_tool_recovery = false;
terminal_try!(self.auto_continue_after_dangling_tool_intent(&mut state, &mut run));
state.output.auto_compaction_blocked_by_recovery = true;
auto_continue_used = true;
terminal_try!(cancellation.check());
}
let (request, request_sequence, request_input_tokens, estimated_input_tokens) =
match self.build_turn_request_and_emit_context_usage(&mut state, &mut run) {
Ok(request) => request,
Err(error) => {
return Err(self.fail_turn_request_build(error, &mut state, &mut run));
}
};
let stream_result = match self.collect_provider_turn(
provider,
request,
&mut state,
&mut run,
&cancellation,
request_sequence,
request_input_tokens,
) {
Ok(stream_result) => stream_result,
Err(failure)
if !failure.cancelled
&& !state.output.auto_compaction_blocked_by_recovery
&& !state.session_persistence.is_degraded()
&& run.continuation_auto_compaction_policy.is_some()
&& failure
.error
.downcast_ref::<crate::providers::error::ProviderError>()
.is_some_and(
crate::providers::error::ProviderError::is_prompt_too_long,
) =>
{
terminal_try!(cancellation.check());
let error = ContextBudgetError {
phase: ContextBudgetPhase::Continuation,
estimated_tokens: request_input_tokens,
threshold_tokens: self.context_budget.threshold_tokens(),
provider_rejected: true,
partial_output: None,
message: failure.error.to_string(),
};
return Err(self.fail_turn_request_build(error.into(), &mut state, &mut run));
}
Err(failure) if !auto_continue_used && failure.eligible_for_auto_continue() => {
terminal_try!(cancellation.check());
let _ = self.record_auto_continue_decision(
&failure,
"continuing",
None,
&mut state,
&mut run,
);
terminal_try!(self.auto_continue_after_incomplete_stream(&mut state, &mut run));
state.output.recovered_incomplete_stream = true;
state.output.auto_compaction_blocked_by_recovery = true;
auto_continue_used = true;
continue;
}
Err(failure) => {
return Err(self.fail_provider_stream(
failure,
auto_continue_used,
&mut state,
&mut run,
));
}
};
if let Err(error) = self.run_reasoning_message_hooks(
&stream_result.reasoning_summaries,
&mut state,
&mut run,
&cancellation,
) {
let _ = self.record_terminal_failure(&mut state, &mut run, &error);
return Err(error);
}
if let Some(input_tokens) = stream_result.provider_input_tokens {
state
.request_projection_cache
.calibrate_anthropic_input_tokens(estimated_input_tokens, input_tokens);
}
let steering_batch = steering.as_ref().and_then(AgentSteering::observe_collapsed);
let reasoning_only = stream_result.tool_calls.is_empty()
&& state.turn_state.assistant_segment().trim().is_empty()
&& !stream_result.reasoning_summaries.is_empty();
if reasoning_only {
if let Some(batch) = steering_batch.clone()
&& !defer_final_steering
{
terminal_try!(cancellation.check());
state.turn_state.finish_text_action_for_continuation();
terminal_try!(inject_steering_update_at_continuation_boundary(
steering.as_ref(),
Some(batch),
&mut state.turn_state,
&state.session_persistence,
&mut run.output_sink,
expand_steering_prompt,
));
tool_execution_expectation.invalidate();
continue;
}
if steering_batch.is_none() && !auto_continue_used {
terminal_try!(
self.auto_continue_after_reasoning_only_turn(&mut state, &mut run)
);
state.output.auto_compaction_blocked_by_recovery = true;
auto_continue_used = true;
continue;
}
}
if stream_result.tool_calls.is_empty()
&& !auto_continue_used
&& !state.turn_state.has_tool_protocol_activity()
&& completion_needs_tool_recovery_for_expectation(
&tool_execution_expectation,
state.turn_state.assistant_segment(),
&state.tool_request_configuration.advertised_tool_names,
)
{
terminal_try!(cancellation.check());
staged_dangling_tool_recovery = true;
continue;
}
if stream_result.tool_calls.is_empty() {
match terminal_try!(self.finish_turn_without_tool_calls(
&mut state,
&mut run,
&cancellation,
&turn_steering,
steering_batch,
&mut tool_execution_expectation,
)) {
TurnLoopStep::Continue => continue,
TurnLoopStep::Break => break,
}
}
let context_window_usage = crate::output::ContextWindowUsage {
current_tokens: stream_result
.provider_input_tokens
.unwrap_or(request_input_tokens),
max_tokens: self.context_budget.max_tokens,
};
terminal_try!(self.run_tool_turn_and_maybe_emit_assistant_complete(
stream_result.tool_calls,
context_window_usage,
&mut state,
&mut run,
&cancellation,
));
if state.output.manual_compaction_request.is_some() {
break;
}
if let Err(error) = cancellation.check() {
let _ = self.record_cancelled_terminal_status(&mut state, &mut run);
return Err(error);
}
let steering_batch = steering.as_ref().and_then(AgentSteering::observe_collapsed);
let steering_provider_text = steering_batch
.as_ref()
.map(|batch| expand_steering_prompt(&batch.text));
self.check_continuation_soft_budget(
&mut state,
&mut run,
steering_provider_text.as_deref(),
)?;
if terminal_try!(inject_steering_update_at_continuation_boundary(
steering.as_ref(),
steering_batch,
&mut state.turn_state,
&state.session_persistence,
&mut run.output_sink,
|text| steering_provider_text.unwrap_or_else(|| text.to_string()),
)) {
tool_execution_expectation.invalidate();
}
}
state.title_guard.finish();
state.output.persistence_degraded = state.session_persistence.is_degraded();
Ok(state.output)
}
pub(super) fn fail_turn_request_build(
&self,
error: anyhow::Error,
state: &mut PrintTurnState<'_>,
run: &mut AgentRunRequest<'_, '_>,
) -> anyhow::Error {
state.output.persistence_degraded = state.session_persistence.is_degraded();
let is_recoverable_continuation_budget = !state.output.persistence_degraded
&& run.continuation_auto_compaction_policy.is_some()
&& error
.downcast_ref::<ContextBudgetError>()
.is_some_and(|budget| budget.phase() == ContextBudgetPhase::Continuation);
if is_recoverable_continuation_budget {
let _ = state.session_persistence.record_terminal_status(
TurnStatus::CompactionRequired,
&state.output.text,
&mut run.output_sink,
);
} else {
let _ = self.record_terminal_failure(state, run, &error);
}
state.output.persistence_degraded = state.session_persistence.is_degraded();
match error.downcast::<ContextBudgetError>() {
Ok(budget) if budget.phase() == ContextBudgetPhase::Continuation => budget
.with_partial_output(std::mem::take(&mut state.output))
.into(),
Ok(budget) => budget.into(),
Err(error) => error,
}
}
pub(super) fn fail_provider_stream(
&self,
failure: ProviderStreamFailure,
auto_continue_used: bool,
state: &mut PrintTurnState<'_>,
run: &mut AgentRunRequest<'_, '_>,
) -> anyhow::Error {
let reason = if auto_continue_used {
Some("auto-continue already used for this turn")
} else {
failure.auto_continue_block_reason()
};
if failure.recovery_diagnostic_relevant() {
let _ = self.record_auto_continue_decision(&failure, "blocked", reason, state, run);
}
self.record_provider_stream_failure(&failure, state, run);
failure.error
}
pub(super) fn finish_turn_without_tool_calls(
&self,
state: &mut PrintTurnState<'_>,
run: &mut AgentRunRequest<'_, '_>,
cancellation: &AgentCancellation,
steering: &TurnSteering<'_>,
steering_batch: Option<SteeringBatch>,
tool_execution_expectation: &mut ToolExecutionExpectation,
) -> anyhow::Result<TurnLoopStep> {
if state.turn_state.assistant_segment().trim().is_empty()
&& let Some(batch) = steering_batch
{
if steering.defer_final {
return Ok(TurnLoopStep::Break);
}
cancellation.check()?;
state.turn_state.finish_text_action_for_continuation();
inject_steering_update_at_continuation_boundary(
steering.steering,
Some(batch),
&mut state.turn_state,
&state.session_persistence,
&mut run.output_sink,
steering.expand,
)?;
tool_execution_expectation.invalidate();
return Ok(TurnLoopStep::Continue);
}
if self.verify_or_continue_terminal_completion(state, run, cancellation)? {
return Ok(TurnLoopStep::Continue);
}
self.complete_assistant_turn_without_tools(state, run, cancellation)?;
cancellation.check()?;
if !steering.defer_final {
let steering_batch = steering.steering.and_then(AgentSteering::observe_collapsed);
if steering_batch.is_some() {
state.turn_state.finish_text_action_for_continuation();
let steering_was_injected = inject_steering_update_at_continuation_boundary(
steering.steering,
steering_batch,
&mut state.turn_state,
&state.session_persistence,
&mut run.output_sink,
steering.expand,
)?;
if steering_was_injected {
tool_execution_expectation.invalidate();
return Ok(TurnLoopStep::Continue);
}
}
}
Ok(TurnLoopStep::Break)
}
pub(super) fn check_continuation_soft_budget(
&self,
state: &mut PrintTurnState<'_>,
run: &mut AgentRunRequest<'_, '_>,
steering_provider_text: Option<&str>,
) -> anyhow::Result<()> {
if state.session_persistence.is_degraded()
|| state.output.auto_compaction_blocked_by_recovery
{
return Ok(());
}
let Some(&(soft_threshold, _)) = run.continuation_auto_compaction_policy.as_ref() else {
return Ok(());
};
let temporary_steering = steering_provider_text
.map(|text| ProviderConversationItem::Message(ChatMessage::user(text)));
let temporary_suffix = temporary_steering.as_slice();
let turn_items = state.turn_state.request_items_slice();
let has_skill_reads = state.base_conversation.iter().chain(turn_items).any(|item| {
matches!(item, ProviderConversationItem::ToolResult(result) if !result.skill_reads.is_empty())
});
let projection = if has_skill_reads {
let turn_with_steering = turn_items
.iter()
.chain(temporary_suffix)
.cloned()
.collect::<Vec<_>>();
let request = self.request_from_shared_conversation(
Arc::clone(&state.base_conversation),
&turn_with_steering,
run.semantic_progress_timeout,
state.prompt_cache_key.as_deref(),
&state.tool_request_configuration,
);
project_provider_request_input_tokens(&self.provider_id, &request)
} else {
state
.request_projection_cache
.project_deterministic_turn_items(turn_items, temporary_suffix)
};
let projection = state
.request_projection_cache
.apply_usage_calibration(projection);
if projection.tokens < soft_threshold {
return Ok(());
}
let _ = state.session_persistence.record_terminal_status(
TurnStatus::CompactionRequired,
&state.output.text,
&mut run.output_sink,
);
state.output.persistence_degraded = state.session_persistence.is_degraded();
Err(ContextBudgetError::new(
ContextBudgetPhase::Continuation,
projection.tokens,
soft_threshold,
self.context_budget.max_tokens,
self.context_budget.reserve_tokens,
)
.with_partial_output(std::mem::take(&mut state.output))
.into())
}
}