mod live_source;
pub mod observed_prefix;
mod placement;
mod pressure;
mod summary_continuation;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
use async_trait::async_trait;
use orchestral_core::agent_protocol::wire::{AgentSessionId, Digest, RunId};
use orchestral_core::agent_session::{
session_range_digest, validate_session_trace, AgentSessionError, AgentSessionEvent,
AgentSessionEventDraft, AgentSessionEventId, AgentSessionJournalStore, AgentSessionRecord,
SessionSourceRange,
};
use orchestral_core::model_protocol::{
ModelContent, ModelContextEstimate, ModelError, ModelMessage, ModelRole, ModelTokenAccounting,
ModelTokenMeter, ModelTokenMeterDescriptor, ModelToolDefinition,
};
use orchestral_core::skill_protocol::SkillId;
use serde::{Deserialize, Serialize};
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum SessionContextError {
#[error(transparent)]
Journal(#[from] AgentSessionError),
#[error("invalid Session Context request: {0}")]
InvalidRequest(String),
#[error("pinned Session Context exceeds the model input budget: used {used}, budget {budget}")]
ContextOverflow { used: u64, budget: u64 },
#[error("Session compaction failed: {0}")]
Compaction(String),
}
pub struct JsonSizeTokenMeter {
bytes_per_token: u64,
}
impl JsonSizeTokenMeter {
pub fn new(bytes_per_token: u64) -> Result<Self, SessionContextError> {
if bytes_per_token != 1 {
return Err(SessionContextError::InvalidRequest(
"the conservative JSON fallback requires exactly one byte per token".to_owned(),
));
}
Ok(Self { bytes_per_token })
}
}
impl Default for JsonSizeTokenMeter {
fn default() -> Self {
Self { bytes_per_token: 1 }
}
}
impl ModelTokenMeter for JsonSizeTokenMeter {
fn meter_descriptor(&self) -> ModelTokenMeterDescriptor {
ModelTokenMeterDescriptor {
strategy: "canonical-json-size".to_owned(),
version: "1".to_owned(),
accounting: ModelTokenAccounting::ConservativeUpperBound,
config_digest: Digest::sha256(format!(
"canonical-json-size/v1\0bytes_per_token={}",
self.bytes_per_token
)),
}
}
fn count_request_input(
&self,
messages: &[ModelMessage],
tools: &[ModelToolDefinition],
) -> Result<u64, ModelError> {
let bytes = serde_jcs::to_vec(&(messages, tools)).map_err(|error| {
ModelError::invalid_request(format!(
"could not serialize model context for metering: {error}"
))
})?;
Ok((bytes.len() as u64).div_ceil(self.bytes_per_token))
}
}
pub struct SessionContextRequest {
pub session_id: AgentSessionId,
pub current_run_id: RunId,
pub through_session_seq: Option<u64>,
pub system_message: Option<ModelMessage>,
pub tools: Vec<ModelToolDefinition>,
pub history_limit: usize,
pub max_context_tokens: u64,
pub reserved_output_tokens: u64,
pub config_digest: Digest,
pub allowed_skill_digests: BTreeMap<SkillId, Digest>,
}
pub struct SessionContextProjection {
pub messages: Vec<ModelMessage>,
pub included_ranges: Vec<SessionSourceRange>,
pub deferred_ranges: Vec<SessionSourceRange>,
pub used_input_tokens: u64,
pub context_estimate: Option<ModelContextEstimate>,
pub planning: Option<observed_prefix::ContextPlanningTrace>,
pub input_budget_tokens: u64,
pub through_session_seq: u64,
pub config_digest: Digest,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ContextTokenPolicy {
UpperBound,
Planning,
}
pub struct AgentSessionContextEngine {
journal: Arc<dyn AgentSessionJournalStore>,
token_meter: Arc<dyn ModelTokenMeter>,
}
impl AgentSessionContextEngine {
pub fn new(
journal: Arc<dyn AgentSessionJournalStore>,
token_meter: Arc<dyn ModelTokenMeter>,
) -> Self {
Self {
journal,
token_meter,
}
}
pub async fn project(
&self,
request: SessionContextRequest,
) -> Result<SessionContextProjection, SessionContextError> {
self.project_with_policy(request, ContextTokenPolicy::UpperBound)
.await
}
pub async fn project_with_policy(
&self,
request: SessionContextRequest,
policy: ContextTokenPolicy,
) -> Result<SessionContextProjection, SessionContextError> {
validate_context_request(&request)?;
let records = self.journal.load_session(&request.session_id).await?;
validate_session_trace(&request.session_id, &records)?;
let records = match request.through_session_seq {
Some(through) if through > records.len() as u64 => {
return Err(SessionContextError::InvalidRequest(format!(
"Session Context cursor {through} is past Journal head {}",
records.len()
)))
}
Some(through) => &records[..through as usize],
None => records.as_slice(),
};
let mut groups = replay_groups(
records,
&request.current_run_id,
&request.allowed_skill_digests,
)?;
let input_budget = request
.max_context_tokens
.saturating_sub(request.reserved_output_tokens);
let mut selected = groups
.values()
.filter(|group| group.pinned)
.map(|group| group.key)
.collect::<BTreeSet<_>>();
let pinned_messages = assemble_messages(&request.system_message, &groups, &selected);
let pinned_tokens = self.context_input_tokens(&pinned_messages, &request.tools, policy)?;
if pinned_tokens > input_budget {
return Err(SessionContextError::ContextOverflow {
used: pinned_tokens,
budget: input_budget,
});
}
let mut retained_anchor = None;
if request.history_limit > 0 {
if let Some(record) = records
.iter()
.rev()
.find(|record| {
record.run_id != request.current_run_id
&& matches!(record.payload, AgentSessionEvent::RunInputCommitted { .. })
})
.filter(|record| !groups.contains_key(&record.session_seq))
{
let AgentSessionEvent::RunInputCommitted { message } = &record.payload else {
unreachable!()
};
let mut candidate_groups = groups.clone();
candidate_groups
.entry(record.session_seq)
.or_insert_with(|| MessageGroup {
key: record.session_seq,
producer_seq: record.session_seq,
source_ranges: vec![single_range(record.session_seq)],
logical_source: single_range(record.session_seq),
messages: vec![message.clone()],
pinned: false,
active_compactable: false,
});
let mut candidate = selected.clone();
candidate.insert(record.session_seq);
if self.context_input_tokens(
&assemble_messages(&request.system_message, &candidate_groups, &candidate),
&request.tools,
policy,
)? <= input_budget
{
groups = candidate_groups;
selected = candidate;
retained_anchor = Some(record.session_seq);
}
}
}
for group in groups
.values()
.rev()
.filter(|group| !group.pinned && Some(group.key) != retained_anchor)
.take(
request
.history_limit
.saturating_sub(usize::from(retained_anchor.is_some())),
)
{
let mut candidate = selected.clone();
candidate.insert(group.key);
let messages = assemble_messages(&request.system_message, &groups, &candidate);
if self.context_input_tokens(&messages, &request.tools, policy)? <= input_budget {
selected = candidate;
}
}
let messages = assemble_messages(&request.system_message, &groups, &selected);
let used_input_tokens = self.count_request_input(&messages, &request.tools)?;
let context_estimate = if policy == ContextTokenPolicy::Planning {
let estimate = self
.token_meter
.estimate_context_input(&messages, &request.tools)
.map_err(|error| SessionContextError::InvalidRequest(error.to_string()))?;
if estimate.tokens > input_budget
|| estimate.tokens > used_input_tokens
|| (estimate.accounting != ModelTokenAccounting::Estimated
&& estimate.tokens != used_input_tokens)
{
return Err(SessionContextError::InvalidRequest(
"context estimate is inconsistent with the input bound or selected budget"
.to_owned(),
));
}
(estimate.accounting == ModelTokenAccounting::Estimated).then_some(estimate)
} else {
None
};
let mut included_ranges = Vec::new();
let mut deferred_ranges = Vec::new();
for group in groups.values() {
let destination = if selected.contains(&group.key) {
&mut included_ranges
} else {
&mut deferred_ranges
};
for source in &group.source_ranges {
if let Some(anchor) = retained_anchor
.filter(|anchor| group.key != *anchor && source.contains(*anchor))
{
if source.first_session_seq < anchor {
destination.push(SessionSourceRange {
first_session_seq: source.first_session_seq,
last_session_seq: anchor - 1,
});
}
if source.last_session_seq > anchor {
destination.push(SessionSourceRange {
first_session_seq: anchor + 1,
last_session_seq: source.last_session_seq,
});
}
} else {
destination.push(source.clone());
}
}
}
Ok(SessionContextProjection {
messages,
included_ranges,
deferred_ranges,
used_input_tokens,
context_estimate,
planning: None,
input_budget_tokens: input_budget,
through_session_seq: records.last().map(|record| record.session_seq).unwrap_or(0),
config_digest: request.config_digest,
})
}
fn context_input_tokens(
&self,
messages: &[ModelMessage],
tools: &[ModelToolDefinition],
policy: ContextTokenPolicy,
) -> Result<u64, SessionContextError> {
if policy == ContextTokenPolicy::UpperBound {
return self.count_request_input(messages, tools);
}
self.token_meter
.estimate_context_input(messages, tools)
.map(|estimate| estimate.tokens)
.map_err(|error| SessionContextError::InvalidRequest(error.to_string()))
}
fn count_request_input(
&self,
messages: &[ModelMessage],
tools: &[ModelToolDefinition],
) -> Result<u64, SessionContextError> {
self.token_meter
.count_request_input(messages, tools)
.map_err(|error| {
SessionContextError::InvalidRequest(format!(
"model token meter rejected Context input: {error}"
))
})
}
}
#[derive(Clone)]
struct MessageGroup {
key: u64,
producer_seq: u64,
source_ranges: Vec<SessionSourceRange>,
logical_source: SessionSourceRange,
messages: Vec<ModelMessage>,
pinned: bool,
active_compactable: bool,
}
fn replay_groups(
records: &[AgentSessionRecord],
current_run_id: &RunId,
allowed_skill_digests: &BTreeMap<SkillId, Digest>,
) -> Result<BTreeMap<u64, MessageGroup>, SessionContextError> {
let mut groups = BTreeMap::new();
let mut loaded_skills = BTreeMap::new();
for record in records {
match &record.payload {
AgentSessionEvent::RunInputCommitted { message } => {
groups.insert(
record.session_seq,
MessageGroup {
key: record.session_seq,
producer_seq: record.session_seq,
source_ranges: vec![single_range(record.session_seq)],
logical_source: single_range(record.session_seq),
messages: vec![message.clone()],
pinned: record.run_id == *current_run_id,
active_compactable: false,
},
);
}
AgentSessionEvent::ToolExchangeCommitted {
assistant,
tool,
retained_artifacts,
..
} => {
groups.insert(
record.session_seq,
MessageGroup {
key: record.session_seq,
producer_seq: record.session_seq,
source_ranges: vec![single_range(record.session_seq)],
logical_source: single_range(record.session_seq),
messages: vec![assistant.clone(), tool.clone()],
pinned: record.run_id == *current_run_id || !retained_artifacts.is_empty(),
active_compactable: retained_artifacts.is_empty(),
},
);
}
AgentSessionEvent::EffectUncertaintyCommitted {
effect_call_id,
model_call_id,
tool_name,
message,
} => {
groups.insert(
record.session_seq,
MessageGroup {
key: record.session_seq,
producer_seq: record.session_seq,
source_ranges: vec![single_range(record.session_seq)],
logical_source: single_range(record.session_seq),
messages: vec![effect_uncertainty_message(
effect_call_id,
model_call_id,
tool_name,
message,
)],
pinned: true,
active_compactable: false,
},
);
}
AgentSessionEvent::RunOutputCommitted { message, .. } => {
groups.insert(
record.session_seq,
MessageGroup {
key: record.session_seq,
producer_seq: record.session_seq,
source_ranges: vec![single_range(record.session_seq)],
logical_source: single_range(record.session_seq),
messages: vec![message.clone()],
pinned: record.run_id == *current_run_id,
active_compactable: false,
},
);
}
AgentSessionEvent::SkillLoaded { load } => {
if record.run_id != *current_run_id {
continue;
}
let descriptor = &load.package.descriptor;
if allowed_skill_digests.get(&descriptor.skill_id) != Some(&descriptor.digest) {
continue;
}
if let Some(previous) =
loaded_skills.insert(descriptor.skill_id.clone(), descriptor.digest.clone())
{
if previous != descriptor.digest {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
format!(
"Skill '{}' changed digest inside one Run without an explicit replacement protocol",
descriptor.skill_id
),
)));
}
continue;
}
groups.insert(
record.session_seq,
MessageGroup {
key: record.session_seq,
producer_seq: record.session_seq,
source_ranges: vec![single_range(record.session_seq)],
logical_source: single_range(record.session_seq),
messages: vec![skill_load_message(load)],
pinned: true,
active_compactable: false,
},
);
}
AgentSessionEvent::CompactionCommitted {
source,
source_digest,
summary,
..
} => {
let observed = session_range_digest(records, source)?;
if observed != *source_digest {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
format!(
"compaction source digest mismatch at session_seq {}",
record.session_seq
),
)));
}
if records.iter().any(|candidate| {
source.contains(candidate.session_seq)
&& matches!(candidate.payload, AgentSessionEvent::SkillLoaded { .. })
}) {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
"compaction cannot shadow durable loaded Skill state".to_owned(),
)));
}
let source_groups = groups
.values()
.filter(|group| source.contains(group.producer_seq))
.collect::<Vec<_>>();
let source_record_count = records
.iter()
.filter(|candidate| source.contains(candidate.session_seq))
.count();
if source_groups.len() != source_record_count {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
"compaction source contains an already-shadowed or unprojected record"
.to_owned(),
)));
}
if source_groups.iter().any(|group| group.pinned) {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
"compaction attempted to shadow the current Run".to_owned(),
)));
}
if summary.role == ModelRole::Assistant
&& !placement::can_replace_source(&groups, records, source)
{
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
"summary placement crosses a surviving Context group".to_owned(),
)));
}
let summary_group = placement::summary_group(
&groups,
source,
record.session_seq,
summary,
false,
false,
)
.map_err(|error| {
SessionContextError::Journal(AgentSessionError::Corrupt(error.to_string()))
})?;
groups.retain(|_, group| !source.contains(group.producer_seq));
groups.insert(record.session_seq, summary_group);
}
AgentSessionEvent::ActiveRunCompactionCommitted {
source,
source_digest,
summary,
..
} => {
let observed = session_range_digest(records, source)?;
if observed != *source_digest {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
format!(
"active-Run compaction source digest mismatch at session_seq {}",
record.session_seq
),
)));
}
if !live_source::valid_active_source(&groups, records, source, &record.run_id) {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
"active-Run compaction may shadow only live, complete Tool exchanges from one Run"
.to_owned(),
)));
}
if (summary.role == ModelRole::Assistant
|| live_source::has_shadowed_records(&groups, source))
&& !placement::can_replace_source(&groups, records, source)
{
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
"summary placement crosses a surviving Context group".to_owned(),
)));
}
let summary_group = placement::summary_group(
&groups,
source,
record.session_seq,
summary,
record.run_id == *current_run_id,
true,
)
.map_err(|error| {
SessionContextError::Journal(AgentSessionError::Corrupt(error.to_string()))
})?;
groups.retain(|_, group| !source.contains(group.producer_seq));
groups.insert(record.session_seq, summary_group);
}
_ => {
return Err(SessionContextError::Journal(AgentSessionError::Corrupt(
"unsupported Session event cannot be projected safely".to_owned(),
)))
}
}
}
Ok(groups)
}
pub(crate) fn skill_load_message(
load: &orchestral_core::skill_protocol::SkillLoad,
) -> ModelMessage {
let descriptor = &load.package.descriptor;
let version = descriptor.version.as_deref().unwrap_or("unversioned");
let resource_base = skill_resource_base(&descriptor.source)
.map(|base| {
format!(
"\nresource_base: {base}\nRelative paths in this Skill's instructions are relative to resource_base."
)
})
.unwrap_or_default();
ModelMessage::text(
ModelRole::System,
format!(
"Loaded Skill (immutable instruction context)\nname: {}\nskill_id: {}\nsource: {:?}:{}{}\nversion: {}\ndigest: {}\n\nInstructions:\n{}",
descriptor.name,
descriptor.skill_id,
descriptor.source.kind,
descriptor.source.locator,
resource_base,
version,
descriptor.digest,
load.package.instructions
),
)
}
pub(crate) fn skill_resource_base(
source: &orchestral_core::skill_protocol::SkillSource,
) -> Option<String> {
let locator = std::path::Path::new(&source.locator);
if !locator.is_absolute() {
return None;
}
locator
.parent()
.map(|parent| parent.to_string_lossy().into_owned())
}
fn effect_uncertainty_message(
effect_call_id: &orchestral_core::tool_protocol::ToolCallId,
model_call_id: &orchestral_core::model_protocol::ModelToolCallId,
tool_name: &str,
message: &str,
) -> ModelMessage {
ModelMessage::text(
ModelRole::System,
format!(
"HOST SAFETY FACT: unresolved Tool effect; never retry automatically.\neffect_call_id: {}\nmodel_call_id: {}\ntool: {}\nobservation: {}\nresolution: explicit Host reconciliation is required",
effect_call_id, model_call_id, tool_name, message
),
)
}
fn assemble_messages(
system_message: &Option<ModelMessage>,
groups: &BTreeMap<u64, MessageGroup>,
selected: &BTreeSet<u64>,
) -> Vec<ModelMessage> {
let mut messages = Vec::new();
if let Some(system) = system_message {
messages.push(system.clone());
}
for group in groups.values() {
if selected.contains(&group.key)
&& group
.messages
.iter()
.all(|message| message.role == ModelRole::System)
{
messages.extend(group.messages.clone());
}
}
let mut conversation = groups.values().collect::<Vec<_>>();
conversation.sort_by_key(|group| (group.logical_source.first_session_seq, group.producer_seq));
for group in conversation {
if selected.contains(&group.key)
&& !group
.messages
.iter()
.all(|message| message.role == ModelRole::System)
{
messages.extend(group.messages.clone());
}
}
messages
}
fn validate_context_request(request: &SessionContextRequest) -> Result<(), SessionContextError> {
if request.session_id.is_empty()
|| request.current_run_id.is_empty()
|| request.history_limit == 0
|| request.max_context_tokens == 0
|| request.reserved_output_tokens >= request.max_context_tokens
|| !request.config_digest.is_sha256()
{
return Err(SessionContextError::InvalidRequest(
"Session/context identities, digest, and token budget are invalid".to_owned(),
));
}
if let Some(system) = &request.system_message {
system.validate().map_err(|error| {
SessionContextError::InvalidRequest(format!("invalid system message: {error}"))
})?;
if system.role != ModelRole::System {
return Err(SessionContextError::InvalidRequest(
"configured system message must have the System role".to_owned(),
));
}
}
for tool in &request.tools {
tool.validate().map_err(|error| {
SessionContextError::InvalidRequest(format!("invalid Tool schema: {error}"))
})?;
}
Ok(())
}
fn single_range(sequence: u64) -> SessionSourceRange {
SessionSourceRange {
first_session_seq: sequence,
last_session_seq: sequence,
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SessionCompactionPolicy {
pub minimum_source_records: usize,
pub keep_recent_records: usize,
}
impl Default for SessionCompactionPolicy {
fn default() -> Self {
Self {
minimum_source_records: 8,
keep_recent_records: 8,
}
}
}
impl SessionCompactionPolicy {
pub fn validate(&self) -> Result<(), SessionContextError> {
if self.minimum_source_records == 0
|| self.keep_recent_records == 0
|| self
.minimum_source_records
.checked_add(self.keep_recent_records)
.is_none()
{
return Err(SessionContextError::InvalidRequest(
"Session compaction limits must be positive and bounded".to_owned(),
));
}
Ok(())
}
pub fn digest(&self) -> Result<Digest, SessionContextError> {
self.validate()?;
serde_jcs::to_vec(self)
.map(Digest::sha256)
.map_err(|error| {
SessionContextError::InvalidRequest(format!(
"could not digest Session compaction policy: {error}"
))
})
}
}
pub fn select_compaction_source(
records: &[AgentSessionRecord],
current_run_id: &RunId,
policy: &SessionCompactionPolicy,
) -> Option<SessionSourceRange> {
policy.validate().ok()?;
let minimum_records = policy
.minimum_source_records
.checked_add(policy.keep_recent_records)?;
if records.len() < minimum_records {
return None;
}
if let Some(previous_compaction) = records.iter().rev().find(|record| {
matches!(
record.payload,
AgentSessionEvent::CompactionCommitted { .. }
| AgentSessionEvent::ActiveRunCompactionCommitted { .. }
)
}) {
let records_since_compaction = records
.len()
.saturating_sub(previous_compaction.session_seq as usize);
if records_since_compaction < minimum_records {
return None;
}
}
let source_len = records.len().saturating_sub(policy.keep_recent_records);
let shadowed = records
.iter()
.filter_map(|record| match &record.payload {
AgentSessionEvent::CompactionCommitted { source, .. }
| AgentSessionEvent::ActiveRunCompactionCommitted { source, .. } => Some(source),
_ => None,
})
.collect::<Vec<_>>();
let mut segment_start = None;
let mut segment_end = 0;
for record in &records[..source_len] {
let is_shadowed = shadowed
.iter()
.any(|source| source.contains(record.session_seq));
let is_barrier = record.run_id == *current_run_id
|| match &record.payload {
AgentSessionEvent::SkillLoaded { .. }
| AgentSessionEvent::EffectUncertaintyCommitted { .. } => true,
AgentSessionEvent::ToolExchangeCommitted {
retained_artifacts, ..
} => !retained_artifacts.is_empty(),
_ => false,
};
if is_shadowed || is_barrier {
if let Some(start) = segment_start {
if segment_end - start + 1 >= policy.minimum_source_records as u64 {
return Some(SessionSourceRange {
first_session_seq: start,
last_session_seq: segment_end,
});
}
}
segment_start = None;
continue;
}
segment_start.get_or_insert(record.session_seq);
segment_end = record.session_seq;
}
segment_start.and_then(|start| {
if segment_end - start + 1 >= policy.minimum_source_records as u64 {
return Some(SessionSourceRange {
first_session_seq: start,
last_session_seq: segment_end,
});
}
None
})
}
pub fn select_active_run_compaction_source(
records: &[AgentSessionRecord],
current_run_id: &RunId,
policy: &SessionCompactionPolicy,
) -> Option<SessionSourceRange> {
policy.validate().ok()?;
let mut live = BTreeSet::new();
for record in records {
match &record.payload {
AgentSessionEvent::ToolExchangeCommitted {
retained_artifacts, ..
} if record.run_id == *current_run_id && retained_artifacts.is_empty() => {
live.insert(record.session_seq);
}
AgentSessionEvent::ActiveRunCompactionCommitted { source, .. } => {
live.retain(|sequence| !source.contains(*sequence));
if record.run_id == *current_run_id {
live.insert(record.session_seq);
}
}
AgentSessionEvent::CompactionCommitted { source, .. } => {
live.retain(|sequence| !source.contains(*sequence));
}
_ => {}
}
}
if let Some(summary_seq) = records.iter().rev().find_map(|record| {
(live.contains(&record.session_seq)
&& record.run_id == *current_run_id
&& matches!(
record.payload,
AgentSessionEvent::ActiveRunCompactionCommitted { .. }
))
.then_some(record.session_seq)
}) {
let summary_index = (summary_seq - 1) as usize;
let mut first_index = summary_index;
while first_index > 0 && live.contains(&records[first_index - 1].session_seq) {
first_index -= 1;
}
let mut last_index = summary_index;
while last_index + 1 < records.len() && live.contains(&records[last_index + 1].session_seq)
{
last_index += 1;
}
let source_last_index = if last_index > summary_index
&& matches!(
records[last_index].payload,
AgentSessionEvent::ToolExchangeCommitted { .. }
) {
last_index - 1
} else {
last_index
};
if source_last_index >= first_index
&& records[first_index..=source_last_index]
.iter()
.any(|record| {
matches!(
record.payload,
AgentSessionEvent::ToolExchangeCommitted { .. }
)
})
{
return Some(SessionSourceRange {
first_session_seq: records[first_index].session_seq,
last_session_seq: records[source_last_index].session_seq,
});
}
}
let live_exchange_count = records
.iter()
.filter(|record| {
live.contains(&record.session_seq)
&& matches!(
record.payload,
AgentSessionEvent::ToolExchangeCommitted { .. }
)
})
.count();
let retain_count = if live_exchange_count > policy.keep_recent_records {
policy.keep_recent_records
} else if live_exchange_count > 1 {
1
} else {
0
};
let retained_recent = records
.iter()
.rev()
.filter(|record| {
live.contains(&record.session_seq)
&& matches!(
record.payload,
AgentSessionEvent::ToolExchangeCommitted { .. }
)
})
.take(retain_count)
.map(|record| record.session_seq)
.collect::<BTreeSet<_>>();
let candidates = live
.difference(&retained_recent)
.copied()
.collect::<BTreeSet<_>>();
let mut start = None;
let mut end = 0;
let mut contains_exchange = false;
for record in records {
if !candidates.contains(&record.session_seq) {
if let Some(first) = start {
if contains_exchange {
return Some(SessionSourceRange {
first_session_seq: first,
last_session_seq: end,
});
}
}
start = None;
contains_exchange = false;
continue;
}
start.get_or_insert(record.session_seq);
end = record.session_seq;
contains_exchange |= matches!(
record.payload,
AgentSessionEvent::ToolExchangeCommitted { .. }
);
}
let contiguous = start.and_then(|first| {
contains_exchange.then_some(SessionSourceRange {
first_session_seq: first,
last_session_seq: end,
})
});
contiguous.or_else(|| {
let groups = replay_groups(records, current_run_id, &BTreeMap::new()).ok()?;
live_source::compactable_segments(&groups, records, current_run_id)
.into_iter()
.find(|source| {
groups
.range(source.first_session_seq..=source.last_session_seq)
.count()
> 1
})
})
}
#[derive(Debug, Clone, PartialEq)]
pub struct SessionCompactionGroup {
pub source: SessionSourceRange,
pub messages: Vec<ModelMessage>,
}
pub struct SessionCompactionInput {
pub session_id: AgentSessionId,
pub source: SessionSourceRange,
pub groups: Vec<SessionCompactionGroup>,
pub focus_messages: Vec<ModelMessage>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SessionSummarizerDescriptor {
pub strategy: String,
pub model: Option<String>,
pub version: String,
pub config_digest: Digest,
}
impl SessionSummarizerDescriptor {
pub fn validate(&self) -> Result<(), SessionContextError> {
if self.strategy.trim().is_empty()
|| self.version.trim().is_empty()
|| self
.model
.as_ref()
.is_some_and(|model| model.trim().is_empty())
|| !self.config_digest.is_sha256()
{
return Err(SessionContextError::InvalidRequest(
"Session summarizer descriptor is invalid".to_owned(),
));
}
Ok(())
}
}
pub struct DeterministicExtractiveSessionSummarizer {
max_summary_chars: usize,
descriptor: SessionSummarizerDescriptor,
}
impl DeterministicExtractiveSessionSummarizer {
pub fn new(max_summary_chars: usize) -> Result<Self, SessionContextError> {
if max_summary_chars < 256 {
return Err(SessionContextError::InvalidRequest(
"deterministic Session summaries require at least 256 characters".to_owned(),
));
}
let config = serde_json::json!({
"contract": "deterministic-extractive-session-summary/v7",
"max_summary_chars": max_summary_chars,
});
let bytes = serde_jcs::to_vec(&config).map_err(|error| {
SessionContextError::InvalidRequest(format!(
"could not digest deterministic Session summarizer config: {error}"
))
})?;
Ok(Self {
max_summary_chars,
descriptor: SessionSummarizerDescriptor {
strategy: "deterministic-extractive".to_owned(),
model: None,
version: "7".to_owned(),
config_digest: Digest::sha256(bytes),
},
})
}
pub fn max_summary_chars(&self) -> usize {
self.max_summary_chars
}
}
struct ExtractiveCandidate {
index: usize,
rendered: String,
terms: BTreeSet<String>,
score: u64,
failed: bool,
latest_tool_exchange: bool,
tool_observation: Option<(String, String)>,
continuation_observation: Option<String>,
}
#[async_trait]
impl AgentSessionSummarizer for DeterministicExtractiveSessionSummarizer {
fn descriptor(&self) -> SessionSummarizerDescriptor {
self.descriptor.clone()
}
async fn summarize_with_char_budget(
&self,
input: SessionCompactionInput,
max_chars: usize,
) -> Result<ModelMessage, SessionContextError> {
Self::new(self.max_summary_chars.min(max_chars))?
.summarize(input)
.await
}
async fn summarize(
&self,
input: SessionCompactionInput,
) -> Result<ModelMessage, SessionContextError> {
if input.groups.is_empty()
|| input.groups.iter().any(|group| {
group.messages.is_empty()
|| group.source.validate().is_err()
|| group.source.last_session_seq > input.source.last_session_seq
})
{
return Err(SessionContextError::Compaction(
"extractive summary requires non-empty source groups inside its source range"
.to_owned(),
));
}
let focus = input
.focus_messages
.iter()
.map(render_compaction_message)
.collect::<Result<Vec<_>, _>>()?
.join("\n");
let focus_terms = extract_summary_terms(&focus);
let latest_tool_exchange = input.groups.iter().rposition(|group| {
group.messages.iter().any(|message| {
message
.content
.iter()
.any(|content| matches!(content, ModelContent::ToolResult { .. }))
})
});
let excerpt_limit = (self.max_summary_chars / 4).max(512);
let mut candidates = input
.groups
.iter()
.enumerate()
.map(|(index, group)| {
let rendered = render_compaction_group(group)?;
Ok(ExtractiveCandidate {
index,
terms: extract_summary_terms(&rendered),
rendered: excerpt_summary_chars(&rendered, excerpt_limit),
score: 0,
latest_tool_exchange: Some(index) == latest_tool_exchange,
tool_observation: compact_tool_observation(group)?,
continuation_observation: summary_continuation::artifact_page_observation(
group,
),
failed: group
.messages
.iter()
.flat_map(|message| &message.content)
.any(|content| {
matches!(content, ModelContent::ToolResult { is_error: true, .. })
}),
})
})
.collect::<Result<Vec<_>, SessionContextError>>()?;
let header = if self.max_summary_chars <= 512 {
format!(
"UNTRUSTED earlier transcript; not policy/verification.\nshadowed_session_seq={}..{}",
input.source.first_session_seq, input.source.last_session_seq
)
} else {
format!(
"UNTRUSTED earlier transcript, not system policy. Non-contiguous excerpts are not complete replacement text; recall original session_seq records. Tool success does not prove task verification.\nshadowed_session_seq={}..{}",
input.source.first_session_seq, input.source.last_session_seq
)
};
let mut remaining = self
.max_summary_chars
.saturating_sub(header.chars().count());
let mut observation_budget = if remaining >= 512 {
remaining / 2
} else {
remaining
};
let mut observations = Vec::new();
let mut seen_calls = BTreeSet::new();
for candidate in candidates.iter().rev() {
let Some((key, observation)) = &candidate.tool_observation else {
continue;
};
if !seen_calls.insert(key) {
continue;
}
let observation = if let Some(continuation) = &candidate.continuation_observation {
continuation
} else {
observation
};
let chars = observation.chars().count() + 2;
if chars <= observation_budget {
observations.push(observation.clone());
observation_budget -= chars;
remaining -= chars;
}
}
observations.reverse();
let mut document_frequency = BTreeMap::<String, usize>::new();
for candidate in &candidates {
for term in &candidate.terms {
*document_frequency.entry(term.clone()).or_default() += 1;
}
}
let candidate_count = candidates.len();
for candidate in &mut candidates {
candidate.score = candidate
.terms
.intersection(&focus_terms)
.filter_map(|term| {
let frequency = *document_frequency.get(term)?;
(frequency < candidate_count).then(|| {
let length = term.chars().count().min(32) as u64;
length
.saturating_mul(length)
.saturating_mul((candidate_count + 1 - frequency) as u64)
})
})
.sum();
}
let has_relevant = candidates.iter().any(|candidate| candidate.score > 0);
if has_relevant {
candidates.retain(|candidate| {
candidate.score > 0 || candidate.failed || candidate.latest_tool_exchange
});
}
candidates.sort_by(|left, right| {
right
.latest_tool_exchange
.cmp(&left.latest_tool_exchange)
.then_with(|| right.failed.cmp(&left.failed))
.then_with(|| right.score.cmp(&left.score))
.then_with(|| right.index.cmp(&left.index))
});
let mut selected = Vec::<(usize, String)>::new();
for candidate in &candidates {
let separator_chars = 2;
if remaining <= separator_chars {
break;
}
let rendered_chars = candidate.rendered.chars().count();
if rendered_chars + separator_chars <= remaining {
selected.push((candidate.index, candidate.rendered.clone()));
remaining -= rendered_chars + separator_chars;
} else {
let available = remaining.saturating_sub(separator_chars);
if available >= 128 {
selected.push((
candidate.index,
excerpt_summary_chars(&candidate.rendered, available),
));
remaining -= available + separator_chars;
}
}
}
if selected.is_empty() {
if let Some(candidate) = candidates.first() {
let available = remaining.saturating_sub(2);
if available > 0 {
selected.push((
candidate.index,
excerpt_summary_chars(&candidate.rendered, available),
));
}
}
}
selected.sort_by_key(|(index, _)| *index);
let mut summary = header;
for observation in observations {
summary.push_str("\n\n");
summary.push_str(&observation);
}
for (_, rendered) in selected {
summary.push_str("\n\n");
summary.push_str(&rendered);
}
debug_assert!(summary.chars().count() <= self.max_summary_chars);
Ok(ModelMessage::text(ModelRole::Assistant, summary))
}
}
fn compact_tool_observation(
group: &SessionCompactionGroup,
) -> Result<Option<(String, String)>, SessionContextError> {
let mut calls = Vec::new();
let mut results = Vec::new();
for message in &group.messages {
for content in &message.content {
match content {
ModelContent::ToolCall {
name, arguments, ..
} => calls.push((name, canonical_summary_json(arguments)?)),
ModelContent::ToolResult {
result, is_error, ..
} => results.push((result, is_error)),
_ => {}
}
}
}
if results.is_empty() {
return Ok(None);
}
let key = if calls.is_empty() {
format!("source:{}", group.source.first_session_seq)
} else {
canonical_summary_json(&serde_json::json!(calls))?
};
let mut observation = format!(
"[recent tool observation session_seq={}..{}; latest identical call]",
group.source.first_session_seq, group.source.last_session_seq
);
for (name, arguments) in calls {
observation.push_str(&format!(
"\n{name} {}",
excerpt_summary_chars(&arguments, 128)
));
}
for (result, is_error) in results {
observation.push_str(&format!(
"\nrecorded_tool_outcome status={}",
if *is_error { "failed" } else { "succeeded" }
));
if let Some(fields) = result.as_object() {
let scalars = fields
.iter()
.filter(|(_, value)| value.is_null() || value.is_boolean() || value.is_number())
.map(|(key, value)| (key.clone(), value.clone()))
.collect::<serde_json::Map<_, _>>();
observation.push_str(" scalar_fields=");
observation.push_str(&excerpt_summary_chars(
&canonical_summary_json(&serde_json::Value::Object(scalars))?,
160,
));
let mut seen_text = BTreeSet::new();
for (key, text) in literal_result_fields(fields) {
if !text.is_empty() && seen_text.insert(text) {
observation.push_str(&format!(
"\n{key} text_excerpt={}",
excerpt_summary_chars(text, 256)
));
}
}
} else {
observation.push_str(" result_excerpt=");
observation.push_str(&excerpt_summary_chars(
&render_compaction_result(result)?,
256,
));
}
}
Ok(Some((key, excerpt_summary_chars(&observation, 512))))
}
fn render_compaction_group(group: &SessionCompactionGroup) -> Result<String, SessionContextError> {
let mut rendered = format!(
"[historical group session_seq={}..{}]",
group.source.first_session_seq, group.source.last_session_seq
);
for message in &group.messages {
for content in &message.content {
if let ModelContent::ToolResult {
call_id,
is_error,
result,
} = content
{
rendered.push_str(&format!(
"\nrecorded_tool_outcome call_id={} status={}",
call_id,
if *is_error { "failed" } else { "succeeded" }
));
if let Some(fields) = result.as_object() {
let scalars = fields
.iter()
.filter(|(_, value)| {
!value.is_string() && !value.is_array() && !value.is_object()
})
.map(|(key, value)| (key.clone(), value.clone()))
.collect::<serde_json::Map<_, _>>();
if !scalars.is_empty() {
rendered.push_str("\nrecorded_result_fields=");
rendered.push_str(&excerpt_summary_chars(
&canonical_summary_json(&serde_json::Value::Object(scalars))?,
512,
));
}
let short_text = fields
.iter()
.filter(|(_, value)| {
value
.as_str()
.is_some_and(|text| text.chars().count() <= 128)
})
.map(|(key, value)| (key.clone(), value.clone()))
.collect::<serde_json::Map<_, _>>();
if !short_text.is_empty() {
rendered.push_str("\nrecorded_result_text=");
rendered.push_str(&excerpt_summary_chars(
&render_compaction_result(&serde_json::Value::Object(short_text))?,
512,
));
}
}
if *is_error {
rendered.push_str("\nerror_result=");
rendered.push_str(&truncate_summary_chars(
&render_compaction_result(result)?,
256,
));
}
}
}
}
for message in &group.messages {
for content in &message.content {
if let ModelContent::ToolCall {
call_id,
name,
arguments,
..
} = content
{
rendered.push_str(&format!(
"\nrecorded_tool_call call_id={call_id} name={name} arguments={}",
excerpt_summary_chars(&canonical_summary_json(arguments)?, 256)
));
}
}
}
for message in &group.messages {
rendered.push('\n');
rendered.push_str(&render_compaction_message(message)?);
}
Ok(rendered)
}
fn render_compaction_message(message: &ModelMessage) -> Result<String, SessionContextError> {
let role = match message.role {
ModelRole::System => "system",
ModelRole::User => "user",
ModelRole::Assistant => "assistant",
ModelRole::Tool => "tool",
_ => {
return Err(SessionContextError::Compaction(
"extractive summary does not support this future model role".to_owned(),
))
}
};
let mut rendered = format!("{role}: ");
for (index, content) in message
.content
.iter()
.filter(|content| !matches!(content, ModelContent::Continuation { .. }))
.enumerate()
{
if index > 0 {
rendered.push_str(" | ");
}
match content {
ModelContent::Text { text } => rendered.push_str(text),
ModelContent::Json { value } => {
rendered.push_str("json=");
rendered.push_str(&canonical_summary_json(value)?);
}
ModelContent::Data { media_type, value } => {
rendered.push_str("data[");
rendered.push_str(media_type);
rendered.push_str("]=");
rendered.push_str(&canonical_summary_json(value)?);
}
ModelContent::ToolCall {
call_id,
name,
arguments,
..
} => {
rendered.push_str("tool_call id=");
rendered.push_str(call_id.as_str());
rendered.push_str(" name=");
rendered.push_str(name);
rendered.push_str(" arguments=");
rendered.push_str(&canonical_summary_json(arguments)?);
}
ModelContent::ToolResult {
call_id,
result,
is_error,
} => {
rendered.push_str("tool_result id=");
rendered.push_str(call_id.as_str());
rendered.push_str(" error=");
rendered.push_str(if *is_error { "true" } else { "false" });
rendered.push_str(" result=");
rendered.push_str(&render_compaction_result(result)?);
}
_ => {
return Err(SessionContextError::Compaction(
"extractive summary does not support this future model content".to_owned(),
))
}
}
}
Ok(rendered)
}
fn render_compaction_result(value: &serde_json::Value) -> Result<String, SessionContextError> {
match value {
serde_json::Value::String(text) => Ok(format!("literal text excerpt:\n{text}")),
serde_json::Value::Object(fields) => {
let metadata = fields
.iter()
.filter(|(_, value)| !value.is_string())
.map(|(key, value)| (key.clone(), value.clone()))
.collect::<serde_json::Map<_, _>>();
let mut rendered = canonical_summary_json(&serde_json::Value::Object(metadata))?;
for (key, text) in literal_result_fields(fields) {
rendered.push_str(&format!(
"\nfield {} literal text excerpt:\n{text}",
canonical_summary_json(&serde_json::Value::String(key.to_owned()))?,
));
}
Ok(rendered)
}
_ => canonical_summary_json(value),
}
}
fn literal_result_fields(
fields: &serde_json::Map<String, serde_json::Value>,
) -> BTreeMap<&str, &str> {
fields
.iter()
.filter_map(|(key, value)| value.as_str().map(|text| (key.as_str(), text)))
.collect()
}
fn canonical_summary_json(value: &serde_json::Value) -> Result<String, SessionContextError> {
let bytes = serde_jcs::to_vec(value).map_err(|error| {
SessionContextError::Compaction(format!(
"could not render canonical summary content: {error}"
))
})?;
String::from_utf8(bytes).map_err(|error| {
SessionContextError::Compaction(format!("canonical summary was not UTF-8: {error}"))
})
}
fn extract_summary_terms(text: &str) -> BTreeSet<String> {
let mut terms = BTreeSet::new();
let mut current = String::new();
let flush = |current: &mut String, terms: &mut BTreeSet<String>| {
if current.chars().count() >= 2 {
terms.insert(std::mem::take(current));
} else {
current.clear();
}
};
for character in text.chars() {
if character.is_alphanumeric() || matches!(character, '_' | '-') {
current.extend(character.to_lowercase());
} else {
flush(&mut current, &mut terms);
}
}
flush(&mut current, &mut terms);
terms
}
fn truncate_summary_chars(value: &str, limit: usize) -> String {
if value.chars().count() <= limit {
return value.to_owned();
}
if limit == 1 {
return "…".to_owned();
}
value.chars().take(limit - 1).chain(['…']).collect()
}
fn excerpt_summary_chars(value: &str, limit: usize) -> String {
const OMITTED: &str = "\n[… excerpt omitted …]\n";
let count = value.chars().count();
if count <= limit {
return value.to_owned();
}
let marker_chars = OMITTED.chars().count();
if limit <= marker_chars {
return truncate_summary_chars(value, limit);
}
let available = limit - marker_chars;
let head = available.div_ceil(2);
let tail = available - head;
value
.chars()
.take(head)
.chain(OMITTED.chars())
.chain(value.chars().skip(count - tail))
.collect()
}
#[async_trait]
pub trait AgentSessionSummarizer: Send + Sync {
fn descriptor(&self) -> SessionSummarizerDescriptor;
async fn summarize(
&self,
input: SessionCompactionInput,
) -> Result<ModelMessage, SessionContextError>;
async fn summarize_with_char_budget(
&self,
input: SessionCompactionInput,
_max_chars: usize,
) -> Result<ModelMessage, SessionContextError> {
self.summarize(input).await
}
}
fn original_compaction_groups(
records: &[AgentSessionRecord],
source: &SessionSourceRange,
) -> Result<Vec<SessionCompactionGroup>, SessionContextError> {
let mut pending = (source.first_session_seq..=source.last_session_seq).collect::<Vec<_>>();
let mut visited = BTreeSet::new();
let mut originals = BTreeMap::new();
while let Some(sequence) = pending.pop() {
if !visited.insert(sequence) {
continue;
}
let record = records.get((sequence - 1) as usize).ok_or_else(|| {
SessionContextError::Compaction(
"original source is outside the Session journal".to_owned(),
)
})?;
let messages = match &record.payload {
AgentSessionEvent::CompactionCommitted { source, .. }
| AgentSessionEvent::ActiveRunCompactionCommitted { source, .. } => {
if source.last_session_seq >= sequence {
return Err(SessionContextError::Compaction(
"compaction references must point backward".to_owned(),
));
}
pending.extend(source.first_session_seq..=source.last_session_seq);
continue;
}
AgentSessionEvent::RunInputCommitted { message }
| AgentSessionEvent::RunOutputCommitted { message, .. } => vec![message.clone()],
AgentSessionEvent::ToolExchangeCommitted {
assistant, tool, ..
} => vec![assistant.clone(), tool.clone()],
_ => {
return Err(SessionContextError::Compaction(
"original source crosses a protected Session fact".to_owned(),
))
}
};
originals.insert(
sequence,
SessionCompactionGroup {
source: single_range(sequence),
messages,
},
);
}
Ok(originals.into_values().collect())
}
pub struct AgentSessionCompactor {
journal: Arc<dyn AgentSessionJournalStore>,
summarizer: Arc<dyn AgentSessionSummarizer>,
summarizer_descriptor: SessionSummarizerDescriptor,
policy: SessionCompactionPolicy,
}
impl AgentSessionCompactor {
pub fn new(
journal: Arc<dyn AgentSessionJournalStore>,
summarizer: Arc<dyn AgentSessionSummarizer>,
policy: SessionCompactionPolicy,
) -> Result<Self, SessionContextError> {
policy.validate()?;
let summarizer_descriptor = summarizer.descriptor();
summarizer_descriptor.validate()?;
Ok(Self {
journal,
summarizer,
summarizer_descriptor,
policy,
})
}
pub fn policy(&self) -> &SessionCompactionPolicy {
&self.policy
}
pub fn summarizer_descriptor(&self) -> &SessionSummarizerDescriptor {
&self.summarizer_descriptor
}
pub async fn compact_if_needed(
&self,
session_id: &AgentSessionId,
current_run_id: &RunId,
) -> Result<Option<AgentSessionRecord>, SessionContextError> {
let records = self.journal.load_session(session_id).await?;
validate_session_trace(session_id, &records)?;
let Some(source) = select_compaction_source(&records, current_run_id, &self.policy) else {
return Ok(None);
};
let source_digest = session_range_digest(&records, &source)?;
let policy_digest = self.policy.digest()?;
let groups = replay_groups(&records, current_run_id, &BTreeMap::new())?;
let source_groups = original_compaction_groups(&records, &source)?;
let focus_messages = groups
.values()
.filter(|group| group.pinned && !source.contains(group.producer_seq))
.flat_map(|group| group.messages.clone())
.collect();
let summary = self
.summarizer
.summarize(SessionCompactionInput {
session_id: session_id.clone(),
source: source.clone(),
groups: source_groups,
focus_messages,
})
.await?;
summary.validate().map_err(|error| {
SessionContextError::Compaction(format!("invalid summary message: {error}"))
})?;
if summary.role == ModelRole::Assistant
&& !placement::can_replace_source(&groups, &records, &source)
{
return Ok(None);
}
let event_id = AgentSessionEventId::new(format!(
"compaction-{}-{}-{}",
session_id.as_str(),
source.first_session_seq,
source.last_session_seq
));
let appended = self
.journal
.append(AgentSessionEventDraft {
event_id,
session_id: session_id.clone(),
run_id: current_run_id.clone(),
payload: AgentSessionEvent::CompactionCommitted {
source,
source_digest,
policy_digest,
summary_config_digest: self.summarizer_descriptor.config_digest.clone(),
summary,
strategy: self.summarizer_descriptor.strategy.clone(),
model: self.summarizer_descriptor.model.clone(),
version: self.summarizer_descriptor.version.clone(),
},
})
.await?;
Ok(Some(appended.record))
}
pub async fn compact_active_run_for_pressure(
&self,
session_id: &AgentSessionId,
current_run_id: &RunId,
) -> Result<Option<AgentSessionRecord>, SessionContextError> {
let records = self.journal.load_session(session_id).await?;
validate_session_trace(session_id, &records)?;
let Some(source) =
select_active_run_compaction_source(&records, current_run_id, &self.policy)
else {
return Ok(None);
};
let groups = replay_groups(&records, current_run_id, &BTreeMap::new())?;
if !live_source::valid_active_source(&groups, &records, &source, current_run_id) {
return Err(SessionContextError::Compaction(
"active-Run compaction source no longer matches live Context groups".to_owned(),
));
}
let source_groups = original_compaction_groups(&records, &source)?;
let focus_messages = groups
.values()
.filter(|group| group.pinned && !source.contains(group.producer_seq))
.flat_map(|group| group.messages.clone())
.collect();
let summary = self
.summarizer
.summarize(SessionCompactionInput {
session_id: session_id.clone(),
source: source.clone(),
groups: source_groups,
focus_messages,
})
.await?;
self.commit_active_run_summary(&records, current_run_id, source, summary)
.await
}
async fn commit_active_run_summary(
&self,
records: &[AgentSessionRecord],
current_run_id: &RunId,
source: SessionSourceRange,
summary: ModelMessage,
) -> Result<Option<AgentSessionRecord>, SessionContextError> {
let session_id = &records[0].session_id;
let groups = replay_groups(records, current_run_id, &BTreeMap::new())?;
if !live_source::valid_active_source(&groups, records, &source, current_run_id) {
return Err(SessionContextError::Compaction(
"active-Run source must contain live Tool producers from one Run".to_owned(),
));
}
if (summary.role == ModelRole::Assistant
|| live_source::has_shadowed_records(&groups, &source))
&& !placement::can_replace_source(&groups, records, &source)
{
return Ok(None);
}
let source_digest = session_range_digest(records, &source)?;
let policy_digest = self.policy.digest()?;
summary.validate().map_err(|error| {
SessionContextError::Compaction(format!("invalid active-Run summary: {error}"))
})?;
let event_id = AgentSessionEventId::new(format!(
"active-run-compaction-{}-{}-{}-{}",
session_id.as_str(),
current_run_id.as_str(),
source.first_session_seq,
source.last_session_seq
));
let appended = self
.journal
.append(AgentSessionEventDraft {
event_id,
session_id: session_id.clone(),
run_id: current_run_id.clone(),
payload: AgentSessionEvent::ActiveRunCompactionCommitted {
source,
source_digest,
policy_digest,
summary_config_digest: self.summarizer_descriptor.config_digest.clone(),
summary,
strategy: self.summarizer_descriptor.strategy.clone(),
model: self.summarizer_descriptor.model.clone(),
version: self.summarizer_descriptor.version.clone(),
},
})
.await?;
Ok(Some(appended.record))
}
}
#[cfg(test)]
mod tests {
use super::*;
use orchestral_core::agent_protocol::wire::{ArtifactRef, ArtifactRefWithDigest};
use orchestral_core::agent_session::{
AgentSessionEvent, AgentSessionEventDraft, InMemoryAgentSessionJournalStore,
};
use orchestral_core::model_protocol::{ModelContent, ModelRequestId, ModelToolCallId};
use orchestral_core::tool_protocol::ToolCallId;
use serde_json::json;
struct FixedSummarizer;
#[async_trait]
impl AgentSessionSummarizer for FixedSummarizer {
fn descriptor(&self) -> SessionSummarizerDescriptor {
SessionSummarizerDescriptor {
strategy: "fixed-test-summary".to_owned(),
model: None,
version: "1".to_owned(),
config_digest: Digest::sha256("fixed-test-summary/v1"),
}
}
async fn summarize(
&self,
input: SessionCompactionInput,
) -> Result<ModelMessage, SessionContextError> {
Ok(ModelMessage::text(
ModelRole::System,
format!(
"summary of {} messages",
input
.groups
.iter()
.map(|group| group.messages.len())
.sum::<usize>()
),
))
}
}
pub(super) async fn append_input(
store: &Arc<InMemoryAgentSessionJournalStore>,
sequence: u64,
run_id: &str,
text: String,
) {
store
.append(AgentSessionEventDraft {
event_id: AgentSessionEventId::new(format!("input-{sequence}")),
session_id: AgentSessionId::new("session-1"),
run_id: RunId::new(run_id),
payload: AgentSessionEvent::RunInputCommitted {
message: ModelMessage::text(ModelRole::User, text),
},
})
.await
.unwrap();
}
async fn append_session_payload(
store: &Arc<InMemoryAgentSessionJournalStore>,
session_id: &AgentSessionId,
run_id: &RunId,
event_id: String,
payload: AgentSessionEvent,
) {
store
.append(AgentSessionEventDraft {
event_id: AgentSessionEventId::new(event_id),
session_id: session_id.clone(),
run_id: run_id.clone(),
payload,
})
.await
.unwrap();
}
pub(super) async fn append_tool_exchange(
store: &Arc<InMemoryAgentSessionJournalStore>,
sequence: u64,
run_id: &str,
width: usize,
) {
let call_id = ModelToolCallId::new(format!("active-call-{sequence}"));
store
.append(AgentSessionEventDraft {
event_id: AgentSessionEventId::new(format!("active-exchange-{sequence}")),
session_id: AgentSessionId::new("session-1"),
run_id: RunId::new(run_id),
payload: AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new(format!("active-request-{sequence}")),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: call_id.clone(),
name: "file_read".to_owned(),
arguments: json!({"path": format!("src/file-{sequence}.rs")}),
extensions: Default::default(),
}],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id,
result: json!({"text": "x".repeat(width)}),
is_error: false,
}],
},
retained_artifacts: Vec::new(),
usage: None,
},
})
.await
.unwrap();
}
fn push_session_record(
records: &mut Vec<AgentSessionRecord>,
session_id: &AgentSessionId,
run_id: &RunId,
event_suffix: impl std::fmt::Display,
payload: AgentSessionEvent,
) {
let session_seq = records.len() as u64 + 1;
records.push(
AgentSessionRecord::seal(
AgentSessionEventDraft {
event_id: AgentSessionEventId::new(format!(
"replay-{event_suffix}-{session_seq}"
)),
session_id: session_id.clone(),
run_id: run_id.clone(),
payload,
},
session_seq,
)
.unwrap(),
);
}
#[tokio::test]
async fn newest_history_cannot_be_evicted_by_older_messages() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
for index in 1..=12 {
append_input(
&store,
index,
if index == 12 { "current" } else { "old" },
format!("message-{index}-{}", "x".repeat(60)),
)
.await;
}
let engine =
AgentSessionContextEngine::new(store, Arc::new(JsonSizeTokenMeter::new(1).unwrap()));
let request = || SessionContextRequest {
session_id: AgentSessionId::new("session-1"),
current_run_id: RunId::new("current"),
through_session_seq: None,
system_message: Some(ModelMessage::text(ModelRole::System, "system")),
tools: Vec::new(),
history_limit: 100,
max_context_tokens: 600,
reserved_output_tokens: 100,
config_digest: Digest::sha256("config"),
allowed_skill_digests: BTreeMap::new(),
};
let projection = engine.project(request()).await.unwrap();
let rendered = serde_json::to_string(&projection.messages).unwrap();
assert!(rendered.contains("message-12"));
assert!(rendered.contains("message-11"));
assert!(!rendered.contains("message-1-"));
assert!(projection.used_input_tokens <= projection.input_budget_tokens);
let planning = engine
.project_with_policy(request(), ContextTokenPolicy::Planning)
.await
.unwrap();
assert_eq!(planning.messages, projection.messages);
assert_eq!(planning.used_input_tokens, projection.used_input_tokens);
assert_eq!(planning.included_ranges, projection.included_ranges);
assert!(planning.context_estimate.is_none());
}
#[tokio::test]
async fn tool_call_and_result_are_selected_as_one_atomic_group() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
append_input(&store, 1, "old", "old".to_owned()).await;
store
.append(AgentSessionEventDraft {
event_id: AgentSessionEventId::new("tool-pair"),
session_id: AgentSessionId::new("session-1"),
run_id: RunId::new("old"),
payload: AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new("request-1"),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: ModelToolCallId::new("call-1"),
name: "echo".to_owned(),
arguments: json!({}),
extensions: Default::default(),
}],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: ModelToolCallId::new("call-1"),
result: json!({ "large": "x".repeat(100) }),
is_error: false,
}],
},
retained_artifacts: Vec::new(),
usage: None,
},
})
.await
.unwrap();
append_input(&store, 3, "current", "current".to_owned()).await;
let engine =
AgentSessionContextEngine::new(store, Arc::new(JsonSizeTokenMeter::new(1).unwrap()));
let projection = engine
.project(SessionContextRequest {
session_id: AgentSessionId::new("session-1"),
current_run_id: RunId::new("current"),
through_session_seq: None,
system_message: None,
tools: Vec::new(),
history_limit: 100,
max_context_tokens: 180,
reserved_output_tokens: 20,
config_digest: Digest::sha256("config"),
allowed_skill_digests: BTreeMap::new(),
})
.await
.unwrap();
let has_call = projection.messages.iter().any(|message| {
message
.content
.iter()
.any(|content| matches!(content, ModelContent::ToolCall { .. }))
});
let has_result = projection.messages.iter().any(|message| {
message
.content
.iter()
.any(|content| matches!(content, ModelContent::ToolResult { .. }))
});
assert_eq!(has_call, has_result);
let prior = engine
.project(SessionContextRequest {
session_id: AgentSessionId::new("session-1"),
current_run_id: RunId::new("current"),
through_session_seq: Some(1),
system_message: None,
tools: Vec::new(),
history_limit: 100,
max_context_tokens: 180,
reserved_output_tokens: 20,
config_digest: Digest::sha256("config"),
allowed_skill_digests: BTreeMap::new(),
})
.await
.unwrap();
assert_eq!(prior.through_session_seq, 1);
assert!(prior.messages.iter().all(|message| {
message.content.iter().all(|content| {
!matches!(
content,
ModelContent::ToolCall { .. } | ModelContent::ToolResult { .. }
)
})
}));
}
#[tokio::test]
async fn compaction_is_traceable_ordered_and_does_not_repeat_without_new_history() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
for index in 1..=6 {
append_input(&store, index, "old", format!("old-{index}")).await;
}
append_input(&store, 7, "current", "current-input".to_owned()).await;
let policy = SessionCompactionPolicy {
minimum_source_records: 3,
keep_recent_records: 2,
};
let policy_digest = policy.digest().unwrap();
let compactor =
AgentSessionCompactor::new(store.clone(), Arc::new(FixedSummarizer), policy).unwrap();
let compacted = compactor
.compact_if_needed(&AgentSessionId::new("session-1"), &RunId::new("current"))
.await
.unwrap()
.expect("old prefix is compacted");
let (source, source_digest, persisted_policy_digest, summary_config_digest) =
match &compacted.payload {
AgentSessionEvent::CompactionCommitted {
source,
source_digest,
policy_digest,
summary_config_digest,
..
} => (source, source_digest, policy_digest, summary_config_digest),
_ => panic!("expected compaction record"),
};
assert_eq!(source.first_session_seq, 1);
assert_eq!(source.last_session_seq, 5);
let records = store
.load_session(&AgentSessionId::new("session-1"))
.await
.unwrap();
assert_eq!(
session_range_digest(&records, source).unwrap(),
*source_digest
);
assert_eq!(*persisted_policy_digest, policy_digest);
assert_eq!(
*summary_config_digest,
Digest::sha256("fixed-test-summary/v1")
);
assert!(compactor
.compact_if_needed(&AgentSessionId::new("session-1"), &RunId::new("current"),)
.await
.unwrap()
.is_none());
let engine =
AgentSessionContextEngine::new(store, Arc::new(JsonSizeTokenMeter::new(1).unwrap()));
let projection = engine
.project(SessionContextRequest {
session_id: AgentSessionId::new("session-1"),
current_run_id: RunId::new("current"),
through_session_seq: None,
system_message: None,
tools: Vec::new(),
history_limit: 100,
max_context_tokens: 10_000,
reserved_output_tokens: 100,
config_digest: Digest::sha256("config"),
allowed_skill_digests: BTreeMap::new(),
})
.await
.unwrap();
let rendered = projection
.messages
.iter()
.map(|message| serde_json::to_string(message).unwrap())
.collect::<Vec<_>>();
assert!(rendered[0].contains("summary of 5 messages"));
assert!(!rendered.join(" ").contains("old-1"));
assert!(rendered.last().unwrap().contains("current-input"));
}
#[tokio::test]
async fn active_run_pressure_compaction_preserves_task_and_atomic_recent_exchanges() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
append_input(&store, 1, "current", "inspect the repository".to_owned()).await;
for sequence in 1..=5 {
append_tool_exchange(&store, sequence, "current", 400).await;
}
let meter = Arc::new(JsonSizeTokenMeter::new(1).unwrap());
let engine = AgentSessionContextEngine::new(store.clone(), meter);
let request = || SessionContextRequest {
session_id: AgentSessionId::new("session-1"),
current_run_id: RunId::new("current"),
through_session_seq: None,
system_message: None,
tools: Vec::new(),
history_limit: 100,
max_context_tokens: 1_800,
reserved_output_tokens: 100,
config_digest: Digest::sha256("config"),
allowed_skill_digests: BTreeMap::new(),
};
assert!(matches!(
engine.project(request()).await,
Err(SessionContextError::ContextOverflow { .. })
));
let compactor = AgentSessionCompactor::new(
store.clone(),
Arc::new(FixedSummarizer),
SessionCompactionPolicy {
minimum_source_records: 3,
keep_recent_records: 2,
},
)
.unwrap();
let compacted = compactor
.compact_active_run_for_pressure(
&AgentSessionId::new("session-1"),
&RunId::new("current"),
)
.await
.unwrap()
.expect("older active exchanges should compact");
assert!(matches!(
compacted.payload,
AgentSessionEvent::ActiveRunCompactionCommitted { .. }
));
let projection = engine.project(request()).await.unwrap();
let rendered = serde_json::to_string(&projection.messages).unwrap();
assert!(rendered.contains("inspect the repository"));
assert!(rendered.contains("summary of 6 messages"));
assert!(!rendered.contains("active-call-1"));
assert!(rendered.contains("active-call-4"));
assert!(rendered.contains("active-call-5"));
let calls = projection
.messages
.iter()
.flat_map(|message| &message.content)
.filter(|content| matches!(content, ModelContent::ToolCall { .. }))
.count();
let results = projection
.messages
.iter()
.flat_map(|message| &message.content)
.filter(|content| matches!(content, ModelContent::ToolResult { .. }))
.count();
assert_eq!((calls, results), (2, 2));
let prior = engine
.project(SessionContextRequest {
through_session_seq: Some(6),
max_context_tokens: 10_000,
..request()
})
.await
.unwrap();
let prior_rendered = serde_json::to_string(&prior.messages).unwrap();
assert!(prior_rendered.contains("active-call-1"));
assert!(!prior_rendered.contains("summary of 6 messages"));
}
#[tokio::test]
async fn repeated_active_run_compaction_merges_the_previous_summary() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
append_input(&store, 1, "current", "continue the task".to_owned()).await;
for sequence in 1..=5 {
append_tool_exchange(&store, sequence, "current", 100).await;
}
let compactor = AgentSessionCompactor::new(
store.clone(),
Arc::new(FixedSummarizer),
SessionCompactionPolicy {
minimum_source_records: 3,
keep_recent_records: 2,
},
)
.unwrap();
compactor
.compact_active_run_for_pressure(
&AgentSessionId::new("session-1"),
&RunId::new("current"),
)
.await
.unwrap()
.unwrap();
append_tool_exchange(&store, 6, "current", 100).await;
append_tool_exchange(&store, 7, "current", 100).await;
let second = compactor
.compact_active_run_for_pressure(
&AgentSessionId::new("session-1"),
&RunId::new("current"),
)
.await
.unwrap()
.expect("the prior summary and newly old exchanges should merge");
let AgentSessionEvent::ActiveRunCompactionCommitted { source, .. } = second.payload else {
panic!("expected active-Run compaction");
};
assert_eq!(
source,
SessionSourceRange {
first_session_seq: 5,
last_session_seq: 8,
}
);
let records = store
.load_session(&AgentSessionId::new("session-1"))
.await
.unwrap();
let originals = original_compaction_groups(&records, &source).unwrap();
assert_eq!(originals.len(), 6);
assert_eq!(originals[0].source, single_range(2));
assert!(originals
.iter()
.flat_map(|group| &group.messages)
.all(|message| message.role != ModelRole::System));
let groups = replay_groups(&records, &RunId::new("current"), &BTreeMap::new()).unwrap();
let selected = groups.keys().copied().collect::<BTreeSet<_>>();
let rendered =
serde_json::to_string(&assemble_messages(&None, &groups, &selected)).unwrap();
assert!(rendered.contains("summary of 12 messages"));
assert!(!rendered.contains("active-call-6"));
assert!(rendered.contains("active-call-7"));
assert!(!rendered.contains("summary of 6 messages"));
}
#[tokio::test]
async fn pressure_compaction_preserves_literal_tool_text_and_replays_without_reencoding() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
let session_id = AgentSessionId::new("session-1");
let run_id = RunId::new("current");
let source = "fn trim_line(text: &str) -> &str {\r\n\t// 保留 `\\n` and \"quotes\"\r\n\ttext.trim_end_matches('\\n')\r\n}\r\n";
let result = json!({
"content": source,
"path": "src/lines.rs",
"range": { "start": 1, "end": 4 },
"truncated": false,
});
append_input(&store, 1, "current", "Inspect line handling".into()).await;
append_session_payload(
&store,
&session_id,
&run_id,
"source-observation".into(),
AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new("source-request"),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: ModelToolCallId::new("source-call"),
name: "inspect".into(),
arguments: json!({"path": "src/lines.rs"}),
extensions: Default::default(),
}],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: ModelToolCallId::new("source-call"),
result: result.clone(),
is_error: false,
}],
},
retained_artifacts: Vec::new(),
usage: None,
},
)
.await;
append_tool_exchange(&store, 2, "current", 10).await;
let originals = store.load_session(&session_id).await.unwrap();
let compactor = AgentSessionCompactor::new(
store.clone(),
Arc::new(DeterministicExtractiveSessionSummarizer::new(4096).unwrap()),
SessionCompactionPolicy {
minimum_source_records: 2,
keep_recent_records: 1,
},
)
.unwrap();
compactor
.compact_active_run_for_pressure(&session_id, &run_id)
.await
.unwrap()
.unwrap();
let records = store.load_session(&session_id).await.unwrap();
assert_eq!(&records[..originals.len()], originals.as_slice());
let AgentSessionEvent::ToolExchangeCommitted { tool, .. } = &records[1].payload else {
panic!("original tool exchange");
};
assert!(
matches!(&tool.content[0], ModelContent::ToolResult { result: saved, .. } if saved == &result)
);
let engine = AgentSessionContextEngine::new(
store.clone(),
Arc::new(JsonSizeTokenMeter::new(1).unwrap()),
);
let request = || SessionContextRequest {
session_id: session_id.clone(),
current_run_id: run_id.clone(),
through_session_seq: None,
system_message: None,
tools: Vec::new(),
history_limit: 100,
max_context_tokens: 20_000,
reserved_output_tokens: 64,
config_digest: Digest::sha256("literal-summary"),
allowed_skill_digests: BTreeMap::new(),
};
let projected = engine.project(request()).await.unwrap();
let summary = projected
.messages
.iter()
.filter(|message| {
message.role == ModelRole::Assistant
&& message
.content
.iter()
.all(|content| matches!(content, ModelContent::Text { .. }))
})
.flat_map(|message| &message.content)
.find_map(|content| match content {
ModelContent::Text { text } => Some(text),
_ => None,
})
.expect("projected literal summary");
assert!(summary.contains(source));
assert!(!summary.contains(&serde_json::to_string(source).unwrap()));
assert!(summary.contains("\"truncated\":false"));
assert!(summary.contains("\"range\":{\"end\":4,\"start\":1}"));
let replayed = engine
.project(SessionContextRequest {
through_session_seq: Some(projected.through_session_seq),
..request()
})
.await
.unwrap();
assert_eq!(replayed.messages, projected.messages);
assert_eq!(store.load_session(&session_id).await.unwrap(), records);
assert!(matches!(
engine
.project(SessionContextRequest {
max_context_tokens: projected.used_input_tokens + 64 - 1,
..request()
})
.await,
Err(SessionContextError::ContextOverflow { .. })
));
}
#[tokio::test]
async fn extractive_literal_text_excerpts_mark_omitted_regions_and_keep_failure_metadata() {
let head = "let start = \"\\n\";\n";
let tail = "let end = '\\n';\n";
let source = format!("{head}{}{tail}", "intermediate source line\n".repeat(300));
let summarizer = DeterministicExtractiveSessionSummarizer::new(2048).unwrap();
let input = |reverse: bool| {
let mut fields = vec![
("content".into(), json!(source)),
("exit_code".into(), json!(2)),
("kind".into(), json!("invalid_text")),
("message".into(), json!("expected one '\\n' boundary")),
("path".into(), json!("src/lines.rs")),
("truncated".into(), json!(true)),
];
if reverse {
fields.reverse();
}
SessionCompactionInput {
session_id: AgentSessionId::new("literal-excerpts"),
source: single_range(1),
focus_messages: Vec::new(),
groups: vec![SessionCompactionGroup {
source: single_range(1),
messages: vec![ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: ModelToolCallId::new("source-error"),
result: serde_json::Value::Object(fields.into_iter().collect()),
is_error: true,
}],
}],
}],
}
};
let summary = summarizer.summarize(input(false)).await.unwrap();
assert_eq!(summary, summarizer.summarize(input(true)).await.unwrap());
let ModelContent::Text { text } = &summary.content[0] else {
panic!("text summary");
};
let head_at = text.find(head).expect("literal source head");
let tail_at = text[head_at..].find(tail).unwrap() + head_at;
assert!(text[head_at..tail_at].contains("\n[… excerpt omitted …]\n"));
assert!(!text.contains(&format!("{head}{tail}")));
assert!(!text.contains(&source));
assert!(text.contains("not complete replacement text"));
assert!(text.contains("status=failed"));
assert!(text.contains("invalid_text"));
assert!(text.contains("expected one '\\n' boundary"));
assert!(text.contains("src/lines.rs"));
assert!(text.contains("\"exit_code\":2"));
assert!(text.contains("\"truncated\":true"));
assert!(text.chars().count() <= 2048);
}
#[tokio::test]
async fn compaction_omits_continuation_but_preserves_recent_and_durable_exchanges() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
let session_id = AgentSessionId::new("session-1");
let run_id = RunId::new("current");
append_input(&store, 1, "current", "Inspect recorded observations".into()).await;
for seq in 1..=5 {
let call_id = ModelToolCallId::new(format!("inspect-{seq}"));
append_session_payload(
&store,
&session_id,
&run_id,
format!("exchange-{seq}"),
AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new(format!("request-{seq}")),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![
ModelContent::Continuation {
namespace: "fixture/continuation".into(),
value: json!({"opaque": format!("private-step-{seq}")}),
},
ModelContent::Data {
media_type: "application/json".into(),
value: json!({"observation": "public-data"}),
},
ModelContent::ToolCall {
call_id: call_id.clone(),
name: "inspect".into(),
arguments: json!({"entry": seq}),
extensions: Default::default(),
},
],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id,
result: json!({"observed": seq}),
is_error: false,
}],
},
retained_artifacts: Vec::new(),
usage: None,
},
)
.await;
}
let originals = store.load_session(&session_id).await.unwrap();
let engine = AgentSessionContextEngine::new(
store.clone(),
Arc::new(JsonSizeTokenMeter::new(1).unwrap()),
);
let request = || SessionContextRequest {
session_id: session_id.clone(),
current_run_id: run_id.clone(),
through_session_seq: None,
system_message: None,
tools: Vec::new(),
history_limit: 100,
max_context_tokens: 20_000,
reserved_output_tokens: 100,
config_digest: Digest::sha256("continuation-compaction"),
allowed_skill_digests: BTreeMap::new(),
};
let before = engine.project(request()).await.unwrap();
let compactor = AgentSessionCompactor::new(
store.clone(),
Arc::new(DeterministicExtractiveSessionSummarizer::new(4096).unwrap()),
SessionCompactionPolicy {
minimum_source_records: 3,
keep_recent_records: 2,
},
)
.unwrap();
compactor
.compact_active_run_for_pressure(&session_id, &run_id)
.await
.unwrap()
.expect("older exchanges compact");
let after = engine.project(request()).await.unwrap();
let summaries = after
.messages
.iter()
.filter(|message| {
message.role == ModelRole::Assistant
&& message
.content
.iter()
.all(|content| matches!(content, ModelContent::Text { .. }))
})
.collect::<Vec<_>>();
assert!(!summaries.is_empty());
let summary = serde_json::to_string(&summaries).unwrap();
assert!(!summary.contains("private-step-"));
assert!(!summary.contains("fixture/continuation"));
assert!(
summary.contains("public-data"),
"ordinary Data remains visible"
);
let continuations = |messages: &[ModelMessage]| {
messages
.iter()
.flat_map(|message| &message.content)
.filter_map(|content| match content {
ModelContent::Continuation { value, .. } => Some(value.clone()),
_ => None,
})
.collect::<Vec<_>>()
};
let original_continuations = continuations(&before.messages);
assert_eq!(original_continuations.len(), 5);
assert_eq!(continuations(&after.messages), original_continuations[3..]);
let records = store.load_session(&session_id).await.unwrap();
assert_eq!(records[..originals.len()], originals);
let replay = engine
.project(SessionContextRequest {
through_session_seq: Some(originals.len() as u64),
..request()
})
.await
.unwrap();
assert_eq!(replay.messages, before.messages);
}
#[tokio::test]
async fn extractive_summary_is_bounded_deterministic_and_retains_referenced_facts() {
const FACTS: usize = 200;
const QUERIES: usize = 100;
const REFERENCED_PER_QUERY: usize = 5;
const MAX_SUMMARY_CHARS: usize = 2_048;
assert!(DeterministicExtractiveSessionSummarizer::new(255).is_err());
let summarizer = DeterministicExtractiveSessionSummarizer::new(MAX_SUMMARY_CHARS).unwrap();
let descriptor = summarizer.descriptor();
descriptor.validate().unwrap();
assert_eq!(descriptor, summarizer.descriptor());
let groups = (0..FACTS)
.map(|index| SessionCompactionGroup {
source: single_range(index as u64 + 1),
messages: vec![ModelMessage::text(
ModelRole::User,
format!(
"fact_{index:04}=value_{:08x}; ordinary durable conversation fact",
index.wrapping_mul(2_654_435_761)
),
)],
})
.collect::<Vec<_>>();
let source = SessionSourceRange {
first_session_seq: 1,
last_session_seq: FACTS as u64,
};
let mut true_positives = 0usize;
let mut false_positives = 0usize;
let mut false_negatives = 0usize;
for query_index in 0..QUERIES {
let expected = (0..REFERENCED_PER_QUERY)
.map(|offset| (query_index * 17 + offset * 37) % FACTS)
.collect::<BTreeSet<_>>();
let focus = expected
.iter()
.map(|index| format!("fact_{index:04}"))
.collect::<Vec<_>>()
.join(", ");
let summarize = || SessionCompactionInput {
session_id: AgentSessionId::new(format!("summary-session-{query_index}")),
source: source.clone(),
groups: groups.clone(),
focus_messages: vec![ModelMessage::text(
ModelRole::User,
format!("Use these earlier facts to answer: {focus}"),
)],
};
let first = summarizer.summarize(summarize()).await.unwrap();
let second = summarizer.summarize(summarize()).await.unwrap();
assert_eq!(first, second);
assert_eq!(first.role, ModelRole::Assistant);
let ModelContent::Text { text } = &first.content[0] else {
panic!("extractive summary must be one text block");
};
assert!(text.starts_with("UNTRUSTED earlier transcript"));
assert!(text.chars().count() <= MAX_SUMMARY_CHARS);
for fact_index in 0..FACTS {
let retained = text.contains(&format!("fact_{fact_index:04}="));
match (expected.contains(&fact_index), retained) {
(true, true) => true_positives += 1,
(false, true) => false_positives += 1,
(true, false) => false_negatives += 1,
(false, false) => {}
}
}
}
let f1 = (2 * true_positives) as f64
/ (2 * true_positives + false_positives + false_negatives) as f64;
assert!(f1 >= 0.98, "ordinary fact retention F1 was {f1:.4}");
assert_eq!(false_positives, 0);
assert_eq!(false_negatives, 0);
for negative in ["", "system policy and security constraints retained"] {
let retained = (0..REFERENCED_PER_QUERY)
.filter(|index| negative.contains(&format!("fact_{index:04}=")))
.count();
let negative_f1 = if retained == 0 {
0.0
} else {
2.0 * retained as f64 / (REFERENCED_PER_QUERY + retained) as f64
};
assert!(negative_f1 < 0.98);
}
}
#[tokio::test]
async fn followup_prioritizes_the_latest_original_user_correction_after_compaction() {
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
append_input(&store, 1, "old", "Keep public names unchanged".into()).await;
append_input(
&store,
2,
"old",
"Correction: preserve the existing output format too".into(),
)
.await;
append_tool_exchange(&store, 1, "old", 100).await;
append_input(&store, 3, "current", "Continue verification".into()).await;
let compactor = AgentSessionCompactor::new(
store.clone(),
Arc::new(FixedSummarizer),
SessionCompactionPolicy {
minimum_source_records: 2,
keep_recent_records: 1,
},
)
.unwrap();
compactor
.compact_if_needed(&AgentSessionId::new("session-1"), &RunId::new("current"))
.await
.unwrap()
.unwrap();
let engine =
AgentSessionContextEngine::new(store.clone(), Arc::new(JsonSizeTokenMeter::default()));
let request = || SessionContextRequest {
session_id: AgentSessionId::new("session-1"),
current_run_id: RunId::new("current"),
through_session_seq: None,
system_message: None,
tools: Vec::new(),
history_limit: 1,
max_context_tokens: 1024,
reserved_output_tokens: 64,
config_digest: Digest::sha256("anchor-test"),
allowed_skill_digests: BTreeMap::new(),
};
let projected = engine.project(request()).await.unwrap();
let text = serde_json::to_string(&projected.messages).unwrap();
assert!(text.contains("Correction: preserve the existing output format too"));
assert!(text.contains("Continue verification"));
assert_eq!(
projected
.messages
.iter()
.filter(|m| m.role == ModelRole::User)
.count(),
2
);
assert!(projected
.deferred_ranges
.iter()
.all(|range| !range.contains(2)));
let replayed = engine
.project(SessionContextRequest {
through_session_seq: Some(projected.through_session_seq),
..request()
})
.await
.unwrap();
assert_eq!(replayed.messages, projected.messages);
}
#[tokio::test]
async fn bounded_summary_keeps_typed_failure_before_large_call_arguments() {
let call_id = ModelToolCallId::new("verification-call");
let summarizer = DeterministicExtractiveSessionSummarizer::new(640).unwrap();
let summary = summarizer
.summarize(SessionCompactionInput {
session_id: AgentSessionId::new("s"),
source: single_range(1),
focus_messages: Vec::new(),
groups: vec![SessionCompactionGroup {
source: single_range(1),
messages: vec![
ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: call_id.clone(),
name: "check".into(),
arguments: json!({"data": "x".repeat(4000)}),
extensions: Default::default(),
}],
},
ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id,
result: json!({"reason": "output contract mismatch"}),
is_error: true,
}],
},
],
}],
})
.await
.unwrap();
let ModelContent::Text { text } = &summary.content[0] else {
panic!("text summary");
};
assert!(text.contains("status=failed"));
assert!(text.contains("output contract mismatch"));
assert!(text.contains("session_seq=1..1"));
assert!(text.chars().count() <= 640);
}
#[tokio::test]
async fn extractive_summary_retains_large_completed_check_alongside_inspection() {
let exchange = |seq, command: &str, result| SessionCompactionGroup {
source: single_range(seq),
messages: vec![
ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: ModelToolCallId::new(format!("call-{seq}")),
name: "run_process".into(),
arguments: json!({"command": command}),
extensions: Default::default(),
}],
},
ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: ModelToolCallId::new(format!("call-{seq}")),
result,
is_error: false,
}],
},
],
};
let summarizer = DeterministicExtractiveSessionSummarizer::new(4096).unwrap();
let input = || SessionCompactionInput {
session_id: AgentSessionId::new("large-check"),
source: SessionSourceRange {
first_session_seq: 1,
last_session_seq: 3,
},
focus_messages: vec![ModelMessage::text(
ModelRole::User,
"Inspect the parser, run its checks and deliver the requested output",
)],
groups: vec![
exchange(1, "inspect parser", json!({"source": "parser code"})),
exchange(
2,
"run checks",
json!({
"output": format!("{}\nCHECK COMPLETE: 243 passed", "checking parser\n".repeat(6000)),
"exit_code": 0,
"alive": false,
}),
),
exchange(
3,
"inspect parser after checks",
json!({
"source": "parser implementation\n".repeat(6000),
}),
),
],
};
let summary = summarizer.summarize(input()).await.unwrap();
assert_eq!(summary, summarizer.summarize(input()).await.unwrap());
let ModelContent::Text { text } = &summary.content[0] else {
panic!("text summary");
};
assert!(text.contains("run checks"));
assert!(text.contains("\"exit_code\":0"));
assert!(text.contains("CHECK COMPLETE: 243 passed"));
assert!(text.contains("session_seq=2..2"));
assert!(text.chars().count() <= 4096);
}
#[tokio::test]
async fn extractive_summary_keeps_current_process_state_among_relevant_old_failures() {
let summarizer = DeterministicExtractiveSessionSummarizer::new(2048).unwrap();
let mut groups = (1..=8)
.map(|seq| SessionCompactionGroup {
source: single_range(seq),
messages: vec![ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: ModelToolCallId::new(format!("old-{seq}")),
result: json!({"reason": "earlier packaging failure".repeat(100)}),
is_error: true,
}],
}],
})
.collect::<Vec<_>>();
groups.push(SessionCompactionGroup {
source: single_range(9),
messages: vec![ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: ModelToolCallId::new("current-process"),
result: json!({"alive": true, "session_id": 71, "output": "进展\n".repeat(4000)}),
is_error: false,
}],
}],
});
let summary = summarizer
.summarize(SessionCompactionInput {
session_id: AgentSessionId::new("process-state"),
source: SessionSourceRange {
first_session_seq: 1,
last_session_seq: 9,
},
groups,
focus_messages: vec![ModelMessage::text(ModelRole::User, "Finish packaging")],
})
.await
.unwrap();
let ModelContent::Text { text } = &summary.content[0] else {
panic!("text summary");
};
assert!(text.contains("session_seq=9..9"));
assert!(text.contains("\"alive\":true"));
assert!(text.contains("\"session_id\":71"));
assert!(text.contains("status=failed"));
assert!(text.contains("Tool success does not prove task verification"));
assert!(text.chars().count() <= 2048);
}
#[tokio::test]
async fn recent_observations_survive_repeated_inspection_without_shared_focus_words() {
let summarizer = DeterministicExtractiveSessionSummarizer::new(4096).unwrap();
let mut groups = Vec::new();
for seq in 1..=30 {
let is_check = seq == 2;
let (name, arguments, result) = if is_check {
(
"execute",
json!({"command": "validate-release"}),
json!({
"exit_code": 0,
"output": format!("{}\nVALIDATION COMPLETE", "progress\n".repeat(3000)),
}),
)
} else {
(
"inspect",
json!({"path": "parser"}),
json!({
"source": "parser implementation\n".repeat(3000),
}),
)
};
let call_id = ModelToolCallId::new(format!("call-{seq}"));
groups.push(SessionCompactionGroup {
source: single_range(seq),
messages: vec![
ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: call_id.clone(),
name: name.into(),
arguments,
extensions: Default::default(),
}],
},
ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id,
result,
is_error: false,
}],
},
],
});
}
let summary = summarizer
.summarize(SessionCompactionInput {
session_id: AgentSessionId::new("repeated-inspection"),
source: SessionSourceRange {
first_session_seq: 1,
last_session_seq: 30,
},
groups,
focus_messages: vec![ModelMessage::text(ModelRole::User, "Inspect parser")],
})
.await
.unwrap();
let ModelContent::Text { text } = &summary.content[0] else {
panic!("text summary");
};
assert!(text.contains("session_seq=2..2"));
assert!(text.contains("validate-release"));
assert!(text.contains("\"exit_code\":0"));
assert!(text.contains("VALIDATION COMPLETE"));
assert!(text.contains("session_seq=30..30"));
assert!(text.chars().count() <= 4096);
}
#[test]
fn ten_thousand_persisted_session_traces_replay_to_online_message_projection() {
const TRACES: usize = 10_000;
let policy_digest = SessionCompactionPolicy {
minimum_source_records: 3,
keep_recent_records: 2,
}
.digest()
.unwrap();
let summary_config_digest = Digest::sha256("session-replay-summary/v1");
let mut compacted_traces = 0usize;
let mut current_tool_traces = 0usize;
for case in 0..TRACES {
let session_id = AgentSessionId::new(format!("replay-session-{case}"));
let old_run_id = RunId::new(format!("replay-old-{case}"));
let current_run_id = RunId::new(format!("replay-current-{case}"));
let old_input = ModelMessage::text(ModelRole::User, format!("old-input-{case}"));
let old_call_id = ModelToolCallId::new(format!("old-call-{case}"));
let old_assistant = ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: old_call_id.clone(),
name: "inspect".to_owned(),
arguments: json!({ "case": case }),
extensions: Default::default(),
}],
};
let old_tool = ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: old_call_id,
result: json!({ "observed": case * 2 }),
is_error: false,
}],
};
let old_output = ModelMessage::text(ModelRole::Assistant, format!("old-output-{case}"));
let old_tail = ModelMessage::text(ModelRole::User, format!("old-tail-{case}"));
let current_input =
ModelMessage::text(ModelRole::User, format!("current-input-{case}"));
let mut records = Vec::with_capacity(9);
push_session_record(
&mut records,
&session_id,
&old_run_id,
case,
AgentSessionEvent::RunInputCommitted {
message: old_input.clone(),
},
);
push_session_record(
&mut records,
&session_id,
&old_run_id,
case,
AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new(format!("old-request-{case}")),
assistant: old_assistant.clone(),
tool: old_tool.clone(),
retained_artifacts: Vec::new(),
usage: None,
},
);
push_session_record(
&mut records,
&session_id,
&old_run_id,
case,
AgentSessionEvent::RunOutputCommitted {
request_id: ModelRequestId::new(format!("old-output-request-{case}")),
message: old_output.clone(),
usage: None,
},
);
push_session_record(
&mut records,
&session_id,
&old_run_id,
case,
AgentSessionEvent::RunInputCommitted {
message: old_tail.clone(),
},
);
push_session_record(
&mut records,
&session_id,
¤t_run_id,
case,
AgentSessionEvent::RunInputCommitted {
message: current_input.clone(),
},
);
let current_exchange = (case % 2 == 0).then(|| {
let call_id = ModelToolCallId::new(format!("current-call-{case}"));
let retained_artifacts = if case % 4 == 0 {
vec![ArtifactRefWithDigest {
artifact_ref: ArtifactRef::new(format!("replay-artifact-{case}")),
digest: Digest::sha256(format!("replay-artifact-bytes-{case}")),
}]
} else {
Vec::new()
};
let result = retained_artifacts.first().map_or_else(
|| json!({ "current_result": case + 1 }),
|artifact| json!({"kind": "artifact", "artifact": artifact}),
);
(
ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: call_id.clone(),
name: "lookup".to_owned(),
arguments: json!({ "current": case }),
extensions: Default::default(),
}],
},
ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id,
result,
is_error: false,
}],
},
retained_artifacts,
)
});
if let Some((assistant, tool, retained_artifacts)) = ¤t_exchange {
push_session_record(
&mut records,
&session_id,
¤t_run_id,
case,
AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new(format!("current-request-{case}")),
assistant: assistant.clone(),
tool: tool.clone(),
retained_artifacts: retained_artifacts.clone(),
usage: None,
},
);
current_tool_traces += 1;
}
let compacted = case % 3 != 0;
let summary =
ModelMessage::text(ModelRole::System, format!("durable-summary-for-{case}"));
if compacted {
let source = SessionSourceRange {
first_session_seq: 1,
last_session_seq: 3,
};
let source_digest = session_range_digest(&records, &source).unwrap();
push_session_record(
&mut records,
&session_id,
¤t_run_id,
case,
AgentSessionEvent::CompactionCommitted {
source,
source_digest,
policy_digest: policy_digest.clone(),
summary_config_digest: summary_config_digest.clone(),
summary: summary.clone(),
strategy: "session-replay-summary".to_owned(),
model: None,
version: "1".to_owned(),
},
);
compacted_traces += 1;
}
let uncertainty = (case % 5 == 0).then(|| {
let effect_call_id = ToolCallId::new(format!("replay-effect-{case}"));
let model_call_id = ModelToolCallId::new(format!("replay-model-call-{case}"));
let tool_name = "replay-effect-tool".to_owned();
let message = "effect acknowledgement missing".to_owned();
push_session_record(
&mut records,
&session_id,
&old_run_id,
case,
AgentSessionEvent::EffectUncertaintyCommitted {
effect_call_id: effect_call_id.clone(),
model_call_id: model_call_id.clone(),
tool_name: tool_name.clone(),
message: message.clone(),
},
);
effect_uncertainty_message(&effect_call_id, &model_call_id, &tool_name, &message)
});
validate_session_trace(&session_id, &records).unwrap();
let persisted_bytes = serde_json::to_vec(&records).unwrap();
let persisted: Vec<AgentSessionRecord> =
serde_json::from_slice(&persisted_bytes).unwrap();
validate_session_trace(&session_id, &persisted).unwrap();
let groups = replay_groups(&persisted, ¤t_run_id, &BTreeMap::new()).unwrap();
let selected = groups.keys().copied().collect::<BTreeSet<_>>();
let replayed = assemble_messages(&None, &groups, &selected);
let mut online = if compacted {
vec![summary, old_tail, current_input]
} else {
vec![
old_input,
old_assistant,
old_tool,
old_output,
old_tail,
current_input,
]
};
if let Some((assistant, tool, _)) = current_exchange {
online.push(assistant);
online.push(tool);
}
if let Some(uncertainty) = uncertainty {
online.insert(usize::from(compacted), uncertainty);
}
assert_eq!(replayed, online);
let replayed_calls = replayed
.iter()
.flat_map(|message| message.content.iter())
.filter(|content| matches!(content, ModelContent::ToolCall { .. }))
.count();
let replayed_results = replayed
.iter()
.flat_map(|message| message.content.iter())
.filter(|content| matches!(content, ModelContent::ToolResult { .. }))
.count();
assert_eq!(replayed_calls, replayed_results);
}
assert_eq!(compacted_traces, 6_666);
assert_eq!(current_tool_traces, 5_000);
}
#[test]
fn ten_thousand_journal_policy_combinations_select_one_deterministic_traceable_source() {
const HISTORIES: usize = 100;
const POLICIES_PER_HISTORY: usize = 100;
let mut seed = 0xC04D_AC71_0A11_CE55u64;
let mut selected = 0usize;
let mut deferred = 0usize;
for history_index in 0..HISTORIES {
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let record_count = 16 + (seed % 33) as usize;
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let current_sequence = 1 + (seed % record_count as u64);
let session_id = AgentSessionId::new(format!("compaction-session-{history_index}"));
let old_run_id = RunId::new(format!("compaction-old-{history_index}"));
let current_run_id = RunId::new(format!("compaction-current-{history_index}"));
let records = (1..=record_count as u64)
.map(|session_seq| {
AgentSessionRecord::seal(
AgentSessionEventDraft {
event_id: AgentSessionEventId::new(format!(
"compaction-event-{history_index}-{session_seq}"
)),
session_id: session_id.clone(),
run_id: if session_seq == current_sequence {
current_run_id.clone()
} else {
old_run_id.clone()
},
payload: AgentSessionEvent::RunInputCommitted {
message: ModelMessage::text(
ModelRole::User,
format!("history-{history_index}-{session_seq}"),
),
},
},
session_seq,
)
.unwrap()
})
.collect::<Vec<_>>();
validate_session_trace(&session_id, &records).unwrap();
for _ in 0..POLICIES_PER_HISTORY {
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let minimum_source_records = 1 + (seed % 12) as usize;
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let keep_recent_records = 1 + (seed % 12) as usize;
let policy = SessionCompactionPolicy {
minimum_source_records,
keep_recent_records,
};
let policy_digest = policy.digest().unwrap();
assert_eq!(policy.digest().unwrap(), policy_digest);
assert_ne!(
SessionCompactionPolicy {
minimum_source_records,
keep_recent_records: keep_recent_records + 1,
}
.digest()
.unwrap(),
policy_digest
);
let first = select_compaction_source(&records, ¤t_run_id, &policy);
let second = select_compaction_source(&records, ¤t_run_id, &policy);
assert_eq!(first, second);
let source_len = record_count.saturating_sub(keep_recent_records);
let expected = if record_count < minimum_source_records + keep_recent_records {
None
} else if current_sequence as usize > source_len {
Some(SessionSourceRange {
first_session_seq: 1,
last_session_seq: source_len as u64,
})
} else if current_sequence.saturating_sub(1) >= minimum_source_records as u64 {
Some(SessionSourceRange {
first_session_seq: 1,
last_session_seq: current_sequence - 1,
})
} else if source_len as u64 - current_sequence >= minimum_source_records as u64 {
Some(SessionSourceRange {
first_session_seq: current_sequence + 1,
last_session_seq: source_len as u64,
})
} else {
None
};
assert_eq!(first, expected);
if let Some(source) = first {
source.validate().unwrap();
let source_records = &records
[(source.first_session_seq - 1) as usize..source.last_session_seq as usize];
assert!(source_records.len() >= minimum_source_records);
assert!(source_records
.iter()
.all(|record| record.run_id != current_run_id));
assert!(source.last_session_seq as usize <= record_count - keep_recent_records);
assert!(session_range_digest(&records, &source).unwrap().is_sha256());
selected += 1;
} else {
deferred += 1;
}
}
}
assert_eq!(selected + deferred, HISTORIES * POLICIES_PER_HISTORY);
assert!(selected > 0);
assert!(deferred > 0);
}
#[test]
fn compaction_advances_around_artifact_and_uncertain_effect_barriers() {
let session_id = AgentSessionId::new("barrier-session");
let old_run_id = RunId::new("barrier-old-run");
let current_run_id = RunId::new("barrier-current-run");
let artifact = ArtifactRefWithDigest {
artifact_ref: ArtifactRef::new("barrier-artifact"),
digest: Digest::sha256("barrier-artifact-bytes"),
};
let artifact_call = ModelToolCallId::new("barrier-artifact-call");
let artifact_exchange = AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new("barrier-artifact-request"),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: artifact_call.clone(),
name: "artifact_source".to_owned(),
arguments: json!({}),
extensions: Default::default(),
}],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: artifact_call,
result: json!({"artifact": artifact.clone()}),
is_error: false,
}],
},
retained_artifacts: vec![artifact],
usage: None,
};
let mut records = Vec::new();
for sequence in 1..=12_u64 {
let payload = match sequence {
4 => artifact_exchange.clone(),
8 => AgentSessionEvent::EffectUncertaintyCommitted {
effect_call_id: ToolCallId::new("barrier-effect"),
model_call_id: ModelToolCallId::new("barrier-model-call"),
tool_name: "dangerous_tool".to_owned(),
message: "effect may have happened".to_owned(),
},
_ => AgentSessionEvent::RunInputCommitted {
message: ModelMessage::text(
ModelRole::User,
format!("barrier-history-{sequence}"),
),
},
};
push_session_record(&mut records, &session_id, &old_run_id, sequence, payload);
}
let policy = SessionCompactionPolicy {
minimum_source_records: 3,
keep_recent_records: 2,
};
let first = select_compaction_source(&records, ¤t_run_id, &policy).unwrap();
assert_eq!(
first,
SessionSourceRange {
first_session_seq: 1,
last_session_seq: 3,
}
);
let source_digest = session_range_digest(&records, &first).unwrap();
push_session_record(
&mut records,
&session_id,
&old_run_id,
"first-compaction",
AgentSessionEvent::CompactionCommitted {
source: first,
source_digest,
policy_digest: policy.digest().unwrap(),
summary_config_digest: Digest::sha256("barrier-summary-config"),
summary: ModelMessage::text(ModelRole::System, "first barrier summary"),
strategy: "barrier-test".to_owned(),
model: None,
version: "1".to_owned(),
},
);
for sequence in 14..=18_u64 {
push_session_record(
&mut records,
&session_id,
&old_run_id,
sequence,
AgentSessionEvent::RunInputCommitted {
message: ModelMessage::text(
ModelRole::User,
format!("post-compaction-{sequence}"),
),
},
);
}
let second = select_compaction_source(&records, ¤t_run_id, &policy).unwrap();
assert_eq!(
second,
SessionSourceRange {
first_session_seq: 5,
last_session_seq: 7,
}
);
}
#[tokio::test]
async fn ten_thousand_generated_long_histories_retain_every_fittable_safety_and_task_anchor() {
const HISTORIES: usize = 100;
const CONFIGS_PER_HISTORY: usize = 100;
let store = Arc::new(InMemoryAgentSessionJournalStore::default());
let meter = Arc::new(JsonSizeTokenMeter::new(1).unwrap());
let engine = AgentSessionContextEngine::new(store.clone(), meter.clone());
let system = ModelMessage::text(ModelRole::System, "stable system policy");
let tools = vec![ModelToolDefinition {
name: "inspect".to_owned(),
description: "Inspect one generated value".to_owned(),
input_schema: json!({
"type": "object",
"required": ["value"],
"properties": { "value": { "type": "string" } },
"additionalProperties": false
}),
}];
let mut seed = 0xA6E7_5E55_D15C_A11Du64;
let mut successful = 0usize;
let mut overflowed = 0usize;
for history_index in 0..HISTORIES {
let session_id = AgentSessionId::new(format!("budget-session-{history_index}"));
let old_run_id = RunId::new(format!("budget-old-{history_index}"));
let current_run_id = RunId::new(format!("budget-current-{history_index}"));
for sequence in 1..=12u64 {
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let width = 12 + (seed % 96) as usize;
let payload = if sequence % 3 == 0 {
let call_id =
ModelToolCallId::new(format!("budget-call-{history_index}-{sequence}"));
AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new(format!(
"budget-request-{history_index}-{sequence}"
)),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: call_id.clone(),
name: "inspect".to_owned(),
arguments: json!({ "value": "x".repeat(width) }),
extensions: Default::default(),
}],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id,
result: json!({ "observed": "y".repeat(width / 2 + 1) }),
is_error: false,
}],
},
retained_artifacts: Vec::new(),
usage: None,
}
} else {
AgentSessionEvent::RunInputCommitted {
message: ModelMessage::text(
ModelRole::User,
format!("history-{history_index}-{sequence}-{}", "z".repeat(width)),
),
}
};
append_session_payload(
&store,
&session_id,
&old_run_id,
format!("budget-event-{history_index}-{sequence}"),
payload,
)
.await;
}
let artifact = ArtifactRefWithDigest {
artifact_ref: ArtifactRef::new(format!("artifact-{history_index}")),
digest: Digest::sha256(format!("artifact-bytes-{history_index}")),
};
let artifact_call_id = ModelToolCallId::new(format!("artifact-call-{history_index}"));
append_session_payload(
&store,
&session_id,
&old_run_id,
format!("budget-artifact-{history_index}"),
AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new(format!("artifact-request-{history_index}")),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: artifact_call_id.clone(),
name: "artifact_source".to_owned(),
arguments: json!({}),
extensions: Default::default(),
}],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id: artifact_call_id,
result: json!({
"kind": "artifact",
"artifact": artifact.clone(),
"summary": "bounded generated Artifact"
}),
is_error: false,
}],
},
retained_artifacts: vec![artifact.clone()],
usage: None,
},
)
.await;
let uncertain_effect_call = format!("uncertain-effect-{history_index}");
let uncertain_model_call = format!("uncertain-model-call-{history_index}");
append_session_payload(
&store,
&session_id,
&old_run_id,
format!("budget-uncertain-effect-{history_index}"),
AgentSessionEvent::EffectUncertaintyCommitted {
effect_call_id: ToolCallId::new(&uncertain_effect_call),
model_call_id: ModelToolCallId::new(&uncertain_model_call),
tool_name: "generated_effect".to_owned(),
message: "dispatch acknowledgement was lost".to_owned(),
},
)
.await;
append_session_payload(
&store,
&session_id,
¤t_run_id,
format!("budget-current-event-{history_index}"),
AgentSessionEvent::RunInputCommitted {
message: ModelMessage::text(
ModelRole::User,
format!("current-task-{history_index}"),
),
},
)
.await;
for (kind, tool_name) in [
("pending-input", "orchestral_request_input"),
("pending-approval", "approval_guarded_tool"),
] {
let call_id = ModelToolCallId::new(format!("{kind}-call-{history_index}"));
append_session_payload(
&store,
&session_id,
¤t_run_id,
format!("budget-{kind}-{history_index}"),
AgentSessionEvent::ToolExchangeCommitted {
request_id: ModelRequestId::new(format!("{kind}-request-{history_index}")),
assistant: ModelMessage {
role: ModelRole::Assistant,
content: vec![ModelContent::ToolCall {
call_id: call_id.clone(),
name: tool_name.to_owned(),
arguments: json!({"marker": kind}),
extensions: Default::default(),
}],
},
tool: ModelMessage {
role: ModelRole::Tool,
content: vec![ModelContent::ToolResult {
call_id,
result: json!({"status": "resolved", "marker": kind}),
is_error: false,
}],
},
retained_artifacts: Vec::new(),
usage: None,
},
)
.await;
}
let records = store.load_session(&session_id).await.unwrap();
let groups = replay_groups(&records, ¤t_run_id, &BTreeMap::new()).unwrap();
let pinned = groups
.values()
.filter(|group| group.pinned)
.map(|group| group.key)
.collect::<BTreeSet<_>>();
let pinned_tokens = meter
.count_request_input(
&assemble_messages(&Some(system.clone()), &groups, &pinned),
&tools,
)
.unwrap();
let old_keys = groups
.values()
.rev()
.filter(|group| !group.pinned)
.map(|group| group.key)
.collect::<Vec<_>>();
let mut recent = pinned.clone();
let mut recent_prefixes = Vec::new();
for key in &old_keys {
recent.insert(*key);
recent_prefixes.push((
recent.clone(),
meter
.count_request_input(
&assemble_messages(&Some(system.clone()), &groups, &recent),
&tools,
)
.unwrap(),
));
}
for config_index in 0..CONFIGS_PER_HISTORY {
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let max_context_tokens = 320 + seed % 4_681;
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let reserved_output_tokens = 1 + seed % (max_context_tokens - 1);
seed = seed.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(1);
let history_limit = 1 + (seed % 12) as usize;
let input_budget = max_context_tokens - reserved_output_tokens;
let config_digest =
Digest::sha256(format!("budget-config-{history_index}-{config_index}"));
let projected = engine
.project(SessionContextRequest {
session_id: session_id.clone(),
current_run_id: current_run_id.clone(),
through_session_seq: None,
system_message: Some(system.clone()),
tools: tools.clone(),
history_limit,
max_context_tokens,
reserved_output_tokens,
config_digest: config_digest.clone(),
allowed_skill_digests: BTreeMap::new(),
})
.await;
if input_budget < pinned_tokens {
assert!(matches!(
projected,
Err(SessionContextError::ContextOverflow { used, budget })
if used == pinned_tokens && budget == input_budget
));
overflowed += 1;
continue;
}
let projection = projected.expect("fittable pinned context projects");
let actual = meter
.count_request_input(&projection.messages, &tools)
.unwrap();
assert_eq!(actual, projection.used_input_tokens);
assert!(actual <= input_budget);
assert_eq!(projection.input_budget_tokens, input_budget);
assert_eq!(projection.config_digest, config_digest);
assert_eq!(projection.through_session_seq, records.len() as u64);
let selected = projection
.included_ranges
.iter()
.map(|range| {
assert_eq!(range.first_session_seq, range.last_session_seq);
range.first_session_seq
})
.collect::<BTreeSet<_>>();
assert!(pinned.is_subset(&selected));
let rendered = serde_json::to_string(&projection.messages).unwrap();
assert!(rendered.contains("stable system policy"));
assert!(rendered.contains(&format!("current-task-{history_index}")));
assert!(rendered.contains(artifact.artifact_ref.as_str()));
assert!(rendered.contains(&uncertain_effect_call));
assert!(rendered.contains(&uncertain_model_call));
assert!(rendered.contains(&format!("pending-input-call-{history_index}")));
assert!(rendered.contains(&format!("pending-approval-call-{history_index}")));
let eligible_history = old_keys
.iter()
.take(history_limit)
.copied()
.collect::<BTreeSet<_>>();
assert!(selected
.difference(&pinned)
.all(|key| eligible_history.contains(key)));
if let Some((required, _)) =
recent_prefixes.iter().rev().find(|(required, tokens)| {
required.len().saturating_sub(pinned.len()) <= history_limit
&& *tokens <= input_budget
})
{
assert!(required.is_subset(&selected));
}
let calls = projection
.messages
.iter()
.flat_map(|message| message.content.iter())
.filter_map(|content| match content {
ModelContent::ToolCall { call_id, .. } => Some(call_id.clone()),
_ => None,
})
.collect::<BTreeSet<_>>();
let results = projection
.messages
.iter()
.flat_map(|message| message.content.iter())
.filter_map(|content| match content {
ModelContent::ToolResult { call_id, .. } => Some(call_id.clone()),
_ => None,
})
.collect::<BTreeSet<_>>();
assert_eq!(calls, results);
successful += 1;
}
}
assert_eq!(successful + overflowed, HISTORIES * CONFIGS_PER_HISTORY);
assert!(successful > 0);
assert!(overflowed > 0);
}
}