use crate::agent::AgentRunError;
use crate::agent::assembly::SessionEvent;
use crate::agent::compaction::compaction::{
SummarizeError, compact_with_model_context, estimate_context_tokens,
};
use crate::agent::messages::compaction_summary;
use crate::agent::session::session::SessionTreeEntry;
use crate::observability::{
ErrorCategory, OperationDetail, OperationOutcome, OperationScope, RuntimeMeasurements,
};
use crate::types::AgentMessage;
impl crate::agent::assembly::AgentHarness {
pub async fn force_compact(
&self,
custom_instructions: Option<String>,
) -> Result<bool, AgentRunError> {
self.start_runtime_extensions().await;
self.do_compact(true, custom_instructions).await
}
pub(crate) async fn run_auto_compaction(&self) -> Result<(), AgentRunError> {
let settings = self.compaction_settings.lock().clone();
if !settings.enabled {
return Ok(());
}
let (context_tokens, context_window) = {
let s = self.agent.state();
let model = match &s.model {
Some(m) => m,
None => return Ok(()),
};
let estimate = estimate_context_tokens(&s.messages);
(estimate.tokens, model.context_window)
};
let algorithm = self.compact_algorithms.algorithm(&settings.algorithm);
if !algorithm
.decide_compact(context_tokens, context_window, &settings)
.await
{
return Ok(());
}
let _ = self.do_compact(false, None).await?;
Ok(())
}
async fn do_compact(
&self,
from_hook: bool,
custom_instructions: Option<String>,
) -> Result<bool, AgentRunError> {
let model = match self.agent.state().model.clone() {
Some(m) => m,
None => return Ok(false),
};
let settings = self.compaction_settings.lock().clone();
self.before_runtime_compaction(&settings.algorithm, from_hook)
.await?;
let scope = OperationScope::start(
self.agent.runtime_observer(),
self.agent.active_run_operation(),
self.agent.observation_context(),
OperationDetail::Compaction {
algorithm: settings.algorithm.clone(),
provider: model.provider.0.clone(),
model: model.id.clone(),
},
);
let entries = match self.session.branch(None).await {
Ok(es) => es,
Err(e) => {
self.emit_harness_event(SessionEvent::Compaction {
from_hook,
summary: format!("compaction skipped: session branch read failed: {e}"),
tokens_before: 0,
});
scope.finish(
OperationOutcome::Failed,
Some(ErrorCategory::Persistence),
RuntimeMeasurements::default(),
);
self.runtime_compaction_failed(
serde_json::json!({
"algorithm": settings.algorithm,
"category": "persistence",
"message": e.to_string(),
}),
false,
)
.await;
return Ok(false);
}
};
let algorithm = self.compact_algorithms.algorithm(&settings.algorithm);
let persistent_model_context = self.runtime_compaction_context_messages();
let result = compact_with_model_context(
algorithm.as_ref(),
model,
&entries,
&persistent_model_context,
&settings,
custom_instructions,
self.stream_fn.clone(),
self.agent.active_token().unwrap_or_default(),
)
.await;
let result = match result {
Ok(r) if !r.summary.is_empty() => r,
Ok(r) => {
scope.finish(
OperationOutcome::Skipped,
None,
RuntimeMeasurements {
input_tokens: r.usage.input,
output_tokens: r.usage.output,
cache_read_tokens: r.usage.cache_read,
cache_write_tokens: r.usage.cache_write,
..Default::default()
},
);
self.runtime_compaction_succeeded(serde_json::json!({
"algorithm": settings.algorithm,
"applied": false,
"tokensBefore": r.tokens_before,
}))
.await;
return Ok(false);
}
Err(SummarizeError::Aborted) => {
scope.finish(
OperationOutcome::Cancelled,
Some(ErrorCategory::Cancellation),
RuntimeMeasurements::default(),
);
self.runtime_compaction_failed(
serde_json::json!({
"algorithm": settings.algorithm,
"category": "cancelled",
}),
true,
)
.await;
return Ok(false);
}
Err(e) => {
scope.finish(
OperationOutcome::Failed,
Some(ErrorCategory::Provider),
RuntimeMeasurements::default(),
);
self.runtime_compaction_failed(
serde_json::json!({
"algorithm": settings.algorithm,
"category": "provider",
"message": e.to_string(),
}),
false,
)
.await;
return Err(AgentRunError::Other(format!("compaction failed: {e}")));
}
};
let first_kept_entry_id = result.first_kept_entry_id.clone().unwrap_or_default();
if let Err(e) = self
.session
.append_compaction(
result.summary.clone(),
first_kept_entry_id.clone(),
result.tokens_before,
None,
from_hook,
)
.await
{
scope.finish(
OperationOutcome::Failed,
Some(ErrorCategory::Persistence),
RuntimeMeasurements {
input_tokens: result.usage.input,
output_tokens: result.usage.output,
cache_read_tokens: result.usage.cache_read,
cache_write_tokens: result.usage.cache_write,
..Default::default()
},
);
self.runtime_compaction_failed(
serde_json::json!({
"algorithm": settings.algorithm,
"category": "persistence",
"message": e.to_string(),
}),
false,
)
.await;
return Err(AgentRunError::Other(format!(
"session append compaction: {e}"
)));
}
self.emit_harness_event(SessionEvent::Compaction {
from_hook,
summary: result.summary.clone(),
tokens_before: result.tokens_before,
});
{
let mut s = self.agent.state();
let mut new_msgs: Vec<AgentMessage> = vec![compaction_summary(result.summary.clone())];
if !first_kept_entry_id.is_empty() {
if let Some(real_idx) = entries.iter().position(|e| e.id() == first_kept_entry_id) {
let kept_in_memory_start = entries[..real_idx]
.iter()
.filter(|e| matches!(e, SessionTreeEntry::Message { .. }))
.count();
if kept_in_memory_start <= s.messages.len() {
new_msgs.extend(s.messages[kept_in_memory_start..].iter().cloned());
}
}
}
s.messages = new_msgs;
}
scope.finish(
OperationOutcome::Succeeded,
None,
RuntimeMeasurements {
input_tokens: result.usage.input,
output_tokens: result.usage.output,
cache_read_tokens: result.usage.cache_read,
cache_write_tokens: result.usage.cache_write,
..Default::default()
},
);
self.runtime_compaction_succeeded(serde_json::json!({
"algorithm": settings.algorithm,
"applied": true,
"tokensBefore": result.tokens_before,
"firstKeptEntryId": first_kept_entry_id,
}))
.await;
Ok(true)
}
}
#[cfg(test)]
tests_bridge_macro::tests_bridge!("agent/compaction/triggers");
#[cfg(test)]
mod triggers_linecov_tests {
tests_bridge_macro::tests_bridge!("agent/compaction/triggers/linecov");
}