use super::*;
impl AgentSession {
pub(super) fn build_turn_request_and_emit_context_usage(
&self,
state: &mut PrintTurnState<'_>,
run: &mut AgentRunRequest<'_, '_>,
) -> anyhow::Result<(ProviderRequest, u64, usize, usize)> {
let mut request = self.request_from_shared_conversation(
Arc::clone(&state.base_conversation),
state.turn_state.request_items_slice(),
run.semantic_progress_timeout,
state.prompt_cache_key.as_deref(),
&state.tool_request_configuration,
);
if let Some(conversation_id) = state.conversation_id.as_deref() {
request = request.with_conversation_id(conversation_id);
}
if self.provider_id == crate::providers::OPENAI_CODEX_PROVIDER {
request.codex_turn_context = Some(state.codex_turn_context.clone());
}
let projection = if request.has_skill_read_provenance() {
project_provider_request_input_tokens(&self.provider_id, &request)
} else {
state
.request_projection_cache
.project_deterministic_turn_items(state.turn_state.request_items_slice(), &[])
};
if self.provider_id == crate::providers::CLAUDE_SUBSCRIPTION_PROVIDER {
self.ensure_request_context_fits(projection, ContextBudgetPhase::Continuation)?;
}
let calibrated = state
.request_projection_cache
.apply_usage_calibration(projection);
let context_projection =
self.ensure_request_context_fits(calibrated, ContextBudgetPhase::Continuation)?;
let current_request_sequence = next_request_sequence(&state.request_sequence);
if let Some(sink) = run.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::ContextUsage {
current_tokens: context_projection.tokens,
max_tokens: self.context_budget.max_tokens,
reasoning_tokens: None,
source: context_projection.source,
request_sequence: current_request_sequence,
})?;
}
Ok((
request,
current_request_sequence,
context_projection.tokens,
projection.tokens,
))
}
#[expect(clippy::too_many_arguments)]
pub(super) fn collect_provider_turn<P: Provider + ?Sized>(
&self,
provider: &P,
request: ProviderRequest,
state: &mut PrintTurnState<'_>,
run: &mut AgentRunRequest<'_, '_>,
cancellation: &AgentCancellation,
request_sequence: u64,
request_input_tokens: usize,
) -> Result<ProviderStreamResult, ProviderStreamFailure> {
collect_provider_events(
provider,
request,
ProviderEventCollector {
cancellation,
output_sink: &mut run.output_sink,
output: &mut state.output,
turn_state: &mut state.turn_state,
assistant_chunk_batch: &mut state.assistant_chunk_batch,
session_persistence: &mut state.session_persistence,
context_budget_max_tokens: self.context_budget.max_tokens,
provider_id: &self.provider_id,
model: &self.model,
request_sequence,
request_input_tokens,
},
)
}
pub(super) fn project_prompt_for_preflight(
&self,
prompt: &str,
session: &Session,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<PreflightRequestProjection> {
let scoped_tools = self.tools_for_task_scope(Some(session), tools)?;
let tools = scoped_tools.as_ref().or(tools);
let initial = self.build_initial_conversation(prompt, Some(session), tools)?;
let request = self.request_from_shared_conversation(
Arc::from(initial.conversation),
&[],
None,
None,
&tool_request_configuration(tools),
);
let projection = project_provider_request_input_tokens(&self.provider_id, &request);
Ok(PreflightRequestProjection {
provider_id: self.provider_id.clone(),
session_id: session.id().to_string(),
session_path: session.path().to_path_buf(),
replay_generation: session.replay_generation(),
request,
projection,
})
}
pub(super) fn project_prompt_input_token_count(
&self,
prompt: &str,
session: Option<&Session>,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<ContextTokenCount> {
let scoped_tools = self.tools_for_task_scope(session, tools)?;
let tools = scoped_tools.as_ref().or(tools);
let conversation = self.build_initial_conversation(prompt, session, tools)?;
let request = self.request_from_conversation(conversation.conversation, None, None, tools);
Ok(project_provider_request_input_tokens(
&self.provider_id,
&request,
))
}
pub(crate) fn project_prompt_input_tokens(
&self,
prompt: &str,
session: &Session,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<usize> {
Ok(self
.project_prompt_input_token_count(prompt, Some(session), tools)?
.tokens)
}
pub(crate) fn project_prompt_input_tokens_without_session(
&self,
prompt: &str,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<usize> {
Ok(self
.project_prompt_input_token_count(prompt, None, tools)?
.tokens)
}
pub(crate) fn ensure_prompt_context_fits(
&self,
prompt: &str,
session: &Session,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<usize> {
Ok(self
.ensure_request_context_fits(
self.project_prompt_input_token_count(prompt, Some(session), tools)?,
ContextBudgetPhase::Initial,
)?
.tokens)
}
pub(crate) fn context_enabled(&self) -> bool {
self.context_budget.enabled
}
pub(crate) fn context_threshold_tokens(&self) -> usize {
self.context_budget.threshold_tokens()
}
pub(crate) fn context_max_tokens(&self) -> usize {
self.context_budget.max_tokens
}
pub(super) fn task_scope_for_run(
&self,
session: Option<&Session>,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<Option<crate::sessions::task_scope::TaskScope>> {
let memory = tools.and_then(|tools| tools.task_scope.as_ref());
if let Some(session) = session {
let scope = self
.replay_cache
.lock()
.map_err(|_| anyhow::anyhow!("session replay cache was poisoned"))?
.task_scope(session)?;
anyhow::ensure!(
scope.is_some() || memory.is_none(),
"missing required delegated task scope; execution blocked"
);
Ok(scope)
} else {
memory
.map(|scope| {
scope
.lock()
.map(|scope| scope.clone())
.map_err(|_| anyhow::anyhow!("task scope lock poisoned"))
})
.transpose()
}
}
pub(super) fn tools_for_task_scope(
&self,
session: Option<&Session>,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<Option<ToolRuntime>> {
let Some(scope) = self.task_scope_for_run(session, tools)? else {
return Ok(None);
};
let tools = tools.ok_or_else(|| {
anyhow::anyhow!("delegated task requires its capability-enforced tool runtime")
})?;
if session.is_none() {
return Ok(Some(tools.clone()));
}
Ok(Some(tools.clone().with_task_scope(scope)?))
}
pub(super) fn build_initial_conversation(
&self,
prompt: &str,
session: Option<&Session>,
tools: Option<&ToolRuntime>,
) -> anyhow::Result<InitialConversation> {
let replay = self
.replay_cache
.lock()
.map_err(|_| anyhow::anyhow!("session replay cache was poisoned"))?
.replay(session)?;
let _legacy_lossy_events = replay.legacy_lossy_events;
let mut conversation = Vec::with_capacity(replay.items.len() + 2);
conversation.push(ProviderConversationItem::Message(ChatMessage::system(
self.system_prompt.clone(),
)));
let scope = if session.is_some() {
replay.task_scope.map_err(anyhow::Error::msg)?
} else {
self.task_scope_for_run(None, tools)?
};
if let Some(scope) = scope {
conversation.push(ProviderConversationItem::Message(ChatMessage::system(
crate::sessions::task_scope::SCOPE_RULES,
)));
conversation.push(ProviderConversationItem::Message(ChatMessage::user(
scope.provider_note()?,
)));
} else if tools.is_some_and(ToolRuntime::is_subagent)
|| session.is_some_and(|session| {
session
.path()
.parent()
.and_then(|path| path.file_name())
.is_some_and(|name| name == "subagents")
})
{
conversation.push(ProviderConversationItem::Message(ChatMessage::system(
crate::sessions::task_scope::LEGACY_SCOPE_NOTE,
)));
}
conversation.extend(replay.items);
if self.reasoning_updates_enabled && session.is_some() {
conversation.push(ProviderConversationItem::ReasoningSelection {
provider: self.provider_id.clone(),
model: self.model.clone(),
effort: self.thinking_level,
});
}
conversation.push(ProviderConversationItem::Message(ChatMessage::user(
prompt.to_string(),
)));
Ok(InitialConversation {
conversation,
session_read_diagnostics: replay.session_read_diagnostics,
replay_warnings: replay.warnings,
})
}
pub(super) fn request_from_conversation(
&self,
conversation: Vec<ProviderConversationItem>,
semantic_progress_timeout: Option<Duration>,
prompt_cache_key: Option<&str>,
tools: Option<&ToolRuntime>,
) -> ProviderRequest {
let tool_configuration = tool_request_configuration(tools);
self.request_from_conversation_with_configuration(
conversation,
semantic_progress_timeout,
prompt_cache_key,
&tool_configuration,
)
}
pub(super) fn request_from_conversation_with_configuration(
&self,
conversation: Vec<ProviderConversationItem>,
semantic_progress_timeout: Option<Duration>,
prompt_cache_key: Option<&str>,
tool_configuration: &ToolRequestConfiguration,
) -> ProviderRequest {
self.configure_request(
ProviderRequest::from_conversation(self.model.clone(), conversation),
semantic_progress_timeout,
prompt_cache_key,
tool_configuration,
)
}
pub(super) fn request_from_shared_conversation(
&self,
base_conversation: Arc<[ProviderConversationItem]>,
turn_items: &[ProviderConversationItem],
semantic_progress_timeout: Option<Duration>,
prompt_cache_key: Option<&str>,
tool_configuration: &ToolRequestConfiguration,
) -> ProviderRequest {
self.configure_request(
ProviderRequest::from_shared_conversation(
self.model.clone(),
base_conversation,
turn_items,
),
semantic_progress_timeout,
prompt_cache_key,
tool_configuration,
)
}
pub(super) fn configure_request(
&self,
request: ProviderRequest,
semantic_progress_timeout: Option<Duration>,
prompt_cache_key: Option<&str>,
tool_configuration: &ToolRequestConfiguration,
) -> ProviderRequest {
let mut disabled_tool_names: Vec<_> = tool_configuration
.disabled_tool_names
.iter()
.cloned()
.collect();
disabled_tool_names.sort_unstable();
let mut request = request
.with_thinking_level(self.thinking_level)
.with_text_verbosity(self.text_verbosity)
.with_default_reasoning_summary(self.send_default_reasoning_summary)
.with_subagents_tool_enabled(tool_configuration.subagents_tool_enabled)
.with_disabled_tool_names(disabled_tool_names)
.with_dynamic_tool_definitions(tool_configuration.dynamic_tool_definitions.clone());
if let Some(catalog) = &tool_configuration.resolved_catalog {
request = request.with_resolved_tool_catalog(catalog.clone());
}
if self.reasoning_updates_enabled {
request = request.with_reasoning_updates(&self.provider_id);
}
if let Some(prompt_cache_key) = prompt_cache_key {
request = request.with_prompt_cache_key(prompt_cache_key);
}
match semantic_progress_timeout {
Some(timeout) => request.with_semantic_progress_timeout(timeout),
None => request,
}
}
pub(super) fn ensure_request_context_fits(
&self,
projection: ContextTokenCount,
phase: ContextBudgetPhase,
) -> anyhow::Result<ContextTokenCount> {
if !self.context_budget.enabled {
return Ok(projection);
}
let estimated_tokens = projection.tokens;
let threshold = self.context_budget.threshold_tokens();
if estimated_tokens <= threshold {
return Ok(projection);
}
Err(ContextBudgetError::new(
phase,
estimated_tokens,
threshold,
self.context_budget.max_tokens,
self.context_budget.reserve_tokens,
)
.into())
}
}