use crate::{
agent::cancellation::AgentCancellation,
agent::load_compact_prompt,
config::{
AuthState, 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, capture_session_snapshot,
record_session_compaction_at_byte_offset_with_snapshot, 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>,
}
#[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) 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>> {
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::agent::runner::provider_from_selection(&provider_config, &selection, &job.cwd)?;
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(),
},
Some(job.settings.provider_stream.semantic_progress_timeout()),
)
}
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::login::refreshed_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_override(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>,
}
#[cfg(test)]
pub(crate) fn compact_session_with_provider(
run: CompactionProviderRun<'_>,
) -> anyhow::Result<Option<CompactionResult>> {
compact_session_with_provider_and_timeout(run, None)
}
fn compact_session_with_provider_and_timeout(
run: CompactionProviderRun<'_>,
semantic_progress_timeout: Option<std::time::Duration>,
) -> anyhow::Result<Option<CompactionResult>> {
let session_snapshot = capture_session_snapshot(run.session.path())?;
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,
)?
.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 mut raw_summary = String::new();
run.provider
.stream_cancellable(request, run.cancellation, &mut |event| {
run.cancellation.check()?;
if let ProviderEvent::TextDelta(delta) = event {
raw_summary.push_str(&delta);
}
run.cancellation.check()?;
Ok(())
})?;
run.cancellation.check()?;
let Some(summary) = sanitize_compaction_summary(&raw_summary) else {
return Ok(None);
};
record_session_compaction_at_byte_offset_with_snapshot(
run.session,
run.cwd,
&summary,
run.provider_id,
run.model,
cutoff_byte_offset,
session_snapshot,
)?;
Ok(Some(CompactionResult {
session_id: run.session.id().to_string(),
summary,
provider: run.provider_id.to_string(),
model: run.model.to_string(),
}))
}
#[cfg(test)]
pub(crate) fn build_compaction_request(
paths: &McPaths,
session: &Session,
model: &str,
custom_instructions: Option<&str>,
) -> anyhow::Result<ProviderRequest> {
let snapshot = session
.read_events_tolerant_bounded(
ConversationReplayLimits::default().max_lines,
ConversationReplayLimits::default().max_bytes,
)?
.events;
build_compaction_request_from_events(paths, session.id(), &snapshot, model, custom_instructions)
}
pub(crate) fn build_compaction_request_from_events(
paths: &McPaths,
session_id: &str,
events: &[SessionEvent],
model: &str,
custom_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);
let 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());
conversation.push(ProviderConversationItem::Message(ChatMessage::user(
final_instruction,
)));
Ok(ProviderRequest::from_conversation_without_tools(
model.to_string(),
conversation,
))
}
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
)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
config::{McPaths, ProviderCredential},
providers::ProviderToolResult,
sessions::{SessionEvent, SessionManager},
};
use serde_json::json;
use std::{
io::{Read, Write},
net::{TcpListener, TcpStream},
sync::{Arc, Mutex, atomic::AtomicBool},
time::Duration,
};
use tempfile::TempDir;
struct CapturingProvider {
request: Mutex<Option<ProviderRequest>>,
deltas: Vec<&'static str>,
error: Option<&'static str>,
}
struct AppendingProvider {
request: Mutex<Option<ProviderRequest>>,
session: Session,
cwd: PathBuf,
}
impl Provider for CapturingProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
*self.request.lock().unwrap() = Some(request);
if let Some(error) = self.error {
anyhow::bail!(error);
}
for delta in &self.deltas {
on_event(ProviderEvent::TextDelta((*delta).to_string()))?;
}
on_event(ProviderEvent::Done)
}
}
impl Provider for AppendingProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
_cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
*self.request.lock().unwrap() = Some(request);
append_event(
&self.session,
&self.cwd,
"user_input",
json!({"text":"CONCURRENT_APPEND_AFTER_SNAPSHOT"}),
);
on_event(ProviderEvent::TextDelta("summary".to_string()))?;
on_event(ProviderEvent::Done)
}
}
fn append_event(session: &Session, cwd: &Path, event_type: &str, payload: serde_json::Value) {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
cwd.to_path_buf(),
payload,
))
.unwrap();
}
fn read_http_request(stream: &mut TcpStream) -> anyhow::Result<()> {
const MAX_REQUEST_BYTES: usize = 1_048_576;
stream.set_read_timeout(Some(Duration::from_secs(2)))?;
let mut buffer = Vec::new();
let header_end = read_until(&mut buffer, stream, b"\r\n\r\n", MAX_REQUEST_BYTES)?;
let headers = std::str::from_utf8(&buffer[..header_end])?;
let mut content_length = None;
let mut chunked = false;
for line in headers.split("\r\n").skip(1) {
let (name, value) = line
.split_once(':')
.ok_or_else(|| anyhow::anyhow!("malformed HTTP header"))?;
if name.eq_ignore_ascii_case("content-length") {
if content_length.is_some() {
anyhow::bail!("duplicate Content-Length header");
}
content_length = Some(
value
.trim()
.parse::<usize>()
.map_err(|_| anyhow::anyhow!("malformed Content-Length header"))?,
);
} else if name.eq_ignore_ascii_case("transfer-encoding") {
chunked = value
.split(',')
.any(|encoding| encoding.trim().eq_ignore_ascii_case("chunked"));
}
}
if !chunked && content_length.is_none() {
anyhow::bail!("missing HTTP request body framing");
}
if chunked {
let mut position = header_end + 4;
loop {
let line_end = read_line_at(&mut buffer, stream, &mut position, MAX_REQUEST_BYTES)?;
let line = std::str::from_utf8(&buffer[position..line_end])?;
let size_text = line.split(';').next().unwrap_or("").trim();
let size = usize::from_str_radix(size_text, 16)
.map_err(|_| anyhow::anyhow!("malformed chunk size"))?;
position = line_end + 2;
if size == 0 {
loop {
let trailer_end =
read_line_at(&mut buffer, stream, &mut position, MAX_REQUEST_BYTES)?;
if trailer_end == position {
return Ok(());
}
let trailer = std::str::from_utf8(&buffer[position..trailer_end])?;
if !trailer.contains(':') {
anyhow::bail!("malformed chunk trailer");
}
position = trailer_end + 2;
}
}
read_range(&mut buffer, stream, &mut position, size, MAX_REQUEST_BYTES)?;
read_range(&mut buffer, stream, &mut position, 2, MAX_REQUEST_BYTES)?;
if &buffer[position - 2..position] != b"\r\n" {
anyhow::bail!("chunk data missing CRLF");
}
}
}
let Some(content_length) = content_length else {
return Ok(());
};
let body_start = header_end + 4;
if buffer.len() < body_start {
anyhow::bail!("HTTP request body starts beyond buffered data");
}
let already_read = buffer.len() - body_start;
if already_read > content_length {
anyhow::bail!("HTTP request contains bytes beyond Content-Length");
}
let mut position = buffer.len();
read_range(
&mut buffer,
stream,
&mut position,
content_length - already_read,
MAX_REQUEST_BYTES,
)?;
Ok(())
}
fn read_until(
buffer: &mut Vec<u8>,
stream: &mut TcpStream,
delimiter: &[u8],
max_bytes: usize,
) -> anyhow::Result<usize> {
loop {
if let Some(position) = buffer
.windows(delimiter.len())
.position(|window| window == delimiter)
{
return Ok(position);
}
read_more(buffer, stream, max_bytes)?;
}
}
fn read_line_at(
buffer: &mut Vec<u8>,
stream: &mut TcpStream,
position: &mut usize,
max_bytes: usize,
) -> anyhow::Result<usize> {
loop {
if let Some(offset) = buffer[*position..]
.windows(2)
.position(|window| window == b"\r\n")
{
return Ok(*position + offset);
}
read_more(buffer, stream, max_bytes)?;
}
}
fn read_range(
buffer: &mut Vec<u8>,
stream: &mut TcpStream,
position: &mut usize,
length: usize,
max_bytes: usize,
) -> anyhow::Result<()> {
let end = position
.checked_add(length)
.ok_or_else(|| anyhow::anyhow!("HTTP request length overflow"))?;
while buffer.len() < end {
read_more(buffer, stream, max_bytes)?;
}
*position = end;
Ok(())
}
fn read_more(
buffer: &mut Vec<u8>,
stream: &mut TcpStream,
max_bytes: usize,
) -> anyhow::Result<()> {
if buffer.len() >= max_bytes {
anyhow::bail!("HTTP request exceeds test helper limit");
}
let mut chunk = [0_u8; 8_192];
let read = stream.read(&mut chunk)?;
if read == 0 {
anyhow::bail!("premature EOF while reading HTTP request");
}
buffer.extend_from_slice(&chunk[..read]);
if buffer.len() > max_bytes {
anyhow::bail!("HTTP request exceeds test helper limit");
}
Ok(())
}
fn start_sse_server(
summary: &'static str,
) -> (String, std::thread::JoinHandle<anyhow::Result<()>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let handle = std::thread::spawn(move || {
let (mut stream, _) = listener.accept()?;
read_http_request(&mut stream)?;
let payload = format!("data: {{\"delta\":\"{summary}\"}}\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())?;
Ok(())
});
(format!("http://{addr}/v1"), handle)
}
#[test]
fn read_http_request_accepts_fragmented_content_length() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let address = listener.local_addr().unwrap();
let server = std::thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
read_http_request(&mut stream)
});
let mut client = TcpStream::connect(address).unwrap();
for part in [
b"POST / HTTP/1.1\r\nContent-Length: 6\r\n\r\n".as_slice(),
b"ab".as_slice(),
b"cdef".as_slice(),
] {
client.write_all(part).unwrap();
}
drop(client);
server.join().unwrap().unwrap();
}
#[test]
fn read_http_request_accepts_chunk_extensions_trailers_and_terminator_bytes() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let address = listener.local_addr().unwrap();
let server = std::thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
read_http_request(&mut stream)
});
let mut client = TcpStream::connect(address).unwrap();
for part in [
b"POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n8;ext=yes\r\n".as_slice(),
b"abc0\r\n\r\n\r\n3\r\n".as_slice(),
b"xyz\r\n0\r\nX-Test: yes\r\n\r\n".as_slice(),
] {
client.write_all(part).unwrap();
}
drop(client);
server.join().unwrap().unwrap();
}
#[test]
fn compaction_context_budget_uses_cached_window_and_preserves_reserve() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut entry = crate::model_catalog::ModelCatalogEntry::new_codex("gpt-test");
entry.context_window = Some(196_000);
crate::model_catalog::write_catalog_cache(
&paths,
crate::providers::OPENAI_CODEX_PROVIDER,
&[entry],
)
.unwrap();
let settings = Settings {
context: Some(ContextBudget {
enabled: true,
max_tokens: 128_000,
reserve_tokens: 40_000,
keep_recent_tokens: 20_000,
..ContextBudget::default()
}),
..Settings::default()
};
let budget = compaction_context_budget(
&settings,
&paths,
crate::providers::OPENAI_CODEX_PROVIDER,
"gpt-test",
);
assert_eq!(budget.max_tokens, 196_000);
assert_eq!(budget.reserve_tokens, 40_000);
assert!(budget.enabled);
assert_eq!(budget.keep_recent_tokens, 20_000);
}
#[test]
fn compaction_context_budget_model_override_wins_after_cached_window() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut entry = crate::model_catalog::ModelCatalogEntry::new_codex("gpt-test");
entry.context_window = Some(196_000);
crate::model_catalog::write_catalog_cache(
&paths,
crate::providers::OPENAI_CODEX_PROVIDER,
&[entry],
)
.unwrap();
let settings = Settings {
context: Some(ContextBudget {
max_tokens: 128_000,
reserve_tokens: 16_384,
model_overrides: std::collections::BTreeMap::from([(
"openai-codex/gpt-test".to_string(),
crate::context::ContextBudgetOverride {
max_tokens: Some(400_000),
reserve_tokens: Some(32_000),
keep_recent_tokens: None,
},
)]),
..ContextBudget::default()
}),
..Settings::default()
};
let budget = compaction_context_budget(
&settings,
&paths,
crate::providers::OPENAI_CODEX_PROVIDER,
"gpt-test",
);
assert_eq!(budget.max_tokens, 400_000);
assert_eq!(budget.reserve_tokens, 32_000);
assert_eq!(
budget.keep_recent_tokens,
ContextBudget::default().keep_recent_tokens
);
}
#[test]
fn compact_session_uses_cached_compaction_model_window_end_to_end() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let (base_url, server) = start_sse_server("cached summary");
let settings = Settings {
compaction: crate::config::CompactionSettings {
provider: Some("local-provider".to_string()),
model: Some("summary-model".to_string()),
..crate::config::CompactionSettings::default()
},
custom_providers: std::collections::BTreeMap::from([(
"local-provider".to_string(),
crate::config::make_custom_provider_config("Local", &base_url, "").unwrap(),
)]),
..Settings::default()
};
crate::config::write_settings(&paths, &settings).unwrap();
let mut entry =
crate::model_catalog::ModelCatalogEntry::new("local-provider", "summary-model");
entry.context_window = Some(196_000);
crate::model_catalog::write_catalog_cache_for_configured_provider(
&paths,
"local-provider",
&[entry],
)
.unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let text = "x".repeat(10_000);
append_event(&session, temp.path(), "user_input", json!({"text": text}));
append_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"ack"}),
);
let request = build_compaction_request(&paths, &session, "summary-model", None).unwrap();
let default_error = ensure_compaction_request_fits(
&request,
"local-provider",
&ContextBudget {
max_tokens: 1,
reserve_tokens: 0,
..ContextBudget::default()
},
)
.unwrap_err()
.to_string();
assert!(
default_error.contains("compaction request is estimated"),
"{default_error}"
);
let budget =
compaction_context_budget(&settings, &paths, "local-provider", "summary-model");
ensure_compaction_request_fits(&request, "local-provider", &budget).unwrap();
let active = EffectiveConfig {
provider: Some(crate::providers::OPENAI_CODEX_PROVIDER.to_string()),
model: Some("active-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: None,
paths: paths.clone(),
};
let result = compact_session(CompactSessionJob {
active_config: active,
settings,
session: session.clone(),
cwd: temp.path().to_path_buf(),
cancellation: AgentCancellation::default(),
custom_instructions: None,
})
.unwrap()
.unwrap();
server.join().unwrap().unwrap();
assert_eq!(result.summary, "cached summary");
assert_eq!(result.provider, "local-provider");
assert_eq!(result.model, "summary-model");
let events = session.read_events().unwrap();
assert_eq!(events.last().unwrap().event_type, "compaction");
assert_eq!(events.last().unwrap().payload["provider"], "local-provider");
assert_eq!(events.last().unwrap().payload["model"], "summary-model");
}
#[test]
fn compaction_request_uses_prompt_replay_tool_items_and_no_tools() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(
&session,
temp.path(),
"user_input",
json!({"text":"inspect"}),
);
append_event(
&session,
temp.path(),
"tool_call",
json!({"id":"call_1","name":"read","arguments":{"path":"a.txt"}}),
);
append_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"file text"}}),
);
append_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"done"}),
);
let request = build_compaction_request(&paths, &session, "compact-model", None).unwrap();
assert_eq!(request.model, "compact-model");
assert!(!request.tools_enabled());
assert!(matches!(
&request.conversation_items()[0],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::System
&& message.content.contains("continuation-ready summary")
));
assert!(request.conversation_items().iter().any(|item| matches!(
item,
ProviderConversationItem::ResponseItem(value)
if value.get("call_id").and_then(serde_json::Value::as_str) == Some("call_1")
)));
assert!(request.conversation_items().iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(ProviderToolResult { output, .. }) if output == "file text"
)));
}
#[test]
fn compaction_provider_run_threads_text_verbosity_to_request() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(&session, temp.path(), "user_input", json!({"text":"work"}));
let provider = CapturingProvider {
request: Mutex::new(None),
deltas: vec!["summary"],
error: None,
};
compact_session_with_provider(CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &provider,
context_budget: &ContextBudget::default(),
text_verbosity: Some(TextVerbosity::Medium),
cancellation: &AgentCancellation::default(),
custom_instructions: None,
})
.unwrap()
.unwrap();
let request = provider.request.lock().unwrap().clone().unwrap();
assert_eq!(request.text_verbosity(), Some(TextVerbosity::Medium));
}
#[test]
fn automatic_compaction_threads_semantic_progress_timeout_to_request() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(&session, temp.path(), "user_input", json!({"text":"work"}));
let provider = CapturingProvider {
request: Mutex::new(None),
deltas: vec!["summary"],
error: None,
};
let timeout = Duration::from_secs(17);
compact_session_with_provider_and_timeout(
CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &provider,
context_budget: &ContextBudget::default(),
text_verbosity: None,
cancellation: &AgentCancellation::default(),
custom_instructions: None,
},
Some(timeout),
)
.unwrap()
.unwrap();
let request = provider.request.lock().unwrap().clone().unwrap();
assert_eq!(request.semantic_progress_timeout(), Some(timeout));
}
#[test]
fn compaction_uses_one_bounded_snapshot_for_request_and_cutoff() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(
&session,
temp.path(),
"user_input",
json!({"text":"SNAPSHOT_ONLY"}),
);
let provider = AppendingProvider {
request: Mutex::new(None),
session: session.clone(),
cwd: temp.path().to_path_buf(),
};
compact_session_with_provider(CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &provider,
context_budget: &ContextBudget::default(),
text_verbosity: Some(TextVerbosity::Low),
cancellation: &AgentCancellation::default(),
custom_instructions: None,
})
.unwrap()
.unwrap();
let request = provider.request.lock().unwrap().clone().unwrap();
let request_items = request.conversation_items();
let material = crate::context::conversation_cache_material(
"local",
"compact-model",
"",
request_items.as_ref(),
);
assert!(material.contains("SNAPSHOT_ONLY"));
assert!(!material.contains("CONCURRENT_APPEND_AFTER_SNAPSHOT"));
let events = session.read_events().unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events.last().unwrap().event_type, "user_input");
assert_eq!(events[0].event_type, "compaction");
assert_eq!(events[0].payload["cutoff_event_count"], 0);
}
#[test]
fn compaction_persists_non_empty_summary_only_after_provider_success() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(&session, temp.path(), "user_input", json!({"text":"work"}));
let provider = CapturingProvider {
request: Mutex::new(None),
deltas: vec![" summary ", "text "],
error: None,
};
let result = compact_session_with_provider(CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &provider,
context_budget: &ContextBudget::default(),
text_verbosity: Some(TextVerbosity::Low),
cancellation: &AgentCancellation::default(),
custom_instructions: None,
})
.unwrap()
.unwrap();
assert_eq!(result.summary, "summary text");
let events = session.read_events().unwrap();
assert_eq!(events.last().unwrap().event_type, "compaction");
assert_eq!(events.last().unwrap().payload["summary"], "summary text");
assert_eq!(events.last().unwrap().payload["provider"], "local");
assert_eq!(events.last().unwrap().payload["model"], "compact-model");
assert_eq!(events.last().unwrap().payload["cutoff_event_count"], 0);
}
#[test]
fn compaction_empty_summary_and_provider_error_do_not_write_checkpoint() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(&session, temp.path(), "user_input", json!({"text":"work"}));
let empty_provider = CapturingProvider {
request: Mutex::new(None),
deltas: vec![" \n\t "],
error: None,
};
assert!(
compact_session_with_provider(CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &empty_provider,
context_budget: &ContextBudget::default(),
text_verbosity: Some(TextVerbosity::Low),
cancellation: &AgentCancellation::default(),
custom_instructions: None,
})
.unwrap()
.is_none()
);
assert_eq!(session.read_events().unwrap().len(), 1);
let error_provider = CapturingProvider {
request: Mutex::new(None),
deltas: vec![],
error: Some("provider failed"),
};
assert!(
compact_session_with_provider(CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &error_provider,
context_budget: &ContextBudget::default(),
text_verbosity: Some(TextVerbosity::Low),
cancellation: &AgentCancellation::default(),
custom_instructions: None,
})
.is_err()
);
assert_eq!(session.read_events().unwrap().len(), 1);
}
#[test]
fn compaction_over_budget_fails_before_provider_and_checkpoint() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(
&session,
temp.path(),
"user_input",
json!({"text":"x".repeat(1000)}),
);
let provider = CapturingProvider {
request: Mutex::new(None),
deltas: vec!["summary"],
error: None,
};
let budget = ContextBudget {
enabled: true,
max_tokens: 4,
reserve_tokens: 1,
keep_recent_tokens: 1,
..ContextBudget::default()
};
let error = compact_session_with_provider(CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &provider,
context_budget: &budget,
text_verbosity: Some(TextVerbosity::Low),
cancellation: &AgentCancellation::default(),
custom_instructions: None,
})
.unwrap_err()
.to_string();
assert!(error.contains("compaction request is estimated"), "{error}");
assert!(provider.request.lock().unwrap().is_none());
assert_eq!(session.read_events().unwrap().len(), 1);
}
#[test]
fn compaction_cancellation_before_stream_does_not_call_provider() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(&session, temp.path(), "user_input", json!({"text":"work"}));
let provider = CapturingProvider {
request: Mutex::new(None),
deltas: vec!["summary"],
error: None,
};
let canceled = Arc::new(AtomicBool::new(true));
let error = compact_session_with_provider(CompactionProviderRun {
paths: &paths,
session: &session,
cwd: temp.path(),
provider_id: "local",
model: "compact-model",
provider: &provider,
context_budget: &ContextBudget::default(),
text_verbosity: Some(TextVerbosity::Low),
cancellation: &AgentCancellation::new(canceled),
custom_instructions: None,
})
.unwrap_err();
assert!(crate::agent::cancellation::is_run_canceled(&error));
assert!(provider.request.lock().unwrap().is_none());
assert_eq!(session.read_events().unwrap().len(), 1);
}
#[test]
fn claude_code_compaction_inherits_sonnet_default() {
let temp = TempDir::new().unwrap();
let active = EffectiveConfig {
provider: Some(crate::providers::CLAUDE_CODE_PROVIDER.to_string()),
model: None,
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: Some(ProviderCredential::NoAuth),
paths: McPaths::from_root(temp.path().join("mc")),
};
assert_eq!(
active_compaction_defaults(&active),
(
crate::providers::CLAUDE_CODE_PROVIDER,
crate::providers::DEFAULT_CLAUDE_CODE_MODEL
)
);
}
#[test]
fn compaction_active_auth_refresh_is_only_required_when_inheriting_active_provider() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let mut active = EffectiveConfig {
provider: Some(crate::providers::OPENAI_CODEX_PROVIDER.to_string()),
model: Some("active-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: None,
paths,
};
let current_auth = AuthState::Missing {
provider: crate::providers::OPENAI_CODEX_PROVIDER.to_string(),
};
let override_settings = Settings {
compaction: crate::config::CompactionSettings {
provider: Some("local-provider".to_string()),
model: Some("summary-model".to_string()),
..crate::config::CompactionSettings::default()
},
..Settings::default()
};
let skipped = refresh_active_auth_for_compaction_if_inherited(
&mut active,
&override_settings,
¤t_auth,
)
.unwrap();
assert!(skipped.is_none());
assert!(active.auth.is_none());
let error = refresh_active_auth_for_compaction_if_inherited(
&mut active,
&Settings::default(),
¤t_auth,
)
.unwrap_err()
.to_string();
assert!(
error.contains("missing OAuth auth for provider 'openai-codex'"),
"{error}"
);
}
#[test]
fn compact_session_uses_compaction_override_without_active_auth() {
let temp = TempDir::new().unwrap();
let paths = McPaths::from_root(temp.path().join("mc"));
let (base_url, server) = start_sse_server("override summary");
let settings = Settings {
compaction: crate::config::CompactionSettings {
provider: Some("local-provider".to_string()),
model: Some("summary-model".to_string()),
..crate::config::CompactionSettings::default()
},
custom_providers: std::collections::BTreeMap::from([(
"local-provider".to_string(),
crate::config::make_custom_provider_config("Local", &base_url, "").unwrap(),
)]),
..Settings::default()
};
crate::config::write_settings(&paths, &settings).unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_event(&session, temp.path(), "user_input", json!({"text":"work"}));
let active = EffectiveConfig {
provider: Some(crate::providers::OPENAI_CODEX_PROVIDER.to_string()),
model: Some("active-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: None,
paths: paths.clone(),
};
let result = compact_session(CompactSessionJob {
active_config: active,
settings,
session: session.clone(),
cwd: temp.path().to_path_buf(),
cancellation: AgentCancellation::default(),
custom_instructions: None,
})
.unwrap()
.unwrap();
server.join().unwrap().unwrap();
assert_eq!(result.summary, "override summary");
assert_eq!(result.provider, "local-provider");
assert_eq!(result.model, "summary-model");
let events = session.read_events().unwrap();
assert_eq!(events.last().unwrap().payload["provider"], "local-provider");
assert_eq!(events.last().unwrap().payload["model"], "summary-model");
}
#[test]
fn compact_session_resolves_override_settings() {
let temp = TempDir::new().unwrap();
let active = EffectiveConfig {
provider: Some("local".to_string()),
model: Some("active-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: Some(ProviderCredential::NoAuth),
paths: McPaths::from_root(temp.path().join("mc")),
};
let settings = Settings::default();
let config = settings
.compaction
.resolve_config(active.provider_id(), active.model.as_deref().unwrap())
.unwrap();
assert_eq!(config.provider, "local");
assert_eq!(config.model, "active-model");
}
}