use crate::types::Message;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
pub const SESSION_COMPACTION_CADENCE_KEY: &str = "session_compaction_cadence";
#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub struct SessionCompactionCadence {
pub session_boundary_index: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_compaction_boundary_index: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_compaction_attempt_boundary_index: Option<u64>,
}
#[derive(Debug, Clone)]
pub struct CompactionContext {
pub last_input_tokens: u64,
pub message_count: usize,
pub estimated_history_tokens: u64,
pub estimated_request_bytes: u64,
pub provider_request_pressure: Option<ProviderRequestPressure>,
pub last_compaction_boundary_index: Option<u64>,
pub session_boundary_index: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ProviderRequestPressure {
pub encoded_bytes: u64,
pub max_bytes: Option<u64>,
}
impl ProviderRequestPressure {
pub fn new(encoded_bytes: u64, max_bytes: Option<u64>) -> Self {
Self {
encoded_bytes,
max_bytes,
}
}
pub fn effective_cap(self, configured_cap: Option<u64>) -> Option<u64> {
match (configured_cap, self.max_bytes) {
(Some(configured), Some(provider)) => Some(configured.min(provider)),
(Some(configured), None) => Some(configured),
(None, Some(provider)) => Some(provider),
(None, None) => None,
}
}
pub fn trigger_threshold(self, configured_cap: Option<u64>) -> Option<u64> {
self.effective_cap(configured_cap).map(|cap| {
cap.saturating_mul(REQUEST_BYTE_TRIGGER_NUMERATOR) / REQUEST_BYTE_TRIGGER_DENOMINATOR
})
}
}
#[derive(Debug, Clone)]
pub struct CompactionResult {
pub messages: Vec<Message>,
pub summary: CompactionSummary,
pub retained: Vec<CompactionRetained>,
pub discarded: Vec<CompactionDiscard>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct CompactionSummary {
pub rebuilt_offset: u64,
pub message: Message,
}
impl CompactionSummary {
pub fn new(rebuilt_offset: u64, message: Message) -> Self {
Self {
rebuilt_offset,
message,
}
}
}
pub const COMPACTION_SUMMARY_PREFIX: &str = "\
[Context compacted] A previous context produced the following summary of work so far. \
The current tool and session state is preserved. Use this summary to continue without \
duplicating work:\n\n";
#[derive(Debug, Clone, PartialEq)]
pub struct CompactionDiscard {
pub source_offset: u64,
pub message: Message,
}
impl CompactionDiscard {
pub fn new(source_offset: u64, message: Message) -> Self {
Self {
source_offset,
message,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct CompactionRetained {
pub source_offset: u64,
pub rebuilt_offset: u64,
pub message: Message,
}
impl CompactionRetained {
pub fn new(source_offset: u64, rebuilt_offset: u64, message: Message) -> Self {
Self {
source_offset,
rebuilt_offset,
message,
}
}
}
const REQUEST_BYTE_TRIGGER_NUMERATOR: u64 = 4;
const REQUEST_BYTE_TRIGGER_DENOMINATOR: u64 = 5;
#[derive(Debug, Clone)]
pub struct CompactionConfig {
pub auto_compact_threshold: u64,
pub max_request_bytes: Option<u64>,
pub recent_turn_budget: usize,
pub max_summary_tokens: u32,
pub min_turns_between_compactions: u32,
}
impl CompactionConfig {
pub fn request_byte_trigger_threshold(&self) -> Option<u64> {
self.max_request_bytes.map(|cap| {
cap.saturating_mul(REQUEST_BYTE_TRIGGER_NUMERATOR) / REQUEST_BYTE_TRIGGER_DENOMINATOR
})
}
}
impl Default for CompactionConfig {
fn default() -> Self {
Self {
auto_compact_threshold: 100_000,
max_request_bytes: None,
recent_turn_budget: 4,
max_summary_tokens: 4096,
min_turns_between_compactions: 3,
}
}
}
pub trait Compactor: Send + Sync {
fn should_compact(&self, ctx: &CompactionContext) -> bool;
fn request_byte_cap(&self, pressure: ProviderRequestPressure) -> Option<u64> {
pressure.max_bytes
}
fn compaction_prompt(&self) -> &str;
fn max_summary_tokens(&self) -> u32;
fn prepare_for_summarization(&self, messages: &[Message]) -> Vec<Message> {
messages.to_vec()
}
fn rebuild_history_under_pressure(
&self,
messages: &[Message],
summary: &str,
_pressure: Option<ProviderRequestPressure>,
) -> CompactionResult {
self.rebuild_history(messages, summary)
}
fn rebuild_history(&self, messages: &[Message], summary: &str) -> CompactionResult;
}
#[derive(Debug, Clone, Copy)]
pub struct CompactionWindow<'a> {
pub messages: &'a [Message],
pub last_input_tokens: u64,
pub session_boundary_index: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CuratedCompactionSummary(String);
impl CuratedCompactionSummary {
pub fn new(text: impl Into<String>) -> Result<Self, CompactionCuratorError> {
let text = text.into();
if text.trim().is_empty() {
return Err(CompactionCuratorError::EmptySummary);
}
Ok(Self(text))
}
pub fn as_str(&self) -> &str {
&self.0
}
pub fn into_string(self) -> String {
self.0
}
}
#[derive(Debug, thiserror::Error)]
pub enum CompactionCuratorError {
#[error("curator produced an empty compaction summary")]
EmptySummary,
#[error("curator failed to produce a compaction summary: {0}")]
Failed(String),
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
pub trait CompactionCurator: Send + Sync {
async fn curate_summary(
&self,
window: CompactionWindow<'_>,
) -> Result<CuratedCompactionSummary, CompactionCuratorError>;
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
#[test]
fn curated_compaction_summary_rejects_empty_text() {
assert!(matches!(
CuratedCompactionSummary::new(""),
Err(CompactionCuratorError::EmptySummary)
));
assert!(matches!(
CuratedCompactionSummary::new(" \n\t"),
Err(CompactionCuratorError::EmptySummary)
));
}
#[test]
fn curated_compaction_summary_round_trips_text() {
let summary = CuratedCompactionSummary::new("curated summary").unwrap();
assert_eq!(summary.as_str(), "curated summary");
assert_eq!(summary.into_string(), "curated summary");
}
#[test]
fn request_byte_trigger_threshold_is_four_fifths_of_the_cap() {
let config = CompactionConfig {
max_request_bytes: Some(10_000_000),
..CompactionConfig::default()
};
assert_eq!(config.request_byte_trigger_threshold(), Some(8_000_000));
assert_eq!(
CompactionConfig::default().request_byte_trigger_threshold(),
None
);
}
#[test]
fn provider_and_configured_request_caps_compose_by_minimum() {
let pressure = ProviderRequestPressure::new(1, Some(8_000_000));
assert_eq!(pressure.effective_cap(Some(10_000_000)), Some(8_000_000));
assert_eq!(pressure.effective_cap(Some(6_000_000)), Some(6_000_000));
assert_eq!(
pressure.trigger_threshold(Some(10_000_000)),
Some(6_400_000)
);
assert_eq!(pressure.trigger_threshold(Some(6_000_000)), Some(4_800_000));
assert_eq!(
ProviderRequestPressure::new(1, None).effective_cap(Some(6_000_000)),
Some(6_000_000)
);
assert_eq!(
ProviderRequestPressure::new(1, Some(8_000_000)).effective_cap(None),
Some(8_000_000)
);
}
}