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