#[cfg(test)]
use crate::cancellation::is_run_canceled;
use crate::{
agent::AgentSession,
agent::steering::AgentSteering,
cancellation::AgentCancellation,
config::{CustomProviderConfig, EffectiveConfig, Settings},
context::ContextBudget,
output::{HookContextMetadata, InvocationMode, OutputEvent, ToolDispatchContext},
primary_agents::{PrimaryAgentProfile, render_primary_agent_prompt_append},
providers::{OPENAI_CODEX_PROVIDER, ProviderSelection, ToolCall},
sessions::Session,
};
use anyhow::Result;
use serde_json::json;
use std::{
collections::HashSet,
path::Path,
sync::{Arc, Mutex, atomic::AtomicU64},
};
mod auto_compaction;
mod preparation;
pub(crate) struct ProviderRunOptions<'a, 'sink> {
pub(crate) settings: Option<Settings>,
pub(crate) prompt: &'a str,
pub(crate) session: Option<&'a Session>,
pub(crate) cwd: &'a Path,
pub(crate) output_sink: Option<&'sink mut dyn crate::agent::AgentOutputSink>,
pub(crate) selected_primary_agent: Option<PrimaryAgentProfile>,
pub(crate) cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
pub(crate) session_title_notifier: Option<crate::session_titles::SessionTitleNotifier>,
pub(crate) herdr_reporter: Option<crate::herdr::HerdrReporter>,
pub(crate) invocation_mode: InvocationMode,
pub(crate) disabled_tools: Option<Arc<Mutex<HashSet<String>>>>,
pub(crate) disabled_subagent_profiles: Option<Arc<Mutex<HashSet<String>>>>,
pub(crate) subagent_profile_discovery:
Option<crate::subagents::profiles::SubagentProfileDiscovery>,
pub(crate) mcp: Option<Arc<Mutex<crate::mcp::manager::McpManager>>>,
}
fn reborrow_output_sink<'a>(
sink: &'a mut Option<&mut dyn crate::agent::AgentOutputSink>,
) -> Option<&'a mut dyn crate::agent::AgentOutputSink> {
match sink {
Some(sink) => Some(&mut **sink),
None => None,
}
}
fn emit_compaction_fast_observation(
output_sink: &mut Option<&mut dyn crate::agent::AgentOutputSink>,
result: &crate::compaction::CompactionResult,
request_sequence: &AtomicU64,
) -> anyhow::Result<()> {
let Some(requested_service_tier) = result.requested_service_tier.as_deref() else {
return Ok(());
};
let Some(sink) = output_sink.as_deref_mut() else {
return Ok(());
};
sink.output_event(OutputEvent::FastObservation {
provider_id: result.provider.clone(),
model: result.model.clone(),
requested_service_tier: requested_service_tier.to_string(),
outcome: result.fast_outcome.clone(),
request_sequence: request_sequence
.load(std::sync::atomic::Ordering::Relaxed)
.saturating_sub(1),
run_order: None,
})
}
pub(crate) const AUTO_COMPACTION_DIAGNOSTIC_MAX_CHARS: usize = 512;
pub(crate) fn auto_compaction_policy(
auto: &crate::config::AutoCompactionSettings,
model_max_tokens: usize,
) -> Option<(usize, String)> {
if !auto.is_enabled() {
return None;
}
let percent_threshold = auto.threshold_percent.map(|percent| {
let numerator = (model_max_tokens as u128).saturating_mul(u128::from(percent));
usize::try_from(numerator.div_ceil(100)).unwrap_or(usize::MAX)
});
let token_threshold = auto
.threshold_tokens
.map(|tokens| usize::try_from(tokens).unwrap_or(usize::MAX));
let threshold = match (percent_threshold, token_threshold) {
(Some(percent), Some(tokens)) => percent.min(tokens),
(Some(threshold), None) | (None, Some(threshold)) => threshold,
(None, None) => return None,
};
Some((threshold, auto.threshold_display()?))
}
pub(crate) fn bounded_auto_compaction_diagnostic(message: impl std::fmt::Display) -> String {
crate::output::redact_sensitive_text(&message.to_string())
.chars()
.take(AUTO_COMPACTION_DIAGNOSTIC_MAX_CHARS)
.collect()
}
fn auto_compaction_eligible(
auto_enabled: bool,
context_enabled: bool,
has_session: bool,
invocation_mode: InvocationMode,
) -> bool {
auto_enabled
&& context_enabled
&& has_session
&& matches!(
invocation_mode,
InvocationMode::Print | InvocationMode::MissionControl
)
}
pub(crate) fn append_primary_agent_to_main_prompt(
agent: AgentSession,
selected: Option<&PrimaryAgentProfile>,
) -> AgentSession {
match selected {
Some(profile) => {
agent.with_appended_system_prompt(&render_primary_agent_prompt_append(profile))
}
None => agent,
}
}
fn apply_cached_context_window_if_eligible(
context_budget: &mut ContextBudget,
config: &EffectiveConfig,
selection: &ProviderSelection,
selected_custom_provider: Option<&CustomProviderConfig>,
) {
if selection.provider != OPENAI_CODEX_PROVIDER
&& selection.provider != crate::providers::ANTHROPIC_PROVIDER
&& 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);
}
pub(crate) fn context_budget_for_selection(
config: &EffectiveConfig,
settings: &Settings,
selection: &ProviderSelection,
) -> ContextBudget {
let mut context_budget = settings.context.clone().unwrap_or_default();
let selected_custom_provider = config.custom_providers.get(&selection.provider);
apply_cached_context_window_if_eligible(
&mut context_budget,
config,
selection,
selected_custom_provider,
);
context_budget
}
pub(crate) fn run_bash_mode_once(
config: &EffectiveConfig,
skills: &crate::skills::SkillDiscovery,
mut options: ProviderRunOptions<'_, '_>,
) -> Result<crate::agent::AgentRunOutput> {
let invocation_mode = options.invocation_mode;
let reporter =
crate::agent::filtered_herdr_reporter(invocation_mode, options.herdr_reporter.clone());
let mut turn = crate::herdr::HerdrTurnReporter::pending(reporter.clone());
options.herdr_reporter = reporter;
let result = run_bash_mode_once_inner(config, skills, options);
turn.finish_bash_result(&result);
result
}
fn run_bash_mode_once_inner(
config: &EffectiveConfig,
skills: &crate::skills::SkillDiscovery,
mut options: ProviderRunOptions<'_, '_>,
) -> Result<crate::agent::AgentRunOutput> {
let Some(bash_prompt) = crate::bash_mode::parse_bash_mode_prompt(options.prompt) else {
anyhow::bail!("bash-mode prompt must start with '!'");
};
let cancellation = options
.cancellation
.take()
.map(AgentCancellation::new)
.unwrap_or_default();
cancellation.check()?;
let settings = options
.settings
.take()
.map(Ok)
.unwrap_or_else(|| crate::config::read_settings(&config.paths))?;
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::SessionHeader {
session_id: options.session.map(|session| session.id().to_string()),
model: config
.model
.clone()
.unwrap_or_else(|| crate::providers::DEFAULT_CODEX_MODEL.to_string()),
cwd: options.cwd.to_path_buf(),
})?;
}
let mut session_persistence =
crate::agent::session_persistence::SessionPersistence::new(None, options.cwd);
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::BashCommand {
command: bash_prompt.command.to_string(),
})?;
}
let hook_settings = settings.hooks.clone();
let tool_settings = settings.tools.clone();
let disabled_tools = options.disabled_tools.clone().unwrap_or_else(|| {
Arc::new(Mutex::new(
crate::config::disabled_tool_names_from_settings(&settings)
.into_iter()
.collect(),
))
});
let tools = crate::tools::ToolRuntime::new_with_full_settings_and_mcp_with_disabled_tools(
options.cwd,
config.paths.clone(),
settings.clone(),
options.mcp.clone(),
disabled_tools,
)?
.with_skills(skills);
let hooks = crate::hooks::HookRuntime::new_with_tool_settings(
options.cwd,
hook_settings,
&tool_settings,
)?;
let call = ToolCall {
id: "bash_mode_1".to_string(),
name: crate::tools::contract::tool_name::BASH.to_string(),
arguments: json!({"command": bash_prompt.command}),
};
let activity =
crate::agent::tool_lifecycle::emit_tool_started(&mut options.output_sink, 0, 0, &call)?;
let dispatch_context = ToolDispatchContext::new_with_hook_context_and_cancellation(
Some(activity.activity_id.clone()),
activity.activity_sender.clone(),
HookContextMetadata {
session_id: options.session.map(|session| session.id().to_string()),
session_path: options.session.map(|session| session.path().to_path_buf()),
provider_id: None,
model_id: None,
agent_id: None,
invocation_mode: options.invocation_mode,
turn_id: Some("turn-0".to_string()),
message_id: Some(format!("tool-{}", activity.activity_id.as_str())),
subagent: false,
},
cancellation.clone(),
);
if let Some(reporter) = options.herdr_reporter.as_ref() {
reporter.report_bash();
}
let outcome = crate::agent::tool_lifecycle::execute_tool_call(
Some(&tools),
Some(&hooks),
&mut session_persistence,
&mut options.output_sink,
None,
call.clone(),
dispatch_context,
)?;
let mut output = crate::agent::AgentRunOutput::default();
match outcome {
crate::agent::tool_lifecycle::ToolExecutionOutcome::Dispatched {
tool_result,
after_hook_failure,
..
} => {
crate::agent::tool_lifecycle::emit_tool_finished(
&mut options.output_sink,
activity.activity_id,
&call,
&tool_result,
)?;
output.tool_results.push((*tool_result).clone());
if let Some(diagnostic) = after_hook_failure {
return Err(crate::hooks::HookPolicyError::new(diagnostic).into());
}
}
crate::agent::tool_lifecycle::ToolExecutionOutcome::BeforeHookFailed {
diagnostic,
tool_result,
} => {
crate::agent::tool_lifecycle::emit_tool_finished(
&mut options.output_sink,
activity.activity_id,
&call,
&tool_result,
)?;
return Err(crate::hooks::HookPolicyError::new(diagnostic).into());
}
}
cancellation.check()?;
Ok(output)
}
pub(crate) fn run_provider_once_streaming(
config: &EffectiveConfig,
instructions: &[crate::instructions::InstructionFile],
skills: &crate::skills::SkillDiscovery,
mut options: ProviderRunOptions<'_, '_>,
) -> Result<crate::agent::AgentRunOutput> {
let invocation_mode = options.invocation_mode;
let reporter = options.herdr_reporter.clone();
crate::agent::run_with_herdr_turn(invocation_mode, reporter, |herdr_reporter| {
options.herdr_reporter = herdr_reporter;
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,
mut options: ProviderRunOptions<'_, '_>,
steering: AgentSteering,
) -> Result<crate::agent::AgentRunOutput> {
let invocation_mode = options.invocation_mode;
let reporter = options.herdr_reporter.clone();
crate::agent::run_with_herdr_turn(invocation_mode, reporter, |herdr_reporter| {
options.herdr_reporter = herdr_reporter;
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 prepared = preparation::prepare(config, instructions, skills, &mut options, cancellation)?;
auto_compaction::run(prepared, instructions, &mut options, steering)
}
#[cfg(test)]
mod primary_agent_tests {
use super::*;
use crate::providers::{
OpenAiCodexProvider, provider_from_selection_with_settings,
provider_from_selection_with_settings_with_codex_exchange,
};
use crate::{
config::{McPaths, ProviderCredential},
context::ContextBudgetOverride,
instructions::InstructionFile,
model_catalog::ModelCatalogEntry,
skills::SkillDiscovery,
};
use std::{collections::BTreeMap, path::PathBuf};
#[derive(Debug, Clone)]
struct NoopTransport;
fn profile() -> PrimaryAgentProfile {
PrimaryAgentProfile {
id: "tars".to_string(),
name: "TARS".to_string(),
description: "Tactical".to_string(),
path: PathBuf::new(),
prompt: "PRIMARY BODY".to_string(),
}
}
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,
auth: Some(ProviderCredential::OAuth {
access: "test-token".to_string(),
account_id: Some("acct".to_string()),
}),
paths,
}
}
#[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_uses_supplied_fast_mode_snapshot() {
let temp = tempfile::TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut entry = ModelCatalogEntry::new_codex("gpt-5.5");
entry.service_tiers = Some(vec![crate::model_catalog::ModelCatalogServiceTier {
id: "priority".to_string(),
name: "fast".to_string(),
description: None,
}]);
crate::model_catalog::write_catalog_cache(&paths, OPENAI_CODEX_PROVIDER, &[entry]).unwrap();
crate::config::set_fast_mode(&paths, false).unwrap();
crate::config::write_auth(
&paths,
&crate::config::Auth {
providers: BTreeMap::from([(
OPENAI_CODEX_PROVIDER.to_string(),
crate::config::AuthProviderRecord::OAuth {
access: "test-token".to_string(),
refresh: Some("test-refresh".to_string()),
expires: Some(4102444800),
account_id: Some("acct".to_string()),
},
)]),
..crate::config::Auth::default()
},
)
.unwrap();
let config = base_config(paths.clone(), OPENAI_CODEX_PROVIDER, "gpt-5.5");
let selection = ProviderSelection {
provider: OPENAI_CODEX_PROVIDER.to_string(),
model: "gpt-5.5".to_string(),
};
let mut settings = Settings::default();
settings.fast.enabled = true;
let provider =
provider_from_selection_with_settings(&config, &selection, temp.path(), &settings)
.unwrap();
assert_eq!(provider.service_tier_for_test(), Some("priority"));
crate::config::set_fast_mode(&paths, true).unwrap();
settings.fast.enabled = false;
let provider =
provider_from_selection_with_settings(&config, &selection, temp.path(), &settings)
.unwrap();
assert_eq!(provider.service_tier_for_test(), None);
let unsupported = ModelCatalogEntry::new_codex("gpt-5.3");
crate::model_catalog::write_catalog_cache(&paths, OPENAI_CODEX_PROVIDER, &[unsupported])
.unwrap();
let unsupported_selection = ProviderSelection {
provider: OPENAI_CODEX_PROVIDER.to_string(),
model: "gpt-5.3".to_string(),
};
settings.fast.enabled = true;
let provider = provider_from_selection_with_settings(
&config,
&unsupported_selection,
temp.path(),
&settings,
)
.unwrap();
assert_eq!(provider.service_tier_for_test(), None);
}
#[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,
auth: Some(ProviderCredential::ApiKey {
key: "anthropic-key".to_string(),
}),
paths: McPaths::from_root(temp.path().join("mc")),
};
let mut auth = crate::config::Auth::default();
auth.providers.insert(
crate::providers::ANTHROPIC_PROVIDER.to_string(),
crate::config::AuthProviderRecord::ApiKey {
key: "anthropic-key".to_string(),
},
);
crate::config::write_auth(&config.paths, &auth).unwrap();
let selection = ProviderSelection::from_config(&config).unwrap();
let _provider = provider_from_selection_with_settings(
&config,
&selection,
temp.path(),
&Settings::default(),
)
.unwrap();
}
#[test]
fn codex_runtime_refresh_completes_before_provider_construction() {
let temp = tempfile::TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut auth = crate::config::Auth::default();
auth.providers.insert(
OPENAI_CODEX_PROVIDER.to_string(),
crate::config::AuthProviderRecord::OAuth {
access: "current-access".into(),
refresh: Some("stored-refresh".into()),
expires: Some(chrono::Utc::now().timestamp() + 3600),
account_id: None,
},
);
crate::config::write_auth(&paths, &auth).unwrap();
let mut config = base_config(paths, OPENAI_CODEX_PROVIDER, "gpt-5.5");
config.auth = None;
assert!(!config.auth_state().is_ready());
let selection = ProviderSelection {
provider: OPENAI_CODEX_PROVIDER.to_string(),
model: "gpt-5.5".to_string(),
};
let calls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let calls_for_exchange = std::sync::Arc::clone(&calls);
let provider = provider_from_selection_with_settings_with_codex_exchange(
&config,
&selection,
temp.path(),
&Settings::default(),
|refresh| {
calls_for_exchange.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
assert_eq!(refresh, "stored-refresh");
Ok(crate::config::NormalizedToken {
access: "refreshed-access".into(),
refresh: None,
expires: Some(4102444800),
account_id: "refreshed-account".into(),
})
},
)
.unwrap();
assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 1);
assert!(provider.service_tier_for_test().is_none());
let stored = crate::config::read_auth(&config.paths).unwrap();
let Some(crate::config::AuthProviderRecord::OAuth {
access, account_id, ..
}) = stored.providers.get(OPENAI_CODEX_PROVIDER)
else {
panic!("missing codex oauth record")
};
assert_eq!(access, "refreshed-access");
assert_eq!(account_id.as_deref(), Some("refreshed-account"));
let request = crate::providers::ProviderRequest::new(
"gpt-5.5",
vec![crate::providers::ChatMessage::user("hello")],
);
let http =
OpenAiCodexProvider::new("gpt-5.5", access.clone(), account_id.clone(), NoopTransport)
.build_http_request(&request)
.unwrap();
assert_eq!(
http.headers.get("chatgpt-account-id").map(String::as_str),
Some("refreshed-account")
);
let pending_temp = tempfile::TempDir::new().unwrap();
let pending_config = base_config(
McPaths::from_root(pending_temp.path().join("mc")),
OPENAI_CODEX_PROVIDER,
"gpt-5.5",
);
let mut pending_config = pending_config;
pending_config.auth = Some(ProviderCredential::OAuth {
access: "unrefreshed-access".into(),
account_id: None,
});
let error = match provider_from_selection_with_settings(
&pending_config,
&selection,
pending_temp.path(),
&Settings::default(),
) {
Ok(_) => panic!("provider construction accepted an unrefreshed OAuth credential"),
Err(error) => error.to_string(),
};
assert!(
error.contains("missing OAuth auth for provider 'openai-codex'"),
"{error}"
);
assert!(
OpenAiCodexProvider::new("gpt-5.5", "unrefreshed-access", None, NoopTransport)
.build_http_request(&request)
.is_err()
);
}
#[test]
fn codex_runtime_refresh_does_not_trust_cached_config_auth() {
let temp = tempfile::TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut auth = crate::config::Auth::default();
auth.providers.insert(
OPENAI_CODEX_PROVIDER.to_string(),
crate::config::AuthProviderRecord::OAuth {
access: "expired-stored-access".into(),
refresh: Some("stored-refresh".into()),
expires: Some(1),
account_id: Some("stored-account".into()),
},
);
crate::config::write_auth(&paths, &auth).unwrap();
let mut config = base_config(paths.clone(), OPENAI_CODEX_PROVIDER, "gpt-5.5");
config.auth = Some(ProviderCredential::OAuth {
access: "cached-access".into(),
account_id: Some("cached-account".into()),
});
assert!(config.auth_state().is_ready());
let selection = ProviderSelection {
provider: OPENAI_CODEX_PROVIDER.to_string(),
model: "gpt-5.5".to_string(),
};
let calls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let calls_for_exchange = std::sync::Arc::clone(&calls);
let provider = provider_from_selection_with_settings_with_codex_exchange(
&config,
&selection,
temp.path(),
&Settings::default(),
|refresh| {
calls_for_exchange.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
assert_eq!(refresh, "stored-refresh");
Ok(crate::config::NormalizedToken {
access: "refreshed-access".into(),
refresh: None,
expires: Some(4102444800),
account_id: "refreshed-account".into(),
})
},
)
.unwrap();
assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 1);
assert!(provider.service_tier_for_test().is_none());
let stored = crate::config::read_auth(&paths).unwrap();
let Some(crate::config::AuthProviderRecord::OAuth {
access, account_id, ..
}) = stored.providers.get(OPENAI_CODEX_PROVIDER)
else {
panic!("missing codex oauth record")
};
assert_eq!(access, "refreshed-access");
assert_eq!(account_id.as_deref(), Some("refreshed-account"));
}
#[test]
fn preparation_reuses_refreshed_credential_when_expiry_is_unknown() {
let temp = tempfile::TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut auth = crate::config::Auth::default();
auth.providers.insert(
OPENAI_CODEX_PROVIDER.to_string(),
crate::config::AuthProviderRecord::OAuth {
access: "expired-access".into(),
refresh: Some("stored-refresh".into()),
expires: Some(1),
account_id: Some("stored-account".into()),
},
);
crate::config::write_auth(&paths, &auth).unwrap();
let mut config = base_config(paths.clone(), OPENAI_CODEX_PROVIDER, "gpt-5.5");
config.auth = None;
let calls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let calls_for_exchange = std::sync::Arc::clone(&calls);
let mut options = ProviderRunOptions {
settings: Some(Settings::default()),
prompt: "hello",
session: None,
cwd: temp.path(),
output_sink: None,
selected_primary_agent: None,
cancellation: None,
session_title_notifier: None,
herdr_reporter: None,
invocation_mode: InvocationMode::Print,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
};
let prepared = preparation::prepare_with_codex_exchange(
&config,
&[],
&SkillDiscovery::default(),
&mut options,
crate::cancellation::AgentCancellation::default(),
|refresh| {
calls_for_exchange.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
assert_eq!(refresh, "stored-refresh");
Ok(crate::config::NormalizedToken {
access: "refreshed-access".into(),
refresh: None,
expires: None,
account_id: "refreshed-account".into(),
})
},
)
.unwrap();
assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 1);
assert_eq!(
prepared.active_config.auth,
Some(ProviderCredential::OAuth {
access: "refreshed-access".into(),
account_id: Some("refreshed-account".into()),
})
);
assert!(prepared.provider.service_tier_for_test().is_none());
let stored = crate::config::read_auth(&paths).unwrap();
let Some(crate::config::AuthProviderRecord::OAuth {
access,
refresh,
expires,
account_id,
}) = stored.providers.get(OPENAI_CODEX_PROVIDER)
else {
panic!("missing codex oauth record")
};
assert_eq!(access, "refreshed-access");
assert_eq!(refresh.as_deref(), Some("stored-refresh"));
assert_eq!(*expires, None);
assert_eq!(account_id.as_deref(), Some("refreshed-account"));
}
#[test]
fn prepared_agent_leaves_existing_disk_context_cache_untouched() {
struct ReplyProvider;
impl crate::providers::Provider for ReplyProvider {
fn stream_cancellable(
&self,
_request: crate::providers::ProviderRequest,
_cancellation: &crate::cancellation::AgentCancellation,
on_event: &mut dyn FnMut(crate::providers::ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
on_event(crate::providers::ProviderEvent::TextDelta("reply".into()))?;
on_event(crate::providers::ProviderEvent::Done)
}
}
let temp = tempfile::TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut auth = crate::config::Auth::default();
auth.providers.insert(
OPENAI_CODEX_PROVIDER.to_string(),
crate::config::AuthProviderRecord::OAuth {
access: "test-access".into(),
refresh: None,
expires: Some(4102444800),
account_id: Some("test-account".into()),
},
);
crate::config::write_auth(&paths, &auth).unwrap();
std::fs::create_dir_all(&paths.cache).unwrap();
let legacy_snapshot = paths.cache.join(format!("{}.json", "a".repeat(64)));
let original = b"legacy diagnostic snapshot: leave unchanged";
std::fs::write(&legacy_snapshot, original).unwrap();
let config = base_config(paths.clone(), OPENAI_CODEX_PROVIDER, "gpt-5.5");
let session = crate::sessions::SessionManager::new(paths.sessions.clone())
.create()
.unwrap();
let mut options = ProviderRunOptions {
settings: Some(Settings::default()),
prompt: "first",
session: Some(&session),
cwd: temp.path(),
output_sink: None,
selected_primary_agent: None,
cancellation: None,
session_title_notifier: None,
herdr_reporter: None,
invocation_mode: InvocationMode::Print,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
};
let prepared = preparation::prepare(
&config,
&[],
&SkillDiscovery::default(),
&mut options,
crate::cancellation::AgentCancellation::default(),
)
.unwrap();
for prompt in ["first", "second with more history"] {
prepared
.parent_agent_for_provider
.run_print_with_tools(&ReplyProvider, prompt, None, Some(&session), temp.path())
.unwrap();
}
let events = session.read_events().unwrap();
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "user_input")
.count(),
2
);
assert!(
!events
.iter()
.any(|event| event.event_type == "context_cache")
);
assert_eq!(std::fs::read(&legacy_snapshot).unwrap(), original);
let cache_files = std::fs::read_dir(&paths.cache)
.unwrap()
.map(|entry| entry.unwrap().path())
.collect::<Vec<_>>();
assert_eq!(cache_files, vec![legacy_snapshot]);
}
#[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 before = base_agent().system_prompt().to_string();
let after = append_primary_agent_to_main_prompt(base_agent(), Some(&profile()));
let expected = format!(
"{before}\n\n<Your-Role Persona=TARS>\nYou must always follow the instructions of your assigned persona:\nPRIMARY BODY\n</Your-Role>"
);
assert_eq!(after.system_prompt(), expected);
}
#[test]
fn auto_compaction_policy_uses_ceil_minimum_and_safe_large_values() {
let percent_first = crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: Some(50),
threshold_tokens: Some(900),
max_compactions_per_run: None,
};
assert_eq!(
auto_compaction_policy(&percent_first, 1_001),
Some((501, "50% or 900 tokens".to_string()))
);
let token_first = crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: Some(80),
threshold_tokens: Some(300),
max_compactions_per_run: None,
};
assert_eq!(
auto_compaction_policy(&token_first, 1_000),
Some((300, "80% or 300 tokens".to_string()))
);
let tie = crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: Some(50),
threshold_tokens: Some(500),
max_compactions_per_run: None,
};
assert_eq!(
auto_compaction_policy(&tie, 1_000),
Some((500, "50% or 500 tokens".to_string()))
);
let percent_only = crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: Some(1),
threshold_tokens: None,
max_compactions_per_run: None,
};
assert_eq!(
auto_compaction_policy(&percent_only, 101),
Some((2, "1%".to_string()))
);
let token_only = crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: None,
threshold_tokens: Some(300),
max_compactions_per_run: None,
};
assert_eq!(
auto_compaction_policy(&token_only, 1_000),
Some((300, "300 tokens".to_string()))
);
let disabled = crate::config::AutoCompactionSettings {
enabled: false,
threshold_percent: Some(50),
threshold_tokens: Some(300),
max_compactions_per_run: None,
};
assert_eq!(auto_compaction_policy(&disabled, 1_000), None);
let enabled_without_thresholds = crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: None,
threshold_tokens: None,
max_compactions_per_run: None,
};
assert_eq!(
auto_compaction_policy(&enabled_without_thresholds, 1_000),
None
);
let large = crate::config::AutoCompactionSettings {
enabled: true,
threshold_percent: Some(100),
threshold_tokens: Some(u64::MAX),
max_compactions_per_run: None,
};
assert_eq!(
auto_compaction_policy(&large, 1_000_001),
Some((1_000_001, format!("100% or {} tokens", u64::MAX)))
);
}
#[test]
fn auto_compaction_eligibility_requires_opt_in_context_session_and_primary_mode() {
for mode in [InvocationMode::Print, 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>>,
#[cfg(unix)]
make_session_read_only_on_completed: Option<PathBuf>,
#[cfg(unix)]
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\"}}],\"usage\":{{\"prompt_tokens\":17,\"completion_tokens\":3,\"total_tokens\":29}}}}\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),
max_compactions_per_run: None,
},
..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,
herdr_reporter: 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_forwards_title_job_before_primary_provider_error() {
let temp = tempfile::TempDir::new().unwrap();
let (base_url, server) = start_auto_compact_server(
vec![Err("primary provider failed"), Err("title provider failed")],
true,
);
let (config, mut settings, _) = prepare_auto_compact_fixture(&temp, &base_url, 400_000);
let provider = config.provider_id().to_string();
settings.session_titles = crate::config::SessionTitleSettings {
enabled: true,
provider: Some(provider),
model: Some("test-model".to_string()),
};
crate::config::write_settings(&config.paths, &settings).unwrap();
let session = crate::sessions::SessionManager::new(config.paths.sessions.clone())
.create()
.unwrap();
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("status 400"), "{error}");
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
let events = loop {
let events = session.read_events().unwrap();
if events.iter().any(|event| {
event.event_type == "diagnostic"
&& event.payload["category"] == "session_title_generation"
&& event.payload["reason"] == "failed"
}) {
break events;
}
assert!(
std::time::Instant::now() < deadline,
"title generation did not start: {events:?}"
);
std::thread::sleep(std::time::Duration::from_millis(1));
};
assert!(events.iter().any(|event| event.event_type == "user_input"));
assert_eq!(
events
.iter()
.filter(|event| {
event.event_type == "diagnostic"
&& event.payload["category"] == "session_title_generation"
})
.count(),
1
);
let requests = server.join().unwrap().unwrap();
assert_eq!(requests.len(), 2);
assert_eq!(
requests
.iter()
.filter(|request| {
String::from_utf8_lossy(request)
.contains("Infer a title from this initial user request")
})
.count(),
1
);
}
fn run_auto_compact_fixture_with_herdr(
temp: &tempfile::TempDir,
config: &EffectiveConfig,
settings: Settings,
session: &Session,
sink: &mut AutoCompactSink,
reporter: crate::herdr::HerdrReporter,
) -> anyhow::Result<crate::agent::AgentRunOutput> {
run_provider_once_streaming(
config,
&[],
&crate::skills::SkillDiscovery::default(),
ProviderRunOptions {
settings: Some(settings),
prompt: "finish current task",
session: Some(session),
cwd: temp.path(),
output_sink: Some(sink),
selected_primary_agent: None,
cancellation: None,
session_title_notifier: None,
herdr_reporter: Some(reporter),
invocation_mode: InvocationMode::Print,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
},
)
}
#[test]
fn herdr_runner_auto_compaction_continuation_reports_one_turn() {
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 (herdr_owner, herdr_lines) = crate::herdr::HerdrOwner::new_for_test();
let mut sink = AutoCompactSink::default();
let output = run_auto_compact_fixture_with_herdr(
&temp,
&config,
settings,
&session,
&mut sink,
herdr_owner.reporter(),
)
.unwrap();
assert!(output.text.contains("first done"));
assert!(output.text.contains("continued done"));
assert_eq!(server.join().unwrap().unwrap().len(), 3);
let statuses = herdr_lines
.lock()
.unwrap()
.iter()
.map(|line| {
let value: serde_json::Value = serde_json::from_str(line).unwrap();
(
value["params"]["state"].as_str().unwrap().to_string(),
value["params"]["message"].as_str().unwrap().to_string(),
)
})
.collect::<Vec<_>>();
assert_eq!(
statuses,
vec![
("working".to_string(), "thinking".to_string()),
("idle".to_string(), "done".to_string()),
]
);
}
#[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_eq!(output.total_tokens, Some(87));
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(
event,
OutputEvent::UsageSnapshot {
final_usage: true,
..
}
))
.count(),
3
);
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 {
cancel_on_started: Some(Arc::clone(&cancelled)),
..AutoCompactSink::default()
};
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_with_both_triggers() {
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 mut settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 300_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
settings.compaction.auto.threshold_percent = Some(25);
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 == "25% or 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);
}
#[test]
fn auto_compaction_zero_allows_more_than_five_compactions_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("__read_large_tool__"),
Ok("third summary"),
Ok("__read_large_tool__"),
Ok("fourth summary"),
Ok("__read_large_tool__"),
Ok("fifth summary"),
Ok("__read_large_tool__"),
Ok("sixth summary"),
Ok("final done"),
],
false,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 70_000, 100);
let mut settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 300_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
settings.compaction.auto.max_compactions_per_run = Some(0);
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(), 13);
assert!(String::from_utf8_lossy(&requests[11]).contains("IN_TURN_CONTEXT_OVERFLOW_MARKER"));
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(event, OutputEvent::CompactionStarted))
.count(),
6
);
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(event, OutputEvent::CompactionCompleted { .. }))
.count(),
6
);
assert_eq!(output.tool_results.len(), 6);
assert!(output.text.contains("final done"));
assert_no_failed_turn_status(&session);
}
#[test]
fn auto_compaction_respects_configured_limit_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("final done"),
],
false,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 70_000, 100);
let mut settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 300_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
settings.compaction.auto.max_compactions_per_run = Some(1);
let mut sink = AutoCompactSink::default();
let output =
run_auto_compact_fixture(&temp, &config, settings, &session, &mut sink, None, None)
.unwrap();
assert_eq!(server.join().unwrap().unwrap().len(), 4);
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(event, OutputEvent::CompactionStarted))
.count(),
1
);
assert_eq!(output.tool_results.len(), 2);
assert!(output.text.contains("final done"));
assert_no_failed_turn_status(&session);
}
#[test]
fn auto_compaction_preflight_consumes_configured_limit() {
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 compaction must not run"),
],
true,
);
let (config, settings, session) =
prepare_auto_compact_fixture_with_history(&temp, &base_url, 390_000, 200_000);
let prompt = format!("PREFLIGHT_LIMIT_MARKER\n{}", "x".repeat(1_500_000));
let mut settings = Settings {
context: Some(crate::context::ContextBudget {
max_tokens: 400_000,
reserve_tokens: 0,
..settings.context.clone().unwrap()
}),
..settings
};
settings.compaction.auto.max_compactions_per_run = Some(1);
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_eq!(server.join().unwrap().unwrap().len(), 2);
assert_eq!(
sink.events
.iter()
.filter(|event| matches!(event, OutputEvent::CompactionStarted))
.count(),
1
);
}
#[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,
herdr_reporter: 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, mut settings, session) =
prepare_auto_compact_fixture(&temp, &base_url, 30_000);
settings.compaction.auto.threshold_percent = Some(90);
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 {
current_tokens,
threshold,
..
} if *current_tokens >= 30_000 && threshold == "90% or 30000 tokens"
)));
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_mission_control_shared_orchestration() {
let mode = 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 { .. }))
);
assert!(
!sink
.events
.iter()
.any(|event| matches!(event, 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 {
cancel_on_started: Some(Arc::clone(&cancelled)),
..AutoCompactSink::default()
};
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 {
cancel_on_completed: Some(Arc::clone(&cancelled)),
..AutoCompactSink::default()
};
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"
}));
}
}