Skip to main content

atman_runtime/
compaction.rs

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/// Result of compacting a messages_handle in place.
969#[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/// Result of applying the automatic compaction policy to an isolated message
978/// handle. The caller must hold that handle's async compaction lock.
979#[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
990/// Apply the root compaction budget, range, summary, and replacement policy to
991/// an isolated message handle. This function does not acquire the async lock
992/// and does not emit session events.
993pub 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
1079/// Compact a messages_handle in place (data-layer primitive, operates on any
1080/// FlowRun's segment). Returns `None` if under `budget` or no compactable
1081/// range. Caller should hold the FlowRun's `compact_lock`.
1082pub 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        // With exactly five users, the fifth recent user is the first message.
1455        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        // The recent-message and recent-user limits meet at the fifth recent user.
1693        // The compacted prefix is replaced by one summary while the tail remains.
1694        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        // 4 messages, all tiny → find_compact_range returns a range but the
1714        // summary message itself is comparable in size, so after >= before.
1715        // Construct a case where find_compact_range returns Some but shrink
1716        // check rejects it: make the range cover near-empty messages so the
1717        // summary overhead exceeds the savings.
1718        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        // budget=1 forces a range, but messages are so small the summary won't help
1728        let result = compact_messages_on_handle(&handle, "x".into(), 1);
1729        // Either no range found (len 6 but tiny), or shrink rejected.
1730        // The key invariant: handle is unchanged if None.
1731        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}