use crate::api::types::{Message, MessageContent};
use crate::api::{ApiClient, ThinkingMode};
use anyhow::Result;
use std::collections::VecDeque;
use tracing::{debug, info, warn};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompressionMethod {
Micro,
Auto,
Full,
}
impl std::fmt::Display for CompressionMethod {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
CompressionMethod::Micro => write!(f, "MicroCompact"),
CompressionMethod::Auto => write!(f, "AutoCompact"),
CompressionMethod::Full => write!(f, "FullCompact"),
}
}
}
#[derive(Debug, Clone)]
pub struct CompressionMetrics {
pub method: CompressionMethod,
pub tokens_before: usize,
pub tokens_after: usize,
pub tokens_saved: usize,
pub messages_before: usize,
pub messages_after: usize,
pub duration_ms: u64,
pub llm_input_tokens: usize,
pub llm_output_tokens: usize,
}
impl CompressionMetrics {
pub fn new(
method: CompressionMethod,
tokens_before: usize,
tokens_after: usize,
messages_before: usize,
messages_after: usize,
duration_ms: u64,
) -> Self {
let tokens_saved = tokens_before.saturating_sub(tokens_after);
Self {
method,
tokens_before,
tokens_after,
tokens_saved,
messages_before,
messages_after,
duration_ms,
llm_input_tokens: 0,
llm_output_tokens: 0,
}
}
pub fn summary(&self) -> String {
format!(
"{}: Saved {} tokens ({} → {}), {} messages → {} (in {}ms)",
self.method,
self.tokens_saved,
self.tokens_before,
self.tokens_after,
self.messages_before,
self.messages_after,
self.duration_ms
)
}
pub fn with_llm_tokens(mut self, input: usize, output: usize) -> Self {
self.llm_input_tokens = input;
self.llm_output_tokens = output;
self
}
}
#[derive(Debug, Clone)]
pub struct AutoCompactConfig {
pub token_threshold: usize,
pub reserve_buffer: usize,
pub max_summary_tokens: usize,
pub max_consecutive_failures: usize,
}
impl Default for AutoCompactConfig {
fn default() -> Self {
Self {
token_threshold: 80, reserve_buffer: 13_000,
max_summary_tokens: 20_000,
max_consecutive_failures: 3,
}
}
}
pub struct AutoCompactManager {
config: AutoCompactConfig,
consecutive_failures: usize,
circuit_open: bool,
last_compression: Option<CompressionMetrics>,
}
impl AutoCompactManager {
pub fn new(config: AutoCompactConfig) -> Self {
Self {
config,
consecutive_failures: 0,
circuit_open: false,
last_compression: None,
}
}
pub fn with_threshold(percentage: usize) -> Self {
let config = AutoCompactConfig {
token_threshold: percentage,
..Default::default()
};
Self::new(config)
}
pub fn should_compress(&self, current_tokens: usize, context_window: usize) -> bool {
if self.circuit_open {
debug!("Circuit breaker open, skipping auto-compression");
return false;
}
let threshold = (context_window * self.config.token_threshold) / 100;
current_tokens >= threshold
}
pub fn record_success(&mut self, metrics: CompressionMetrics) {
self.consecutive_failures = 0;
self.last_compression = Some(metrics);
}
pub fn record_failure(&mut self) {
self.consecutive_failures += 1;
if self.consecutive_failures >= self.config.max_consecutive_failures {
self.circuit_open = true;
warn!(
"AutoCompact circuit breaker opened after {} consecutive failures",
self.consecutive_failures
);
}
}
pub fn reset_circuit(&mut self) {
self.circuit_open = false;
self.consecutive_failures = 0;
}
pub fn is_circuit_open(&self) -> bool {
self.circuit_open
}
pub fn last_compression(&self) -> Option<&CompressionMetrics> {
self.last_compression.as_ref()
}
}
#[derive(Debug, Clone)]
pub struct FileAccessTracker {
recent_accesses: VecDeque<(String, std::time::Instant)>,
max_tracked_calls: usize,
}
impl Default for FileAccessTracker {
fn default() -> Self {
Self::new(20)
}
}
impl FileAccessTracker {
pub fn new(max_tracked_calls: usize) -> Self {
Self {
recent_accesses: VecDeque::with_capacity(max_tracked_calls),
max_tracked_calls,
}
}
pub fn record_access(&mut self, path: &str) {
let now = std::time::Instant::now();
self.recent_accesses.retain(|(p, _)| p != path);
self.recent_accesses.push_back((path.to_string(), now));
while self.recent_accesses.len() > self.max_tracked_calls {
self.recent_accesses.pop_front();
}
}
pub fn get_recent_files(&self, n: usize) -> Vec<String> {
self.recent_accesses
.iter()
.rev()
.take(n)
.map(|(p, _)| p.clone())
.collect()
}
pub fn get_tracked_files(&self) -> Vec<(String, std::time::Duration)> {
let now = std::time::Instant::now();
self.recent_accesses
.iter()
.map(|(p, t)| (p.clone(), now.duration_since(*t)))
.collect()
}
pub fn clear(&mut self) {
self.recent_accesses.clear();
}
pub fn len(&self) -> usize {
self.recent_accesses.len()
}
pub fn is_empty(&self) -> bool {
self.recent_accesses.is_empty()
}
}
fn estimate_tokens(msg: &Message) -> usize {
let content_tokens = msg.content.text().len() / 4;
let reasoning_tokens = msg
.reasoning_content
.as_ref()
.map(|r| r.len() / 4)
.unwrap_or(0);
let tool_call_tokens = msg.tool_calls.as_ref().map(|tc| tc.len() * 50).unwrap_or(0);
content_tokens + reasoning_tokens + tool_call_tokens + 4 }
fn truncate_chars(s: &str, max_chars: usize) -> String {
let mut chars = s.chars();
let truncated: String = chars.by_ref().take(max_chars).collect();
if chars.next().is_some() {
format!("{}...[truncated]", truncated)
} else {
s.to_string()
}
}
fn truncate_chars_with_total(s: &str, max_chars: usize) -> String {
let total_chars = s.chars().count();
if total_chars <= max_chars {
return s.to_string();
}
let truncated: String = s.chars().take(max_chars).collect();
format!(
"{}\n...[{} chars truncated by micro_compact]",
truncated,
total_chars.saturating_sub(max_chars)
)
}
pub fn micro_compact(messages: &mut Vec<Message>) -> CompressionMetrics {
let start = std::time::Instant::now();
let tokens_before = messages.iter().map(estimate_tokens).sum::<usize>();
let messages_before = messages.len();
if messages.len() <= 12 {
return CompressionMetrics::new(
CompressionMethod::Micro,
tokens_before,
tokens_before,
messages_before,
messages_before,
start.elapsed().as_millis() as u64,
);
}
const MIN_MESSAGES_TO_KEEP: usize = 22;
let system_msg = messages.first().cloned();
let keep_start = messages.len().saturating_sub(MIN_MESSAGES_TO_KEEP);
if keep_start <= 1 {
return CompressionMetrics::new(
CompressionMethod::Micro,
tokens_before,
tokens_before,
messages_before,
messages_before,
start.elapsed().as_millis() as u64,
);
}
let mut compressed = Vec::new();
if let Some(sys) = system_msg {
compressed.push(sys);
}
let to_compress = if keep_start > 1 {
&messages[1..keep_start]
} else {
&[]
};
let recent = &messages[keep_start..];
if !to_compress.is_empty() {
let old_user_count = to_compress.iter().filter(|m| m.role == "user").count();
let old_assistant_count = to_compress.iter().filter(|m| m.role == "assistant").count();
let old_tool_count = to_compress.iter().filter(|m| m.role == "tool").count();
let summary = format!(
"[MICRO-COMPACT: {} earlier messages compressed ({} user, {} assistant, {} tool results). \
Key decisions and file edits preserved.]",
to_compress.len(),
old_user_count,
old_assistant_count,
old_tool_count
);
compressed.push(Message::user(summary));
}
for (i, msg) in recent.iter().enumerate() {
let mut cleaned = msg.clone();
let is_recent = i >= recent.len().saturating_sub(4);
if !is_recent {
cleaned.reasoning_content = None;
}
if msg.role == "tool" && !is_recent {
if let MessageContent::Text(text) = &msg.content {
if text.chars().count() > 500 {
let truncated = truncate_chars_with_total(text, 500);
cleaned.content = MessageContent::Text(truncated);
}
}
}
compressed.push(cleaned);
}
let tokens_after = compressed.iter().map(estimate_tokens).sum::<usize>();
let messages_after = compressed.len();
let metrics = CompressionMetrics::new(
CompressionMethod::Micro,
tokens_before,
tokens_after,
messages_before,
messages_after,
start.elapsed().as_millis() as u64,
);
*messages = compressed;
metrics
}
pub async fn auto_compact(
client: &ApiClient,
messages: &mut Vec<Message>,
_config: &AutoCompactConfig,
) -> Result<CompressionMetrics> {
let start = std::time::Instant::now();
let tokens_before = messages.iter().map(estimate_tokens).sum::<usize>();
let messages_before = messages.len();
if messages.len() <= 8 {
return Ok(CompressionMetrics::new(
CompressionMethod::Auto,
tokens_before,
tokens_before,
messages_before,
messages_before,
start.elapsed().as_millis() as u64,
));
}
let system_msg = messages.first().cloned();
let keep_recent = 6;
let recent_start = messages.len().saturating_sub(keep_recent);
let recent_msgs: Vec<Message> = messages[recent_start..].to_vec();
let to_summarize = &messages[1..recent_start];
if to_summarize.is_empty() {
return Ok(CompressionMetrics::new(
CompressionMethod::Auto,
tokens_before,
tokens_before,
messages_before,
messages_before,
start.elapsed().as_millis() as u64,
));
}
let summary_content = format!(
"Summarize the following conversation history concisely. \
Preserve key facts, decisions, file paths, and action items. \
Omit routine tool outputs unless they contain errors or important results.\n\n{}",
to_summarize
.iter()
.enumerate()
.map(|(i, m)| {
let content = truncate_chars(m.content.text(), 800);
format!("[{}] {}: {}", i, m.role, content)
})
.collect::<Vec<_>>()
.join("\n\n")
);
let summary_request = vec![
Message::system("You are a context summarizer. Compress conversation history while preserving critical information for task completion. Be concise."),
Message::user(summary_content),
];
let response = tokio::time::timeout(
std::time::Duration::from_secs(60),
client.chat(summary_request, None, ThinkingMode::Disabled),
)
.await
.map_err(|_| anyhow::anyhow!("AutoCompact API call timed out after 60s"))??;
let summary = response
.choices
.first()
.map(|c| c.message.content.text().to_string())
.unwrap_or_else(|| "[Context compression: conversation history]".to_string());
let mut compressed = Vec::new();
if let Some(sys) = system_msg {
compressed.push(sys);
}
compressed.push(Message::user(format!(
"[AUTO-COMPACT SUMMARY — {} earlier messages]:\n{}",
to_summarize.len(),
summary
)));
compressed.extend(recent_msgs);
let tokens_after = compressed.iter().map(estimate_tokens).sum::<usize>();
let messages_after = compressed.len();
let metrics = CompressionMetrics::new(
CompressionMethod::Auto,
tokens_before,
tokens_after,
messages_before,
messages_after,
start.elapsed().as_millis() as u64,
)
.with_llm_tokens(
response.usage.prompt_tokens,
response.usage.completion_tokens,
);
*messages = compressed;
Ok(metrics)
}
pub async fn full_compact(
client: &ApiClient,
messages: &mut Vec<Message>,
file_tracker: &FileAccessTracker,
_target_budget: usize,
) -> Result<CompressionMetrics> {
let start = std::time::Instant::now();
let tokens_before = messages.iter().map(estimate_tokens).sum::<usize>();
let messages_before = messages.len();
let system_msg = messages.first().cloned();
let last_user_msg = messages.iter().rev().find(|m| m.role == "user").cloned();
if messages.len() <= 4 {
return Ok(CompressionMetrics::new(
CompressionMethod::Full,
tokens_before,
tokens_before,
messages_before,
messages_before,
start.elapsed().as_millis() as u64,
));
}
let summary_prompt = "Create a comprehensive but concise summary of this entire conversation. \
Include: 1) Task goal and current status, 2) Key files modified/accessed, \
3) Important decisions made, 4) Current blockers or next steps, \
5) Any active plans or schemas in use."
.to_string();
let summary_request = vec![
Message::system("You are a comprehensive context summarizer. Preserve all critical information for continuing the task."),
Message::user(format!("{}\n\nConversation:\n{}",
summary_prompt,
messages.iter().enumerate().map(|(i, m)| {
let content = truncate_chars(m.content.text(), 600);
format!("[{}] {}: {}", i, m.role, content)
}).collect::<Vec<_>>().join("\n")
)),
];
let response = tokio::time::timeout(
std::time::Duration::from_secs(90),
client.chat(summary_request, None, ThinkingMode::Disabled),
)
.await
.map_err(|_| anyhow::anyhow!("FullCompact API call timed out after 90s"))??;
let summary = response
.choices
.first()
.map(|c| c.message.content.text().to_string())
.unwrap_or_else(|| "[Full context compression applied]".to_string());
let mut compressed = Vec::new();
if let Some(sys) = system_msg {
compressed.push(sys);
}
compressed.push(Message::user(format!(
"[FULL-COMPACT — Complete conversation summary]:\n{}\n\n[Recent files and context follow]",
summary
)));
let recent_files = file_tracker.get_recent_files(5);
if !recent_files.is_empty() {
let mut file_context = String::from("\n## Recently Accessed Files:\n");
for path in recent_files {
if let Ok(content) = tokio::fs::read_to_string(&path).await {
let total_chars = content.chars().count();
let truncated = if total_chars > 20_000 {
format!(
"{}\n...[truncated, {} total chars]",
content.chars().take(20_000).collect::<String>(),
total_chars
)
} else {
content
};
file_context.push_str(&format!("\n### {}\n```\n{}\n```\n", path, truncated));
} else {
file_context.push_str(&format!("\n### {} (unavailable)\n", path));
}
}
compressed.push(Message::user(file_context));
}
if let Some(last_user) = last_user_msg {
if compressed.last().map(|m| m.content.text()) != Some(last_user.content.text()) {
compressed.push(last_user);
}
}
let tokens_after = compressed.iter().map(estimate_tokens).sum::<usize>();
let messages_after = compressed.len();
let metrics = CompressionMetrics::new(
CompressionMethod::Full,
tokens_before,
tokens_after,
messages_before,
messages_after,
start.elapsed().as_millis() as u64,
)
.with_llm_tokens(
response.usage.prompt_tokens,
response.usage.completion_tokens,
);
*messages = compressed;
Ok(metrics)
}
pub struct CompressionOrchestrator {
auto_manager: AutoCompactManager,
file_tracker: FileAccessTracker,
metrics_history: Vec<CompressionMetrics>,
}
impl Default for CompressionOrchestrator {
fn default() -> Self {
Self::new()
}
}
impl CompressionOrchestrator {
pub fn new() -> Self {
Self {
auto_manager: AutoCompactManager::new(AutoCompactConfig::default()),
file_tracker: FileAccessTracker::default(),
metrics_history: Vec::new(),
}
}
pub fn with_config(config: AutoCompactConfig) -> Self {
Self {
auto_manager: AutoCompactManager::new(config),
file_tracker: FileAccessTracker::default(),
metrics_history: Vec::new(),
}
}
pub fn run_micro(&mut self, messages: &mut Vec<Message>) -> CompressionMetrics {
let metrics = micro_compact(messages);
self.metrics_history.push(metrics.clone());
metrics
}
pub async fn run_auto(
&mut self,
client: &ApiClient,
messages: &mut Vec<Message>,
) -> Result<CompressionMetrics> {
match auto_compact(client, messages, &self.auto_manager.config).await {
Ok(metrics) => {
self.auto_manager.record_success(metrics.clone());
self.metrics_history.push(metrics.clone());
Ok(metrics)
}
Err(e) => {
self.auto_manager.record_failure();
Err(e)
}
}
}
pub async fn run_full(
&mut self,
client: &ApiClient,
messages: &mut Vec<Message>,
) -> Result<CompressionMetrics> {
let metrics = full_compact(
client,
messages,
&self.file_tracker,
50_000, )
.await?;
self.metrics_history.push(metrics.clone());
Ok(metrics)
}
pub async fn check_and_compress(
&mut self,
client: &ApiClient,
messages: &mut Vec<Message>,
current_tokens: usize,
context_window: usize,
) -> Option<CompressionMetrics> {
if !self
.auto_manager
.should_compress(current_tokens, context_window)
{
return None;
}
match self.run_auto(client, messages).await {
Ok(metrics) => {
info!("AutoCompact triggered: {}", metrics.summary());
Some(metrics)
}
Err(e) => {
warn!("AutoCompact failed, falling back to MicroCompact: {}", e);
let metrics = self.run_micro(messages);
info!("MicroCompact fallback: {}", metrics.summary());
Some(metrics)
}
}
}
pub fn record_file_access(&mut self, path: &str) {
self.file_tracker.record_access(path);
}
pub fn file_tracker(&self) -> &FileAccessTracker {
&self.file_tracker
}
pub fn metrics_history(&self) -> &[CompressionMetrics] {
&self.metrics_history
}
pub fn total_tokens_saved(&self) -> usize {
self.metrics_history.iter().map(|m| m.tokens_saved).sum()
}
pub fn reset(&mut self) {
self.auto_manager.reset_circuit();
self.file_tracker.clear();
self.metrics_history.clear();
}
}
#[cfg(test)]
#[path = "../../tests/unit/agent/compression/compression_test.rs"]
mod tests;