use crate::{
agent::load_compact_prompt,
auth::AuthState,
cancellation::AgentCancellation,
config::{
EffectiveConfig, McPaths, Settings, TextVerbosity, load_effective_provider_selection,
},
context::{
ContextBudget, ConversationReplayLimits, build_conversation_replay_from_events,
project_provider_request_input_tokens,
},
providers::{
ChatMessage, Provider, ProviderConversationItem, ProviderEvent, ProviderRequest,
ProviderSelection,
},
sessions::{Session, SessionEvent, sanitize_compaction_summary},
};
use std::path::{Path, PathBuf};
#[derive(Clone)]
pub(crate) struct CompactSessionJob {
pub(crate) active_config: EffectiveConfig,
pub(crate) settings: Settings,
pub(crate) session: Session,
pub(crate) cwd: PathBuf,
pub(crate) cancellation: AgentCancellation,
pub(crate) custom_instructions: Option<String>,
pub(crate) additional_instructions: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct CompactionResult {
pub(crate) session_id: String,
pub(crate) summary: String,
pub(crate) provider: String,
pub(crate) model: String,
pub(crate) requested_service_tier: Option<String>,
pub(crate) fast_outcome: crate::fast::FastOutcome,
pub(crate) rotation_warning: Option<String>,
}
pub(crate) fn compaction_success_message(result: &CompactionResult) -> String {
format!(
"compaction complete: session {}; older turns remain in JSONL history; future provider requests use summary boundary",
result.session_id
)
}
pub(crate) fn compact_session(job: CompactSessionJob) -> anyhow::Result<Option<CompactionResult>> {
compact_session_observed(job, &mut |_, _| Ok(()))
}
pub(crate) enum CompactionObservation<'a> {
RequestStarted,
ProviderEvent(&'a ProviderEvent),
}
pub(crate) fn compact_session_observed(
job: CompactSessionJob,
observe: &mut dyn FnMut(&str, CompactionObservation<'_>) -> anyhow::Result<()>,
) -> anyhow::Result<Option<CompactionResult>> {
let (active_provider, active_model) = active_compaction_defaults(&job.active_config);
let compaction_config = job
.settings
.compaction
.resolve_config(active_provider, active_model)
.map_err(|message| anyhow::anyhow!(crate::output::redact_sensitive_text(&message)))?;
let provider_config = compaction_effective_config(&job.active_config, &compaction_config)?;
let selection = ProviderSelection::from_config(&provider_config)?;
let provider = crate::providers::provider_from_selection_with_settings_for_workload(
&provider_config,
&selection,
&job.cwd,
&job.settings,
crate::fast::FastWorkload::Compaction,
)?;
let budget = compaction_context_budget(
&job.settings,
&provider_config.paths,
&selection.provider,
&selection.model,
);
compact_session_with_provider_and_timeout(
CompactionProviderRun {
paths: &provider_config.paths,
session: &job.session,
cwd: &job.cwd,
provider_id: &selection.provider,
model: &selection.model,
provider: provider.as_ref(),
context_budget: &budget,
text_verbosity: job.settings.text_verbosity_for(&selection.provider),
cancellation: &job.cancellation,
custom_instructions: job.custom_instructions.as_deref(),
additional_instructions: job.additional_instructions.as_deref(),
},
Some(job.settings.provider_stream.semantic_progress_timeout()),
observe,
)
}
pub(crate) fn refresh_active_auth_for_compaction_if_inherited(
active_config: &mut EffectiveConfig,
settings: &Settings,
current_auth_state: &AuthState,
) -> anyhow::Result<Option<AuthState>> {
if !compaction_inherits_active_provider(settings) {
return Ok(None);
}
let auth_state =
crate::auth::refreshed_codex_auth_state(&active_config.paths, current_auth_state)?;
active_config.auth = auth_state.credential().cloned();
Ok(Some(auth_state))
}
fn compaction_inherits_active_provider(settings: &Settings) -> bool {
settings.compaction.provider.is_none() && settings.compaction.model.is_none()
}
fn active_compaction_defaults(active_config: &EffectiveConfig) -> (&str, &str) {
(
active_config.provider_id(),
active_config.model.as_deref().unwrap_or_else(|| {
crate::providers::default_model_for_provider(active_config.provider_id())
}),
)
}
fn compaction_effective_config(
active_config: &EffectiveConfig,
compaction_config: &crate::config::CompactionConfig,
) -> anyhow::Result<EffectiveConfig> {
let (active_provider, active_model) = active_compaction_defaults(active_config);
if compaction_config.provider == active_provider && compaction_config.model == active_model {
return Ok(active_config.clone());
}
load_effective_provider_selection(
&active_config.paths,
&compaction_config.provider,
&compaction_config.model,
)
}
fn compaction_context_budget(
settings: &Settings,
paths: &McPaths,
provider: &str,
model: &str,
) -> ContextBudget {
let mut budget = settings.context.clone().unwrap_or_default();
if let Some(context_window) =
crate::model_catalog::cached_model_context_window(paths, provider, model)
{
budget.max_tokens = context_window;
}
budget.apply_model_limits(provider, model);
budget
}
pub(crate) struct CompactionProviderRun<'a> {
pub(crate) paths: &'a McPaths,
pub(crate) session: &'a Session,
pub(crate) cwd: &'a Path,
pub(crate) provider_id: &'a str,
pub(crate) model: &'a str,
pub(crate) provider: &'a dyn Provider,
pub(crate) context_budget: &'a ContextBudget,
pub(crate) text_verbosity: Option<TextVerbosity>,
pub(crate) cancellation: &'a AgentCancellation,
pub(crate) custom_instructions: Option<&'a str>,
pub(crate) additional_instructions: Option<&'a str>,
}
fn compact_session_with_provider_and_timeout(
run: CompactionProviderRun<'_>,
semantic_progress_timeout: Option<std::time::Duration>,
observe: &mut dyn FnMut(&str, CompactionObservation<'_>) -> anyhow::Result<()>,
) -> anyhow::Result<Option<CompactionResult>> {
let session_snapshot = run.session.prepare_compaction_snapshot()?;
let snapshot = run.session.read_events_tolerant_bounded(
ConversationReplayLimits::default().max_lines,
ConversationReplayLimits::default().max_bytes,
)?;
let cutoff_byte_offset = snapshot.cutoff_bytes;
let snapshot_events = snapshot.events;
let mut request = build_compaction_request_from_events(
run.paths,
run.session.id(),
&snapshot_events,
run.model,
run.custom_instructions,
run.additional_instructions,
)?
.with_text_verbosity(run.text_verbosity);
if let Some(timeout) = semantic_progress_timeout {
request = request.with_semantic_progress_timeout(timeout);
}
ensure_compaction_request_fits(&request, run.provider_id, run.context_budget)?;
run.cancellation.check()?;
let requested_service_tier = run.provider.requested_service_tier().map(str::to_string);
let mut returned_service_tier = None;
let mut raw_summary = String::new();
observe(run.provider_id, CompactionObservation::RequestStarted)?;
run.provider
.stream_cancellable(request, run.cancellation, &mut |event| {
run.cancellation.check()?;
observe(
run.provider_id,
CompactionObservation::ProviderEvent(&event),
)?;
match event {
ProviderEvent::TextDelta(delta) => raw_summary.push_str(&delta),
ProviderEvent::ServiceTier(tier) => returned_service_tier = Some(tier.clone()),
_ => {}
}
run.cancellation.check()?;
Ok(())
})?;
run.cancellation.check()?;
let Some(summary) = sanitize_compaction_summary(&raw_summary) else {
return Ok(None);
};
let rotation_warning =
match crate::sessions::record_session_compaction_at_byte_offset_with_snapshot(
run.session,
run.cwd,
&summary,
run.provider_id,
run.model,
cutoff_byte_offset,
session_snapshot,
) {
Ok(()) => None,
Err(error) if error.committed_summary() == Some(summary.as_str()) => {
Some(crate::output::sanitize_display_text(&format!(
"compaction checkpoint committed; rotation cleanup warning: {error}"
)))
}
Err(error) => return Err(error.into()),
};
let fast_outcome = crate::fast::fast_outcome(
run.provider_id,
requested_service_tier.as_deref(),
returned_service_tier.as_deref(),
);
Ok(Some(CompactionResult {
session_id: run.session.id().to_string(),
summary,
provider: run.provider_id.to_string(),
model: run.model.to_string(),
requested_service_tier,
fast_outcome,
rotation_warning,
}))
}
pub(crate) fn build_compaction_request_from_events(
paths: &McPaths,
session_id: &str,
events: &[SessionEvent],
model: &str,
custom_instructions: Option<&str>,
additional_instructions: Option<&str>,
) -> anyhow::Result<ProviderRequest> {
let prompt = load_compact_prompt(Some(&paths.prompts))?;
let replay = build_conversation_replay_from_events(session_id, events);
let mut conversation = Vec::with_capacity(replay.items.len() + 2);
conversation.push(ProviderConversationItem::Message(ChatMessage::system(
prompt,
)));
conversation.extend(replay.items);
if let Some(note) =
crate::sessions::effects::ExecutionEffects::from_events(session_id, events).provider_note()
{
conversation.push(ProviderConversationItem::Message(ChatMessage::user(note)));
}
let mut final_instruction = custom_instructions
.map(str::trim)
.filter(|instructions| !instructions.is_empty())
.map(|instructions| {
format!(
"Produce the compacted continuation summary now. Apply these user instructions for this compaction: {instructions}"
)
})
.unwrap_or_else(|| "Produce the compacted continuation summary now.".to_string());
if let Some(additional_instructions) = additional_instructions
.map(str::trim)
.filter(|instructions| !instructions.is_empty())
{
final_instruction.push_str("\n\n");
final_instruction.push_str(additional_instructions);
}
conversation.push(ProviderConversationItem::Message(ChatMessage::user(
final_instruction,
)));
Ok(
ProviderRequest::from_conversation_without_tools(model.to_string(), conversation)
.with_conversation_id(crate::providers::conversation_id_for_session(session_id)),
)
}
fn ensure_compaction_request_fits(
request: &ProviderRequest,
provider_id: &str,
context_budget: &ContextBudget,
) -> anyhow::Result<()> {
if !context_budget.enabled {
return Ok(());
}
let estimated_tokens = project_provider_request_input_tokens(provider_id, request).tokens;
let threshold = context_budget.threshold_tokens();
if estimated_tokens <= threshold {
return Ok(());
}
anyhow::bail!(
"compaction request is estimated at {estimated_tokens} tokens, exceeding threshold {threshold} tokens (max_tokens={}, reserve_tokens={}); no checkpoint was written; choose a larger compaction model or start /new",
context_budget.max_tokens,
context_budget.reserve_tokens
)
}