use crate::{
agent::cancellation::{AgentCancellation, is_run_canceled},
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>>>,
}
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,
}
}
const AUTO_COMPACTION_DIAGNOSTIC_MAX_CHARS: usize = 512;
const MAX_AUTO_COMPACTIONS_PER_RUN: usize = 4;
fn auto_compaction_policy(
auto: &crate::config::AutoCompactionSettings,
model_max_tokens: usize,
) -> Option<(usize, String)> {
if !auto.is_enabled() {
return None;
}
let threshold = match (auto.threshold_percent, auto.threshold_tokens) {
(Some(percent), None) => {
let numerator = (model_max_tokens as u128).saturating_mul(u128::from(percent));
usize::try_from(numerator.div_ceil(100)).unwrap_or(usize::MAX)
}
(None, Some(tokens)) => usize::try_from(tokens).unwrap_or(usize::MAX),
_ => return None,
};
Some((threshold, auto.threshold_display()?))
}
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::Shell | InvocationMode::MissionControl
)
}
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.clone())
.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 auto = settings.compaction.auto.clone();
let auto_eligible = auto_compaction_eligible(
auto.is_enabled(),
context_budget.enabled,
options.session.is_some(),
options.invocation_mode,
);
let mut prompt = options.prompt.to_string();
let mut prompt_origin = crate::output::UserPromptOrigin::User;
let mut effective_prompt: Option<String> = None;
let mut title_job = title_job;
let mut combined = crate::agent::AgentRunOutput::default();
let mut auto_compaction_count = 0_usize;
let mut preflight_compaction_used = false;
let mut projected_candidate: Option<(String, String)> = None;
let auto_policy = auto_eligible
.then(|| auto_compaction_policy(&auto, parent_agent_for_provider.context_max_tokens()))
.flatten();
if auto_eligible {
cancellation.check()?;
let session = options
.session
.expect("auto-compaction eligibility requires session");
let effective_current_prompt = AgentSession::effective_prompt_for_projection(
&prompt,
Some(&tools),
options.invocation_mode,
);
let hard_threshold = context_budget.threshold_tokens();
let static_tokens = parent_agent_for_provider
.project_prompt_input_tokens_without_session(&effective_current_prompt, Some(&tools))?;
effective_prompt = Some(effective_current_prompt.clone());
if static_tokens <= hard_threshold {
let current_tokens = parent_agent_for_provider.project_prompt_input_tokens(
&effective_current_prompt,
session,
Some(&tools),
)?;
if current_tokens > hard_threshold {
preflight_compaction_used = true;
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionTriggered {
current_tokens,
max_tokens: parent_agent_for_provider.context_max_tokens(),
threshold: format!("preflight hard budget ({hard_threshold} tokens)"),
})?;
sink.output_event(OutputEvent::CompactionStarted)?;
}
let compaction =
crate::compaction::compact_session(crate::compaction::CompactSessionJob {
active_config: config.clone(),
settings: settings.clone(),
session: session.clone(),
cwd: options.cwd.to_path_buf(),
cancellation: cancellation.clone(),
custom_instructions: None,
});
match compaction {
Ok(Some(_)) => {}
Ok(None) => {
let message = bounded_auto_compaction_diagnostic(
"automatic preflight compaction produced no usable summary; old primary remains authoritative; no provider request was sent",
);
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.clone(),
})?;
}
return Err(anyhow::anyhow!(message));
}
Err(error) if is_run_canceled(&error) => {
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: "automatic preflight compaction canceled; old primary remains authoritative; no provider request was sent".to_string(),
})?;
}
return Err(error);
}
Err(error) => {
let authority = error
.downcast_ref::<crate::sessions::CompactionRotationError>()
.map(|_| "")
.unwrap_or(" old primary remains authoritative;");
let message = bounded_auto_compaction_diagnostic(format!(
"automatic preflight compaction failed;{authority} no provider request was sent: {error}"
));
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.clone(),
})?;
}
return Err(anyhow::anyhow!(message));
}
}
let compacted_tokens = match parent_agent_for_provider.ensure_prompt_context_fits(
&effective_current_prompt,
session,
Some(&tools),
) {
Ok(tokens) => tokens,
Err(error) => {
let message = bounded_auto_compaction_diagnostic(format!(
"automatic preflight compaction completed, but current prompt still exceeds hard context budget; new checkpoint remains authoritative; no provider request was sent: {error}"
));
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.clone(),
})?;
}
return Err(anyhow::anyhow!(message));
}
};
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionCompleted {
current_tokens: compacted_tokens,
max_tokens: parent_agent_for_provider.context_max_tokens(),
})?;
}
auto_compaction_count = auto_compaction_count.saturating_add(1);
}
}
}
loop {
let request = crate::agent::AgentRunRequest {
prompt: &prompt,
prompt_origin,
effective_prompt: effective_prompt.as_deref(),
tools: Some(&tools),
hooks: Some(&hooks),
session: options.session,
cwd: options.cwd,
output_sink: reborrow_output_sink(&mut options.output_sink),
cancellation: cancellation.clone(),
session_title_job: if auto_eligible {
None
} else {
title_job.take()
},
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: herdr_reporter.clone(),
continuation_auto_compaction_policy: (auto_eligible
&& auto_compaction_count < MAX_AUTO_COMPACTIONS_PER_RUN)
.then(|| auto_policy.clone())
.flatten(),
};
let output_result = match steering.as_ref() {
Some(steering) if auto_eligible => parent_agent_for_provider
.run_print_with_tools_streaming_output_cancellable_with_deferred_steering(
provider.as_ref(),
request,
steering.clone(),
),
Some(steering) => parent_agent_for_provider
.run_print_with_tools_streaming_output_cancellable_with_steering(
provider.as_ref(),
request,
steering.clone(),
),
None => parent_agent_for_provider
.run_print_with_tools_streaming_output_cancellable(provider.as_ref(), request),
};
let mut continuation_context_overflow = None;
let output = match output_result {
Ok(output) => output,
Err(error)
if auto_eligible
&& auto_compaction_count < MAX_AUTO_COMPACTIONS_PER_RUN
&& error
.downcast_ref::<crate::agent::ContextBudgetError>()
.is_some_and(|budget| {
budget.phase() == crate::agent::ContextBudgetPhase::Continuation
}) =>
{
let budget = error
.downcast::<crate::agent::ContextBudgetError>()
.expect("context budget error checked above");
let threshold = budget.threshold_tokens();
let threshold_display = auto_policy
.as_ref()
.filter(|(soft_threshold, _)| *soft_threshold == threshold)
.map(|(_, display)| display.clone())
.unwrap_or_else(|| format!("continuation hard budget ({threshold} tokens)"));
continuation_context_overflow = Some((
budget.estimated_tokens(),
threshold_display,
budget.to_string(),
));
budget.into_partial_output().unwrap_or_default()
}
Err(error)
if preflight_compaction_used
&& error
.downcast_ref::<crate::agent::RequiredUserInputPersistenceError>()
.is_some() =>
{
let message = bounded_auto_compaction_diagnostic(format!(
"current user input persistence failed after preflight compaction; new checkpoint remains authoritative; no provider request was sent: {error}"
));
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.clone(),
})?;
}
return Err(anyhow::anyhow!(message));
}
Err(error)
if auto_compaction_count > 0
&& error
.downcast_ref::<crate::agent::RequiredUserInputPersistenceError>()
.is_some() =>
{
let message = bounded_auto_compaction_diagnostic(format!(
"automatic continuation persistence failed; new checkpoint remains authoritative; no provider continuation was sent: {error}"
));
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.clone(),
})?;
}
return Err(anyhow::anyhow!(message));
}
Err(error) => return Err(error),
};
if auto_eligible
&& let Some(mut job) = title_job.take()
&& let Some(session) = options.session
&& let Ok(metadata) =
crate::sessions::titles::should_start_title_generation_metadata(session)
&& metadata.latest_title.is_none()
&& metadata.user_input_count == 1
{
if let Some(first_prompt) = metadata.first_user_input_text {
job.first_prompt = first_prompt;
}
let _ = crate::sessions::titles::spawn_background(job).join();
}
let persistence_degraded = output.persistence_degraded;
let recovered_incomplete_stream = output.recovered_incomplete_stream;
let auto_compaction_blocked_by_recovery = output.auto_compaction_blocked_by_recovery;
if !combined.text.is_empty() && !output.text.is_empty() {
combined.text.push('\n');
}
combined.text.push_str(&output.text);
combined.usage = output.usage;
combined.total_tokens = match (combined.total_tokens, output.total_tokens) {
(Some(total), Some(next)) => Some(total.saturating_add(next)),
(None, next) => next,
(total, None) => total,
};
combined.tool_results.extend(output.tool_results);
combined.persistence_degraded |= persistence_degraded;
combined.recovered_incomplete_stream |= recovered_incomplete_stream;
combined.auto_compaction_blocked_by_recovery |= auto_compaction_blocked_by_recovery;
if !auto_eligible || auto_compaction_count >= MAX_AUTO_COMPACTIONS_PER_RUN {
return Ok(combined);
}
if persistence_degraded || auto_compaction_blocked_by_recovery {
if let Some((_, _, message)) = continuation_context_overflow {
return Err(anyhow::anyhow!(message));
}
return Ok(combined);
}
cancellation.check()?;
let session = options
.session
.expect("auto-compaction eligibility requires session");
let observed_steering = steering.as_ref().and_then(AgentSteering::observe_collapsed);
let candidate = observed_steering
.as_ref()
.map(|batch| batch.text.as_str())
.unwrap_or("continue");
let effective_candidate = match projected_candidate.as_ref() {
Some((raw, effective)) if raw == candidate => effective.clone(),
_ => {
let effective = AgentSession::effective_prompt_for_projection(
candidate,
Some(&tools),
options.invocation_mode,
);
projected_candidate = Some((candidate.to_string(), effective.clone()));
effective
}
};
let current_tokens = parent_agent_for_provider.project_prompt_input_tokens(
&effective_candidate,
session,
Some(&tools),
)?;
let max_tokens = parent_agent_for_provider.context_max_tokens();
if continuation_context_overflow.is_none()
&& !auto.triggered(current_tokens as u64, max_tokens as u64)
{
if let Some(batch) = observed_steering {
prompt = batch.text;
effective_prompt = Some(effective_candidate);
prompt_origin = crate::output::UserPromptOrigin::Steering;
continue;
}
return Ok(combined);
}
let (trigger_tokens, trigger_threshold) = continuation_context_overflow
.as_ref()
.map(|(estimated, threshold, _)| (*estimated, threshold.clone()))
.unwrap_or_else(|| {
(
current_tokens,
auto.threshold_display()
.unwrap_or_else(|| "invalid threshold".to_string()),
)
});
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionTriggered {
current_tokens: trigger_tokens,
max_tokens,
threshold: trigger_threshold,
})?;
sink.output_event(OutputEvent::CompactionStarted)?;
}
let compaction = crate::compaction::compact_session(crate::compaction::CompactSessionJob {
active_config: config.clone(),
settings: settings.clone(),
session: session.clone(),
cwd: options.cwd.to_path_buf(),
cancellation: cancellation.clone(),
custom_instructions: None,
});
let result = match compaction {
Ok(Some(result)) => result,
Ok(None) => {
let message = "automatic compaction produced no usable summary; old primary remains authoritative; no continuation was sent";
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.to_string(),
})?;
}
anyhow::bail!(message)
}
Err(error) => {
if is_run_canceled(&error) {
let message = "automatic compaction canceled; old primary remains authoritative; no continuation was sent";
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: message.to_string(),
})?;
}
return Err(error);
}
let message = if let Some(rotation_error) =
error.downcast_ref::<crate::sessions::CompactionRotationError>()
{
bounded_auto_compaction_diagnostic(format!(
"automatic compaction failed; no continuation was sent: {rotation_error}"
))
} else {
bounded_auto_compaction_diagnostic(format!(
"automatic compaction failed; old primary remains authoritative; no continuation was sent: {error}"
))
};
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.clone(),
})?;
}
return Err(anyhow::anyhow!(message));
}
};
drop(result);
auto_compaction_count = auto_compaction_count.saturating_add(1);
let observed_steering = steering.as_ref().and_then(AgentSteering::observe_collapsed);
let candidate = observed_steering
.as_ref()
.map(|batch| batch.text.as_str())
.unwrap_or("continue");
let effective_candidate = match projected_candidate.as_ref() {
Some((raw, effective)) if raw == candidate => effective.clone(),
_ => {
let effective = AgentSession::effective_prompt_for_projection(
candidate,
Some(&tools),
options.invocation_mode,
);
projected_candidate = Some((candidate.to_string(), effective.clone()));
effective
}
};
let compacted_tokens = match parent_agent_for_provider.ensure_prompt_context_fits(
&effective_candidate,
session,
Some(&tools),
) {
Ok(tokens) => tokens,
Err(error) => {
let message = bounded_auto_compaction_diagnostic(format!(
"automatic compaction completed, but projected continuation exceeds normal context budget; new checkpoint remains authoritative; no provider continuation was sent: {error}"
));
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.clone(),
})?;
}
return Err(anyhow::anyhow!(message));
}
};
if auto.triggered(compacted_tokens as u64, max_tokens as u64) {
let message = "automatic compaction completed, but projected context remains at or above threshold; new checkpoint remains authoritative; no continuation was sent";
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: message.to_string(),
})?;
}
anyhow::bail!(message)
}
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionCompleted {
current_tokens: compacted_tokens,
max_tokens,
})?;
}
if let Err(error) = cancellation.check() {
let message = "automatic continuation canceled; new checkpoint remains authoritative; no provider continuation was sent";
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: message.to_string(),
})?;
}
return Err(error);
}
match observed_steering {
Some(batch) => {
prompt = batch.text;
effective_prompt = Some(effective_candidate);
prompt_origin = crate::output::UserPromptOrigin::Steering;
}
None => {
prompt = "continue".to_string();
effective_prompt = None;
prompt_origin = crate::output::UserPromptOrigin::AutomaticCompaction;
}
}
}
}
#[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 auto_compaction_eligibility_requires_opt_in_context_session_and_primary_mode() {
for mode in [
InvocationMode::Print,
InvocationMode::Shell,
InvocationMode::MissionControl,
] {
assert!(auto_compaction_eligible(true, true, true, mode));
}
assert!(!auto_compaction_eligible(
false,
true,
true,
InvocationMode::Print
));
assert!(!auto_compaction_eligible(
true,
false,
true,
InvocationMode::Print
));
assert!(!auto_compaction_eligible(
true,
true,
false,
InvocationMode::Print
));
assert!(!auto_compaction_eligible(
true,
true,
true,
InvocationMode::Subagent
));
}
#[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);
}
#[derive(Default)]
struct AutoCompactSink {
events: Vec<OutputEvent>,
cancel_on_completed: Option<Arc<std::sync::atomic::AtomicBool>>,
cancel_on_started: Option<Arc<std::sync::atomic::AtomicBool>>,
make_session_read_only_on_completed: Option<PathBuf>,
make_session_read_only_on_tool_result: Option<PathBuf>,
enqueue_steering_on_completed: Option<(AgentSteering, String)>,
}
impl crate::agent::AgentOutputSink for AutoCompactSink {
fn assistant_delta(&mut self, _text: &str) -> anyhow::Result<()> {
Ok(())
}
fn output_event(&mut self, event: OutputEvent) -> anyhow::Result<()> {
if matches!(event, OutputEvent::CompactionStarted)
&& let Some(cancelled) = self.cancel_on_started.as_ref()
{
cancelled.store(true, std::sync::atomic::Ordering::Release);
}
if matches!(event, OutputEvent::CompactionCompleted { .. })
&& let Some(cancelled) = self.cancel_on_completed.as_ref()
{
cancelled.store(true, std::sync::atomic::Ordering::Release);
}
if matches!(event, OutputEvent::ToolResult { .. }) {
#[cfg(unix)]
if let Some(path) = self.make_session_read_only_on_tool_result.as_ref() {
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o400))?;
}
}
if matches!(event, OutputEvent::CompactionCompleted { .. }) {
#[cfg(unix)]
if let Some(path) = self.make_session_read_only_on_completed.as_ref() {
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o400))?;
}
if let Some((steering, text)) = self.enqueue_steering_on_completed.take() {
steering
.try_enqueue(text)
.map_err(|error| anyhow::anyhow!("enqueue test steering: {error:?}"))?;
}
}
self.events.push(event);
Ok(())
}
fn tool_block(&mut self, _block: &str) -> anyhow::Result<()> {
Ok(())
}
}
fn start_auto_compact_server(
responses: Vec<Result<&'static str, &'static str>>,
allow_missing_last_request: bool,
) -> (
String,
std::thread::JoinHandle<anyhow::Result<Vec<Vec<u8>>>>,
) {
use std::io::{Read, Write};
use std::net::TcpListener;
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let address = listener.local_addr().unwrap();
let handle = std::thread::spawn(move || {
let mut requests = Vec::with_capacity(responses.len());
let response_count = responses.len();
for (index, response_body) in responses.into_iter().enumerate() {
let optional = allow_missing_last_request && index + 1 == response_count;
let max_idle_attempts = if optional { 50 } else { 1_000 };
let mut idle_attempts = 0;
let (mut stream, _) = loop {
match listener.accept() {
Ok(connection) => {
connection.0.set_nonblocking(false)?;
break connection;
}
Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
idle_attempts += 1;
if idle_attempts >= max_idle_attempts {
if optional {
return Ok(requests);
}
anyhow::bail!("timed out waiting for test provider request");
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
Err(error) => return Err(error.into()),
}
};
let mut request = Vec::new();
let mut buffer = [0_u8; 8192];
let (header_end, content_length) = loop {
let count = stream.read(&mut buffer)?;
if count == 0 {
anyhow::bail!("test server received premature request EOF");
}
request.extend_from_slice(&buffer[..count]);
if let Some(header_end) =
request.windows(4).position(|part| part == b"\r\n\r\n")
{
let headers = std::str::from_utf8(&request[..header_end])?;
let content_length = headers
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().ok())
.flatten()
})
.ok_or_else(|| {
anyhow::anyhow!("test request missing Content-Length")
})?;
break (header_end, content_length);
}
};
while request.len() < header_end + 4 + content_length {
let count = stream.read(&mut buffer)?;
if count == 0 {
anyhow::bail!("test server received truncated request body");
}
request.extend_from_slice(&buffer[..count]);
}
requests.push(request);
match response_body {
Ok(text) => {
let payload = if text == "__read_large_tool__" {
concat!(
"data: {\"choices\":[{\"delta\":{\"content\":\"first done\"}}]}\n\n",
"data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"call_large_read\",\"function\":{\"name\":\"read\",\"arguments\":\"{\\\"paths\\\":[\\\"large.txt\\\"]}\"}}]}}]}\n\n",
"data: {\"choices\":[{\"finish_reason\":\"tool_calls\"}]}\n\n",
"data: [DONE]\n\n"
)
.to_string()
} else {
format!(
"data: {{\"choices\":[{{\"delta\":{{\"content\":\"{text}\"}}}}]}}\n\ndata: {{\"choices\":[{{\"finish_reason\":\"stop\"}}]}}\n\ndata: [DONE]\n\n"
)
};
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
payload.len(),
payload
);
stream.write_all(response.as_bytes())?;
}
Err(message) => {
let response = format!(
"HTTP/1.1 400 Bad Request\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
message.len(),
message
);
stream.write_all(response.as_bytes())?;
}
}
stream.flush()?;
}
Ok(requests)
});
(format!("http://{address}/v1"), handle)
}
fn prepare_auto_compact_fixture(
temp: &tempfile::TempDir,
base_url: &str,
threshold_tokens: u64,
) -> (EffectiveConfig, Settings, Session) {
prepare_auto_compact_fixture_with_history(temp, base_url, threshold_tokens, 200_000)
}
fn prepare_auto_compact_fixture_with_history(
temp: &tempfile::TempDir,
base_url: &str,
threshold_tokens: u64,
history_chars: usize,
) -> (EffectiveConfig, Settings, Session) {
let paths = McPaths::from_root(temp.path().join("mc"));
let provider_id = format!(
"local-auto-{}",
base_url.trim_end_matches("/v1").rsplit(':').next().unwrap()
);
let mut config = base_config(paths.clone(), &provider_id, "test-model");
config.auth = Some(ProviderCredential::NoAuth);
config.custom_providers.insert(
provider_id.clone(),
crate::config::make_custom_provider_config("Local", base_url, "").unwrap(),
);
let settings = Settings {
custom_providers: config.custom_providers.clone(),
context: Some(crate::context::ContextBudget {
max_tokens: 500_000,
reserve_tokens: 0,
..crate::context::ContextBudget::default()
}),
compaction: crate::config::CompactionSettings {
auto: crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: None,
threshold_tokens: Some(threshold_tokens),
},
..crate::config::CompactionSettings::default()
},
..Settings::default()
};
crate::config::write_settings(&paths, &settings).unwrap();
let mut catalog_entry = ModelCatalogEntry::new(&provider_id, "test-model");
catalog_entry.supports_reasoning = Some(false);
crate::model_catalog::write_catalog_cache_for_configured_provider(
&paths,
&provider_id,
&[catalog_entry],
)
.unwrap();
let manager = crate::sessions::SessionManager::new(paths.sessions.clone());
let session = manager.create().unwrap();
crate::sessions::record_session_event(
Some(&session),
temp.path(),
crate::sessions::SessionEventKind::UserInput,
serde_json::json!({"text": "x".repeat(history_chars)}),
)
.unwrap();
crate::sessions::record_session_event(
Some(&session),
temp.path(),
crate::sessions::SessionEventKind::AssistantOutput,
serde_json::json!({"text": "old answer"}),
)
.unwrap();
(config, settings, session)
}
fn run_auto_compact_fixture(
temp: &tempfile::TempDir,
config: &EffectiveConfig,
settings: Settings,
session: &Session,
sink: &mut AutoCompactSink,
steering: Option<AgentSteering>,
cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
) -> anyhow::Result<crate::agent::AgentRunOutput> {
run_auto_compact_fixture_in_mode(
temp,
config,
settings,
session,
sink,
steering,
cancellation,
InvocationMode::Print,
)
}
fn run_auto_compact_fixture_with_prompt(
temp: &tempfile::TempDir,
config: &EffectiveConfig,
settings: Settings,
session: &Session,
sink: &mut AutoCompactSink,
prompt: &str,
cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
) -> anyhow::Result<crate::agent::AgentRunOutput> {
let options = ProviderRunOptions {
settings: Some(settings),
prompt,
session: Some(session),
cwd: temp.path(),
output_sink: Some(sink),
selected_primary_agent: None,
cancellation,
session_title_notifier: None,
invocation_mode: InvocationMode::Print,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
};
run_provider_once_streaming(
config,
&[],
&crate::skills::SkillDiscovery::default(),
options,
)
}
#[test]
fn auto_compaction_preflights_oversized_current_input_without_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) =
start_auto_compact_server(vec![Ok("summary"), Ok("answer")], false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 400_000);
let prompt = format!("PREFLIGHT_CURRENT_INPUT_MARKER\n{}", "x".repeat(1_500_000));
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 400_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let mut sink = AutoCompactSink::default();
let output = run_auto_compact_fixture_with_prompt(
&temp, &config, settings, &session, &mut sink, &prompt, None,
)
.unwrap();
assert!(output.text.contains("answer"));
let requests = server.join().unwrap().unwrap();
assert_eq!(requests.len(), 2);
let compaction_request = String::from_utf8_lossy(&requests[0]);
let user_request = String::from_utf8_lossy(&requests[1]);
assert!(compaction_request.contains("Produce the compacted continuation summary now."));
assert!(!compaction_request.contains("PREFLIGHT_CURRENT_INPUT_MARKER"));
assert!(user_request.contains("PREFLIGHT_CURRENT_INPUT_MARKER"));
assert!(
sink.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionTriggered { .. }))
);
assert!(
sink.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionStarted))
);
assert!(
sink.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionCompleted { .. }))
);
assert!(
!sink
.events
.iter()
.any(|event| matches!(event, OutputEvent::AutomaticUserPrompt { .. }))
);
let events = session.read_events().unwrap();
assert_eq!(events[0].event_type, "compaction");
let persisted = events
.iter()
.filter(|event| event.event_type == "user_input")
.filter_map(|event| event.payload["text"].as_str())
.filter(|text| *text == prompt)
.count();
assert_eq!(persisted, 1);
}
#[test]
fn auto_compaction_preflight_leaves_budget_for_post_tool_compaction() {
let temp = tempfile::TempDir::new().unwrap();
write_large_context_fixture(&temp);
let (base_url, server) = start_auto_compact_server(
vec![
Ok("preflight summary"),
Ok("__read_large_tool__"),
Ok("post-tool summary"),
Ok("continued done"),
],
false,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 390_000, 200_000);
let prompt = format!("PREFLIGHT_TOOL_GROWTH_MARKER\n{}", "x".repeat(1_500_000));
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 400_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let mut sink = AutoCompactSink::default();
let output = run_auto_compact_fixture_with_prompt(
&temp, &config, settings, &session, &mut sink, &prompt, None,
)
.unwrap();
let requests = server.join().unwrap().unwrap();
assert_eq!(requests.len(), 4);
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(event, OutputEvent::CompactionStarted))
.count(),
2
);
assert!(output.text.contains("continued done"));
assert_eq!(output.tool_results.len(), 1);
assert!(String::from_utf8_lossy(&requests[2]).contains("IN_TURN_CONTEXT_OVERFLOW_MARKER"));
assert_no_failed_turn_status(&session);
}
#[test]
fn auto_compaction_preflight_cancellation_preserves_identity_and_session() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(Vec::new(), false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let prompt = format!("PREFLIGHT_CANCEL_MARKER\n{}", "x".repeat(1_500_000));
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 400_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false));
let mut sink = AutoCompactSink {
events: Vec::new(),
cancel_on_completed: None,
cancel_on_started: Some(Arc::clone(&cancelled)),
make_session_read_only_on_completed: None,
make_session_read_only_on_tool_result: None,
enqueue_steering_on_completed: None,
};
let error = run_auto_compact_fixture_with_prompt(
&temp,
&config,
settings,
&session,
&mut sink,
&prompt,
Some(cancelled),
)
.unwrap_err();
assert!(is_run_canceled(&error), "{error}");
assert!(server.join().unwrap().unwrap().is_empty());
assert!(
!session
.read_events()
.unwrap()
.iter()
.any(|event| event.event_type == "compaction")
);
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "warning" && message.contains("preflight compaction canceled")
)));
}
#[test]
fn auto_compaction_skips_preflight_when_prompt_without_history_exceeds_budget() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(Vec::new(), false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let prompt = format!("STATIC_OVERFLOW_MARKER\n{}", "x".repeat(200_000));
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 40_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let mut sink = AutoCompactSink::default();
let error = run_auto_compact_fixture_with_prompt(
&temp, &config, settings, &session, &mut sink, &prompt, None,
)
.unwrap_err()
.to_string();
assert!(
error.contains("full session history request is estimated"),
"{error}"
);
assert!(server.join().unwrap().unwrap().is_empty());
assert!(
!sink
.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionTriggered { .. }))
);
assert!(
!session
.read_events()
.unwrap()
.iter()
.any(|event| event.event_type == "compaction")
);
}
fn write_large_context_fixture(temp: &tempfile::TempDir) {
std::fs::write(
temp.path().join("large.txt"),
format!(
"IN_TURN_CONTEXT_OVERFLOW_MARKER\n{}",
"0123456789abcdef".repeat(56_250)
),
)
.unwrap();
}
fn assert_no_failed_turn_status(session: &Session) {
let mut paths = vec![session.path().to_path_buf()];
let history = session
.path()
.parent()
.unwrap()
.join(".history")
.join(session.id());
if let Ok(entries) = std::fs::read_dir(history) {
paths.extend(entries.flatten().map(|entry| entry.path()));
}
for path in paths {
let Ok(contents) = std::fs::read_to_string(&path) else {
continue;
};
for line in contents.lines() {
let event: serde_json::Value = serde_json::from_str(line).unwrap();
assert!(
event["event_type"] != "turn_status" || event["payload"]["status"] != "failed",
"recoverable compaction boundary persisted failed status in {}",
path.display()
);
}
}
}
#[test]
fn auto_compaction_triggers_at_soft_threshold_inside_tool_continuation() {
let temp = tempfile::TempDir::new().unwrap();
write_large_context_fixture(&temp);
let (base_url, server) = start_auto_compact_server(
vec![
Ok("__read_large_tool__"),
Ok("summary"),
Ok("continued done"),
Ok("recursive compaction must not run"),
],
true,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 70_000, 100);
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 300_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let mut sink = AutoCompactSink::default();
let output =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap();
let requests = server.join().unwrap().unwrap();
assert_eq!(requests.len(), 3);
assert!(output.text.contains("first done"));
assert!(output.text.contains("continued done"));
assert_eq!(output.tool_results.len(), 1);
let compaction_request = String::from_utf8_lossy(&requests[1]);
assert!(compaction_request.contains("IN_TURN_CONTEXT_OVERFLOW_MARKER"));
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::CompactionTriggered {
current_tokens,
threshold,
..
} if *current_tokens >= 70_000 && *current_tokens < 300_000 && threshold == "70000 tokens"
)));
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::AutomaticUserPrompt { text } if text == "continue"
)));
assert_no_failed_turn_status(&session);
}
#[test]
fn auto_compaction_recovers_from_hard_in_turn_context_overflow_without_false_failure() {
let temp = tempfile::TempDir::new().unwrap();
write_large_context_fixture(&temp);
let (base_url, server) = start_auto_compact_server(
vec![
Ok("__read_large_tool__"),
Ok("summary"),
Ok("continued done"),
],
false,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 300_000, 100);
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 80_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let mut sink = AutoCompactSink::default();
let output =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap();
assert!(output.text.contains("continued done"));
assert_eq!(server.join().unwrap().unwrap().len(), 3);
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::CompactionTriggered { current_tokens, threshold, .. }
if *current_tokens > 80_000
&& threshold == "continuation hard budget (80000 tokens)"
)));
assert_no_failed_turn_status(&session);
}
#[test]
fn auto_compaction_can_repeat_after_new_tool_growth() {
let temp = tempfile::TempDir::new().unwrap();
write_large_context_fixture(&temp);
let (base_url, server) = start_auto_compact_server(
vec![
Ok("__read_large_tool__"),
Ok("first summary"),
Ok("__read_large_tool__"),
Ok("second summary"),
Ok("final done"),
],
false,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 70_000, 100);
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 300_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let mut sink = AutoCompactSink::default();
let output =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap();
let requests = server.join().unwrap().unwrap();
assert_eq!(requests.len(), 5);
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(event, OutputEvent::CompactionStarted))
.count(),
2
);
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(event, OutputEvent::CompactionCompleted { .. }))
.count(),
2
);
assert_eq!(output.tool_results.len(), 2);
assert!(output.text.contains("final done"));
assert_no_failed_turn_status(&session);
}
#[cfg(unix)]
#[test]
fn in_turn_context_overflow_does_not_compact_degraded_persistence() {
use std::os::unix::fs::PermissionsExt;
let temp = tempfile::TempDir::new().unwrap();
write_large_context_fixture(&temp);
let (base_url, server) = start_auto_compact_server(
vec![Ok("__read_large_tool__"), Ok("compaction must not run")],
true,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 70_000, 100);
let settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 80_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
let mut sink = AutoCompactSink {
make_session_read_only_on_tool_result: Some(session.path().to_path_buf()),
..AutoCompactSink::default()
};
let error =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap_err()
.to_string();
std::fs::set_permissions(session.path(), std::fs::Permissions::from_mode(0o600)).unwrap();
assert!(error.contains("full session history request is estimated"));
assert_eq!(server.join().unwrap().unwrap().len(), 1);
assert!(
!sink
.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionTriggered { .. }))
);
}
#[allow(clippy::too_many_arguments)]
fn run_auto_compact_fixture_in_mode(
temp: &tempfile::TempDir,
config: &EffectiveConfig,
settings: Settings,
session: &Session,
sink: &mut AutoCompactSink,
steering: Option<AgentSteering>,
cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
invocation_mode: InvocationMode,
) -> anyhow::Result<crate::agent::AgentRunOutput> {
let options = ProviderRunOptions {
settings: Some(settings),
prompt: "finish current task",
session: Some(session),
cwd: temp.path(),
output_sink: Some(sink),
selected_primary_agent: None,
cancellation,
session_title_notifier: None,
invocation_mode,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
};
match steering {
Some(steering) => run_provider_once_streaming_with_steering(
config,
&[],
&crate::skills::SkillDiscovery::default(),
options,
steering,
),
None => run_provider_once_streaming(
config,
&[],
&crate::skills::SkillDiscovery::default(),
options,
),
}
}
#[test]
fn auto_compaction_rotates_then_persists_and_sends_labeled_continue() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(
vec![
Ok("first done"),
Ok("summary"),
Ok("continued done"),
Ok("recursive compaction must not run"),
],
true,
);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let mut sink = AutoCompactSink::default();
let output =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap();
let requests = server.join().unwrap().unwrap();
assert_eq!(requests.len(), 3);
assert!(output.text.contains("first done"));
assert!(output.text.contains("continued done"));
assert!(
sink.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionTriggered { .. }))
);
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::AutomaticUserPrompt { text } if text == "continue"
)));
let events = session.read_events().unwrap();
let automatic = events.iter().find(|event| {
event.event_type == "user_input" && event.payload["origin"] == "automatic_compaction"
});
assert_eq!(
automatic.and_then(|event| event.payload["text"].as_str()),
Some("continue")
);
}
#[test]
fn auto_compaction_pending_steering_replaces_continue_and_is_acknowledged_after_persist() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(
vec![Ok("first done"), Ok("summary"), Ok("steered done")],
false,
);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let steering = AgentSteering::new();
steering
.try_enqueue("prioritize deterministic tests #git-status".to_string())
.unwrap();
let expected = steering.observe_collapsed().unwrap().text;
let mut sink = AutoCompactSink::default();
let output = run_auto_compact_fixture(
&temp,
&config,
settings,
&session,
&mut sink,
Some(steering.clone()),
None,
)
.unwrap();
let requests = server.join().unwrap().unwrap();
assert_eq!(requests.len(), 3);
let request = String::from_utf8_lossy(&requests[2]);
assert!(request.contains("prioritize deterministic tests #git-status"));
assert!(
request.contains("<context_injections>")
|| request.contains("<context_injections>")
);
assert!(output.text.contains("steered done"));
assert_eq!(steering.pending_count(), 0);
assert!(
!sink
.events
.iter()
.any(|event| matches!(event, OutputEvent::AutomaticUserPrompt { .. }))
);
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::UserPrompt { text } if text == &expected
)));
let events = session.read_events().unwrap();
let persisted = events
.iter()
.find(|event| event.event_type == "user_input" && event.payload["origin"] == "steering")
.and_then(|event| event.payload["text"].as_str())
.expect("persisted steering");
assert!(persisted.starts_with(expected.as_str()));
assert!(persisted.contains("<context_injections>"));
}
#[test]
fn auto_compaction_runs_through_shell_and_mission_control_shared_orchestration() {
for mode in [InvocationMode::Shell, InvocationMode::MissionControl] {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(
vec![Ok("first done"), Ok("summary"), Ok("continued done")],
false,
);
let (config, settings, session) =
prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let mut sink = AutoCompactSink::default();
let output = run_auto_compact_fixture_in_mode(
&temp, &config, settings, &session, &mut sink, None, None, mode,
)
.unwrap();
assert_eq!(server.join().unwrap().unwrap().len(), 3, "{mode:?}");
assert!(output.text.contains("continued done"), "{mode:?}");
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::AutomaticUserPrompt { text } if text == "continue"
)));
}
}
#[test]
fn auto_compaction_is_excluded_for_subagent_invocation_mode() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(vec![Ok("done")], false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let mut sink = AutoCompactSink::default();
run_auto_compact_fixture_in_mode(
&temp,
&config,
settings,
&session,
&mut sink,
None,
None,
InvocationMode::Subagent,
)
.unwrap();
assert_eq!(server.join().unwrap().unwrap().len(), 1);
assert!(
!sink
.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionTriggered { .. }))
);
}
#[cfg(unix)]
#[test]
fn auto_compaction_continuation_append_failure_keeps_checkpoint_and_steering_queued() {
use std::os::unix::fs::PermissionsExt;
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(
vec![Ok("first done"), Ok("summary"), Ok("must not continue")],
true,
);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let steering = AgentSteering::new();
steering.try_enqueue("first steering".to_string()).unwrap();
let mut sink = AutoCompactSink {
events: Vec::new(),
cancel_on_completed: None,
cancel_on_started: None,
make_session_read_only_on_completed: Some(session.path().to_path_buf()),
make_session_read_only_on_tool_result: None,
enqueue_steering_on_completed: Some((steering.clone(), "later steering".to_string())),
};
let error = run_auto_compact_fixture(
&temp,
&config,
settings,
&session,
&mut sink,
Some(steering.clone()),
None,
)
.unwrap_err()
.to_string();
std::fs::set_permissions(session.path(), std::fs::Permissions::from_mode(0o600)).unwrap();
assert!(
error.contains("automatic continuation persistence failed"),
"{error}"
);
assert!(error.contains("new checkpoint remains authoritative"));
assert!(error.contains("no provider continuation was sent"));
assert_eq!(server.join().unwrap().unwrap().len(), 2);
assert_eq!(steering.pending_count(), 2);
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "error"
&& message.contains("new checkpoint remains authoritative")
)));
let events = session.read_events().unwrap();
assert_eq!(events[0].event_type, "compaction");
assert!(!events.iter().any(|event| {
event.event_type == "user_input" && event.payload["origin"] == "steering"
}));
}
#[test]
fn auto_compaction_provider_failure_emits_error_and_sends_no_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) =
start_auto_compact_server(vec![Ok("first done"), Err("compaction unavailable")], false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let mut sink = AutoCompactSink::default();
let error =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap_err()
.to_string();
assert!(
error.starts_with("automatic compaction failed; old primary remains authoritative; no continuation was sent:"),
"{error}"
);
assert_eq!(server.join().unwrap().unwrap().len(), 2);
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "error"
&& message.starts_with("automatic compaction failed; old primary remains authoritative; no continuation was sent:")
)));
assert!(!session.read_events().unwrap().iter().any(|event| {
event.event_type == "user_input" && event.payload["origin"] == "automatic_compaction"
}));
}
#[test]
fn auto_compaction_resumed_provider_failure_does_not_send_another_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(
vec![
Ok("first done"),
Ok("summary"),
Err("continued request failed"),
Ok("must not recurse"),
],
true,
);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let mut sink = AutoCompactSink::default();
let error =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap_err()
.to_string();
assert!(error.contains("continued request failed"), "{error}");
assert_eq!(server.join().unwrap().unwrap().len(), 3);
assert!(session.read_events().unwrap().iter().any(|event| {
event.event_type == "user_input"
&& event.payload["origin"] == "automatic_compaction"
&& event.payload["text"] == "continue"
}));
}
#[test]
fn auto_compaction_ineffective_reduction_emits_error_and_sends_no_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) =
start_auto_compact_server(vec![Ok("first done"), Ok("summary")], false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 1);
let mut sink = AutoCompactSink::default();
let error =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap_err()
.to_string();
let expected = "automatic compaction completed, but projected context remains at or above threshold; new checkpoint remains authoritative; no continuation was sent";
assert_eq!(error, expected);
assert_eq!(server.join().unwrap().unwrap().len(), 2);
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "error" && message == expected
)));
assert!(!session.read_events().unwrap().iter().any(|event| {
event.event_type == "user_input" && event.payload["origin"] == "automatic_compaction"
}));
}
#[test]
fn auto_compaction_post_checkpoint_normal_context_overflow_sends_no_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let large_summary: &'static str = Box::leak("summary ".repeat(30_000).into_boxed_str());
let (base_url, server) = start_auto_compact_server(
vec![Ok("first done"), Ok(large_summary), Ok("must not continue")],
true,
);
let (config, mut settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 1, 1);
let context = settings.context.as_mut().expect("context settings");
context.max_tokens = 40_000;
context.reserve_tokens = 0;
let mut sink = AutoCompactSink::default();
let error =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap_err()
.to_string();
assert!(
error.starts_with("automatic compaction completed, but projected continuation exceeds normal context budget;"),
"{error}"
);
assert!(error.contains("new checkpoint remains authoritative"));
assert!(error.contains("no provider continuation was sent"));
assert!(
error.contains("exceeding threshold 40000 tokens"),
"{error}"
);
assert!(error.chars().count() <= AUTO_COMPACTION_DIAGNOSTIC_MAX_CHARS);
assert_eq!(server.join().unwrap().unwrap().len(), 2);
assert!(
sink.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionTriggered { .. }))
);
assert!(
sink.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionStarted))
);
assert!(!sink.events.iter().any(|event| matches!(
event,
OutputEvent::CompactionCompleted { .. } | OutputEvent::AutomaticUserPrompt { .. }
)));
assert!(sink.events.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "error" && message == &error
)));
let events = session.read_events().unwrap();
assert_eq!(events[0].event_type, "compaction");
assert!(!events.iter().any(|event| {
event.event_type == "user_input" && event.payload["origin"] == "automatic_compaction"
}));
}
#[test]
fn auto_compaction_below_threshold_delivers_pending_steering_without_compaction() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) =
start_auto_compact_server(vec![Ok("first done"), Ok("steered done")], false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 400_000);
let steering = AgentSteering::new();
steering.try_enqueue("keep going".to_string()).unwrap();
let mut sink = AutoCompactSink::default();
let output = run_auto_compact_fixture(
&temp,
&config,
settings,
&session,
&mut sink,
Some(steering.clone()),
None,
)
.unwrap();
assert_eq!(server.join().unwrap().unwrap().len(), 2);
assert!(output.text.contains("steered done"));
assert_eq!(steering.pending_count(), 0);
assert!(
!sink
.events
.iter()
.any(|event| matches!(event, OutputEvent::CompactionTriggered { .. }))
);
}
#[test]
fn auto_compaction_cancellation_before_summary_preserves_cancellation_identity() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(vec![Ok("first done")], false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false));
let mut sink = AutoCompactSink {
events: Vec::new(),
cancel_on_completed: None,
cancel_on_started: Some(Arc::clone(&cancelled)),
make_session_read_only_on_completed: None,
make_session_read_only_on_tool_result: None,
enqueue_steering_on_completed: None,
};
let error = run_auto_compact_fixture(
&temp,
&config,
settings,
&session,
&mut sink,
None,
Some(cancelled),
)
.unwrap_err();
assert!(is_run_canceled(&error), "{error}");
assert_eq!(server.join().unwrap().unwrap().len(), 1);
assert!(!sink.events.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "error" && message.contains("automatic compaction failed")
)));
}
#[test]
fn auto_compaction_cancellation_after_rotation_sends_no_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) =
start_auto_compact_server(vec![Ok("first done"), Ok("summary")], false);
let (config, settings, session) = prepare_auto_compact_fixture(&temp, &base_url, 30_000);
let cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false));
let mut sink = AutoCompactSink {
events: Vec::new(),
cancel_on_completed: Some(Arc::clone(&cancelled)),
cancel_on_started: None,
make_session_read_only_on_completed: None,
make_session_read_only_on_tool_result: None,
enqueue_steering_on_completed: None,
};
let error = run_auto_compact_fixture(
&temp,
&config,
settings,
&session,
&mut sink,
None,
Some(cancelled),
)
.unwrap_err()
.to_string();
assert!(error.contains("cancel"), "{error}");
assert_eq!(server.join().unwrap().unwrap().len(), 2);
assert!(!session.read_events().unwrap().iter().any(|event| {
event.event_type == "user_input" && event.payload["origin"] == "automatic_compaction"
}));
}
}