use crate::{
agent::cancellation::AgentCancellation,
agent::steering::AgentSteering,
agent::{AgentSession, AgentSessionConfig},
config::{
AnthropicCacheTtl, CustomProviderConfig, EffectiveConfig, ProviderCredential, Settings,
},
context::ContextBudget,
output::{HookContextMetadata, InvocationMode, OutputEvent, ToolDispatchContext},
primary_agents::{PrimaryAgentProfile, render_primary_agent_prompt_append},
providers::{
AnthropicProvider, CLAUDE_CODE_PROVIDER, ClaudeCodeProvider, CodexRefreshedAuth,
OPENAI_CODEX_PROVIDER, OpenAiCodexProvider, OpenAiCompatibleProvider, Provider,
ProviderSelection, ReqwestHttpTransport, ToolCall, claude_code::auth::ClaudeCodeAuth,
},
sessions::Session,
};
use anyhow::Result;
use serde_json::json;
use std::{
collections::HashSet,
path::Path,
sync::{Arc, Mutex},
};
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) session_title_notifier: Option<crate::sessions::titles::SessionTitleNotifier>,
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>>>,
}
pub(crate) fn new_parent_agent_with_subagent_profiles(
model: impl Into<String>,
prompt_dir: Option<&Path>,
instructions: &[crate::instructions::InstructionFile],
skills: &crate::skills::SkillDiscovery,
discovery: &crate::subagents::profiles::SubagentProfileDiscovery,
) -> Result<AgentSession> {
let subagents_section =
crate::subagents::profiles::render_subagent_profiles_prompt(prompt_dir, discovery)?;
new_parent_agent_with_rendered_subagent_profiles(
model,
prompt_dir,
instructions,
skills,
subagents_section.as_deref(),
)
}
fn new_parent_agent_with_rendered_subagent_profiles(
model: impl Into<String>,
prompt_dir: Option<&Path>,
instructions: &[crate::instructions::InstructionFile],
skills: &crate::skills::SkillDiscovery,
subagents_section: Option<&str>,
) -> Result<AgentSession> {
AgentSession::new_with_prompt_dir_and_subagents(
model,
prompt_dir,
instructions,
skills,
subagents_section,
)
}
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 supported_custom_provider<'a>(
config: &'a EffectiveConfig,
selection: &ProviderSelection,
) -> Result<Option<&'a CustomProviderConfig>> {
let selected_custom_provider = config.custom_providers.get(&selection.provider);
if selection.provider.as_str() != OPENAI_CODEX_PROVIDER
&& selection.provider.as_str() != crate::providers::ANTHROPIC_PROVIDER
&& selection.provider.as_str() != CLAUDE_CODE_PROVIDER
&& selected_custom_provider.is_none()
{
anyhow::bail!(
"unsupported provider '{}'; supported providers: '{}', '{}', '{}' and configured custom providers",
selection.provider,
OPENAI_CODEX_PROVIDER,
crate::providers::ANTHROPIC_PROVIDER,
CLAUDE_CODE_PROVIDER
);
}
Ok(selected_custom_provider)
}
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 != CLAUDE_CODE_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_override(&selection.provider, &selection.model);
}
fn custom_provider_from_config<T>(
paths: &crate::config::McPaths,
provider_id: &str,
model: &str,
api_key: Option<String>,
custom: &CustomProviderConfig,
transport: T,
) -> OpenAiCompatibleProvider<T> {
OpenAiCompatibleProvider::custom(
provider_id.to_string(),
model.to_string(),
api_key,
custom.base_url.clone(),
custom.use_responses_endpoint,
transport,
)
.with_text_verbosity_support(custom.supports_text_verbosity)
.with_reasoning_protocol(
custom.reasoning_protocol,
crate::model_catalog::cached_model_max_output_tokens(paths, provider_id, model),
)
}
pub(crate) fn provider_from_selection(
config: &EffectiveConfig,
selection: &ProviderSelection,
cwd: &Path,
) -> Result<Arc<dyn Provider>> {
let cache_ttl = crate::config::read_settings(&config.paths)?.anthropic_cache_ttl;
provider_from_selection_with_anthropic_cache_ttl(config, selection, cwd, cache_ttl)
}
fn provider_from_selection_with_anthropic_cache_ttl(
config: &EffectiveConfig,
selection: &ProviderSelection,
_cwd: &Path,
anthropic_cache_ttl: Option<AnthropicCacheTtl>,
) -> Result<Arc<dyn Provider>> {
let selected_custom_provider = supported_custom_provider(config, selection)?;
let auth = config.require_auth()?;
match (selection.provider.as_str(), auth) {
(OPENAI_CODEX_PROVIDER, ProviderCredential::OAuth { access, account_id }) => {
let paths = config.paths.clone();
Ok(Arc::new(
OpenAiCodexProvider::new(
selection.model.clone(),
access,
account_id,
ReqwestHttpTransport,
)
.with_auth_refresh(move |cancellation| {
cancellation.check()?;
let ProviderCredential::OAuth { access, account_id } =
crate::login::force_refresh_codex_credential_from_store(&paths)?
else {
unreachable!("openai-codex forced refresh returns OAuth credentials")
};
cancellation.check()?;
Ok(CodexRefreshedAuth {
access_token: access,
account_id,
})
}),
))
}
(OPENAI_CODEX_PROVIDER, ProviderCredential::ApiKey { .. }) => anyhow::bail!(
"provider 'openai-codex' requires provider-keyed OAuth auth; refusing to send API-key/runtime token to Codex transport"
),
(OPENAI_CODEX_PROVIDER, ProviderCredential::NoAuth) => anyhow::bail!(
"provider 'openai-codex' requires provider-keyed OAuth auth; refusing no-auth configuration"
),
(crate::providers::ANTHROPIC_PROVIDER, ProviderCredential::ApiKey { key }) => {
let max_output_tokens = crate::model_catalog::cached_model_max_output_tokens(
&config.paths,
&selection.provider,
&selection.model,
);
let thinking_level = crate::thinking::resolve_thinking_level(
&crate::thinking::available_thinking_levels(
&selection.provider,
&selection.model,
crate::model_catalog::cached_model_thinking_metadata(
&config.paths,
&selection.provider,
&selection.model,
)
.as_ref(),
crate::thinking::ThinkingCapabilityScope::BuiltIn,
),
config.thinking_level,
);
Ok(Arc::new(
AnthropicProvider::new(selection.model.clone(), key, ReqwestHttpTransport)
.with_cache_ttl(anthropic_cache_ttl)
.with_max_output_tokens(max_output_tokens)
.with_thinking_level(thinking_level),
))
}
(crate::providers::ANTHROPIC_PROVIDER, ProviderCredential::OAuth { .. }) => {
anyhow::bail!("provider 'anthropic' requires Anthropic API-key auth, not OAuth token")
}
(crate::providers::ANTHROPIC_PROVIDER, ProviderCredential::NoAuth) => anyhow::bail!(
"provider 'anthropic' requires Anthropic API-key auth; refusing no-auth configuration"
),
(CLAUDE_CODE_PROVIDER, ProviderCredential::OAuth { .. }) => {
let max_output_tokens = crate::model_catalog::cached_model_max_output_tokens(
&config.paths,
&selection.provider,
&selection.model,
);
Ok(Arc::new(
ClaudeCodeProvider::new(
selection.model.clone(),
ClaudeCodeAuth::resolve_oauth()?,
ReqwestHttpTransport,
)
.with_cache_ttl(anthropic_cache_ttl)
.with_max_output_tokens(max_output_tokens),
))
}
(CLAUDE_CODE_PROVIDER, ProviderCredential::ApiKey { key }) => {
let max_output_tokens = crate::model_catalog::cached_model_max_output_tokens(
&config.paths,
&selection.provider,
&selection.model,
);
Ok(Arc::new(
ClaudeCodeProvider::new(
selection.model.clone(),
ClaudeCodeAuth::ApiKey { key },
ReqwestHttpTransport,
)
.with_cache_ttl(anthropic_cache_ttl)
.with_max_output_tokens(max_output_tokens),
))
}
(provider_id, ProviderCredential::ApiKey { key }) => {
let Some(custom) = selected_custom_provider else {
unreachable!("provider support checked above")
};
Ok(Arc::new(custom_provider_from_config(
&config.paths,
provider_id,
&selection.model,
Some(key),
custom,
ReqwestHttpTransport,
)))
}
(provider_id, ProviderCredential::NoAuth) => {
let Some(custom) = selected_custom_provider else {
unreachable!("provider support checked above")
};
Ok(Arc::new(custom_provider_from_config(
&config.paths,
provider_id,
&selection.model,
None,
custom,
ReqwestHttpTransport,
)))
}
(provider_id, ProviderCredential::OAuth { .. }) => anyhow::bail!(
"provider '{provider_id}' requires custom-provider auth mode, not OAuth token"
),
}
}
pub(crate) fn run_bash_mode_once(
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 tools = crate::tools::ToolRuntime::new_with_full_settings_and_mcp(
options.cwd,
config.paths.clone(),
settings.clone(),
options.mcp.clone(),
)?
.with_disabled_tools(options.disabled_tools.clone().unwrap_or_else(|| {
Arc::new(Mutex::new(
crate::config::disabled_tool_names_from_settings(&settings)
.into_iter()
.collect(),
))
}))
.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}),
};
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(),
);
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 } => {
let local_result = crate::agent::tool_lifecycle::before_hook_failure_activity_result(
&call,
&diagnostic,
);
crate::agent::tool_lifecycle::emit_tool_finished(
&mut options.output_sink,
activity.activity_id,
&call,
&local_result,
)?;
return Err(crate::hooks::HookPolicyError::new(diagnostic).into());
}
}
cancellation.check()?;
Ok(output)
}
pub(crate) fn run_provider_once_streaming(
config: &EffectiveConfig,
instructions: &[crate::instructions::InstructionFile],
skills: &crate::skills::SkillDiscovery,
options: ProviderRunOptions<'_, '_>,
) -> Result<crate::agent::AgentRunOutput> {
run_provider_once_streaming_inner(config, instructions, skills, options, None)
}
pub(crate) fn run_provider_once_streaming_with_steering(
config: &EffectiveConfig,
instructions: &[crate::instructions::InstructionFile],
skills: &crate::skills::SkillDiscovery,
options: ProviderRunOptions<'_, '_>,
steering: AgentSteering,
) -> Result<crate::agent::AgentRunOutput> {
run_provider_once_streaming_inner(config, instructions, skills, options, Some(steering))
}
fn run_provider_once_streaming_inner(
config: &EffectiveConfig,
instructions: &[crate::instructions::InstructionFile],
skills: &crate::skills::SkillDiscovery,
mut options: ProviderRunOptions<'_, '_>,
steering: Option<AgentSteering>,
) -> Result<crate::agent::AgentRunOutput> {
let cancellation = options
.cancellation
.take()
.map(AgentCancellation::new)
.unwrap_or_default();
cancellation.check()?;
let mut refreshed_config;
let config = if config.provider_id() == OPENAI_CODEX_PROVIDER {
refreshed_config = config.clone();
let credential = crate::login::codex_credential_from_store(&config.paths)?;
refreshed_config.auth = Some(credential);
&refreshed_config
} else {
config
};
let selection = ProviderSelection::from_config(config)?;
let selected_custom_provider = supported_custom_provider(config, &selection)?;
let settings = options
.settings
.take()
.map(Ok)
.unwrap_or_else(|| crate::config::read_settings(&config.paths))?;
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 disabled_tool_names = disabled_tools
.lock()
.map(|disabled| disabled.clone())
.map_err(|_| anyhow::anyhow!("disabled tools lock poisoned"))?;
let herdr_reporter = crate::herdr::HerdrReporter::from_env(&settings.integrations.herdr);
let mut context_budget = settings.context.clone().unwrap_or_default();
apply_cached_context_window_if_eligible(
&mut context_budget,
config,
&selection,
selected_custom_provider,
);
let model = selection.model.clone();
let mut cached_thinking_metadata = crate::model_catalog::cached_model_thinking_metadata(
&config.paths,
&selection.provider,
&selection.model,
);
if selected_custom_provider.is_some_and(|_| {
cached_thinking_metadata.as_ref().is_none_or(|metadata| {
metadata.supports_reasoning.is_none() && metadata.reasoning_efforts.is_none()
})
}) {
let refresh_result = if crate::model_catalog::automatic_refresh_allowed(&selection.provider)
{
crate::model_catalog::refresh_catalog_for_provider(
&config.paths,
&selection.provider,
&selection.model,
)
} else {
Ok(crate::model_catalog::CatalogRefreshOutcome::Updated)
};
if matches!(
&refresh_result,
Ok(crate::model_catalog::CatalogRefreshOutcome::NoMatch)
) {
crate::model_catalog::automatic_refresh_suppressed(&selection.provider);
}
if let Err(error) = refresh_result {
crate::model_catalog::automatic_refresh_failed(&selection.provider);
let message = format!(
"metadata refresh failed for provider '{}' (category: {}); automatic retry temporarily suppressed",
selection.provider,
error.category().as_str()
);
let _ = crate::sessions::try_record_session_event(
options.session,
options.cwd,
crate::sessions::SessionEventKind::Diagnostic,
serde_json::json!({"level":"warning", "message": message.clone()}),
);
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message,
})?;
}
}
cached_thinking_metadata = crate::model_catalog::cached_model_thinking_metadata(
&config.paths,
&selection.provider,
&selection.model,
);
}
let capability_scope = match selected_custom_provider {
Some(custom) => crate::thinking::ThinkingCapabilityScope::Custom(custom.reasoning_protocol),
None => crate::thinking::ThinkingCapabilityScope::BuiltIn,
};
let thinking_levels = crate::thinking::available_thinking_levels(
&selection.provider,
&selection.model,
cached_thinking_metadata.as_ref(),
capability_scope,
);
let thinking_level =
crate::thinking::resolve_thinking_level(&thinking_levels, config.thinking_level);
let send_default_reasoning_summary = thinking_levels.len() > 1;
let subagent_profile_discovery =
options
.subagent_profile_discovery
.clone()
.unwrap_or_else(|| {
crate::subagents::profiles::discover_subagent_profiles(&config.paths.subagents)
});
let disabled_subagent_profiles = options
.disabled_subagent_profiles
.as_ref()
.map(|disabled| {
disabled
.lock()
.map(|disabled| disabled.clone())
.map_err(|_| anyhow::anyhow!("disabled subagent profiles lock poisoned"))
})
.transpose()?
.unwrap_or_else(|| {
crate::config::disabled_subagent_profile_names_from_settings(&settings)
.into_iter()
.collect()
});
let enabled_subagent_profile_discovery = crate::subagents::profiles::filter_enabled_profiles(
&subagent_profile_discovery,
&disabled_subagent_profiles,
);
let subagent_profiles_prompt =
if disabled_tool_names.contains(crate::tools::contract::tool_name::SUBAGENTS) {
None
} else {
crate::subagents::profiles::render_subagent_profiles_prompt(
Some(&config.paths.prompts),
&enabled_subagent_profile_discovery,
)?
};
let agent_config = AgentSessionConfig::new(selection.provider.clone(), model.clone())
.with_thinking_level(thinking_level)
.with_text_verbosity(settings.text_verbosity_for(&selection.provider))
.with_thinking_levels(thinking_levels)
.with_default_reasoning_summary(send_default_reasoning_summary)
.with_context_budget(context_budget)
.with_context_cache_dir(config.paths.cache.clone());
let agent = AgentSession::new_with_prompt_dir_and_subagents_and_disabled(
model.clone(),
Some(&config.paths.prompts),
instructions,
skills,
None,
&disabled_tool_names,
)?
.with_config(agent_config.clone());
let parent_agent_for_provider = AgentSession::new_with_prompt_dir_and_subagents_and_disabled(
model.clone(),
Some(&config.paths.prompts),
instructions,
skills,
subagent_profiles_prompt.as_deref(),
&disabled_tool_names,
)?
.with_config(agent_config);
let parent_agent_for_provider = append_primary_agent_to_main_prompt(
parent_agent_for_provider,
options.selected_primary_agent.as_ref(),
);
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,
cwd: options.cwd.to_path_buf(),
})?;
}
let provider = provider_from_selection_with_anthropic_cache_ttl(
config,
&selection,
options.cwd,
settings.anthropic_cache_ttl,
)?;
let title_job = match settings.session_titles.eligible_config() {
Ok(Some(title_config)) => {
options
.session
.cloned()
.map(|session| crate::sessions::titles::SessionTitleJob {
paths: config.paths.clone(),
session,
cwd: options.cwd.to_path_buf(),
config: title_config,
first_prompt: options.prompt.to_string(),
cancellation: cancellation.clone(),
notifier: options.session_title_notifier.clone(),
})
}
Ok(None) => None,
Err(message) => {
let sanitized = crate::output::redact_sensitive_text(&message);
if let Err(persistence_error) = crate::sessions::try_record_session_event(
options.session,
options.cwd,
crate::sessions::SessionEventKind::Diagnostic,
serde_json::json!({"level":"warning", "message": sanitized}),
) && let Some(sink) = options.output_sink.as_deref_mut()
{
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: format!(
"session persistence failed; future resume may be incomplete: {persistence_error}",
),
})?;
}
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: sanitized,
})?;
}
None
}
};
let hook_settings = settings.hooks.clone();
let tool_settings = settings.tools.clone();
let base_tools = crate::tools::ToolRuntime::new_with_full_settings_and_mcp(
options.cwd,
config.paths.clone(),
settings.clone(),
options.mcp.clone(),
)?
.with_disabled_tools(Arc::clone(&disabled_tools))
.with_skills(skills);
let hooks = crate::hooks::HookRuntime::new_with_tool_settings(
options.cwd,
hook_settings,
&tool_settings,
)?;
let inherited_hooks = if hooks.is_inert() {
None
} else {
Some(hooks.clone())
};
let subagent_provider_override = crate::subagents::SubagentProviderOverride::new(
config.paths.clone(),
Arc::new({
let paths = config.paths.clone();
move |selection: &ProviderSelection, cwd: &Path| {
let provider_config = crate::config::load_effective_provider_selection(
&paths,
&selection.provider,
&selection.model,
)?;
let provider_selection = ProviderSelection::from_config(&provider_config)?;
let provider = provider_from_selection(&provider_config, &provider_selection, cwd)?;
let scope = crate::thinking::capability_scope_for_provider(
&provider_config.custom_providers,
&provider_selection.provider,
);
Ok(crate::subagents::ResolvedProviderOverride { provider, scope })
}
}),
);
let subagent_config = crate::subagents::SubagentRunConfig {
parent_agent: agent.clone(),
provider: Arc::clone(&provider),
provider_override: Some(subagent_provider_override),
parent_tools: base_tools.clone(),
parent_cwd: options.cwd.to_path_buf(),
cancellation: cancellation.clone(),
profiles: enabled_subagent_profile_discovery.profiles,
subagent_profiles_prompt: subagent_profiles_prompt.clone(),
sessions_root: options
.session
.and_then(|active| active.path().parent().map(std::path::Path::to_path_buf)),
depth: 0,
parent_activity_id: None,
activity_sender: None,
inherited_hooks,
semantic_progress_timeout: settings
.provider_stream
.subagent_semantic_progress_timeout(),
schema_validation_max_retries: settings.subagents.schema_validation_max_retries(),
};
let tools = base_tools.with_subagents(move |arguments, context| {
let mut config = subagent_config.clone();
config.parent_activity_id = context.parent_activity_id;
config.activity_sender = context.activity_sender;
crate::subagents::dispatch_subagents(arguments, config)
});
let request = crate::agent::AgentRunRequest {
prompt: options.prompt,
tools: Some(&tools),
hooks: Some(&hooks),
session: options.session,
cwd: options.cwd,
output_sink: options.output_sink,
cancellation,
session_title_job: title_job,
semantic_progress_timeout: Some(settings.provider_stream.semantic_progress_timeout()),
invocation_mode: options.invocation_mode,
agent_id: options
.selected_primary_agent
.as_ref()
.map(|profile| profile.id.clone()),
initial_instructions: instructions,
ttsr: settings.ttsr.clone(),
herdr_reporter,
};
match steering {
Some(steering) => parent_agent_for_provider
.run_print_with_tools_streaming_output_cancellable_with_steering(
provider.as_ref(),
request,
steering,
),
None => parent_agent_for_provider
.run_print_with_tools_streaming_output_cancellable(provider.as_ref(), request),
}
}
#[cfg(test)]
mod primary_agent_tests {
use super::*;
use crate::{
config::{McPaths, ProviderCredential},
context::ContextBudgetOverride,
instructions::InstructionFile,
model_catalog::ModelCatalogEntry,
skills::SkillDiscovery,
};
use std::{collections::BTreeMap, path::PathBuf};
fn profile() -> PrimaryAgentProfile {
PrimaryAgentProfile {
id: "tars".to_string(),
name: "TARS".to_string(),
description: "Tactical".to_string(),
path: PathBuf::new(),
prompt: "PRIMARY BODY".to_string(),
}
}
#[derive(Debug, Clone)]
struct NoopTransport;
fn base_agent() -> AgentSession {
AgentSession::new(
"gpt-test",
&[] as &[InstructionFile],
&SkillDiscovery::default(),
)
}
fn base_config(paths: McPaths, provider: &str, model: &str) -> EffectiveConfig {
EffectiveConfig {
provider: Some(provider.to_string()),
model: Some(model.to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: Some(ProviderCredential::OAuth {
access: "test-token".to_string(),
account_id: Some("acct".to_string()),
}),
paths,
}
}
#[test]
fn unsupported_provider_validation_preserves_error_text() {
let temp = tempfile::TempDir::new().unwrap();
let config = base_config(
McPaths::from_root(temp.path().join("mc")),
"unsupported",
"model-a",
);
let selection = ProviderSelection {
provider: "unsupported".to_string(),
model: "model-a".to_string(),
};
let error = supported_custom_provider(&config, &selection)
.unwrap_err()
.to_string();
assert_eq!(
error,
"unsupported provider 'unsupported'; supported providers: 'openai-codex', 'anthropic', 'claude-code' and configured custom providers"
);
}
#[test]
fn custom_provider_from_config_uses_selected_endpoint_mode() {
let chat = CustomProviderConfig {
label: "Chat Provider".to_string(),
base_url: "https://chat.example/v1".to_string(),
api_key_env_var: None,
models_dev_provider: None,
use_responses_endpoint: false,
supports_text_verbosity: false,
reasoning_protocol: crate::config::CustomReasoningProtocol::default(),
extra_models: Vec::new(),
};
let responses = CustomProviderConfig {
label: "Responses Provider".to_string(),
base_url: "https://responses.example/v1".to_string(),
api_key_env_var: None,
models_dev_provider: None,
use_responses_endpoint: true,
supports_text_verbosity: false,
reasoning_protocol: crate::config::CustomReasoningProtocol::default(),
extra_models: Vec::new(),
};
let request = crate::providers::ProviderRequest::new(
"model-a",
vec![crate::providers::ChatMessage::user("hello")],
);
let chat_http = custom_provider_from_config(
&crate::config::McPaths::from_root(std::env::temp_dir()),
"chat-provider",
"model-a",
None,
&chat,
NoopTransport,
)
.build_http_request(&request);
let responses_http = custom_provider_from_config(
&crate::config::McPaths::from_root(std::env::temp_dir()),
"responses-provider",
"model-a",
None,
&responses,
NoopTransport,
)
.build_http_request(&request);
assert_eq!(chat_http.url, "https://chat.example/v1/chat/completions");
assert_eq!(responses_http.url, "https://responses.example/v1/responses");
assert!(chat_http.body.get("messages").is_some());
assert!(responses_http.body.get("input").is_some());
}
#[test]
fn custom_provider_from_config_uses_selected_reasoning_protocol() {
let temp = tempfile::TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut entry = ModelCatalogEntry::new("anthropic-custom", "model-a");
entry.max_output_tokens = Some(32_768);
let custom = CustomProviderConfig {
label: "Anthropic-like".to_string(),
base_url: "https://provider.example/v1".to_string(),
api_key_env_var: None,
models_dev_provider: None,
use_responses_endpoint: false,
supports_text_verbosity: false,
reasoning_protocol: crate::config::CustomReasoningProtocol::AnthropicLike,
extra_models: Vec::new(),
};
crate::config::write_settings(
&paths,
&crate::config::Settings {
custom_providers: BTreeMap::from([(
"anthropic-custom".to_string(),
custom.clone(),
)]),
..Default::default()
},
)
.unwrap();
crate::model_catalog::write_catalog_cache_for_configured_provider(
&paths,
"anthropic-custom",
&[entry],
)
.unwrap();
let request = crate::providers::ProviderRequest::new(
"model-a",
vec![crate::providers::ChatMessage::user("hello")],
)
.with_thinking_level(crate::thinking::ThinkingLevel::High)
.with_default_reasoning_summary(true);
for use_responses_endpoint in [false, true] {
let custom = CustomProviderConfig {
use_responses_endpoint,
..custom.clone()
};
let body = custom_provider_from_config(
&paths,
"anthropic-custom",
"model-a",
None,
&custom,
NoopTransport,
)
.build_http_request(&request)
.body;
assert_eq!(
body["thinking"],
serde_json::json!({
"type": "enabled",
"budget_tokens": 16_384,
})
);
assert_eq!(body["max_tokens"], 32_768);
assert!(body.get("reasoning_effort").is_none());
assert!(body.get("reasoning").is_none());
}
}
#[test]
fn cached_context_window_applies_only_for_supported_catalog_providers() {
let temp = tempfile::TempDir::new().unwrap();
let mut config = base_config(
McPaths::from_root(temp.path().join("mc")),
OPENAI_CODEX_PROVIDER,
"gpt-test",
);
let mut entry = ModelCatalogEntry::new_codex("gpt-test");
entry.context_window = Some(196_000);
crate::model_catalog::write_catalog_cache(&config.paths, OPENAI_CODEX_PROVIDER, &[entry])
.unwrap();
let selection = ProviderSelection {
provider: OPENAI_CODEX_PROVIDER.to_string(),
model: "gpt-test".to_string(),
};
let mut budget = ContextBudget {
max_tokens: 10,
..ContextBudget::default()
};
apply_cached_context_window_if_eligible(&mut budget, &config, &selection, None);
assert_eq!(budget.max_tokens, 196_000);
config.provider = Some("unsupported".to_string());
let unsupported = ProviderSelection {
provider: "unsupported".to_string(),
model: "gpt-test".to_string(),
};
let mut budget = ContextBudget {
max_tokens: 10,
..ContextBudget::default()
};
apply_cached_context_window_if_eligible(&mut budget, &config, &unsupported, None);
assert_eq!(budget.max_tokens, 10);
}
#[test]
fn model_context_override_wins_after_cached_window_for_provider_runs() {
let temp = tempfile::TempDir::new().unwrap();
let config = base_config(
McPaths::from_root(temp.path().join("mc")),
OPENAI_CODEX_PROVIDER,
"gpt-test",
);
let mut entry = ModelCatalogEntry::new_codex("gpt-test");
entry.context_window = Some(196_000);
crate::model_catalog::write_catalog_cache(&config.paths, OPENAI_CODEX_PROVIDER, &[entry])
.unwrap();
let selection = ProviderSelection {
provider: OPENAI_CODEX_PROVIDER.to_string(),
model: "gpt-test".to_string(),
};
let mut budget = ContextBudget {
max_tokens: 10,
reserve_tokens: 2,
keep_recent_tokens: 3,
model_overrides: BTreeMap::from([(
"openai-codex/gpt-test".to_string(),
ContextBudgetOverride {
max_tokens: Some(400_000),
reserve_tokens: Some(32_000),
keep_recent_tokens: Some(50_000),
},
)]),
..ContextBudget::default()
};
apply_cached_context_window_if_eligible(&mut budget, &config, &selection, None);
assert_eq!(budget.max_tokens, 400_000);
assert_eq!(budget.reserve_tokens, 32_000);
assert_eq!(budget.keep_recent_tokens, 50_000);
}
#[test]
fn model_context_override_requires_exact_provider_model_match() {
let temp = tempfile::TempDir::new().unwrap();
let config = base_config(
McPaths::from_root(temp.path().join("mc")),
OPENAI_CODEX_PROVIDER,
"gpt-other",
);
let selection = ProviderSelection {
provider: OPENAI_CODEX_PROVIDER.to_string(),
model: "gpt-other".to_string(),
};
let mut budget = ContextBudget {
max_tokens: 10,
model_overrides: BTreeMap::from([(
"openai-codex/gpt-test".to_string(),
ContextBudgetOverride {
max_tokens: Some(400_000),
..ContextBudgetOverride::default()
},
)]),
..ContextBudget::default()
};
apply_cached_context_window_if_eligible(&mut budget, &config, &selection, None);
assert_eq!(budget.max_tokens, 10);
}
#[test]
fn provider_from_selection_constructs_anthropic_provider() {
let temp = tempfile::TempDir::new().unwrap();
let config = EffectiveConfig {
provider: Some(crate::providers::ANTHROPIC_PROVIDER.to_string()),
model: Some(crate::providers::DEFAULT_ANTHROPIC_MODEL.to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: Some("anthropic-key".to_string()),
auth: Some(ProviderCredential::ApiKey {
key: "anthropic-key".to_string(),
}),
paths: McPaths::from_root(temp.path().join("mc")),
};
let selection = ProviderSelection::from_config(&config).unwrap();
let _provider = provider_from_selection(&config, &selection, temp.path()).unwrap();
}
#[test]
fn provider_from_selection_constructs_claude_code_provider_with_api_key() {
let temp = tempfile::TempDir::new().unwrap();
let config = EffectiveConfig {
provider: Some(crate::providers::CLAUDE_CODE_PROVIDER.to_string()),
model: Some(crate::providers::DEFAULT_CLAUDE_CODE_MODEL.to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: Some(ProviderCredential::ApiKey {
key: "api-key".to_string(),
}),
paths: McPaths::from_root(temp.path().join("mc")),
};
let selection = ProviderSelection::from_config(&config).unwrap();
let _provider = provider_from_selection(&config, &selection, temp.path()).unwrap();
}
#[test]
fn provider_selection_threads_anthropic_cache_ttl_into_claude_code_provider() {
let temp = tempfile::TempDir::new().unwrap();
let config = EffectiveConfig {
provider: Some(crate::providers::CLAUDE_CODE_PROVIDER.to_string()),
model: Some(crate::providers::DEFAULT_CLAUDE_CODE_MODEL.to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: Some(ProviderCredential::ApiKey {
key: "api-key".to_string(),
}),
paths: McPaths::from_root(temp.path().join("mc")),
};
let selection = ProviderSelection::from_config(&config).unwrap();
let provider = provider_from_selection_with_anthropic_cache_ttl(
&config,
&selection,
temp.path(),
Some(AnthropicCacheTtl::OneHour),
)
.unwrap();
assert_eq!(
provider.anthropic_cache_ttl_for_test(),
Some(AnthropicCacheTtl::OneHour)
);
}
#[test]
fn claude_code_rejects_no_auth_credential() {
let temp = tempfile::TempDir::new().unwrap();
let config = EffectiveConfig {
provider: Some(crate::providers::CLAUDE_CODE_PROVIDER.to_string()),
model: Some(crate::providers::DEFAULT_CLAUDE_CODE_MODEL.to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: Some(ProviderCredential::NoAuth),
paths: McPaths::from_root(temp.path().join("mc")),
};
let selection = ProviderSelection {
provider: crate::providers::CLAUDE_CODE_PROVIDER.to_string(),
model: crate::providers::DEFAULT_CLAUDE_CODE_MODEL.to_string(),
};
let error = match provider_from_selection(&config, &selection, temp.path()) {
Ok(_) => panic!("claude-code accepted no-auth credential"),
Err(error) => error.to_string(),
};
assert!(
error.contains("requires Claude Code OAuth credentials"),
"{error}"
);
}
#[test]
fn primary_agent_none_preserves_main_prompt() {
let agent = base_agent();
let before = agent.system_prompt().to_string();
let after = append_primary_agent_to_main_prompt(agent, None);
assert_eq!(after.system_prompt(), before);
}
#[test]
fn selected_primary_agent_appends_to_provider_parent_prompt() {
let agent = base_agent();
let after = append_primary_agent_to_main_prompt(agent, Some(&profile()));
assert!(
after
.system_prompt()
.contains("Primary agent profile (TARS):")
);
assert!(after.system_prompt().contains("PRIMARY BODY"));
}
#[test]
fn selected_primary_agent_is_excluded_from_subagent_base_prompt() {
let agent = base_agent();
let plain = agent.clone();
let provider = append_primary_agent_to_main_prompt(agent, Some(&profile()));
assert!(provider.system_prompt().contains("PRIMARY BODY"));
assert!(!plain.system_prompt().contains("PRIMARY BODY"));
}
#[test]
fn system_prompt_modal_includes_selected_primary_agent_for_main_prompt() {
let agent = append_primary_agent_to_main_prompt(base_agent(), Some(&profile()));
assert!(agent.system_prompt().contains("PRIMARY BODY"));
}
#[test]
fn system_prompt_modal_none_preserves_existing_prompt() {
let agent = base_agent();
let before = agent.system_prompt().to_string();
let agent = append_primary_agent_to_main_prompt(agent, None);
assert_eq!(agent.system_prompt(), before);
}
}