1use crate::message::{Message, MessagePart, MessageRole};
2
3pub const KEEP_RECENT_MESSAGES: usize = 10;
4pub const KEEP_RECENT_USER_TURNS: usize = 5;
5const KEEP_RECENT_TOKEN_FRACTION: f64 = 0.05;
6const COMPACTION_SAFETY_MARGIN_MIN: u64 = 2_000;
7const COMPACTION_SAFETY_MARGIN_MAX_RATIO: f64 = 0.05;
8
9#[derive(Debug, Clone, Copy, Default)]
10pub struct CompactionBudgetContext {
11 pub fixed_input_tokens: Option<u64>,
12}
13
14impl CompactionBudgetContext {
15 pub fn estimated_input_tokens(self, message_tokens: u64) -> u64 {
16 message_tokens.saturating_add(self.fixed_input_tokens.unwrap_or(0))
17 }
18
19 pub fn history_budget(self, info: &crate::model_registry::ModelInfo) -> Option<u64> {
20 let fixed_input_tokens = self.fixed_input_tokens?;
21 let output_cap = (info.context_budget as f64 * 0.20) as u64;
22 let output_floor = 8_000_u64.min(output_cap);
23 let output_reserve = (info.max_output_tokens.unwrap_or(32_000) as u64)
24 .max(output_floor)
25 .min(output_cap);
26 let safety_cap = (info.context_budget as f64 * COMPACTION_SAFETY_MARGIN_MAX_RATIO) as u64;
27 let safety_margin = ((info.context_budget as f64 * 0.02) as u64)
28 .max(COMPACTION_SAFETY_MARGIN_MIN.min(safety_cap))
29 .min(safety_cap);
30 Some(
31 info.context_budget
32 .saturating_sub(output_reserve)
33 .saturating_sub(safety_margin)
34 .saturating_sub(fixed_input_tokens),
35 )
36 }
37}
38
39pub fn estimate_tokens_for_message(msg: &Message) -> u64 {
40 let mut chars = 0usize;
41 let mut fixed_tokens = 0u64;
42 for part in &msg.parts {
43 chars += match part {
44 MessagePart::ContextRecord(record) => record.render_for_model().len(),
45 MessagePart::CompactSummary { summary, .. } => summary.len(),
46 MessagePart::Text { text } => text.len(),
47 MessagePart::Thinking { thinking, .. } => thinking.len(),
48 MessagePart::ToolResult { content, .. } => content.len(),
49 MessagePart::Image { source } => {
50 fixed_tokens = fixed_tokens.saturating_add(match source.detail {
51 crate::provider::ImageDetail::Low => 85,
52 crate::provider::ImageDetail::Auto => 1_024,
53 crate::provider::ImageDetail::High => 1_536,
54 crate::provider::ImageDetail::Original => 2_048,
55 });
56 0
57 }
58 MessagePart::ToolUse {
59 name,
60 input,
61 intent,
62 ..
63 } => {
64 name.len()
65 + input.to_string().len()
66 + intent.as_ref().map_or(0, |intent| {
67 crate::message::TOOL_CALL_INTENT_FIELD.len() + intent.as_str().len() + 5
68 })
69 }
70 };
71 }
72 chars = chars.saturating_add(estimate_role_overhead(msg.role));
73 (chars as f64 / 3.5).ceil() as u64 + fixed_tokens
74}
75
76fn estimate_role_overhead(role: MessageRole) -> usize {
77 match role {
78 MessageRole::System => 12,
79 MessageRole::User => 8,
80 MessageRole::Assistant => 8,
81 MessageRole::Tool => 16,
82 }
83}
84
85pub fn estimate_tokens_for_messages(messages: &[Message]) -> u64 {
86 messages.iter().map(estimate_tokens_for_message).sum()
87}
88
89#[derive(Debug, Clone, PartialEq, Eq)]
90pub struct CompactRange {
91 pub start: usize,
92 pub end: usize,
93 pub tokens_saved_estimate: u64,
94}
95
96pub fn is_plan_related(msg: &Message) -> bool {
97 for part in &msg.parts {
98 match part {
99 MessagePart::ToolUse { name, .. } if name.starts_with("plan.") => return true,
100 MessagePart::ToolResult { content, .. } if content.starts_with("# Plan:") => {
101 return true;
102 }
103 _ => {}
104 }
105 }
106 false
107}
108
109pub fn is_compaction_summary(msg: &Message) -> bool {
110 if !matches!(msg.role, MessageRole::System) {
111 return false;
112 }
113 msg.parts
114 .iter()
115 .any(|part| matches!(part, MessagePart::CompactSummary { .. }))
116}
117
118fn find_kth_recent_user(messages: &[Message], k: usize) -> usize {
119 let mut user_count = 0;
120 for (index, message) in messages.iter().enumerate().rev() {
121 if message.role == MessageRole::User {
122 user_count += 1;
123 if user_count == k {
124 return index;
125 }
126 }
127 }
128 0
129}
130
131fn align_compact_end_to_tool_transactions(messages: &[Message], end: usize) -> usize {
132 let mut tool_use_messages = std::collections::HashMap::new();
133 for (index, message) in messages.iter().take(end).enumerate() {
134 for part in &message.parts {
135 if let MessagePart::ToolUse { id, .. } = part {
136 tool_use_messages.entry(id.as_str()).or_insert(index);
137 }
138 }
139 }
140
141 let mut aligned_end = end;
142 for message in messages.iter().skip(end) {
143 for part in &message.parts {
144 if let MessagePart::ToolResult { tool_use_id, .. } = part
145 && let Some(&use_index) = tool_use_messages.get(tool_use_id.as_str())
146 {
147 aligned_end = aligned_end.min(use_index);
148 }
149 }
150 }
151 aligned_end
152}
153
154pub fn find_compact_range(messages: &[Message], budget: u64) -> Option<CompactRange> {
155 let total = estimate_tokens_for_messages(messages);
156 if total <= budget || messages.len() < 4 {
157 return None;
158 }
159
160 let start = messages
161 .iter()
162 .rposition(is_compaction_summary)
163 .unwrap_or(0);
164 let keep_recent_tokens = (budget as f64 * KEEP_RECENT_TOKEN_FRACTION).ceil() as u64;
165 let mut recent_tokens = 0u64;
166 let mut token_end = messages.len();
167 for (index, message) in messages.iter().enumerate().rev() {
168 recent_tokens = recent_tokens.saturating_add(estimate_tokens_for_message(message));
169 token_end = index;
170 if recent_tokens >= keep_recent_tokens {
171 break;
172 }
173 }
174 let message_end = messages.len().saturating_sub(KEEP_RECENT_MESSAGES);
175 let end = message_end
176 .min(token_end)
177 .min(find_kth_recent_user(messages, KEEP_RECENT_USER_TURNS));
178 let end = align_compact_end_to_tool_transactions(messages, end);
179 if end < start + 2 {
180 return None;
181 }
182
183 let tokens_saved_estimate = messages[start..end]
184 .iter()
185 .map(estimate_tokens_for_message)
186 .sum();
187 Some(CompactRange {
188 start,
189 end,
190 tokens_saved_estimate,
191 })
192}
193
194pub fn estimate_compacted_message_tokens(
195 messages: &[Message],
196 range: &CompactRange,
197 summary: &str,
198) -> u64 {
199 let turn_id = messages
200 .get(range.start)
201 .map(|m| m.turn_id.clone())
202 .unwrap_or_else(crate::event::TurnId::now);
203 let after = replace_range_with_summary(messages, range, summary.to_string(), turn_id);
204 estimate_tokens_for_messages(&after)
205}
206
207pub fn filter_orphan_tool_messages(messages: &mut Vec<Message>) {
208 crate::message::retain_complete_tool_pairs(messages);
209}
210
211pub fn find_compact_summaries(messages: &[Message]) -> Vec<CompactSummary> {
212 let mut out = Vec::new();
213 for (idx, msg) in messages.iter().enumerate() {
214 if let Some(summary) = compact_summary(msg) {
215 out.push(CompactSummary {
216 message_index: idx,
217 seq_start: summary.seq_start,
218 seq_end: summary.seq_end,
219 count: summary.count,
220 });
221 }
222 }
223 out
224}
225
226#[derive(Debug, Clone, PartialEq, Eq)]
227pub struct CompactSummary {
228 pub message_index: usize,
229 pub seq_start: u64,
230 pub seq_end: u64,
231 pub count: usize,
232}
233
234struct CompactSummaryPart {
235 seq_start: u64,
236 seq_end: u64,
237 count: usize,
238}
239
240fn extract_anchor(messages: &[Message]) -> Option<(String, &[Message])> {
241 let first = messages.first()?;
242 let summary = first.parts.iter().find_map(|part| match part {
243 MessagePart::CompactSummary { summary, .. } => Some(summary.clone()),
244 _ => None,
245 })?;
246 Some((summary, &messages[1..]))
247}
248
249fn compact_summary(msg: &Message) -> Option<CompactSummaryPart> {
250 if msg.role != MessageRole::System {
251 return None;
252 }
253 msg.parts.iter().find_map(|part| match part {
254 MessagePart::CompactSummary {
255 seq_start,
256 seq_end,
257 count,
258 ..
259 } => Some(CompactSummaryPart {
260 seq_start: *seq_start,
261 seq_end: *seq_end,
262 count: *count,
263 }),
264 _ => None,
265 })
266}
267
268pub async fn maybe_auto_compact(
269 session: &crate::session::Session,
270 model: &str,
271 providers: &crate::provider::ProviderRegistry,
272) {
273 maybe_auto_compact_with_budget(
274 session,
275 model,
276 providers,
277 CompactionBudgetContext::default(),
278 )
279 .await;
280}
281
282pub async fn maybe_auto_compact_with_budget(
283 session: &crate::session::Session,
284 model: &str,
285 providers: &crate::provider::ProviderRegistry,
286 budget_context: CompactionBudgetContext,
287) {
288 let _compact_guard = session.acquire_compact_lock().await;
289 maybe_auto_compact_locked(session, model, providers, budget_context).await;
290}
291
292pub fn spawn_auto_compact(
293 session: std::sync::Arc<crate::session::Session>,
294 model: String,
295 providers: crate::provider::ProviderRegistry,
296) {
297 tokio::task::spawn_blocking(move || {
298 let Ok(rt) = tokio::runtime::Builder::new_current_thread()
299 .enable_all()
300 .build()
301 else {
302 session.push_system_note("compaction skipped: background runtime init failed".into());
303 return;
304 };
305 rt.block_on(async move {
306 maybe_auto_compact(&session, &model, &providers).await;
307 });
308 });
309}
310
311pub async fn start_auto_compact(
312 session: std::sync::Arc<crate::session::Session>,
313 model: String,
314 providers: crate::provider::ProviderRegistry,
315) {
316 start_auto_compact_with_budget(
317 session,
318 model,
319 providers,
320 CompactionBudgetContext::default(),
321 )
322 .await;
323}
324
325pub async fn start_auto_compact_with_budget(
326 session: std::sync::Arc<crate::session::Session>,
327 model: String,
328 providers: crate::provider::ProviderRegistry,
329 budget_context: CompactionBudgetContext,
330) {
331 let compact_guard = session.acquire_compact_lock_owned().await;
332 tokio::task::spawn_blocking(move || {
333 let Ok(rt) = tokio::runtime::Builder::new_current_thread()
334 .enable_all()
335 .build()
336 else {
337 drop(compact_guard);
338 session.push_system_note("compaction skipped: background runtime init failed".into());
339 return;
340 };
341 rt.block_on(async move {
342 maybe_auto_compact_locked(&session, &model, &providers, budget_context).await;
343 drop(compact_guard);
344 });
345 });
346}
347
348async fn maybe_auto_compact_locked(
349 session: &crate::session::Session,
350 model: &str,
351 providers: &crate::provider::ProviderRegistry,
352 budget_context: CompactionBudgetContext,
353) {
354 let forced = session.take_manual_compact_request();
355 let info = crate::model_registry::model_info(model);
356 let trigger = info.compaction_trigger_threshold();
357 let target = budget_context
358 .history_budget(&info)
359 .map(|budget| budget.min(info.compaction_target_after()))
360 .unwrap_or_else(|| info.compaction_target_after());
361 let msgs = session.messages();
362 let window_tokens = estimate_tokens_for_messages(&msgs);
363 let current = budget_context.estimated_input_tokens(window_tokens);
364 if !forced && current <= trigger {
365 return;
366 }
367 if !forced && !session.approval_cooldown_ok_for_compact() {
368 return;
369 }
370 let Some(range) = find_compact_range(&msgs, target) else {
371 let (replacement, rewritten_count) =
372 build_budgeted_turn_rewrite(&msgs, target, model, providers).await;
373 let after_tokens = estimate_tokens_for_messages(&replacement);
374 if rewritten_count == 0 || after_tokens >= window_tokens || after_tokens > target {
375 session.emit_compact_warning(
376 model,
377 current,
378 trigger,
379 info.context_budget,
380 "no compactible span — retained user content cannot fit the history budget",
381 );
382 return;
383 }
384 match session.commit_rewritten_window(
385 replacement,
386 window_tokens,
387 window_tokens,
388 rewritten_count,
389 ) {
390 Some(_) => {}
391 None => {
392 session.emit_compact_warning(
393 model,
394 current,
395 trigger,
396 info.context_budget,
397 "retained turn output rewrite did not shrink the transcript",
398 );
399 }
400 }
401 return;
402 };
403 let _ = session
404 .stream_tx()
405 .send(crate::stream::StreamFrame::CompactionSummary {
406 phase: crate::stream::CompactionPhase::Running,
407 range_start: range.start,
408 range_end: range.end.saturating_sub(1),
409 summary: String::new(),
410 before_tokens: current,
411 after_tokens: 0,
412 compacted_count: range.end - range.start,
413 });
414 let send_failed = |session: &crate::session::Session, reason: &str| {
415 let _ = session
416 .stream_tx()
417 .send(crate::stream::StreamFrame::CompactionSummary {
418 phase: crate::stream::CompactionPhase::Failed,
419 range_start: range.start,
420 range_end: range.end.saturating_sub(1),
421 summary: reason.to_string(),
422 before_tokens: current,
423 after_tokens: current,
424 compacted_count: range.end - range.start,
425 });
426 };
427 let mut filtered: Vec<Message> = msgs[range.start..range.end].to_vec();
428 filter_orphan_tool_messages(&mut filtered);
429 let (anchor, new_messages) = extract_anchor(&filtered)
430 .map(|(anchor, remaining)| (Some(anchor), remaining.to_vec()))
431 .unwrap_or_else(|| (None, filtered.clone()));
432 let summary =
433 match generate_llm_summary(anchor.as_deref(), &new_messages, model, providers).await {
434 Ok(text) => text,
435 Err(err) => {
436 session.emit_compact_warning(
437 model,
438 current,
439 trigger,
440 info.context_budget,
441 &format!("LLM summary failed: {err}. Degraded to placeholder."),
442 );
443 format!(
444 "[atman: compacted {} messages, LLM summary unavailable at {}]",
445 range.end - range.start,
446 chrono::Utc::now().to_rfc3339()
447 )
448 }
449 };
450 let final_summary =
451 match request_review_if_enabled(session, forced, &filtered, &range, current, summary).await
452 {
453 ReviewOutcome::Commit(s) => s,
454 ReviewOutcome::Rejected => {
455 send_failed(
456 session,
457 "compaction rejected by user; keeping full transcript",
458 );
459 session.push_system_note(
460 "compaction rejected by user; keeping full transcript".into(),
461 );
462 return;
463 }
464 };
465 let replacement =
466 build_budgeted_replacement(&msgs, &range, &final_summary, target, model, providers).await;
467 let after_tokens = estimate_tokens_for_messages(&replacement);
468 if after_tokens >= window_tokens {
469 send_failed(
470 session,
471 &format!(
472 "compaction skipped: replacement would not shrink transcript ({} >= {} tokens)",
473 after_tokens, window_tokens
474 ),
475 );
476 session.push_system_note(format!(
477 "compaction skipped: replacement would not shrink transcript ({} >= {} tokens)",
478 after_tokens, window_tokens
479 ));
480 return;
481 }
482 match session.commit_compacted_window(
483 final_summary,
484 replacement,
485 range,
486 window_tokens,
487 window_tokens,
488 ) {
489 Some(result) => {
490 session.push_system_note(format!(
491 "auto-compacted {}..{} — {} → {} tokens",
492 result.compacted_start,
493 result.compacted_end,
494 result.before_tokens,
495 result.after_tokens
496 ));
497 }
498 None => {
499 session.emit_compact_warning(
500 model,
501 current,
502 trigger,
503 info.context_budget,
504 "no compactible span — history too short or already fully compacted",
505 );
506 }
507 }
508}
509
510enum ReviewOutcome {
511 Commit(String),
512 Rejected,
513}
514
515async fn request_review_if_enabled(
516 session: &crate::session::Session,
517 forced: bool,
518 slice: &[Message],
519 range: &CompactRange,
520 tokens_before: u64,
521 summary: String,
522) -> ReviewOutcome {
523 if !session.compact_review_mode().should_review(forced) {
524 return ReviewOutcome::Commit(summary);
525 }
526 let reviews = session.compact_reviews();
527 if reviews.subscriber_count() == 0 {
528 return ReviewOutcome::Commit(summary);
529 }
530 let pending = crate::session::PendingCompactReview {
531 review_id: uuid::Uuid::now_v7().to_string(),
532 summary: summary.clone(),
533 slice_preview: format_slice_for_preview(slice),
534 slice_count: slice.len(),
535 range_start: range.start,
536 range_end: range.end,
537 tokens_before,
538 emitted_at: chrono::Utc::now(),
539 };
540 let rx = reviews.request(pending);
541 match rx.await {
542 Ok(crate::session::CompactReviewDecision::AcceptAsIs) => ReviewOutcome::Commit(summary),
543 Ok(crate::session::CompactReviewDecision::AcceptEdited { summary: edited }) => {
544 ReviewOutcome::Commit(edited)
545 }
546 Ok(crate::session::CompactReviewDecision::Reject) | Err(_) => ReviewOutcome::Rejected,
547 }
548}
549
550fn format_slice_for_preview(slice: &[Message]) -> String {
551 let mut out = String::new();
552 for (i, msg) in slice.iter().enumerate() {
553 let role = msg.role.as_str();
554 let body = serialize_message_for_summary(msg);
555 let truncated: String = body.chars().take(400).collect();
556 out.push_str(&format!("[{i}] {role}: {truncated}\n"));
557 }
558 out.chars().take(16_000).collect()
559}
560
561const SUMMARY_SYSTEM_PROMPT: &str =
562 "You are an anchored context summarization assistant for coding sessions.";
563
564const SUMMARY_INSTRUCTIONS: &str = r#"You are an anchored context summarization assistant.
565
566Below is:
5671. <current-anchor>: the existing handoff state, which is authoritative and must be preserved.
5682. <new-messages>: only the messages that arrived since the anchor was written.
569
570Merge the NEW facts from <new-messages> INTO the current anchor, producing an upgraded full anchor.
571
572STRUCTURAL RULES (data model, not optional style):
573- ## Objective: unchanged unless the new messages show the user explicitly redirected.
574- ### Completed: ONLY ADD newly completed items. Never remove or re-evaluate an existing completed item. If a completed item is now in question, add it to ### Active or ### Blocked instead. NEVER delete from Completed.
575- ### Active: update based on new messages; move newly-done items to Completed.
576- ### Blocked: update based on new messages; remove resolved ones.
577- ## Decisions: only add new decisions. Never remove old ones.
578- ## Next Move: replace based on current end state.
579- Keep every section, even when empty.
580- Preserve exact file paths, symbols, commands, error strings, identifiers.
581
582Output exactly this Markdown structure:
583## Objective
584## Important Details
585## Work State
586### Completed
587### Active
588### Blocked
589## Decisions
590## Next Move
591## Relevant Files
592
593Do not mention the summary process or that context was compacted.
594Respond in the same language as the conversation."#;
595
596async fn generate_llm_summary(
597 anchor: Option<&str>,
598 slice: &[Message],
599 model: &str,
600 providers: &crate::provider::ProviderRegistry,
601) -> Result<String, crate::error::RuntimeError> {
602 let provider = providers.resolve(model).ok_or_else(|| {
603 crate::error::RuntimeError::ToolFailed(format!("no provider for {model}"))
604 })?;
605 let payload = format_slice_for_summary(slice);
606 let (messages, dump_user) = if let Some(anchor) = anchor {
607 let anchor_user = format!("<current-anchor>\n{anchor}\n</current-anchor>");
608 let new_user =
609 format!("<new-messages>\n{payload}\n</new-messages>\n\n{SUMMARY_INSTRUCTIONS}");
610 (
611 vec![
612 Message::user_text(crate::event::TurnId::now(), anchor_user.clone()),
613 Message::user_text(crate::event::TurnId::now(), new_user.clone()),
614 ],
615 format!("{anchor_user}\n\n{new_user}"),
616 )
617 } else {
618 let user = format!(
619 "<conversation_history>\n{payload}\n</conversation_history>\n\n{SUMMARY_INSTRUCTIONS}"
620 );
621 (
622 vec![Message::user_text(
623 crate::event::TurnId::now(),
624 user.clone(),
625 )],
626 user,
627 )
628 };
629 if let Ok(dir) = std::env::var("ATMAN_COMPACT_DUMP") {
630 let _ = std::fs::write(
631 format!("{dir}/compact_request.txt"),
632 format!("=== SYSTEM ===\n{SUMMARY_SYSTEM_PROMPT}\n\n=== USER ===\n{dump_user}"),
633 );
634 }
635 let req = crate::provider::LlmRequest {
636 model: model.into(),
637 messages,
638 system: Some(SUMMARY_SYSTEM_PROMPT.into()),
639 input: crate::value::Value::Unit,
640 schema: None,
641 cache_prompt: false,
642 prompt_cache_key: None,
643 tools: Vec::new(),
644 reasoning: crate::provider::ReasoningSelection::ProviderDefault,
645 stall_timeout_secs: 0,
646 };
647 let outcome = provider.call(req).await?;
648 let text = outcome.text_concat();
649 if text.trim().is_empty() {
650 return Err(crate::error::RuntimeError::ToolFailed(
651 "empty summary from provider".into(),
652 ));
653 }
654 Ok(text)
655}
656
657async fn build_budgeted_turn_rewrite(
658 messages: &[Message],
659 history_budget: u64,
660 model: &str,
661 providers: &crate::provider::ProviderRegistry,
662) -> (Vec<Message>, usize) {
663 let mut replacement = messages.to_vec();
664 let mut group_index = 0;
665 let mut rewritten_count = 0;
666 while estimate_tokens_for_messages(&replacement) > history_budget {
667 let groups = user_turn_ranges(&replacement);
668 let Some((start, end)) = groups.get(group_index).copied() else {
669 break;
670 };
671 let output = replacement[start + 1..end].to_vec();
672 if output.is_empty() {
673 group_index += 1;
674 continue;
675 }
676 let output_tokens = estimate_tokens_for_messages(&output);
677 let summary = generate_llm_summary(None, &output, model, providers)
678 .await
679 .unwrap_or_else(|_| deterministic_turn_omission(&output));
680 let mut summary_message = Message::assistant_text(
681 replacement[start].turn_id.clone(),
682 format!("[atman: compacted turn output]\n{summary}\n[/atman: compacted turn output]"),
683 );
684 if estimate_tokens_for_message(&summary_message) >= output_tokens {
685 summary_message = Message::assistant_text(
686 replacement[start].turn_id.clone(),
687 deterministic_turn_omission(&output),
688 );
689 }
690 rewritten_count += output.len();
691 replacement.splice(start + 1..end, [summary_message]);
692 group_index += 1;
693 }
694 filter_orphan_tool_messages(&mut replacement);
695 (replacement, rewritten_count)
696}
697
698async fn build_budgeted_replacement(
699 messages: &[Message],
700 range: &CompactRange,
701 anchor_summary: &str,
702 history_budget: u64,
703 model: &str,
704 providers: &crate::provider::ProviderRegistry,
705) -> Vec<Message> {
706 let turn_id = messages
707 .get(range.start)
708 .map(|message| message.turn_id.clone())
709 .unwrap_or_else(crate::event::TurnId::now);
710 let mut replacement =
711 replace_range_with_summary(messages, range, anchor_summary.to_string(), turn_id);
712 filter_orphan_tool_messages(&mut replacement);
713 if estimate_tokens_for_messages(&replacement) <= history_budget {
714 return replacement;
715 }
716
717 let mut group_index = 0;
718 loop {
719 let groups = user_turn_ranges(&replacement);
720 if group_index >= groups.len()
721 || estimate_tokens_for_messages(&replacement) <= history_budget
722 {
723 break;
724 }
725 let (start, end) = groups[group_index];
726 let output: Vec<Message> = replacement[start + 1..end].to_vec();
727 if output.is_empty() {
728 group_index += 1;
729 continue;
730 }
731 let output_tokens = estimate_tokens_for_messages(&output);
732 let summary = generate_llm_summary(None, &output, model, providers)
733 .await
734 .unwrap_or_else(|_| deterministic_turn_omission(&output));
735 let mut summary_message = Message::assistant_text(
736 replacement[start].turn_id.clone(),
737 format!("[atman: compacted turn output]\n{summary}\n[/atman: compacted turn output]"),
738 );
739 if estimate_tokens_for_message(&summary_message) >= output_tokens {
740 summary_message = Message::assistant_text(
741 replacement[start].turn_id.clone(),
742 deterministic_turn_omission(&output),
743 );
744 }
745 replacement.splice(start + 1..end, [summary_message]);
746 group_index += 1;
747 }
748
749 if estimate_tokens_for_messages(&replacement) > history_budget {
750 let groups = user_turn_ranges(&replacement);
751 for (start, end) in groups.into_iter().rev() {
752 let output = replacement[start + 1..end].to_vec();
753 if !output.is_empty() {
754 replacement.splice(
755 start + 1..end,
756 [Message::assistant_text(
757 replacement[start].turn_id.clone(),
758 deterministic_turn_omission(&output),
759 )],
760 );
761 }
762 }
763 }
764
765 if estimate_tokens_for_messages(&replacement) > history_budget {
766 replacement = compaction_floor(&replacement);
767 }
768
769 filter_orphan_tool_messages(&mut replacement);
770 replacement
771}
772
773fn user_turn_ranges(messages: &[Message]) -> Vec<(usize, usize)> {
774 let starts: Vec<usize> = messages
775 .iter()
776 .enumerate()
777 .filter_map(|(index, message)| (message.role == MessageRole::User).then_some(index))
778 .collect();
779 starts
780 .iter()
781 .enumerate()
782 .map(|(index, start)| {
783 (
784 *start,
785 starts.get(index + 1).copied().unwrap_or(messages.len()),
786 )
787 })
788 .collect()
789}
790
791fn deterministic_turn_omission(messages: &[Message]) -> String {
792 format!(
793 "[atman: omitted {} oversized assistant/system/tool messages during persistent compaction]",
794 messages.len()
795 )
796}
797
798fn compaction_floor(messages: &[Message]) -> Vec<Message> {
799 let mut out = Vec::new();
800 if let Some(anchor) = messages
801 .iter()
802 .find(|message| is_compaction_summary(message))
803 {
804 out.push(anchor.clone());
805 }
806
807 let mut latest =
808 std::collections::HashMap::<&str, (&crate::context_plan::ContextRecord, &Message)>::new();
809 for message in messages {
810 for part in &message.parts {
811 if let MessagePart::ContextRecord(record) = part {
812 latest
813 .entry(record.key())
814 .and_modify(|(current, source)| {
815 if record.revision() >= current.revision() {
816 *current = record;
817 *source = message;
818 }
819 })
820 .or_insert((record, message));
821 }
822 }
823 }
824 let mut records: Vec<_> = latest.into_values().collect();
825 records.sort_by(|(left, _), (right, _)| left.key().cmp(right.key()));
826 out.extend(
827 records.into_iter().map(|(record, source)| {
828 Message::context_record(source.turn_id.clone(), record.clone())
829 }),
830 );
831
832 out.extend(messages.iter().filter_map(|message| {
833 if message.role != MessageRole::User {
834 return None;
835 }
836 let mut user = message.clone();
837 user.parts.retain(|part| {
838 !matches!(
839 part,
840 MessagePart::CompactSummary { .. } | MessagePart::ContextRecord(_)
841 )
842 });
843 (!user.parts.is_empty()).then_some(user)
844 }));
845 out
846}
847
848fn format_slice_for_summary(slice: &[Message]) -> String {
849 let mut out = String::new();
850 for (i, msg) in slice.iter().enumerate() {
851 let role = msg.role.as_str();
852 let body = serialize_message_for_summary(msg);
853 let truncated: String = body.chars().take(4000).collect();
854 out.push_str(&format!("[{i}] {role}: {truncated}\n\n"));
855 }
856 out.chars().take(120_000).collect()
857}
858
859fn serialize_message_for_summary(msg: &Message) -> String {
860 let mut parts = Vec::new();
861 for part in &msg.parts {
862 match part {
863 MessagePart::ContextRecord(record) => {
864 if record.retention() == crate::context_plan::ContextRecordRetention::Timeline {
865 parts.push(record.render_for_model());
866 }
867 }
868 MessagePart::CompactSummary { summary, .. } => {
869 parts.push(summary.clone());
870 }
871 MessagePart::Text { text } => {
872 parts.push(text.clone());
873 }
874 MessagePart::Thinking { thinking, .. } => {
875 let truncated: String = thinking.chars().take(1000).collect();
876 parts.push(format!("[thinking: {truncated}]"));
877 }
878 MessagePart::ToolUse {
879 name,
880 input,
881 intent,
882 ..
883 } => {
884 let input_str = if input.is_null() {
885 String::new()
886 } else {
887 input.to_string()
888 };
889 let truncated: String = input_str.chars().take(2000).collect();
890 let purpose = intent
891 .as_ref()
892 .map(|intent| format!(" purpose={}", intent.as_str()))
893 .unwrap_or_default();
894 parts.push(format!("[tool_call: {name}{purpose}({truncated})]"));
895 }
896 MessagePart::ToolResult {
897 content,
898 is_error,
899 tool_use_id,
900 } => {
901 let truncated: String = content.chars().take(3000).collect();
902 let marker = if *is_error { "ERROR" } else { "ok" };
903 let id_short: String = tool_use_id.chars().take(12).collect();
904 parts.push(format!("[tool_result {id_short}… {marker}: {truncated}]"));
905 }
906 MessagePart::Image { .. } => {
907 parts.push("[image]".into());
908 }
909 }
910 }
911 parts.join(" ")
912}
913
914pub fn replace_range_with_summary(
915 messages: &[Message],
916 range: &CompactRange,
917 summary: String,
918 turn_id: crate::event::TurnId,
919) -> Vec<Message> {
920 let retained_records = latest_context_records_before(messages, range.end);
921 let mut out =
922 Vec::with_capacity(1 + retained_records.len() + messages.len().saturating_sub(range.end));
923 out.push(Message::system_compact_summary(
924 turn_id,
925 summary,
926 range.start as u64,
927 range.end.saturating_sub(1) as u64,
928 range.end - range.start,
929 ));
930 out.extend(retained_records);
931 out.extend_from_slice(&messages[range.end..]);
932 out
933}
934
935fn latest_context_records_before(messages: &[Message], end: usize) -> Vec<Message> {
936 let suffix_keys: std::collections::HashSet<&str> = messages[end..]
937 .iter()
938 .flat_map(|message| &message.parts)
939 .filter_map(|part| match part {
940 MessagePart::ContextRecord(record) => Some(record.key()),
941 _ => None,
942 })
943 .collect();
944 let mut latest = std::collections::HashMap::<&str, (usize, &Message, &MessagePart)>::new();
945 for (index, message) in messages[..end].iter().enumerate() {
946 for part in &message.parts {
947 if let MessagePart::ContextRecord(record) = part
948 && record.retention() == crate::context_plan::ContextRecordRetention::Latest
949 && !suffix_keys.contains(record.key())
950 {
951 latest.insert(record.key(), (index, message, part));
952 }
953 }
954 }
955 let mut retained: Vec<_> = latest.into_values().collect();
956 retained.sort_by_key(|(index, _, _)| *index);
957 retained
958 .into_iter()
959 .map(|(_, message, part)| Message {
960 role: MessageRole::System,
961 parts: vec![part.clone()],
962 turn_id: message.turn_id.clone(),
963 origin: crate::message::MessageOrigin::Internal,
964 })
965 .collect()
966}
967
968#[derive(Debug, Clone, PartialEq, Eq)]
970pub struct HandleCompactResult {
971 pub before_tokens: u64,
972 pub after_tokens: u64,
973 pub compacted_start: usize,
974 pub compacted_end: usize,
975}
976
977#[derive(Debug, Clone, PartialEq)]
980pub struct HandleAutoCompactResult {
981 pub before_tokens: u64,
982 pub after_tokens: u64,
983 pub compacted_start: usize,
984 pub compacted_end: usize,
985 pub compacted_count: usize,
986 pub summary: String,
987 pub checkpoint_messages: Vec<Message>,
988}
989
990pub async fn maybe_auto_compact_handle_locked(
994 handle: &std::sync::Arc<std::sync::Mutex<Vec<Message>>>,
995 model: &str,
996 providers: &crate::provider::ProviderRegistry,
997 budget_context: CompactionBudgetContext,
998 forced: bool,
999) -> Option<HandleAutoCompactResult> {
1000 let snapshot = handle.lock().unwrap().clone();
1001 let info = crate::model_registry::model_info(model);
1002 let trigger = info.compaction_trigger_threshold();
1003 let target = budget_context
1004 .history_budget(&info)
1005 .map(|budget| budget.min(info.compaction_target_after()))
1006 .unwrap_or_else(|| info.compaction_target_after());
1007 let before_tokens = estimate_tokens_for_messages(&snapshot);
1008 let current = budget_context.estimated_input_tokens(before_tokens);
1009 if !forced && current <= trigger {
1010 return None;
1011 }
1012
1013 let (replacement, summary, compacted_start, compacted_end, compacted_count, must_fit_target) =
1014 if let Some(range) = find_compact_range(&snapshot, target) {
1015 let mut filtered = snapshot[range.start..range.end].to_vec();
1016 filter_orphan_tool_messages(&mut filtered);
1017 let (anchor, new_messages) = extract_anchor(&filtered)
1018 .map(|(anchor, remaining)| (Some(anchor), remaining.to_vec()))
1019 .unwrap_or_else(|| (None, filtered));
1020 let summary = generate_llm_summary(anchor.as_deref(), &new_messages, model, providers)
1021 .await
1022 .unwrap_or_else(|_| {
1023 format!(
1024 "[atman: compacted {} messages; summary unavailable]",
1025 range.end - range.start
1026 )
1027 });
1028 let replacement =
1029 build_budgeted_replacement(&snapshot, &range, &summary, target, model, providers)
1030 .await;
1031 let compacted_end = range.end.saturating_sub(1);
1032 let compacted_count = range.end - range.start;
1033 (
1034 replacement,
1035 summary,
1036 range.start,
1037 compacted_end,
1038 compacted_count,
1039 false,
1040 )
1041 } else {
1042 let (replacement, rewritten_count) =
1043 build_budgeted_turn_rewrite(&snapshot, target, model, providers).await;
1044 if rewritten_count == 0 {
1045 return None;
1046 }
1047 (
1048 replacement,
1049 format!(
1050 "[atman: persistently compacted output from {rewritten_count} retained messages]"
1051 ),
1052 0,
1053 0,
1054 rewritten_count,
1055 true,
1056 )
1057 };
1058 let after_tokens = estimate_tokens_for_messages(&replacement);
1059 if after_tokens >= before_tokens || (must_fit_target && after_tokens > target) {
1060 return None;
1061 }
1062
1063 let mut messages = handle.lock().unwrap();
1064 if *messages != snapshot {
1065 return None;
1066 }
1067 *messages = replacement.clone();
1068 Some(HandleAutoCompactResult {
1069 before_tokens,
1070 after_tokens,
1071 compacted_start,
1072 compacted_end,
1073 compacted_count,
1074 summary,
1075 checkpoint_messages: replacement,
1076 })
1077}
1078
1079pub fn compact_messages_on_handle(
1083 handle: &std::sync::Arc<std::sync::Mutex<Vec<Message>>>,
1084 summary: String,
1085 budget: u64,
1086) -> Option<HandleCompactResult> {
1087 let mut msgs = handle.lock().unwrap();
1088 let before_tokens = estimate_tokens_for_messages(&msgs);
1089 let range = find_compact_range(&msgs, budget)?;
1090 let turn_id = msgs
1091 .get(range.start)
1092 .map(|m| m.turn_id.clone())
1093 .unwrap_or_else(crate::event::TurnId::now);
1094 let after = replace_range_with_summary(&msgs, &range, summary, turn_id);
1095 let after_tokens = estimate_tokens_for_messages(&after);
1096 if after_tokens >= before_tokens {
1097 return None;
1098 }
1099 let result = HandleCompactResult {
1100 before_tokens,
1101 after_tokens,
1102 compacted_start: range.start,
1103 compacted_end: range.end.saturating_sub(1),
1104 };
1105 *msgs = after;
1106 Some(result)
1107}
1108
1109#[cfg(test)]
1110mod tests {
1111 use super::*;
1112 use crate::event::TurnId;
1113 use crate::message::MessageOrigin;
1114
1115 fn user(text: &str) -> Message {
1116 Message::user_text(TurnId::now(), text)
1117 }
1118 fn assistant(text: &str) -> Message {
1119 Message::assistant_text(TurnId::now(), text)
1120 }
1121 fn system(text: &str) -> Message {
1122 Message::system_text(TurnId::now(), text)
1123 }
1124
1125 fn context_record(key: &str, revision: u64, text: &str) -> Message {
1126 Message::context_record(
1127 TurnId::now(),
1128 crate::context_plan::ContextRecord::new(
1129 key,
1130 revision,
1131 crate::context_plan::ContextRecordAuthority::Runtime,
1132 crate::context_plan::ContextRecordRetention::Latest,
1133 crate::context_plan::ContextRecordBody::text(text),
1134 ),
1135 )
1136 }
1137
1138 fn context_tombstone(key: &str, revision: u64) -> Message {
1139 Message::context_record(
1140 TurnId::now(),
1141 crate::context_plan::ContextRecord::new(
1142 key,
1143 revision,
1144 crate::context_plan::ContextRecordAuthority::Runtime,
1145 crate::context_plan::ContextRecordRetention::Latest,
1146 crate::context_plan::ContextRecordBody::tombstone(),
1147 ),
1148 )
1149 }
1150
1151 #[test]
1152 fn replacement_keeps_only_the_latest_live_record_per_key() {
1153 let messages = vec![
1154 user("old"),
1155 context_record("session.goal", 1, "first"),
1156 context_record("session.goal", 2, "second"),
1157 assistant("old answer"),
1158 user("current"),
1159 ];
1160 let replacement = replace_range_with_summary(
1161 &messages,
1162 &CompactRange {
1163 start: 0,
1164 end: 4,
1165 tokens_saved_estimate: 1,
1166 },
1167 "summary".into(),
1168 TurnId::now(),
1169 );
1170
1171 assert_eq!(replacement.len(), 3);
1172 assert!(is_compaction_summary(&replacement[0]));
1173 assert!(matches!(
1174 replacement[1].parts.as_slice(),
1175 [MessagePart::ContextRecord(record)]
1176 if record.key() == "session.goal" && record.revision() == 2
1177 ));
1178 assert_eq!(replacement[2].text_concat(), "current");
1179 }
1180
1181 #[test]
1182 fn compaction_budget_reserves_output_safety_and_fixed_input_only_at_compact_time() {
1183 let info = crate::model_registry::ModelInfo {
1184 name: "test".into(),
1185 context_budget: 100_000,
1186 compact_threshold_ratio: 0.8,
1187 reasoning: crate::provider::ReasoningSelection::ProviderDefault,
1188 capabilities: crate::provider::ModelCapabilities::default(),
1189 image_detail: crate::provider::ImageDetail::Auto,
1190 max_output_tokens: Some(10_000),
1191 };
1192 let budget = CompactionBudgetContext {
1193 fixed_input_tokens: Some(5_000),
1194 }
1195 .history_budget(&info);
1196 assert_eq!(budget, Some(100_000 - 10_000 - 2_000 - 5_000));
1197 }
1198
1199 #[test]
1200 fn compaction_budget_saturates_when_fixed_input_exceeds_context() {
1201 let info = crate::model_registry::ModelInfo {
1202 name: "test".into(),
1203 context_budget: 20_000,
1204 compact_threshold_ratio: 0.8,
1205 reasoning: crate::provider::ReasoningSelection::ProviderDefault,
1206 capabilities: crate::provider::ModelCapabilities::default(),
1207 image_detail: crate::provider::ImageDetail::Auto,
1208 max_output_tokens: None,
1209 };
1210 assert_eq!(
1211 CompactionBudgetContext {
1212 fixed_input_tokens: Some(100_000)
1213 }
1214 .history_budget(&info),
1215 Some(0)
1216 );
1217 }
1218
1219 #[test]
1220 fn compaction_preflight_estimate_uses_current_messages_and_fixed_prefix() {
1221 let budget = CompactionBudgetContext {
1222 fixed_input_tokens: Some(7_000),
1223 };
1224 assert_eq!(budget.estimated_input_tokens(11_000), 18_000);
1225 assert_eq!(
1226 CompactionBudgetContext::default().estimated_input_tokens(11_000),
1227 11_000
1228 );
1229 }
1230
1231 #[tokio::test]
1232 async fn budgeted_replacement_groups_by_user_boundary_and_persists_omission() {
1233 let first_turn = TurnId::now();
1234 let second_turn = TurnId::now();
1235 let mut messages = vec![system(&"old".repeat(20_000)), assistant("old answer")];
1236 messages.push(Message::user_text(first_turn.clone(), "first user"));
1237 messages.push(Message::assistant_text(
1238 second_turn.clone(),
1239 "first output".repeat(4_000),
1240 ));
1241 messages.push(assistant_with_tool_use(
1242 "calling tool",
1243 "fs.read",
1244 serde_json::json!({"path": "/tmp/example"}),
1245 ));
1246 messages.push(tool_result(
1247 "call_test",
1248 &"tool output".repeat(4_000),
1249 false,
1250 ));
1251 messages.push(Message::user_text(second_turn.clone(), "current user"));
1252 messages.push(Message::assistant_text(
1253 first_turn,
1254 "current output".repeat(4_000),
1255 ));
1256 let range = CompactRange {
1257 start: 0,
1258 end: 2,
1259 tokens_saved_estimate: 1,
1260 };
1261
1262 let replacement = build_budgeted_replacement(
1263 &messages,
1264 &range,
1265 "anchor",
1266 500,
1267 "missing-provider",
1268 &crate::provider::ProviderRegistry::default(),
1269 )
1270 .await;
1271
1272 let texts: Vec<String> = replacement.iter().map(Message::text_concat).collect();
1273 assert!(texts.iter().any(|text| text == "current user"));
1274 assert!(
1275 texts
1276 .iter()
1277 .any(|text| text.contains("omitted 3 oversized"))
1278 );
1279 assert!(!replacement.iter().any(|message| {
1280 message.parts.iter().any(|part| {
1281 matches!(
1282 part,
1283 MessagePart::ToolUse { .. } | MessagePart::ToolResult { .. }
1284 )
1285 })
1286 }));
1287 assert_eq!(user_turn_ranges(&replacement).len(), 2);
1288 }
1289
1290 #[tokio::test]
1291 async fn budgeted_replacement_floor_keeps_anchor_records_and_user_inputs() {
1292 let messages = vec![
1293 system("old system"),
1294 context_record("session.goal", 1, "old goal"),
1295 context_record("session.goal", 2, "current goal"),
1296 context_tombstone("session.workspace", 3),
1297 user("first user"),
1298 assistant("first output"),
1299 user("current user"),
1300 assistant("current output"),
1301 ];
1302 let replacement = build_budgeted_replacement(
1303 &messages,
1304 &CompactRange {
1305 start: 0,
1306 end: 4,
1307 tokens_saved_estimate: 1,
1308 },
1309 "anchor",
1310 1,
1311 "missing-provider",
1312 &crate::provider::ProviderRegistry::default(),
1313 )
1314 .await;
1315
1316 assert!(is_compaction_summary(&replacement[0]));
1317 let records: Vec<_> = replacement
1318 .iter()
1319 .flat_map(|message| &message.parts)
1320 .filter_map(|part| match part {
1321 MessagePart::ContextRecord(record) => Some(record),
1322 _ => None,
1323 })
1324 .collect();
1325 assert_eq!(records.len(), 2);
1326 assert_eq!(records[0].key(), "session.goal");
1327 assert_eq!(records[0].revision(), 2);
1328 assert_eq!(records[1].key(), "session.workspace");
1329 assert!(records[1].body().is_tombstone());
1330 assert_eq!(
1331 replacement
1332 .iter()
1333 .filter(|message| message.role == MessageRole::User)
1334 .map(Message::text_concat)
1335 .collect::<Vec<_>>(),
1336 ["first user", "current user"]
1337 );
1338 assert!(
1339 !replacement
1340 .iter()
1341 .any(|message| message.role == MessageRole::Assistant)
1342 );
1343 }
1344
1345 #[tokio::test]
1346 async fn turn_rewrite_compacts_oversized_tool_output_without_dropping_recent_users() {
1347 let first = TurnId::now();
1348 let current = TurnId::now();
1349 let messages = vec![
1350 Message::user_text(first, "first user"),
1351 assistant_with_tool_use(
1352 &"calling tool".repeat(2_000),
1353 "fs.read",
1354 serde_json::json!({"path": "/tmp/example"}),
1355 ),
1356 tool_result("call_test", &"tool output".repeat(8_000), false),
1357 Message::user_text(current, "current user"),
1358 ];
1359
1360 assert!(find_compact_range(&messages, 500).is_none());
1361 let (replacement, rewritten_count) = build_budgeted_turn_rewrite(
1362 &messages,
1363 500,
1364 "missing-provider",
1365 &crate::provider::ProviderRegistry::default(),
1366 )
1367 .await;
1368
1369 assert_eq!(rewritten_count, 2);
1370 assert!(estimate_tokens_for_messages(&replacement) <= 500);
1371 let users: Vec<String> = replacement
1372 .iter()
1373 .filter(|message| message.role == MessageRole::User)
1374 .map(Message::text_concat)
1375 .collect();
1376 assert_eq!(users, vec!["first user", "current user"]);
1377 assert!(replacement.iter().any(|message| {
1378 message
1379 .text_concat()
1380 .contains("omitted 2 oversized assistant/system/tool messages")
1381 }));
1382 assert!(!replacement.iter().any(|message| {
1383 message.parts.iter().any(|part| {
1384 matches!(
1385 part,
1386 MessagePart::ToolUse { .. } | MessagePart::ToolResult { .. }
1387 )
1388 })
1389 }));
1390 }
1391
1392 #[test]
1393 fn user_turn_ranges_ignore_misanchored_turn_ids() {
1394 let first = TurnId::now();
1395 let second = TurnId::now();
1396 let messages = vec![
1397 Message::user_text(first.clone(), "u1"),
1398 Message::assistant_text(second.clone(), "a1"),
1399 Message::user_text(second, "u2"),
1400 Message::assistant_text(first, "a2"),
1401 ];
1402 assert_eq!(user_turn_ranges(&messages), vec![(0, 2), (2, 4)]);
1403 }
1404
1405 #[test]
1406 fn summary_instructions_keep_decisions_before_next_move() {
1407 let objective = SUMMARY_INSTRUCTIONS
1408 .find("## Objective")
1409 .expect("objective");
1410 let decisions = SUMMARY_INSTRUCTIONS
1411 .find("## Decisions")
1412 .expect("decisions");
1413 let next_move = SUMMARY_INSTRUCTIONS
1414 .find("## Next Move")
1415 .expect("next move");
1416 assert!(objective < decisions);
1417 assert!(decisions < next_move);
1418 }
1419
1420 #[test]
1421 fn estimate_scales_with_char_length() {
1422 let short = user("hi");
1423 let long = user(&"x".repeat(3500));
1424 assert!(estimate_tokens_for_message(&long) > estimate_tokens_for_message(&short) * 100);
1425 }
1426
1427 #[test]
1428 fn find_compact_returns_none_when_under_budget() {
1429 let msgs = vec![user("a"), assistant("b"), user("c"), assistant("d")];
1430 assert!(find_compact_range(&msgs, 1000).is_none());
1431 }
1432
1433 #[test]
1434 fn find_compact_returns_none_for_short_history() {
1435 let msgs = vec![user(&"x".repeat(9000))];
1436 assert!(find_compact_range(&msgs, 100).is_none());
1437 }
1438
1439 #[test]
1440 fn find_kth_recent_user_handles_exact_excess_and_mixed_history() {
1441 let exact = vec![
1442 user("u0"),
1443 assistant("a0"),
1444 system("s0"),
1445 tool_result("call-0", "result", false),
1446 user("u1"),
1447 assistant("a1"),
1448 user("u2"),
1449 system("s1"),
1450 user("u3"),
1451 assistant("a3"),
1452 user("u4"),
1453 ];
1454 assert_eq!(find_kth_recent_user(&exact, KEEP_RECENT_USER_TURNS), 0);
1456
1457 let mut excess = exact.clone();
1458 excess.push(user("u5"));
1459 assert_eq!(find_kth_recent_user(&excess, KEEP_RECENT_USER_TURNS), 4);
1460
1461 let too_few = vec![user("a"), assistant("b"), assistant("c")];
1462 assert_eq!(find_kth_recent_user(&too_few, KEEP_RECENT_USER_TURNS), 0);
1463 }
1464
1465 #[test]
1466 fn find_compact_range_preserves_minimum_recent_messages_without_anchor() {
1467 let mut msgs = vec![system("head")];
1468 msgs.extend((0..25).map(|index| assistant(&format!("old {index}"))));
1469 msgs.extend(
1470 (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
1471 );
1472 let range = find_compact_range(&msgs, 1).expect("range");
1473 assert_eq!(range.start, 0);
1474 assert_eq!(range.end, msgs.len() - KEEP_RECENT_MESSAGES);
1475 }
1476
1477 #[test]
1478 fn find_compact_range_keeps_tool_use_with_later_result() {
1479 let mut msgs = (0..11)
1480 .map(|index| assistant(&format!("old {index}")))
1481 .collect::<Vec<_>>();
1482 msgs.push(assistant_with_tool_use(
1483 "calling tool",
1484 "fs.read",
1485 serde_json::json!({"path": "/tmp/example"}),
1486 ));
1487 msgs.push(tool_result("call_test", "result", false));
1488 msgs.extend((0..9).map(|index| {
1489 if index % 2 == 0 {
1490 user(&format!("recent user {index}"))
1491 } else {
1492 assistant(&format!("recent assistant {index}"))
1493 }
1494 }));
1495
1496 let range = find_compact_range(&msgs, 1).expect("range");
1497 assert_eq!(range.end, 11);
1498 assert!(matches!(
1499 msgs[range.end].parts.as_slice(),
1500 [MessagePart::Text { .. }, MessagePart::ToolUse { id, .. }] if id == "call_test"
1501 ));
1502 }
1503
1504 #[test]
1505 fn find_compact_range_keeps_parallel_tool_batch_together() {
1506 let mut msgs = (0..10)
1507 .map(|index| assistant(&format!("old {index}")))
1508 .collect::<Vec<_>>();
1509 msgs.push(assistant_with_tool_uses(&["call_a", "call_b"]));
1510 msgs.push(tool_result("call_a", "first result", false));
1511 msgs.push(tool_result("call_b", "second result", false));
1512 msgs.extend((0..9).map(|index| {
1513 if index % 2 == 0 {
1514 user(&format!("recent user {index}"))
1515 } else {
1516 assistant(&format!("recent assistant {index}"))
1517 }
1518 }));
1519
1520 let range = find_compact_range(&msgs, 1).expect("range");
1521 assert_eq!(range.end, 10);
1522 assert_eq!(
1523 msgs[range.end]
1524 .parts
1525 .iter()
1526 .filter(|part| matches!(part, MessagePart::ToolUse { .. }))
1527 .count(),
1528 2
1529 );
1530 }
1531
1532 #[test]
1533 fn find_compact_range_keeps_boundary_after_closed_tool_batch() {
1534 let mut msgs = (0..9)
1535 .map(|index| assistant(&format!("old {index}")))
1536 .collect::<Vec<_>>();
1537 msgs.push(assistant_with_tool_uses(&["call_a", "call_b"]));
1538 msgs.push(Message {
1539 role: MessageRole::Tool,
1540 parts: vec![
1541 MessagePart::ToolResult {
1542 tool_use_id: "call_a".into(),
1543 content: "first result".into(),
1544 is_error: false,
1545 },
1546 MessagePart::ToolResult {
1547 tool_use_id: "call_b".into(),
1548 content: "second result".into(),
1549 is_error: false,
1550 },
1551 ],
1552 turn_id: TurnId::now(),
1553 origin: MessageOrigin::User,
1554 });
1555 msgs.push(assistant("batch complete"));
1556 msgs.extend((0..10).map(|index| {
1557 if index % 2 == 0 {
1558 user(&format!("recent user {index}"))
1559 } else {
1560 assistant(&format!("recent assistant {index}"))
1561 }
1562 }));
1563
1564 let range = find_compact_range(&msgs, 1).expect("range");
1565 assert_eq!(range.end, 12);
1566 assert_eq!(msgs[range.end].role, MessageRole::User);
1567 }
1568
1569 #[test]
1570 fn find_compact_range_preserves_minimum_window_for_large_tool_result() {
1571 let mut msgs = vec![user(&"h".repeat(500_000))];
1572 for index in 1..31 {
1573 if matches!(index, 20 | 22 | 24 | 26 | 28 | 30) {
1574 msgs.push(user(&"u".repeat(800)));
1575 } else {
1576 msgs.push(assistant(&"a".repeat(800)));
1577 }
1578 }
1579 msgs.push(tool_result("call-large", &"t".repeat(22_000), false));
1580
1581 let budget = 120_000;
1582 let minimum_recent_tokens = (budget as f64 * KEEP_RECENT_TOKEN_FRACTION).ceil() as u64;
1583 assert!(estimate_tokens_for_message(msgs.last().unwrap()) > minimum_recent_tokens);
1584
1585 let range = find_compact_range(&msgs, budget).expect("range");
1586 assert!(
1587 range.end <= msgs.len() - KEEP_RECENT_MESSAGES,
1588 "range was {range:?}"
1589 );
1590 assert!(msgs.len() - range.end >= KEEP_RECENT_MESSAGES);
1591 assert!(estimate_tokens_for_messages(&msgs[range.end..]) >= minimum_recent_tokens);
1592 }
1593
1594 #[test]
1595 fn find_compact_range_handles_four_to_twenty_one_message_histories() {
1596 for len in 4..=21 {
1597 let msgs = (0..len)
1598 .map(|_| user(&"x".repeat(5000)))
1599 .collect::<Vec<_>>();
1600 assert_eq!(
1601 find_compact_range(&msgs, 1).is_some(),
1602 len >= KEEP_RECENT_MESSAGES + 2,
1603 "len={len}"
1604 );
1605 }
1606 }
1607
1608 #[test]
1609 fn find_compact_range_recent_users_limit_mixed_history() {
1610 let msgs = vec![
1611 system("head"),
1612 assistant("a0"),
1613 user("u0"),
1614 tool_result("call-0", "r0", false),
1615 assistant("a1"),
1616 user("u1"),
1617 assistant("a2"),
1618 tool_result("call-1", "r1", false),
1619 user("u2"),
1620 assistant("a3"),
1621 system("note"),
1622 user("u3"),
1623 tool_result("call-2", "r2", false),
1624 assistant("a4"),
1625 user("u4"),
1626 assistant("a5"),
1627 tool_result("call-3", "r3", false),
1628 user("u5"),
1629 assistant("a6"),
1630 system("tail"),
1631 assistant("a7"),
1632 ];
1633
1634 let range = find_compact_range(&msgs, 1).expect("range");
1635 assert_eq!(range.end, 5);
1636 assert_eq!(msgs[range.end].role, MessageRole::User);
1637 }
1638
1639 #[test]
1640 fn find_compact_range_preserves_recent_user_turns_after_anchor() {
1641 let mut msgs = vec![system("head"), compaction_summary("summary")];
1642 msgs.extend((0..6).flat_map(|index| {
1643 [
1644 user(&format!("user {index}")),
1645 assistant("assistant"),
1646 assistant("tool fragment"),
1647 ]
1648 }));
1649 msgs.extend((0..12).map(|_| assistant("recent fragment")));
1650 let range = find_compact_range(&msgs, 1).expect("range");
1651 assert_eq!(range.start, 1);
1652 assert_eq!(range.end, 5);
1653 assert_eq!(msgs[range.end].role, MessageRole::User);
1654 }
1655
1656 #[test]
1657 fn find_compact_range_returns_none_when_end_cannot_cover_two_messages() {
1658 let msgs = vec![
1659 system("head"),
1660 compaction_summary("summary"),
1661 assistant("tail"),
1662 user("tail"),
1663 ];
1664 assert!(find_compact_range(&msgs, 1).is_none());
1665 }
1666
1667 #[test]
1668 fn extract_anchor_removes_leading_compact_summary() {
1669 let messages = vec![compaction_summary("anchor"), user("new")];
1670 let (anchor, remaining) = extract_anchor(&messages).expect("anchor");
1671 assert_eq!(anchor, "anchor");
1672 assert_eq!(remaining, &messages[1..]);
1673 }
1674
1675 #[test]
1676 fn extract_anchor_returns_none_without_leading_summary() {
1677 let messages = vec![user("new")];
1678 assert!(extract_anchor(&messages).is_none());
1679 }
1680
1681 #[test]
1682 fn compact_messages_on_handle_replaces_range_in_place() {
1683 let mut messages = vec![system("head")];
1684 messages.extend((0..9).map(|index| assistant(&format!("old {index}"))));
1685 messages.extend(
1686 (0..6).flat_map(|index| [user(&format!("user {index}")), assistant(&"x".repeat(4000))]),
1687 );
1688 messages.extend((0..10).map(|index| assistant(&format!("recent {index}"))));
1689 messages.push(user("tail"));
1690 let handle: std::sync::Arc<std::sync::Mutex<Vec<Message>>> =
1691 std::sync::Arc::new(std::sync::Mutex::new(messages));
1692 let result = compact_messages_on_handle(&handle, "gist".into(), 100);
1695 let result = result.expect("should compact");
1696 assert!(result.after_tokens < result.before_tokens);
1697 let msgs = handle.lock().unwrap();
1698 assert_eq!(result.compacted_start, 0);
1699 assert_eq!(result.compacted_end, 13);
1700 assert!(is_compaction_summary(&msgs[0]));
1701 assert_eq!(msgs.last().unwrap().text_concat(), "tail");
1702 }
1703
1704 #[test]
1705 fn compact_messages_on_handle_none_when_under_budget() {
1706 let handle: std::sync::Arc<std::sync::Mutex<Vec<Message>>> =
1707 std::sync::Arc::new(std::sync::Mutex::new(vec![user("short")]));
1708 assert!(compact_messages_on_handle(&handle, "g".into(), 100_000).is_none());
1709 }
1710
1711 #[test]
1712 fn compact_messages_on_handle_none_when_summary_would_not_shrink() {
1713 let handle: std::sync::Arc<std::sync::Mutex<Vec<Message>>> =
1719 std::sync::Arc::new(std::sync::Mutex::new(vec![
1720 system("h"),
1721 user("."),
1722 assistant("."),
1723 user("."),
1724 assistant("."),
1725 user("t"),
1726 ]));
1727 let result = compact_messages_on_handle(&handle, "x".into(), 1);
1729 let before_len = handle.lock().unwrap().len();
1732 if result.is_none() {
1733 assert_eq!(handle.lock().unwrap().len(), before_len);
1734 }
1735 }
1736
1737 #[test]
1738 fn replace_range_puts_summary_system_message_in_place() {
1739 let msgs = vec![
1740 system("head"),
1741 user("m1"),
1742 assistant("m2"),
1743 user("m3"),
1744 assistant("m4"),
1745 user("tail"),
1746 ];
1747 let range = CompactRange {
1748 start: 1,
1749 end: 5,
1750 tokens_saved_estimate: 100,
1751 };
1752 let out = replace_range_with_summary(
1753 &msgs,
1754 &range,
1755 "gist: talked about m1..m4".into(),
1756 TurnId::now(),
1757 );
1758 assert_eq!(out.len(), 2, "summary + tail");
1759 assert_eq!(out[0].role, MessageRole::System);
1760 assert!(out[0].text_concat().contains("gist: talked about"));
1761 assert!(matches!(
1762 out[0].parts.as_slice(),
1763 [MessagePart::CompactSummary {
1764 seq_start: 1,
1765 seq_end: 4,
1766 count: 4,
1767 ..
1768 }]
1769 ));
1770 assert_eq!(out[1].role, MessageRole::User);
1771 assert_eq!(out[1].text_concat(), "tail");
1772 }
1773
1774 #[test]
1775 fn find_compact_range_anchors_on_latest_structured_summary() {
1776 let mut msgs = vec![
1777 system("head"),
1778 Message::system_compact_summary(TurnId::now(), "old", 0, 1, 2),
1779 ];
1780 msgs.extend(
1781 (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
1782 );
1783 msgs.extend((0..12).map(|index| assistant(&format!("recent {index}"))));
1784 let range = find_compact_range(&msgs, 1).expect("range");
1785 assert_eq!(range.start, 1);
1786 assert_eq!(range.end, 4);
1787 }
1788
1789 fn assistant_with_tool_use(text: &str, tool_name: &str, input: serde_json::Value) -> Message {
1790 Message {
1791 role: MessageRole::Assistant,
1792 parts: vec![
1793 MessagePart::Text { text: text.into() },
1794 MessagePart::ToolUse {
1795 id: "call_test".into(),
1796 name: tool_name.into(),
1797 input,
1798 intent: None,
1799 },
1800 ],
1801 turn_id: TurnId::now(),
1802 origin: MessageOrigin::User,
1803 }
1804 }
1805
1806 fn assistant_with_tool_uses(ids: &[&str]) -> Message {
1807 Message {
1808 role: MessageRole::Assistant,
1809 parts: ids
1810 .iter()
1811 .map(|id| MessagePart::ToolUse {
1812 id: (*id).into(),
1813 name: "fs.read".into(),
1814 input: serde_json::json!({"path": format!("/tmp/{id}")}),
1815 intent: None,
1816 })
1817 .collect(),
1818 turn_id: TurnId::now(),
1819 origin: MessageOrigin::User,
1820 }
1821 }
1822
1823 fn tool_result(id: &str, content: &str, is_error: bool) -> Message {
1824 Message {
1825 role: MessageRole::Tool,
1826 parts: vec![MessagePart::ToolResult {
1827 tool_use_id: id.into(),
1828 content: content.into(),
1829 is_error,
1830 }],
1831 turn_id: TurnId::now(),
1832 origin: MessageOrigin::User,
1833 }
1834 }
1835
1836 fn thinking(text: &str) -> Message {
1837 Message {
1838 role: MessageRole::Assistant,
1839 parts: vec![
1840 MessagePart::Thinking {
1841 thinking: text.into(),
1842 signature: None,
1843 },
1844 MessagePart::Text {
1845 text: "after thinking".into(),
1846 },
1847 ],
1848 turn_id: TurnId::now(),
1849 origin: MessageOrigin::User,
1850 }
1851 }
1852
1853 #[test]
1854 fn format_slice_for_summary_includes_tool_use() {
1855 let slice = vec![
1856 user("read the file"),
1857 assistant_with_tool_use(
1858 "let me check",
1859 "fs.read",
1860 serde_json::json!({"path": "/tmp/foo.rs"}),
1861 ),
1862 tool_result("call_test", "fn main() {}", false),
1863 ];
1864 let out = format_slice_for_summary(&slice);
1865 assert!(out.contains("fs.read"), "missing tool name: {out}");
1866 assert!(out.contains("/tmp/foo.rs"), "missing tool input: {out}");
1867 assert!(
1868 out.contains("fn main()"),
1869 "missing tool_result content: {out}"
1870 );
1871 assert!(out.contains("tool_call"), "missing tool_call marker: {out}");
1872 assert!(
1873 out.contains("tool_result"),
1874 "missing tool_result marker: {out}"
1875 );
1876 }
1877
1878 #[test]
1879 fn format_slice_for_summary_includes_thinking() {
1880 let slice = vec![thinking("I should consider the edge case")];
1881 let out = format_slice_for_summary(&slice);
1882 assert!(out.contains("thinking"), "missing thinking marker: {out}");
1883 assert!(out.contains("edge case"), "missing thinking content: {out}");
1884 }
1885
1886 #[test]
1887 fn format_slice_for_summary_marks_error_tool_results() {
1888 let slice = vec![tool_result("call_1", "permission denied", true)];
1889 let out = format_slice_for_summary(&slice);
1890 assert!(out.contains("ERROR"), "missing ERROR marker: {out}");
1891 }
1892
1893 #[test]
1894 fn format_slice_for_summary_truncates_long_tool_input() {
1895 let long_input = serde_json::json!({"content": "x".repeat(5000)});
1896 let slice = vec![assistant_with_tool_use("check", "fs.write", long_input)];
1897 let out = format_slice_for_summary(&slice);
1898 let tool_call_line = out
1899 .lines()
1900 .find(|l| l.contains("tool_call"))
1901 .unwrap_or_else(|| panic!("no tool_call line in {out}"));
1902 assert!(
1903 tool_call_line.chars().count() < 2200,
1904 "tool_call line not truncated: {tool_call_line}"
1905 );
1906 }
1907
1908 fn compaction_summary(text: &str) -> Message {
1909 Message::system_compact_summary(TurnId::now(), text, 1, 5, 5)
1910 }
1911
1912 #[test]
1913 fn is_compaction_summary_detects_structured_variant() {
1914 assert!(is_compaction_summary(&compaction_summary("gist")));
1915 assert!(!is_compaction_summary(&system("plain system msg")));
1916 assert!(!is_compaction_summary(&user("user msg")));
1917 }
1918
1919 #[test]
1920 fn find_compact_range_spans_across_compaction_summaries() {
1921 let mut msgs = vec![
1922 system("head"),
1923 user(&"x".repeat(3000)),
1924 assistant(&"y".repeat(3000)),
1925 user(&"z".repeat(3000)),
1926 compaction_summary("first compaction summary"),
1927 ];
1928 msgs.extend(
1929 (0..6).flat_map(|index| [user(&format!("user {index}")), assistant(&"x".repeat(3000))]),
1930 );
1931 msgs.extend((0..10).map(|index| assistant(&format!("recent {index}"))));
1932 let range = find_compact_range(&msgs, 500).expect("expected range across summary");
1933 assert_eq!(
1934 range.start, 4,
1935 "range should anchor at the structured summary"
1936 );
1937 assert!(
1938 range.end > 4,
1939 "range should include later work, got {range:?}"
1940 );
1941 assert!(
1942 range.end - range.start >= 3,
1943 "range must cover >= 3 msgs, got {}",
1944 range.end - range.start
1945 );
1946 }
1947
1948 #[test]
1949 fn find_compact_starts_from_summary() {
1950 let mut msgs = vec![user("a"), assistant("b"), compaction_summary("summary 1")];
1951 msgs.extend(
1952 (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
1953 );
1954 msgs.extend((0..12).map(|index| assistant(&format!("recent {index}"))));
1955 let range = find_compact_range(&msgs, 10).expect("expected range");
1956 assert_eq!(
1957 range.start, 2,
1958 "should start from the compact summary anchor"
1959 );
1960 assert_eq!(range.end, 5, "the fifth recent user is retained");
1961 }
1962
1963 #[test]
1964 fn find_compact_range_includes_older_compaction_summaries() {
1965 let mut msgs = vec![compaction_summary("summary 0")];
1966 msgs.extend((0..6).flat_map(|index| {
1967 [
1968 user(&format!("old user {index}")),
1969 assistant(&"x".repeat(2000)),
1970 ]
1971 }));
1972 msgs.push(compaction_summary("summary 1"));
1973 msgs.extend((0..6).flat_map(|index| {
1974 [
1975 user(&format!("new user {index}")),
1976 assistant(&"z".repeat(2000)),
1977 ]
1978 }));
1979 msgs.extend((0..10).map(|index| assistant(&format!("tail {index}"))));
1980 let range = find_compact_range(&msgs, 500).expect("expected range");
1981 assert_eq!(range.start, 13, "should compact from the latest summary");
1982 assert!(
1983 range.end > range.start,
1984 "should include work after the latest summary"
1985 );
1986 }
1987
1988 #[test]
1989 fn compacted_message_tokens_detects_growth() {
1990 let msgs = vec![compaction_summary("summary 0"), user("a"), assistant("b")];
1991 let range = CompactRange {
1992 start: 1,
1993 end: 3,
1994 tokens_saved_estimate: 0,
1995 };
1996 let before = estimate_tokens_for_messages(&msgs);
1997 let after = estimate_compacted_message_tokens(
1998 &msgs,
1999 &range,
2000 "a very long summary that expands the transcript a lot",
2001 );
2002 assert!(after > before, "expected growth to be detectable");
2003 }
2004
2005 #[test]
2006 fn find_compact_starts_from_zero_without_summary() {
2007 let mut msgs = (0..26)
2008 .map(|index| assistant(&format!("old {index}")))
2009 .collect::<Vec<_>>();
2010 msgs.extend(
2011 (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
2012 );
2013 let range = find_compact_range(&msgs, 10).expect("expected range");
2014 assert_eq!(range.start, 0, "should start from 0 without summary");
2015 assert_eq!(range.end, 28, "the recent-message limit is retained");
2016 }
2017
2018 #[test]
2019 fn filter_orphan_tool_messages_removes_orphan_results() {
2020 use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
2021 let turn = TurnId::now();
2022 let msgs = vec![
2023 Message {
2024 role: MessageRole::Tool,
2025 parts: vec![MessagePart::ToolResult {
2026 tool_use_id: "orphan".into(),
2027 content: "no matching use".into(),
2028 is_error: false,
2029 }],
2030 turn_id: turn.clone(),
2031 origin: MessageOrigin::User,
2032 },
2033 Message {
2034 role: MessageRole::Assistant,
2035 parts: vec![MessagePart::ToolUse {
2036 id: "call_1".into(),
2037 name: "fs.read".into(),
2038 input: serde_json::json!({}),
2039 intent: None,
2040 }],
2041 turn_id: turn.clone(),
2042 origin: MessageOrigin::User,
2043 },
2044 Message {
2045 role: MessageRole::Tool,
2046 parts: vec![MessagePart::ToolResult {
2047 tool_use_id: "call_1".into(),
2048 content: "ok".into(),
2049 is_error: false,
2050 }],
2051 turn_id: turn,
2052 origin: MessageOrigin::User,
2053 },
2054 ];
2055 let mut filtered = msgs;
2056 filter_orphan_tool_messages(&mut filtered);
2057 assert_eq!(filtered.len(), 2, "orphan result should be removed");
2058 }
2059
2060 #[test]
2061 fn filter_orphan_tool_parts_preserves_valid_mixed_message_content() {
2062 use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
2063 let turn = TurnId::now();
2064 let mut messages = vec![
2065 Message {
2066 role: MessageRole::Assistant,
2067 parts: vec![
2068 MessagePart::Text {
2069 text: "keep assistant text".into(),
2070 },
2071 MessagePart::ToolUse {
2072 id: "valid".into(),
2073 name: "fs.read".into(),
2074 input: serde_json::json!({}),
2075 intent: None,
2076 },
2077 MessagePart::ToolUse {
2078 id: "orphan-use".into(),
2079 name: "fs.read".into(),
2080 input: serde_json::json!({}),
2081 intent: None,
2082 },
2083 ],
2084 turn_id: turn.clone(),
2085 origin: MessageOrigin::User,
2086 },
2087 Message {
2088 role: MessageRole::Tool,
2089 parts: vec![
2090 MessagePart::Text {
2091 text: "keep tool text".into(),
2092 },
2093 MessagePart::ToolResult {
2094 tool_use_id: "valid".into(),
2095 content: "ok".into(),
2096 is_error: false,
2097 },
2098 MessagePart::ToolResult {
2099 tool_use_id: "orphan-result".into(),
2100 content: "drop".into(),
2101 is_error: false,
2102 },
2103 ],
2104 turn_id: turn,
2105 origin: MessageOrigin::User,
2106 },
2107 ];
2108
2109 filter_orphan_tool_messages(&mut messages);
2110
2111 assert_eq!(messages.len(), 2);
2112 assert!(
2113 matches!(&messages[0].parts[..], [MessagePart::Text { .. }, MessagePart::ToolUse { id, .. }] if id == "valid")
2114 );
2115 assert!(
2116 matches!(&messages[1].parts[..], [MessagePart::Text { .. }, MessagePart::ToolResult { tool_use_id, .. }] if tool_use_id == "valid")
2117 );
2118 }
2119}