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::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/// Result of compacting a messages_handle in place.
1026#[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/// Result of applying the automatic compaction policy to an isolated message
1035/// handle. The caller must hold that handle's async compaction lock.
1036#[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
1047/// Apply the root compaction budget, range, summary, and replacement policy to
1048/// an isolated message handle. This function does not acquire the async lock
1049/// and does not emit session events.
1050pub 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
1136/// Compact a messages_handle in place (data-layer primitive, operates on any
1137/// FlowRun's segment). Returns `None` if under `budget` or no compactable
1138/// range. Caller should hold the FlowRun's `compact_lock`.
1139pub 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        // With exactly five users, the fifth recent user is the first message.
1512        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        // The recent-message and recent-user limits meet at the fifth recent user.
1750        // The compacted prefix is replaced by one summary while the tail remains.
1751        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        // 4 messages, all tiny → find_compact_range returns a range but the
1771        // summary message itself is comparable in size, so after >= before.
1772        // Construct a case where find_compact_range returns Some but shrink
1773        // check rejects it: make the range cover near-empty messages so the
1774        // summary overhead exceeds the savings.
1775        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        // budget=1 forces a range, but messages are so small the summary won't help
1785        let result = compact_messages_on_handle(&handle, "x".into(), 1);
1786        // Either no range found (len 6 but tiny), or shrink rejected.
1787        // The key invariant: handle is unchanged if None.
1788        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}