use async_trait::async_trait;
use std::sync::Arc;
use crate::event::HarnessUsage;
use crate::model::{
collect_model_response, ChatMessage, ModelClient, ModelClientError, ModelResponse,
ModelTurnInput,
};
use crate::tools::ToolSpec;
pub const DEFAULT_TRIGGER_FRACTION: f64 = 0.90;
pub const PRUNE_PROTECT_TURNS: usize = 2;
pub const PRUNE_MINIMUM_TOKENS: u64 = 2_000;
pub const PRUNE_MIN_CONTENT_TOKENS: u64 = 100;
pub const PRUNED_CONTENT_STUB: &str =
"[pruned — output removed to free context; re-run the tool if needed]";
pub const DEFAULT_TAIL_MIN_MESSAGES: usize = 4;
pub const DEFAULT_SUMMARY_MAX_TOKENS: i32 = 2_000;
pub const DEFAULT_USER_MESSAGE_TOKEN_BUDGET: u64 = 20_000;
pub fn estimate_tokens(s: &str) -> u64 {
if s.trim().is_empty() {
return 0;
}
let mut ascii: u64 = 0;
let mut non_ascii: u64 = 0;
for c in s.chars() {
if c.is_ascii() {
ascii += 1;
} else {
non_ascii += 1;
}
}
ascii.div_ceil(4) + non_ascii
}
pub fn estimate_chat_message_tokens(m: &ChatMessage) -> u64 {
let tokens = match m {
ChatMessage::User { content, .. } => estimate_tokens(content),
ChatMessage::Assistant {
text,
tool_calls,
thinking,
usage: _,
} => {
let text_tokens = text.as_deref().map(estimate_tokens).unwrap_or(0);
let tc_tokens: u64 = tool_calls
.iter()
.map(|tc| estimate_tokens(&tc.input.to_string()) + estimate_tokens(&tc.name) + 8)
.sum();
let thinking_tokens = thinking
.as_ref()
.map(|t| {
estimate_tokens(&t.text)
+ t.signature.as_deref().map(estimate_tokens).unwrap_or(0)
})
.unwrap_or(0);
text_tokens + tc_tokens + thinking_tokens
}
ChatMessage::Tool { content, .. } => estimate_tokens(content) + 16,
};
tokens.max(1)
}
pub fn estimate_messages_tokens(messages: &[ChatMessage]) -> u64 {
messages.iter().map(estimate_chat_message_tokens).sum()
}
fn usage_context_tokens(u: &HarnessUsage) -> u64 {
u.input_tokens
}
pub fn estimate_context_tokens_anchored(messages: &[ChatMessage]) -> u64 {
let anchor = messages
.iter()
.enumerate()
.rev()
.find_map(|(idx, m)| match m {
ChatMessage::Assistant { usage: Some(u), .. } => Some((idx, usage_context_tokens(u))),
_ => None,
});
match anchor {
Some((idx, measured)) => {
let tail: u64 = messages[idx + 1..]
.iter()
.map(estimate_chat_message_tokens)
.sum();
measured + tail
}
None => estimate_messages_tokens(messages),
}
}
pub struct CompactionContext {
pub system_prompt: Option<String>,
pub model_client: Arc<dyn ModelClient>,
pub context_window_tokens: u64,
pub tools: Vec<ToolSpec>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct CompactionOutcome {
pub messages: Vec<ChatMessage>,
pub usage: Option<crate::event::HarnessUsage>,
}
#[derive(Debug, thiserror::Error)]
pub enum CompactionError {
#[error("compaction model call failed: {0}")]
ModelCall(#[from] ModelClientError),
#[error("model produced empty summary; refusing to fold history")]
EmptySummary,
}
#[async_trait]
pub trait CompactionStrategy: Send + Sync {
fn should_compact(&self, messages: &[ChatMessage], context_window_tokens: u64) -> bool;
async fn compact(
&self,
messages: Vec<ChatMessage>,
ctx: &CompactionContext,
) -> Result<CompactionOutcome, CompactionError>;
}
pub type DefaultCompactionStrategy = SummarizeCompactionStrategy;
pub struct SummarizeCompactionStrategy {
pub trigger_fraction: f64,
pub tail_min_messages: usize,
pub summary_max_tokens: i32,
pub summary_prompt: String,
pub user_message_token_budget: u64,
}
impl Default for SummarizeCompactionStrategy {
fn default() -> Self {
Self {
trigger_fraction: DEFAULT_TRIGGER_FRACTION,
tail_min_messages: DEFAULT_TAIL_MIN_MESSAGES,
summary_max_tokens: DEFAULT_SUMMARY_MAX_TOKENS,
summary_prompt: DEFAULT_SUMMARY_PROMPT.into(),
user_message_token_budget: DEFAULT_USER_MESSAGE_TOKEN_BUDGET,
}
}
}
impl SummarizeCompactionStrategy {
pub fn with_trigger_fraction(mut self, fraction: f64) -> Self {
self.trigger_fraction = fraction;
self
}
pub fn with_tail_min_messages(mut self, n: usize) -> Self {
self.tail_min_messages = n;
self
}
pub fn with_summary_max_tokens(mut self, n: i32) -> Self {
self.summary_max_tokens = n;
self
}
pub fn with_user_message_token_budget(mut self, budget: u64) -> Self {
self.user_message_token_budget = budget;
self
}
}
pub const DEFAULT_SUMMARY_PROMPT: &str = "You are performing a CONTEXT CHECKPOINT COMPACTION. \
Create a handoff summary for another agent instance that will resume this task.\n\n\
Include:\n\
- Current progress and key decisions made\n\
- Important context, constraints, or user preferences that must be respected\n\
- What remains to be done (clear next steps)\n\
- Any critical data, file paths, command outputs, or references needed to continue\n\n\
If a prior <conversation-summary> block exists in this conversation, produce an UPDATED \
summary that supersedes it (incorporating all activity since). \
Output only the summary text — no preamble, no closing remarks.";
#[async_trait]
impl CompactionStrategy for SummarizeCompactionStrategy {
fn should_compact(&self, messages: &[ChatMessage], context_window_tokens: u64) -> bool {
if messages.len() <= self.tail_min_messages {
return false;
}
let tokens = estimate_context_tokens_anchored(messages);
let threshold = ((context_window_tokens as f64) * self.trigger_fraction).round() as u64;
tokens > threshold
}
async fn compact(
&self,
messages: Vec<ChatMessage>,
ctx: &CompactionContext,
) -> Result<CompactionOutcome, CompactionError> {
if messages.len() <= self.tail_min_messages {
return Ok(CompactionOutcome {
messages,
usage: None,
});
}
let threshold = ((ctx.context_window_tokens as f64) * self.trigger_fraction).round() as u64;
let (messages, freed_tokens) = prune_tool_outputs(messages, PRUNE_PROTECT_TURNS);
if freed_tokens > 0 {
tracing::debug!(
target: "harness::compaction",
freed_tokens,
"prune pass freed tokens"
);
let tokens_after_prune =
estimate_context_tokens_anchored(&messages).saturating_sub(freed_tokens);
if tokens_after_prune <= threshold {
tracing::info!(
target: "harness::compaction",
freed_tokens,
tokens_after_prune,
"prune sufficient — summarise skipped"
);
return Ok(CompactionOutcome {
messages,
usage: None,
});
}
}
let mut summarize_messages = messages.clone();
summarize_messages.push(ChatMessage::User {
content: self.summary_prompt.clone(),
attachments: vec![],
});
let request = ModelTurnInput {
system_prompt: ctx.system_prompt.clone(),
messages: summarize_messages,
tools: ctx.tools.clone(),
hosted_tools: vec![],
tool_choice: crate::model::ToolChoice::Auto,
parallel_tool_calls: None,
};
let stream = ctx.model_client.stream(request).await?;
let response = collect_model_response(stream).await?;
let (summary_text, usage) = match response {
ModelResponse::Message { text, usage, .. } => (text, usage),
ModelResponse::ToolCall { .. } => return Err(CompactionError::EmptySummary),
};
if summary_text.trim().is_empty() {
return Err(CompactionError::EmptySummary);
}
let user_texts = collect_user_message_texts(&messages);
if user_texts.is_empty() {
return Ok(CompactionOutcome { messages, usage });
}
let out =
build_compacted_history(&user_texts, &summary_text, self.user_message_token_budget);
Ok(CompactionOutcome {
messages: out,
usage,
})
}
}
fn serialize_summary(summary: &str) -> String {
format!("<conversation-summary>\n{summary}\n</conversation-summary>")
}
fn collect_user_message_texts(messages: &[ChatMessage]) -> Vec<String> {
messages
.iter()
.filter_map(|m| match m {
ChatMessage::User { content, .. } if !is_summary_message(content) => {
Some(content.clone())
}
_ => None,
})
.collect()
}
fn is_summary_message(content: &str) -> bool {
content.trim_start().starts_with("<conversation-summary>")
}
fn build_compacted_history(
user_texts: &[String],
summary_text: &str,
token_budget: u64,
) -> Vec<ChatMessage> {
let mut selected: Vec<String> = Vec::new();
let mut remaining = token_budget;
for text in user_texts.iter().rev() {
if remaining == 0 {
break;
}
let tokens = estimate_tokens(text);
if tokens <= remaining {
selected.push(text.clone());
remaining -= tokens;
} else {
selected.push(truncate_to_token_budget(text, remaining));
break;
}
}
selected.reverse(); let mut out = Vec::with_capacity(selected.len() + 1);
for text in selected {
out.push(ChatMessage::User {
content: text,
attachments: vec![],
});
}
out.push(ChatMessage::User {
content: serialize_summary(summary_text),
attachments: vec![],
});
out
}
fn truncate_to_token_budget(s: &str, budget: u64) -> String {
if budget == 0 {
return String::new();
}
let mut ascii: u64 = 0;
let mut non_ascii: u64 = 0;
let mut end = 0usize;
for (byte_pos, c) in s.char_indices() {
let (na, nn) = if c.is_ascii() {
(ascii + 1, non_ascii)
} else {
(ascii, non_ascii + 1)
};
if na.div_ceil(4) + nn > budget {
break;
}
ascii = na;
non_ascii = nn;
end = byte_pos + c.len_utf8();
}
s[..end].to_string()
}
pub fn prune_tool_outputs(
messages: Vec<ChatMessage>,
protect_turns: usize,
) -> (Vec<ChatMessage>, u64) {
let stub_tokens = estimate_tokens(PRUNED_CONTENT_STUB);
let mut user_turns_seen: usize = 0;
let mut candidates: Vec<(usize, u64)> = Vec::new();
let mut gross_freed: u64 = 0;
for (i, msg) in messages.iter().enumerate().rev() {
match msg {
ChatMessage::User { content, .. } => {
if is_summary_message(content) {
break;
}
user_turns_seen += 1;
}
ChatMessage::Tool { content, .. } => {
if user_turns_seen < protect_turns {
continue;
}
let tokens = estimate_tokens(content);
if tokens >= PRUNE_MIN_CONTENT_TOKENS {
candidates.push((i, tokens));
gross_freed += tokens;
}
}
_ => {}
}
}
let replacements = candidates.len() as u64;
let net_freed = gross_freed.saturating_sub(stub_tokens * replacements);
if net_freed < PRUNE_MINIMUM_TOKENS {
return (messages, 0);
}
let mut out = messages;
for (i, _) in &candidates {
if let ChatMessage::Tool { content, .. } = &mut out[*i] {
*content = PRUNED_CONTENT_STUB.to_string();
}
}
(out, net_freed)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::model::{ModelChunk, ModelClient};
use crate::tools::ToolInvocation;
use async_trait::async_trait;
use futures::stream::{BoxStream, StreamExt};
#[derive(Clone)]
struct FixedSummaryClient {
summary: String,
}
#[async_trait]
impl ModelClient for FixedSummaryClient {
fn hosted_capability(
&self,
_capability: crate::model::HostedCapability,
) -> crate::model::CapabilitySupport {
crate::model::CapabilitySupport::Unsupported
}
async fn stream(
&self,
_input: ModelTurnInput,
) -> Result<BoxStream<'static, Result<ModelChunk, ModelClientError>>, ModelClientError>
{
let chunks = vec![
Ok(ModelChunk::TextDelta {
msg_id: "sum".into(),
delta: self.summary.clone(),
}),
Ok(ModelChunk::Done {
stop_reason: "end_turn".into(),
usage: None,
}),
];
Ok(futures::stream::iter(chunks).boxed())
}
}
fn user(s: &str) -> ChatMessage {
ChatMessage::User {
content: s.into(),
attachments: vec![],
}
}
fn assistant_text(s: &str) -> ChatMessage {
ChatMessage::Assistant {
text: Some(s.into()),
tool_calls: vec![],
thinking: None,
usage: None,
}
}
fn tool_msg(id: &str, content: &str) -> ChatMessage {
ChatMessage::Tool {
tool_call_id: id.into(),
content: content.into(),
is_error: false,
attachments: vec![],
}
}
#[test]
fn token_estimate_grows_with_content_size() {
let small = user("hi");
let big = user(&"x".repeat(8000));
assert!(estimate_chat_message_tokens(&big) > estimate_chat_message_tokens(&small));
}
#[test]
fn estimate_tokens_splits_ascii_and_cjk() {
assert_eq!(estimate_tokens(""), 0);
assert_eq!(estimate_tokens(" \n"), 0);
assert_eq!(estimate_tokens("abcd"), 1);
assert_eq!(estimate_tokens("abcde"), 2); assert_eq!(estimate_tokens("你好世界"), 4); assert_eq!(estimate_tokens("hi你好"), 3); }
#[test]
fn token_estimate_counts_cjk_near_one_per_char() {
let cjk = user(&"汉".repeat(1000));
let estimate = estimate_chat_message_tokens(&cjk);
assert!(
estimate >= 1000,
"CJK undercounted: got {estimate}, want >= 1000"
);
}
#[test]
fn token_estimate_includes_tool_call_input() {
let bare = assistant_text("done");
let with_tool = ChatMessage::Assistant {
text: Some("done".into()),
tool_calls: vec![ToolInvocation {
id: "tc".into(),
name: "bash".into(),
input: serde_json::json!({"command": "echo lots of bytes here for sure"}),
raw_emitted_args: None,
}],
thinking: None,
usage: None,
};
assert!(estimate_chat_message_tokens(&with_tool) > estimate_chat_message_tokens(&bare));
}
fn assistant_with_usage(text: &str, input: u64, cache_read: u64) -> ChatMessage {
ChatMessage::Assistant {
text: Some(text.into()),
tool_calls: vec![],
thinking: None,
usage: Some(HarnessUsage {
input_tokens: input,
cache_read_input_tokens: cache_read,
..Default::default()
}),
}
}
#[test]
fn anchored_estimate_uses_reported_usage_for_prefix() {
let msgs = vec![
user("hi"),
assistant_with_usage("ok", 40_000, 0),
user(&"x".repeat(4000)), ];
let tail = estimate_chat_message_tokens(&msgs[2]);
assert_eq!(estimate_context_tokens_anchored(&msgs), 40_000 + tail);
}
#[test]
fn anchored_estimate_does_not_double_count_cache_reads() {
let msgs = vec![assistant_with_usage("ok", 10_000, 25_000)];
assert_eq!(estimate_context_tokens_anchored(&msgs), 10_000);
}
#[test]
fn anchored_estimate_picks_newest_usage_anchor() {
let msgs = vec![
user("a"),
assistant_with_usage("first", 10_000, 0),
user("b"),
assistant_with_usage("second", 30_000, 0),
user("tail"),
];
let tail = estimate_chat_message_tokens(&msgs[4]);
assert_eq!(estimate_context_tokens_anchored(&msgs), 30_000 + tail);
}
#[test]
fn anchored_estimate_falls_back_without_usage() {
let msgs = vec![user("hi"), assistant_text("there"), user("more")];
assert_eq!(
estimate_context_tokens_anchored(&msgs),
estimate_messages_tokens(&msgs)
);
}
#[test]
fn should_compact_skips_when_below_threshold() {
let strat = SummarizeCompactionStrategy::default();
let messages = vec![user("hello"), assistant_text("hi")];
assert!(!strat.should_compact(&messages, 200_000));
}
#[test]
fn should_compact_fires_when_above_threshold() {
let strat = SummarizeCompactionStrategy::default();
let messages = vec![
user(&"x".repeat(8000)),
assistant_text(&"y".repeat(8000)),
user(&"x".repeat(8000)),
assistant_text(&"y".repeat(8000)),
user(&"x".repeat(8000)),
];
assert!(strat.should_compact(&messages, 11_000));
}
#[test]
fn should_compact_respects_tail_min_floor() {
let strat = SummarizeCompactionStrategy::default();
let messages = vec![
user(&"x".repeat(100_000)),
assistant_text(&"y".repeat(100_000)),
];
assert!(!strat.should_compact(&messages, 1_000));
}
#[tokio::test]
async fn compact_folds_history_into_summary_plus_tail() {
let strat = SummarizeCompactionStrategy::default().with_tail_min_messages(2);
let ctx = CompactionContext {
system_prompt: None,
model_client: Arc::new(FixedSummaryClient {
summary: "we ran ls and grep".into(),
}),
context_window_tokens: 10_000,
tools: vec![],
};
let messages = vec![
user("first user"),
assistant_text("response 1"),
user("second user"),
tool_msg("tc1", "tool result"),
user("third user"),
assistant_text("final response"),
];
let outcome = strat.compact(messages, &ctx).await.unwrap();
let out = outcome.messages;
assert_eq!(out.len(), 4, "3 user messages + 1 summary");
match &out[0] {
ChatMessage::User { content, .. } => assert_eq!(content, "first user"),
other => panic!("expected User at [0], got {other:?}"),
}
match &out[1] {
ChatMessage::User { content, .. } => assert_eq!(content, "second user"),
other => panic!("expected User at [1], got {other:?}"),
}
match &out[2] {
ChatMessage::User { content, .. } => assert_eq!(content, "third user"),
other => panic!("expected User at [2], got {other:?}"),
}
assert!(matches!(&out[3], ChatMessage::User { content, .. }
if content.contains("<conversation-summary>") && content.contains("we ran ls and grep")));
assert!(out.len() < 6);
assert!(outcome.usage.is_none());
}
#[tokio::test]
async fn compact_returns_empty_summary_error_on_blank_response() {
let strat = SummarizeCompactionStrategy::default().with_tail_min_messages(2);
let ctx = CompactionContext {
system_prompt: None,
model_client: Arc::new(FixedSummaryClient { summary: "".into() }),
context_window_tokens: 10_000,
tools: vec![],
};
let messages = vec![
user("a"),
assistant_text("b"),
user("c"),
assistant_text("d"),
];
let err = strat.compact(messages, &ctx).await.unwrap_err();
assert!(matches!(err, CompactionError::EmptySummary));
}
#[tokio::test]
async fn compact_skips_when_messages_at_or_below_tail_min() {
let strat = SummarizeCompactionStrategy::default().with_tail_min_messages(4);
let ctx = CompactionContext {
system_prompt: None,
model_client: Arc::new(FixedSummaryClient {
summary: "irrelevant".into(),
}),
context_window_tokens: 1_000,
tools: vec![],
};
let messages = vec![
user("1"),
assistant_text("2"),
user("3"),
assistant_text("4"),
];
let outcome = strat.compact(messages.clone(), &ctx).await.unwrap();
assert_eq!(outcome.messages, messages);
assert!(outcome.usage.is_none());
}
fn big_tool(id: &str) -> ChatMessage {
tool_msg(id, &"x".repeat(8_400))
}
fn small_tool(id: &str) -> ChatMessage {
tool_msg(id, &"x".repeat(40))
}
#[test]
fn prune_replaces_old_large_tool_outputs() {
let messages = vec![
user("turn 0"),
big_tool("t0"),
user("turn 1"),
big_tool("t1"),
user("turn 2"),
big_tool("t2"),
];
let (pruned, freed) = prune_tool_outputs(messages, 2);
assert!(freed > 0, "expected tokens to be freed");
assert_eq!(pruned[1], tool_msg("t0", PRUNED_CONTENT_STUB));
assert_ne!(pruned[3], tool_msg("t1", PRUNED_CONTENT_STUB));
assert_ne!(pruned[5], tool_msg("t2", PRUNED_CONTENT_STUB));
}
#[test]
fn prune_does_not_touch_small_tool_outputs() {
let messages = vec![
user("turn 0"),
small_tool("t0"),
user("turn 1"),
big_tool("t1"),
user("turn 2"),
big_tool("t2"),
];
let (pruned, _freed) = prune_tool_outputs(messages.clone(), 2);
assert_eq!(pruned[1], messages[1]);
}
#[test]
fn prune_no_op_when_savings_below_minimum() {
let messages = vec![
user("turn 0"),
small_tool("t0"),
user("turn 1"),
big_tool("t1"),
user("turn 2"),
];
let (out, freed) = prune_tool_outputs(messages.clone(), 2);
assert_eq!(freed, 0);
assert_eq!(out, messages);
}
#[test]
fn prune_stops_at_summary_boundary() {
let messages = vec![
user("turn 0"),
big_tool("t_before"),
user("<conversation-summary>\nprevious context\n</conversation-summary>"),
user("turn 1"),
big_tool("t1"),
user("turn 2"),
big_tool("t2"),
];
let (out, freed) = prune_tool_outputs(messages.clone(), 2);
assert_eq!(
freed, 0,
"summary boundary must block pruning of earlier content"
);
assert_eq!(out, messages);
}
#[tokio::test]
async fn compact_skips_summarise_when_prune_sufficient() {
struct PanicSummaryClient;
#[async_trait::async_trait]
impl ModelClient for PanicSummaryClient {
fn hosted_capability(
&self,
_capability: crate::model::HostedCapability,
) -> crate::model::CapabilitySupport {
crate::model::CapabilitySupport::Unsupported
}
async fn stream(
&self,
_: ModelTurnInput,
) -> Result<BoxStream<'static, Result<ModelChunk, ModelClientError>>, ModelClientError>
{
panic!("summarise should not be called when prune is sufficient");
}
}
let big_content = "y".repeat(200_000); let messages = vec![
user("turn 0"),
tool_msg("t0", &big_content),
user("turn 1"),
assistant_text("ok"),
user("turn 2"),
assistant_text("done"),
];
let ctx = CompactionContext {
system_prompt: None,
model_client: Arc::new(PanicSummaryClient),
context_window_tokens: 60_000,
tools: vec![],
};
let strat = SummarizeCompactionStrategy {
trigger_fraction: 0.9,
tail_min_messages: 2,
..Default::default()
};
let outcome = strat.compact(messages, &ctx).await.unwrap();
assert!(outcome.usage.is_none());
assert_eq!(outcome.messages[1], tool_msg("t0", PRUNED_CONTENT_STUB));
}
}