use crate::{
agent::AgentSession,
agent::steering::AgentSteering,
cancellation::AgentCancellation,
config::{CustomProviderConfig, EffectiveConfig, Settings},
context::ContextBudget,
output::{HookContextMetadata, InvocationMode, OutputEvent, ToolDispatchContext},
primary_agents::{PrimaryAgentProfile, render_primary_agent_prompt_append},
providers::{OPENAI_CODEX_PROVIDER, ProviderSelection, ToolCall},
sessions::Session,
};
use anyhow::Result;
use serde_json::json;
use std::{
collections::HashSet,
path::Path,
sync::{Arc, Mutex, atomic::AtomicU64},
};
mod auto_compaction;
mod preparation;
pub(crate) struct ProviderRunOptions<'a, 'sink> {
pub(crate) settings: Option<Settings>,
pub(crate) prompt: &'a str,
pub(crate) session: Option<&'a Session>,
pub(crate) cwd: &'a Path,
pub(crate) output_sink: Option<&'sink mut dyn crate::agent::AgentOutputSink>,
pub(crate) selected_primary_agent: Option<PrimaryAgentProfile>,
pub(crate) cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
pub(crate) subagent_controls: Option<crate::subagents::control::SubagentControls>,
pub(crate) session_title_notifier: Option<crate::session_titles::SessionTitleNotifier>,
pub(crate) herdr_reporter: Option<crate::herdr::HerdrReporter>,
pub(crate) invocation_mode: InvocationMode,
pub(crate) disabled_tools: Option<Arc<Mutex<HashSet<String>>>>,
pub(crate) disabled_subagent_profiles: Option<Arc<Mutex<HashSet<String>>>>,
pub(crate) subagent_profile_discovery:
Option<crate::subagents::profiles::SubagentProfileDiscovery>,
pub(crate) mcp: Option<Arc<Mutex<crate::mcp::manager::McpManager>>>,
}
fn reborrow_output_sink<'a>(
sink: &'a mut Option<&mut dyn crate::agent::AgentOutputSink>,
) -> Option<&'a mut dyn crate::agent::AgentOutputSink> {
match sink {
Some(sink) => Some(&mut **sink),
None => None,
}
}
fn emit_compaction_fast_observation(
output_sink: &mut Option<&mut dyn crate::agent::AgentOutputSink>,
result: &crate::compaction::CompactionResult,
request_sequence: &AtomicU64,
) -> anyhow::Result<()> {
let Some(requested_service_tier) = result.requested_service_tier.as_deref() else {
return Ok(());
};
let Some(sink) = output_sink.as_deref_mut() else {
return Ok(());
};
sink.output_event(OutputEvent::FastObservation {
provider_id: result.provider.clone(),
model: result.model.clone(),
requested_service_tier: requested_service_tier.to_string(),
outcome: result.fast_outcome.clone(),
request_sequence: request_sequence
.load(std::sync::atomic::Ordering::Relaxed)
.saturating_sub(1),
run_order: None,
})
}
pub(crate) const AUTO_COMPACTION_DIAGNOSTIC_MAX_CHARS: usize = 512;
pub(crate) fn auto_compaction_policy(
auto: &crate::config::AutoCompactionSettings,
model_max_tokens: usize,
) -> Option<(usize, String)> {
if !auto.is_enabled() {
return None;
}
let percent_threshold = auto.threshold_percent.map(|percent| {
let numerator = (model_max_tokens as u128).saturating_mul(u128::from(percent));
usize::try_from(numerator.div_ceil(100)).unwrap_or(usize::MAX)
});
let token_threshold = auto
.threshold_tokens
.map(|tokens| usize::try_from(tokens).unwrap_or(usize::MAX));
let threshold = match (percent_threshold, token_threshold) {
(Some(percent), Some(tokens)) => percent.min(tokens),
(Some(threshold), None) | (None, Some(threshold)) => threshold,
(None, None) => return None,
};
Some((threshold, auto.threshold_display()?))
}
pub(crate) fn bounded_auto_compaction_diagnostic(message: impl std::fmt::Display) -> String {
crate::output::redact_sensitive_text(&message.to_string())
.chars()
.take(AUTO_COMPACTION_DIAGNOSTIC_MAX_CHARS)
.collect()
}
fn auto_compaction_eligible(
auto_enabled: bool,
context_enabled: bool,
has_session: bool,
invocation_mode: InvocationMode,
) -> bool {
auto_enabled
&& context_enabled
&& has_session
&& matches!(
invocation_mode,
InvocationMode::Print
| InvocationMode::MissionControl
| InvocationMode::LocalAgent
| InvocationMode::Side
)
}
pub(crate) fn append_primary_agent_to_main_prompt(
agent: AgentSession,
selected: Option<&PrimaryAgentProfile>,
) -> AgentSession {
match selected {
Some(profile) => {
agent.with_appended_system_prompt(&render_primary_agent_prompt_append(profile))
}
None => agent,
}
}
fn apply_cached_context_window_if_eligible(
context_budget: &mut ContextBudget,
config: &EffectiveConfig,
selection: &ProviderSelection,
selected_custom_provider: Option<&CustomProviderConfig>,
) {
if selection.provider != OPENAI_CODEX_PROVIDER
&& selection.provider != crate::providers::ANTHROPIC_PROVIDER
&& selection.provider != crate::providers::CLAUDE_SUBSCRIPTION_PROVIDER
&& selected_custom_provider.is_none()
{
return;
}
if let Some(context_window) = crate::model_catalog::cached_model_context_window(
&config.paths,
&selection.provider,
&selection.model,
) {
context_budget.max_tokens = context_window;
}
context_budget.apply_model_limits(&selection.provider, &selection.model);
}
pub(crate) fn context_budget_for_selection(
config: &EffectiveConfig,
settings: &Settings,
selection: &ProviderSelection,
) -> ContextBudget {
let mut context_budget = settings.context.clone().unwrap_or_default();
let selected_custom_provider = config.custom_providers.get(&selection.provider);
apply_cached_context_window_if_eligible(
&mut context_budget,
config,
selection,
selected_custom_provider,
);
settings
.compaction
.apply_conversation_limit(&mut context_budget, &selection.provider);
context_budget
}
pub(crate) fn run_bash_mode_once(
config: &EffectiveConfig,
skills: &crate::skills::SkillDiscovery,
mut options: ProviderRunOptions<'_, '_>,
) -> Result<crate::agent::AgentRunOutput> {
let invocation_mode = options.invocation_mode;
let reporter =
crate::agent::filtered_herdr_reporter(invocation_mode, options.herdr_reporter.clone());
let mut turn = crate::herdr::HerdrTurnReporter::pending(reporter.clone());
options.herdr_reporter = reporter;
let result = run_bash_mode_once_inner(config, skills, options);
turn.finish_bash_result(&result);
result
}
fn run_bash_mode_once_inner(
config: &EffectiveConfig,
skills: &crate::skills::SkillDiscovery,
mut options: ProviderRunOptions<'_, '_>,
) -> Result<crate::agent::AgentRunOutput> {
let Some(bash_prompt) = crate::bash_mode::parse_bash_mode_prompt(options.prompt) else {
anyhow::bail!("bash-mode prompt must start with '!'");
};
let cancellation = options
.cancellation
.take()
.map(AgentCancellation::new)
.unwrap_or_default();
cancellation.check()?;
let settings = options
.settings
.take()
.map(Ok)
.unwrap_or_else(|| crate::config::read_settings(&config.paths))?;
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::SessionHeader {
session_id: options.session.map(|session| session.id().to_string()),
model: config
.model
.clone()
.unwrap_or_else(|| crate::providers::DEFAULT_CODEX_MODEL.to_string()),
cwd: options.cwd.to_path_buf(),
})?;
}
let mut session_persistence =
crate::agent::session_persistence::SessionPersistence::new(None, options.cwd);
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::BashCommand {
command: bash_prompt.command.to_string(),
})?;
}
let hook_settings = settings.hooks.clone();
let tool_settings = settings.tools.clone();
let disabled_tools = options.disabled_tools.clone().unwrap_or_else(|| {
Arc::new(Mutex::new(
crate::config::disabled_tool_names_from_settings(&settings)
.into_iter()
.collect(),
))
});
let tools = crate::tools::ToolRuntime::new_with_full_settings_and_mcp_with_disabled_tools(
options.cwd,
config.paths.clone(),
settings.clone(),
options.mcp.clone(),
disabled_tools,
)?
.with_skills(skills);
let hooks = crate::hooks::HookRuntime::new_with_tool_settings(
options.cwd,
hook_settings,
&tool_settings,
)?;
let call = ToolCall {
id: "bash_mode_1".to_string(),
name: crate::tools::contract::tool_name::BASH.to_string(),
arguments: json!({"command": bash_prompt.command, "intent": "Run the user's shell command"}),
};
let activity =
crate::agent::tool_lifecycle::emit_tool_started(&mut options.output_sink, 0, 0, &call)?;
let dispatch_context = ToolDispatchContext::new_with_hook_context_and_cancellation(
Some(activity.activity_id.clone()),
activity.activity_sender.clone(),
HookContextMetadata {
session_id: options.session.map(|session| session.id().to_string()),
session_path: options.session.map(|session| session.path().to_path_buf()),
provider_id: None,
model_id: None,
agent_id: None,
invocation_mode: options.invocation_mode,
turn_id: Some("turn-0".to_string()),
message_id: Some(format!("tool-{}", activity.activity_id.as_str())),
subagent: false,
},
cancellation.clone(),
);
if let Some(reporter) = options.herdr_reporter.as_ref() {
reporter.report_bash();
}
let outcome = crate::agent::tool_lifecycle::execute_tool_call(
Some(&tools),
Some(&hooks),
&mut session_persistence,
&mut options.output_sink,
None,
call.clone(),
dispatch_context,
)?;
let mut output = crate::agent::AgentRunOutput::default();
match outcome {
crate::agent::tool_lifecycle::ToolExecutionOutcome::Dispatched {
tool_result,
after_hook_failure,
..
} => {
crate::agent::tool_lifecycle::emit_tool_finished(
&mut options.output_sink,
activity.activity_id,
&call,
&tool_result,
)?;
output.tool_results.push((*tool_result).clone());
if let Some(diagnostic) = after_hook_failure {
return Err(crate::hooks::HookPolicyError::new(diagnostic).into());
}
}
crate::agent::tool_lifecycle::ToolExecutionOutcome::BeforeHookFailed {
diagnostic,
tool_result,
} => {
crate::agent::tool_lifecycle::emit_tool_finished(
&mut options.output_sink,
activity.activity_id,
&call,
&tool_result,
)?;
return Err(crate::hooks::HookPolicyError::new(diagnostic).into());
}
}
cancellation.check()?;
Ok(output)
}
pub(crate) fn run_provider_once_streaming_with_steering(
config: &EffectiveConfig,
instructions: &[crate::instructions::InstructionFile],
skills: &crate::skills::SkillDiscovery,
mut options: ProviderRunOptions<'_, '_>,
steering: AgentSteering,
) -> Result<crate::agent::AgentRunOutput> {
let invocation_mode = options.invocation_mode;
let reporter = options.herdr_reporter.clone();
crate::agent::run_with_herdr_turn(invocation_mode, reporter, |herdr_reporter| {
options.herdr_reporter = herdr_reporter;
let cancellation = options
.cancellation
.take()
.map(AgentCancellation::new)
.unwrap_or_default();
cancellation.check()?;
let prepared =
preparation::prepare(config, instructions, skills, &mut options, cancellation)?;
auto_compaction::run(prepared, instructions, &mut options, Some(steering))
})
}