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 history_budget(self, info: &crate::model_registry::ModelInfo) -> Option<u64> {
16        let fixed_input_tokens = self.fixed_input_tokens?;
17        let output_cap = (info.context_budget as f64 * 0.20) as u64;
18        let output_floor = 8_000_u64.min(output_cap);
19        let output_reserve = (info.max_output_tokens.unwrap_or(32_000) as u64)
20            .max(output_floor)
21            .min(output_cap);
22        let safety_cap = (info.context_budget as f64 * COMPACTION_SAFETY_MARGIN_MAX_RATIO) as u64;
23        let safety_margin = ((info.context_budget as f64 * 0.02) as u64)
24            .max(COMPACTION_SAFETY_MARGIN_MIN.min(safety_cap))
25            .min(safety_cap);
26        Some(
27            info.context_budget
28                .saturating_sub(output_reserve)
29                .saturating_sub(safety_margin)
30                .saturating_sub(fixed_input_tokens),
31        )
32    }
33}
34
35pub fn estimate_tokens_for_message(msg: &Message) -> u64 {
36    let mut chars = 0usize;
37    for part in &msg.parts {
38        chars += match part {
39            MessagePart::CompactSummary { summary, .. } => summary.len(),
40            MessagePart::Text { text } => text.len(),
41            MessagePart::Thinking { thinking, .. } => thinking.len(),
42            MessagePart::ToolResult { content, .. } => content.len(),
43            MessagePart::Image { .. } => 512,
44            MessagePart::ToolUse { name, input, .. } => name.len() + input.to_string().len(),
45        };
46    }
47    chars = chars.saturating_add(estimate_role_overhead(msg.role));
48    (chars as f64 / 3.5).ceil() as u64
49}
50
51fn estimate_role_overhead(role: MessageRole) -> usize {
52    match role {
53        MessageRole::System => 12,
54        MessageRole::User => 8,
55        MessageRole::Assistant => 8,
56        MessageRole::Tool => 16,
57    }
58}
59
60pub fn estimate_tokens_for_messages(messages: &[Message]) -> u64 {
61    messages.iter().map(estimate_tokens_for_message).sum()
62}
63
64#[derive(Debug, Clone, PartialEq, Eq)]
65pub struct CompactRange {
66    pub start: usize,
67    pub end: usize,
68    pub tokens_saved_estimate: u64,
69}
70
71pub fn is_plan_related(msg: &Message) -> bool {
72    for part in &msg.parts {
73        match part {
74            MessagePart::ToolUse { name, .. } if name.starts_with("plan.") => return true,
75            MessagePart::ToolResult { content, .. } if content.starts_with("# Plan:") => {
76                return true;
77            }
78            _ => {}
79        }
80    }
81    false
82}
83
84pub fn is_compaction_summary(msg: &Message) -> bool {
85    if !matches!(msg.role, MessageRole::System) {
86        return false;
87    }
88    msg.parts
89        .iter()
90        .any(|part| matches!(part, MessagePart::CompactSummary { .. }))
91}
92
93fn find_kth_recent_user(messages: &[Message], k: usize) -> usize {
94    let mut user_count = 0;
95    for (index, message) in messages.iter().enumerate().rev() {
96        if message.role == MessageRole::User {
97            user_count += 1;
98            if user_count == k {
99                return index;
100            }
101        }
102    }
103    0
104}
105
106pub fn find_compact_range(messages: &[Message], budget: u64) -> Option<CompactRange> {
107    let total = estimate_tokens_for_messages(messages);
108    if total <= budget || messages.len() < 4 {
109        return None;
110    }
111
112    let start = messages
113        .iter()
114        .rposition(is_compaction_summary)
115        .unwrap_or(0);
116    let keep_recent_tokens = (budget as f64 * KEEP_RECENT_TOKEN_FRACTION).ceil() as u64;
117    let mut recent_tokens = 0u64;
118    let mut token_end = messages.len();
119    for (index, message) in messages.iter().enumerate().rev() {
120        recent_tokens = recent_tokens.saturating_add(estimate_tokens_for_message(message));
121        token_end = index;
122        if recent_tokens >= keep_recent_tokens {
123            break;
124        }
125    }
126    let message_end = messages.len().saturating_sub(KEEP_RECENT_MESSAGES);
127    let end = message_end
128        .min(token_end)
129        .min(find_kth_recent_user(messages, KEEP_RECENT_USER_TURNS));
130    if end < start + 2 {
131        return None;
132    }
133
134    let tokens_saved_estimate = messages[start..end]
135        .iter()
136        .map(estimate_tokens_for_message)
137        .sum();
138    Some(CompactRange {
139        start,
140        end,
141        tokens_saved_estimate,
142    })
143}
144
145pub fn estimate_compacted_message_tokens(
146    messages: &[Message],
147    range: &CompactRange,
148    summary: &str,
149) -> u64 {
150    let turn_id = messages
151        .get(range.start)
152        .map(|m| m.turn_id.clone())
153        .unwrap_or_else(crate::event::TurnId::now);
154    let after = replace_range_with_summary(messages, range, summary.to_string(), turn_id);
155    estimate_tokens_for_messages(&after)
156}
157
158pub fn filter_orphan_tool_messages(messages: &mut Vec<Message>) {
159    let use_ids: std::collections::HashSet<String> = messages
160        .iter()
161        .flat_map(|m| {
162            m.parts.iter().filter_map(|p| match p {
163                MessagePart::ToolUse { id, .. } => Some(id.clone()),
164                _ => None,
165            })
166        })
167        .collect();
168    let result_ids: std::collections::HashSet<String> = messages
169        .iter()
170        .flat_map(|m| {
171            m.parts.iter().filter_map(|p| match p {
172                MessagePart::ToolResult { tool_use_id, .. } => Some(tool_use_id.clone()),
173                _ => None,
174            })
175        })
176        .collect();
177    let mut seen_results: std::collections::HashSet<String> = std::collections::HashSet::new();
178    messages.retain(|m| {
179        for p in &m.parts {
180            match p {
181                MessagePart::ToolUse { id, .. } if !result_ids.contains(id) => return false,
182                MessagePart::ToolResult { tool_use_id, .. } => {
183                    if !use_ids.contains(tool_use_id) || !seen_results.insert(tool_use_id.clone()) {
184                        return false;
185                    }
186                }
187                _ => {}
188            }
189        }
190        true
191    });
192}
193
194pub fn find_compact_summaries(messages: &[Message]) -> Vec<CompactSummary> {
195    let mut out = Vec::new();
196    for (idx, msg) in messages.iter().enumerate() {
197        if let Some(summary) = compact_summary(msg) {
198            out.push(CompactSummary {
199                message_index: idx,
200                seq_start: summary.seq_start,
201                seq_end: summary.seq_end,
202                count: summary.count,
203            });
204        }
205    }
206    out
207}
208
209#[derive(Debug, Clone, PartialEq, Eq)]
210pub struct CompactSummary {
211    pub message_index: usize,
212    pub seq_start: u64,
213    pub seq_end: u64,
214    pub count: usize,
215}
216
217struct CompactSummaryPart {
218    seq_start: u64,
219    seq_end: u64,
220    count: usize,
221}
222
223fn extract_anchor(messages: &[Message]) -> Option<(String, &[Message])> {
224    let first = messages.first()?;
225    let summary = first.parts.iter().find_map(|part| match part {
226        MessagePart::CompactSummary { summary, .. } => Some(summary.clone()),
227        _ => None,
228    })?;
229    Some((summary, &messages[1..]))
230}
231
232fn compact_summary(msg: &Message) -> Option<CompactSummaryPart> {
233    if msg.role != MessageRole::System {
234        return None;
235    }
236    msg.parts.iter().find_map(|part| match part {
237        MessagePart::CompactSummary {
238            seq_start,
239            seq_end,
240            count,
241            ..
242        } => Some(CompactSummaryPart {
243            seq_start: *seq_start,
244            seq_end: *seq_end,
245            count: *count,
246        }),
247        _ => None,
248    })
249}
250
251pub async fn maybe_auto_compact(
252    session: &crate::session::Session,
253    model: &str,
254    providers: &crate::provider::ProviderRegistry,
255) {
256    maybe_auto_compact_with_budget(
257        session,
258        model,
259        providers,
260        CompactionBudgetContext::default(),
261    )
262    .await;
263}
264
265pub async fn maybe_auto_compact_with_budget(
266    session: &crate::session::Session,
267    model: &str,
268    providers: &crate::provider::ProviderRegistry,
269    budget_context: CompactionBudgetContext,
270) {
271    let _compact_guard = session.acquire_compact_lock().await;
272    maybe_auto_compact_locked(session, model, providers, budget_context).await;
273}
274
275pub fn spawn_auto_compact(
276    session: std::sync::Arc<crate::session::Session>,
277    model: String,
278    providers: crate::provider::ProviderRegistry,
279) {
280    tokio::task::spawn_blocking(move || {
281        let Ok(rt) = tokio::runtime::Builder::new_current_thread()
282            .enable_all()
283            .build()
284        else {
285            session.push_system_note("compaction skipped: background runtime init failed".into());
286            return;
287        };
288        rt.block_on(async move {
289            maybe_auto_compact(&session, &model, &providers).await;
290        });
291    });
292}
293
294pub async fn start_auto_compact(
295    session: std::sync::Arc<crate::session::Session>,
296    model: String,
297    providers: crate::provider::ProviderRegistry,
298) {
299    start_auto_compact_with_budget(
300        session,
301        model,
302        providers,
303        CompactionBudgetContext::default(),
304    )
305    .await;
306}
307
308pub async fn start_auto_compact_with_budget(
309    session: std::sync::Arc<crate::session::Session>,
310    model: String,
311    providers: crate::provider::ProviderRegistry,
312    budget_context: CompactionBudgetContext,
313) {
314    let compact_guard = session.acquire_compact_lock_owned().await;
315    tokio::task::spawn_blocking(move || {
316        let Ok(rt) = tokio::runtime::Builder::new_current_thread()
317            .enable_all()
318            .build()
319        else {
320            drop(compact_guard);
321            session.push_system_note("compaction skipped: background runtime init failed".into());
322            return;
323        };
324        rt.block_on(async move {
325            maybe_auto_compact_locked(&session, &model, &providers, budget_context).await;
326            drop(compact_guard);
327        });
328    });
329}
330
331async fn maybe_auto_compact_locked(
332    session: &crate::session::Session,
333    model: &str,
334    providers: &crate::provider::ProviderRegistry,
335    budget_context: CompactionBudgetContext,
336) {
337    let forced = session.take_manual_compact_request();
338    let info = crate::model_registry::model_info(model);
339    let trigger = info.compaction_trigger_threshold();
340    let target = budget_context
341        .history_budget(&info)
342        .map(|budget| budget.min(info.compaction_target_after()))
343        .unwrap_or_else(|| info.compaction_target_after());
344    let msgs = session.messages();
345    let window_tokens = estimate_tokens_for_messages(&msgs);
346    let provider_tokens = session.last_input_tokens();
347    let current = if provider_tokens > 0 {
348        provider_tokens
349    } else {
350        window_tokens
351    };
352    if !forced && current <= trigger {
353        return;
354    }
355    if !forced && !session.approval_cooldown_ok_for_compact() {
356        return;
357    }
358    let Some(range) = find_compact_range(&msgs, target) else {
359        let (replacement, rewritten_count) =
360            build_budgeted_turn_rewrite(&msgs, target, model, providers).await;
361        let after_tokens = estimate_tokens_for_messages(&replacement);
362        if rewritten_count == 0 || after_tokens >= window_tokens || after_tokens > target {
363            session.emit_compact_warning(
364                model,
365                current,
366                trigger,
367                info.context_budget,
368                "no compactible span — retained user content cannot fit the history budget",
369            );
370            return;
371        }
372        match session.commit_rewritten_window(
373            replacement,
374            window_tokens,
375            window_tokens,
376            rewritten_count,
377        ) {
378            Some(_) => {}
379            None => {
380                session.emit_compact_warning(
381                    model,
382                    current,
383                    trigger,
384                    info.context_budget,
385                    "retained turn output rewrite did not shrink the transcript",
386                );
387            }
388        }
389        return;
390    };
391    let _ = session
392        .stream_tx()
393        .send(crate::stream::StreamFrame::CompactionSummary {
394            phase: crate::stream::CompactionPhase::Running,
395            range_start: range.start,
396            range_end: range.end.saturating_sub(1),
397            summary: String::new(),
398            before_tokens: current,
399            after_tokens: 0,
400            compacted_count: range.end - range.start,
401        });
402    let send_failed = |session: &crate::session::Session, reason: &str| {
403        let _ = session
404            .stream_tx()
405            .send(crate::stream::StreamFrame::CompactionSummary {
406                phase: crate::stream::CompactionPhase::Failed,
407                range_start: range.start,
408                range_end: range.end.saturating_sub(1),
409                summary: reason.to_string(),
410                before_tokens: current,
411                after_tokens: current,
412                compacted_count: range.end - range.start,
413            });
414    };
415    let mut filtered: Vec<Message> = msgs[range.start..range.end].to_vec();
416    filter_orphan_tool_messages(&mut filtered);
417    let (anchor, new_messages) = extract_anchor(&filtered)
418        .map(|(anchor, remaining)| (Some(anchor), remaining.to_vec()))
419        .unwrap_or_else(|| (None, filtered.clone()));
420    let summary =
421        match generate_llm_summary(anchor.as_deref(), &new_messages, model, providers).await {
422            Ok(text) => text,
423            Err(err) => {
424                session.emit_compact_warning(
425                    model,
426                    current,
427                    trigger,
428                    info.context_budget,
429                    &format!("LLM summary failed: {err}. Degraded to placeholder."),
430                );
431                format!(
432                    "[atman: compacted {} messages, LLM summary unavailable at {}]",
433                    range.end - range.start,
434                    chrono::Utc::now().to_rfc3339()
435                )
436            }
437        };
438    let final_summary =
439        match request_review_if_enabled(session, forced, &filtered, &range, current, summary).await
440        {
441            ReviewOutcome::Commit(s) => s,
442            ReviewOutcome::Rejected => {
443                send_failed(
444                    session,
445                    "compaction rejected by user; keeping full transcript",
446                );
447                session.push_system_note(
448                    "compaction rejected by user; keeping full transcript".into(),
449                );
450                return;
451            }
452        };
453    let replacement =
454        build_budgeted_replacement(&msgs, &range, &final_summary, target, model, providers).await;
455    let after_tokens = estimate_tokens_for_messages(&replacement);
456    if after_tokens >= window_tokens {
457        send_failed(
458            session,
459            &format!(
460                "compaction skipped: replacement would not shrink transcript ({} >= {} tokens)",
461                after_tokens, window_tokens
462            ),
463        );
464        session.push_system_note(format!(
465            "compaction skipped: replacement would not shrink transcript ({} >= {} tokens)",
466            after_tokens, window_tokens
467        ));
468        return;
469    }
470    match session.commit_compacted_window(
471        final_summary,
472        replacement,
473        range,
474        window_tokens,
475        window_tokens,
476    ) {
477        Some(result) => {
478            session.push_system_note(format!(
479                "auto-compacted {}..{} — {} → {} tokens",
480                result.compacted_start,
481                result.compacted_end,
482                result.before_tokens,
483                result.after_tokens
484            ));
485        }
486        None => {
487            session.emit_compact_warning(
488                model,
489                current,
490                trigger,
491                info.context_budget,
492                "no compactible span — history too short or already fully compacted",
493            );
494        }
495    }
496}
497
498enum ReviewOutcome {
499    Commit(String),
500    Rejected,
501}
502
503async fn request_review_if_enabled(
504    session: &crate::session::Session,
505    forced: bool,
506    slice: &[Message],
507    range: &CompactRange,
508    tokens_before: u64,
509    summary: String,
510) -> ReviewOutcome {
511    if !session.compact_review_mode().should_review(forced) {
512        return ReviewOutcome::Commit(summary);
513    }
514    let reviews = session.compact_reviews();
515    if reviews.subscriber_count() == 0 {
516        return ReviewOutcome::Commit(summary);
517    }
518    let pending = crate::session::PendingCompactReview {
519        review_id: uuid::Uuid::now_v7().to_string(),
520        summary: summary.clone(),
521        slice_preview: format_slice_for_preview(slice),
522        slice_count: slice.len(),
523        range_start: range.start,
524        range_end: range.end,
525        tokens_before,
526        emitted_at: chrono::Utc::now(),
527    };
528    let rx = reviews.request(pending);
529    match rx.await {
530        Ok(crate::session::CompactReviewDecision::AcceptAsIs) => ReviewOutcome::Commit(summary),
531        Ok(crate::session::CompactReviewDecision::AcceptEdited { summary: edited }) => {
532            ReviewOutcome::Commit(edited)
533        }
534        Ok(crate::session::CompactReviewDecision::Reject) | Err(_) => ReviewOutcome::Rejected,
535    }
536}
537
538fn format_slice_for_preview(slice: &[Message]) -> String {
539    let mut out = String::new();
540    for (i, msg) in slice.iter().enumerate() {
541        let role = msg.role.as_str();
542        let body = serialize_message_for_summary(msg);
543        let truncated: String = body.chars().take(400).collect();
544        out.push_str(&format!("[{i}] {role}: {truncated}\n"));
545    }
546    out.chars().take(16_000).collect()
547}
548
549const SUMMARY_SYSTEM_PROMPT: &str =
550    "You are an anchored context summarization assistant for coding sessions.";
551
552const SUMMARY_INSTRUCTIONS: &str = r#"You are an anchored context summarization assistant.
553
554Below is:
5551. <current-anchor>: the existing handoff state, which is authoritative and must be preserved.
5562. <new-messages>: only the messages that arrived since the anchor was written.
557
558Merge the NEW facts from <new-messages> INTO the current anchor, producing an upgraded full anchor.
559
560STRUCTURAL RULES (data model, not optional style):
561- ## Objective: unchanged unless the new messages show the user explicitly redirected.
562- ### 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.
563- ### Active: update based on new messages; move newly-done items to Completed.
564- ### Blocked: update based on new messages; remove resolved ones.
565- ## Decisions: only add new decisions. Never remove old ones.
566- ## Next Move: replace based on current end state.
567- Keep every section, even when empty.
568- Preserve exact file paths, symbols, commands, error strings, identifiers.
569
570Output exactly this Markdown structure:
571## Objective
572## Important Details
573## Work State
574### Completed
575### Active
576### Blocked
577## Decisions
578## Next Move
579## Relevant Files
580
581Do not mention the summary process or that context was compacted.
582Respond in the same language as the conversation."#;
583
584async fn generate_llm_summary(
585    anchor: Option<&str>,
586    slice: &[Message],
587    model: &str,
588    providers: &crate::provider::ProviderRegistry,
589) -> Result<String, crate::error::RuntimeError> {
590    let provider = providers.resolve(model).ok_or_else(|| {
591        crate::error::RuntimeError::ToolFailed(format!("no provider for {model}"))
592    })?;
593    let payload = format_slice_for_summary(slice);
594    let (messages, dump_user) = if let Some(anchor) = anchor {
595        let anchor_user = format!("<current-anchor>\n{anchor}\n</current-anchor>");
596        let new_user =
597            format!("<new-messages>\n{payload}\n</new-messages>\n\n{SUMMARY_INSTRUCTIONS}");
598        (
599            vec![
600                Message::user_text(crate::event::TurnId::now(), anchor_user.clone()),
601                Message::user_text(crate::event::TurnId::now(), new_user.clone()),
602            ],
603            format!("{anchor_user}\n\n{new_user}"),
604        )
605    } else {
606        let user = format!(
607            "<conversation_history>\n{payload}\n</conversation_history>\n\n{SUMMARY_INSTRUCTIONS}"
608        );
609        (
610            vec![Message::user_text(
611                crate::event::TurnId::now(),
612                user.clone(),
613            )],
614            user,
615        )
616    };
617    if let Ok(dir) = std::env::var("ATMAN_COMPACT_DUMP") {
618        let _ = std::fs::write(
619            format!("{dir}/compact_request.txt"),
620            format!("=== SYSTEM ===\n{SUMMARY_SYSTEM_PROMPT}\n\n=== USER ===\n{dump_user}"),
621        );
622    }
623    let req = crate::provider::LlmRequest {
624        model: model.into(),
625        messages,
626        system: Some(SUMMARY_SYSTEM_PROMPT.into()),
627        input: crate::value::Value::Unit,
628        schema: None,
629        cache_prompt: false,
630        tools: Vec::new(),
631        thinking_enabled: false,
632        stall_timeout_secs: 0,
633    };
634    let outcome = provider.call(req).await?;
635    let text = outcome.text_concat();
636    if text.trim().is_empty() {
637        return Err(crate::error::RuntimeError::ToolFailed(
638            "empty summary from provider".into(),
639        ));
640    }
641    Ok(text)
642}
643
644async fn build_budgeted_turn_rewrite(
645    messages: &[Message],
646    history_budget: u64,
647    model: &str,
648    providers: &crate::provider::ProviderRegistry,
649) -> (Vec<Message>, usize) {
650    let mut replacement = messages.to_vec();
651    let mut group_index = 0;
652    let mut rewritten_count = 0;
653    while estimate_tokens_for_messages(&replacement) > history_budget {
654        let groups = user_turn_ranges(&replacement);
655        let Some((start, end)) = groups.get(group_index).copied() else {
656            break;
657        };
658        let output = replacement[start + 1..end].to_vec();
659        if output.is_empty() {
660            group_index += 1;
661            continue;
662        }
663        let output_tokens = estimate_tokens_for_messages(&output);
664        let summary = generate_llm_summary(None, &output, model, providers)
665            .await
666            .unwrap_or_else(|_| deterministic_turn_omission(&output));
667        let mut summary_message = Message::assistant_text(
668            replacement[start].turn_id.clone(),
669            format!("[atman: compacted turn output]\n{summary}\n[/atman: compacted turn output]"),
670        );
671        if estimate_tokens_for_message(&summary_message) >= output_tokens {
672            summary_message = Message::assistant_text(
673                replacement[start].turn_id.clone(),
674                deterministic_turn_omission(&output),
675            );
676        }
677        rewritten_count += output.len();
678        replacement.splice(start + 1..end, [summary_message]);
679        group_index += 1;
680    }
681    filter_orphan_tool_messages(&mut replacement);
682    (replacement, rewritten_count)
683}
684
685async fn build_budgeted_replacement(
686    messages: &[Message],
687    range: &CompactRange,
688    anchor_summary: &str,
689    history_budget: u64,
690    model: &str,
691    providers: &crate::provider::ProviderRegistry,
692) -> Vec<Message> {
693    let turn_id = messages
694        .get(range.start)
695        .map(|message| message.turn_id.clone())
696        .unwrap_or_else(crate::event::TurnId::now);
697    let mut replacement =
698        replace_range_with_summary(messages, range, anchor_summary.to_string(), turn_id);
699    filter_orphan_tool_messages(&mut replacement);
700    if estimate_tokens_for_messages(&replacement) <= history_budget {
701        return replacement;
702    }
703
704    let mut group_index = 0;
705    loop {
706        let groups = user_turn_ranges(&replacement);
707        if group_index >= groups.len()
708            || estimate_tokens_for_messages(&replacement) <= history_budget
709        {
710            break;
711        }
712        let (start, end) = groups[group_index];
713        let output: Vec<Message> = replacement[start + 1..end].to_vec();
714        if output.is_empty() {
715            group_index += 1;
716            continue;
717        }
718        let output_tokens = estimate_tokens_for_messages(&output);
719        let summary = generate_llm_summary(None, &output, model, providers)
720            .await
721            .unwrap_or_else(|_| deterministic_turn_omission(&output));
722        let mut summary_message = Message::assistant_text(
723            replacement[start].turn_id.clone(),
724            format!("[atman: compacted turn output]\n{summary}\n[/atman: compacted turn output]"),
725        );
726        if estimate_tokens_for_message(&summary_message) >= output_tokens {
727            summary_message = Message::assistant_text(
728                replacement[start].turn_id.clone(),
729                deterministic_turn_omission(&output),
730            );
731        }
732        replacement.splice(start + 1..end, [summary_message]);
733        group_index += 1;
734    }
735
736    if estimate_tokens_for_messages(&replacement) > history_budget {
737        let groups = user_turn_ranges(&replacement);
738        for (start, end) in groups.into_iter().rev() {
739            let output = replacement[start + 1..end].to_vec();
740            if !output.is_empty() {
741                replacement.splice(
742                    start + 1..end,
743                    [Message::assistant_text(
744                        replacement[start].turn_id.clone(),
745                        deterministic_turn_omission(&output),
746                    )],
747                );
748            }
749        }
750    }
751
752    if estimate_tokens_for_messages(&replacement) > history_budget {
753        replacement.retain(|message| message.role == MessageRole::User);
754    }
755
756    filter_orphan_tool_messages(&mut replacement);
757    replacement
758}
759
760fn user_turn_ranges(messages: &[Message]) -> Vec<(usize, usize)> {
761    let starts: Vec<usize> = messages
762        .iter()
763        .enumerate()
764        .filter_map(|(index, message)| (message.role == MessageRole::User).then_some(index))
765        .collect();
766    starts
767        .iter()
768        .enumerate()
769        .map(|(index, start)| {
770            (
771                *start,
772                starts.get(index + 1).copied().unwrap_or(messages.len()),
773            )
774        })
775        .collect()
776}
777
778fn deterministic_turn_omission(messages: &[Message]) -> String {
779    format!(
780        "[atman: omitted {} oversized assistant/system/tool messages during persistent compaction]",
781        messages.len()
782    )
783}
784
785fn format_slice_for_summary(slice: &[Message]) -> String {
786    let mut out = String::new();
787    for (i, msg) in slice.iter().enumerate() {
788        let role = msg.role.as_str();
789        let body = serialize_message_for_summary(msg);
790        let truncated: String = body.chars().take(4000).collect();
791        out.push_str(&format!("[{i}] {role}: {truncated}\n\n"));
792    }
793    out.chars().take(120_000).collect()
794}
795
796fn serialize_message_for_summary(msg: &Message) -> String {
797    let mut parts = Vec::new();
798    for part in &msg.parts {
799        match part {
800            MessagePart::CompactSummary { summary, .. } => {
801                parts.push(summary.clone());
802            }
803            MessagePart::Text { text } => {
804                parts.push(text.clone());
805            }
806            MessagePart::Thinking { thinking, .. } => {
807                let truncated: String = thinking.chars().take(1000).collect();
808                parts.push(format!("[thinking: {truncated}]"));
809            }
810            MessagePart::ToolUse { name, input, .. } => {
811                let input_str = if input.is_null() {
812                    String::new()
813                } else {
814                    input.to_string()
815                };
816                let truncated: String = input_str.chars().take(2000).collect();
817                parts.push(format!("[tool_call: {name}({truncated})]"));
818            }
819            MessagePart::ToolResult {
820                content,
821                is_error,
822                tool_use_id,
823            } => {
824                let truncated: String = content.chars().take(3000).collect();
825                let marker = if *is_error { "ERROR" } else { "ok" };
826                let id_short: String = tool_use_id.chars().take(12).collect();
827                parts.push(format!("[tool_result {id_short}… {marker}: {truncated}]"));
828            }
829            MessagePart::Image { .. } => {
830                parts.push("[image]".into());
831            }
832        }
833    }
834    parts.join(" ")
835}
836
837pub fn replace_range_with_summary(
838    messages: &[Message],
839    range: &CompactRange,
840    summary: String,
841    turn_id: crate::event::TurnId,
842) -> Vec<Message> {
843    let mut out = Vec::with_capacity(1 + messages.len().saturating_sub(range.end));
844    out.push(Message::system_compact_summary(
845        turn_id,
846        summary,
847        range.start as u64,
848        range.end.saturating_sub(1) as u64,
849        range.end - range.start,
850    ));
851    out.extend_from_slice(&messages[range.end..]);
852    out
853}
854
855/// Result of compacting a messages_handle in place.
856#[derive(Debug, Clone, PartialEq, Eq)]
857pub struct HandleCompactResult {
858    pub before_tokens: u64,
859    pub after_tokens: u64,
860    pub compacted_start: usize,
861    pub compacted_end: usize,
862}
863
864/// Compact a messages_handle in place (data-layer primitive, operates on any
865/// FlowRun's segment). Returns `None` if under `budget` or no compactable
866/// range. Caller should hold the FlowRun's `compact_lock`.
867pub fn compact_messages_on_handle(
868    handle: &std::sync::Arc<std::sync::Mutex<Vec<Message>>>,
869    summary: String,
870    budget: u64,
871) -> Option<HandleCompactResult> {
872    let mut msgs = handle.lock().unwrap();
873    let before_tokens = estimate_tokens_for_messages(&msgs);
874    let range = find_compact_range(&msgs, budget)?;
875    let turn_id = msgs
876        .get(range.start)
877        .map(|m| m.turn_id.clone())
878        .unwrap_or_else(crate::event::TurnId::now);
879    let after = replace_range_with_summary(&msgs, &range, summary, turn_id);
880    let after_tokens = estimate_tokens_for_messages(&after);
881    if after_tokens >= before_tokens {
882        return None;
883    }
884    let result = HandleCompactResult {
885        before_tokens,
886        after_tokens,
887        compacted_start: range.start,
888        compacted_end: range.end.saturating_sub(1),
889    };
890    *msgs = after;
891    Some(result)
892}
893
894#[cfg(test)]
895mod tests {
896    use super::*;
897    use crate::event::TurnId;
898    use crate::message::MessageOrigin;
899
900    fn user(text: &str) -> Message {
901        Message::user_text(TurnId::now(), text)
902    }
903    fn assistant(text: &str) -> Message {
904        Message::assistant_text(TurnId::now(), text)
905    }
906    fn system(text: &str) -> Message {
907        Message::system_text(TurnId::now(), text)
908    }
909
910    #[test]
911    fn compaction_budget_reserves_output_safety_and_fixed_input_only_at_compact_time() {
912        let info = crate::model_registry::ModelInfo {
913            name: "test".into(),
914            context_budget: 100_000,
915            compact_threshold_ratio: 0.8,
916            thinking_enabled: false,
917            max_output_tokens: Some(10_000),
918        };
919        let budget = CompactionBudgetContext {
920            fixed_input_tokens: Some(5_000),
921        }
922        .history_budget(&info);
923        assert_eq!(budget, Some(100_000 - 10_000 - 2_000 - 5_000));
924    }
925
926    #[test]
927    fn compaction_budget_saturates_when_fixed_input_exceeds_context() {
928        let info = crate::model_registry::ModelInfo {
929            name: "test".into(),
930            context_budget: 20_000,
931            compact_threshold_ratio: 0.8,
932            thinking_enabled: false,
933            max_output_tokens: None,
934        };
935        assert_eq!(
936            CompactionBudgetContext {
937                fixed_input_tokens: Some(100_000)
938            }
939            .history_budget(&info),
940            Some(0)
941        );
942    }
943
944    #[tokio::test]
945    async fn budgeted_replacement_groups_by_user_boundary_and_persists_omission() {
946        let first_turn = TurnId::now();
947        let second_turn = TurnId::now();
948        let mut messages = vec![system(&"old".repeat(20_000)), assistant("old answer")];
949        messages.push(Message::user_text(first_turn.clone(), "first user"));
950        messages.push(Message::assistant_text(
951            second_turn.clone(),
952            "first output".repeat(4_000),
953        ));
954        messages.push(assistant_with_tool_use(
955            "calling tool",
956            "fs.read",
957            serde_json::json!({"path": "/tmp/example"}),
958        ));
959        messages.push(tool_result(
960            "call_test",
961            &"tool output".repeat(4_000),
962            false,
963        ));
964        messages.push(Message::user_text(second_turn.clone(), "current user"));
965        messages.push(Message::assistant_text(
966            first_turn,
967            "current output".repeat(4_000),
968        ));
969        let range = CompactRange {
970            start: 0,
971            end: 2,
972            tokens_saved_estimate: 1,
973        };
974
975        let replacement = build_budgeted_replacement(
976            &messages,
977            &range,
978            "anchor",
979            500,
980            "missing-provider",
981            &crate::provider::ProviderRegistry::default(),
982        )
983        .await;
984
985        let texts: Vec<String> = replacement.iter().map(Message::text_concat).collect();
986        assert!(texts.iter().any(|text| text == "current user"));
987        assert!(
988            texts
989                .iter()
990                .any(|text| text.contains("omitted 3 oversized"))
991        );
992        assert!(!replacement.iter().any(|message| {
993            message.parts.iter().any(|part| {
994                matches!(
995                    part,
996                    MessagePart::ToolUse { .. } | MessagePart::ToolResult { .. }
997                )
998            })
999        }));
1000        assert_eq!(user_turn_ranges(&replacement).len(), 2);
1001    }
1002
1003    #[tokio::test]
1004    async fn turn_rewrite_compacts_oversized_tool_output_without_dropping_recent_users() {
1005        let first = TurnId::now();
1006        let current = TurnId::now();
1007        let messages = vec![
1008            Message::user_text(first, "first user"),
1009            assistant_with_tool_use(
1010                &"calling tool".repeat(2_000),
1011                "fs.read",
1012                serde_json::json!({"path": "/tmp/example"}),
1013            ),
1014            tool_result("call_test", &"tool output".repeat(8_000), false),
1015            Message::user_text(current, "current user"),
1016        ];
1017
1018        assert!(find_compact_range(&messages, 500).is_none());
1019        let (replacement, rewritten_count) = build_budgeted_turn_rewrite(
1020            &messages,
1021            500,
1022            "missing-provider",
1023            &crate::provider::ProviderRegistry::default(),
1024        )
1025        .await;
1026
1027        assert_eq!(rewritten_count, 2);
1028        assert!(estimate_tokens_for_messages(&replacement) <= 500);
1029        let users: Vec<String> = replacement
1030            .iter()
1031            .filter(|message| message.role == MessageRole::User)
1032            .map(Message::text_concat)
1033            .collect();
1034        assert_eq!(users, vec!["first user", "current user"]);
1035        assert!(replacement.iter().any(|message| {
1036            message
1037                .text_concat()
1038                .contains("omitted 2 oversized assistant/system/tool messages")
1039        }));
1040        assert!(!replacement.iter().any(|message| {
1041            message.parts.iter().any(|part| {
1042                matches!(
1043                    part,
1044                    MessagePart::ToolUse { .. } | MessagePart::ToolResult { .. }
1045                )
1046            })
1047        }));
1048    }
1049
1050    #[test]
1051    fn user_turn_ranges_ignore_misanchored_turn_ids() {
1052        let first = TurnId::now();
1053        let second = TurnId::now();
1054        let messages = vec![
1055            Message::user_text(first.clone(), "u1"),
1056            Message::assistant_text(second.clone(), "a1"),
1057            Message::user_text(second, "u2"),
1058            Message::assistant_text(first, "a2"),
1059        ];
1060        assert_eq!(user_turn_ranges(&messages), vec![(0, 2), (2, 4)]);
1061    }
1062
1063    #[test]
1064    fn summary_instructions_keep_decisions_before_next_move() {
1065        let objective = SUMMARY_INSTRUCTIONS
1066            .find("## Objective")
1067            .expect("objective");
1068        let decisions = SUMMARY_INSTRUCTIONS
1069            .find("## Decisions")
1070            .expect("decisions");
1071        let next_move = SUMMARY_INSTRUCTIONS
1072            .find("## Next Move")
1073            .expect("next move");
1074        assert!(objective < decisions);
1075        assert!(decisions < next_move);
1076    }
1077
1078    #[test]
1079    fn estimate_scales_with_char_length() {
1080        let short = user("hi");
1081        let long = user(&"x".repeat(3500));
1082        assert!(estimate_tokens_for_message(&long) > estimate_tokens_for_message(&short) * 100);
1083    }
1084
1085    #[test]
1086    fn find_compact_returns_none_when_under_budget() {
1087        let msgs = vec![user("a"), assistant("b"), user("c"), assistant("d")];
1088        assert!(find_compact_range(&msgs, 1000).is_none());
1089    }
1090
1091    #[test]
1092    fn find_compact_returns_none_for_short_history() {
1093        let msgs = vec![user(&"x".repeat(9000))];
1094        assert!(find_compact_range(&msgs, 100).is_none());
1095    }
1096
1097    #[test]
1098    fn find_kth_recent_user_handles_exact_excess_and_mixed_history() {
1099        let exact = vec![
1100            user("u0"),
1101            assistant("a0"),
1102            system("s0"),
1103            tool_result("call-0", "result", false),
1104            user("u1"),
1105            assistant("a1"),
1106            user("u2"),
1107            system("s1"),
1108            user("u3"),
1109            assistant("a3"),
1110            user("u4"),
1111        ];
1112        // With exactly five users, the fifth recent user is the first message.
1113        assert_eq!(find_kth_recent_user(&exact, KEEP_RECENT_USER_TURNS), 0);
1114
1115        let mut excess = exact.clone();
1116        excess.push(user("u5"));
1117        assert_eq!(find_kth_recent_user(&excess, KEEP_RECENT_USER_TURNS), 4);
1118
1119        let too_few = vec![user("a"), assistant("b"), assistant("c")];
1120        assert_eq!(find_kth_recent_user(&too_few, KEEP_RECENT_USER_TURNS), 0);
1121    }
1122
1123    #[test]
1124    fn find_compact_range_preserves_minimum_recent_messages_without_anchor() {
1125        let mut msgs = vec![system("head")];
1126        msgs.extend((0..25).map(|index| assistant(&format!("old {index}"))));
1127        msgs.extend(
1128            (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
1129        );
1130        let range = find_compact_range(&msgs, 1).expect("range");
1131        assert_eq!(range.start, 0);
1132        assert_eq!(range.end, msgs.len() - KEEP_RECENT_MESSAGES);
1133    }
1134
1135    #[test]
1136    fn find_compact_range_preserves_minimum_window_for_large_tool_result() {
1137        let mut msgs = vec![user(&"h".repeat(500_000))];
1138        for index in 1..31 {
1139            if matches!(index, 20 | 22 | 24 | 26 | 28 | 30) {
1140                msgs.push(user(&"u".repeat(800)));
1141            } else {
1142                msgs.push(assistant(&"a".repeat(800)));
1143            }
1144        }
1145        msgs.push(tool_result("call-large", &"t".repeat(22_000), false));
1146
1147        let budget = 120_000;
1148        let minimum_recent_tokens = (budget as f64 * KEEP_RECENT_TOKEN_FRACTION).ceil() as u64;
1149        assert!(estimate_tokens_for_message(msgs.last().unwrap()) > minimum_recent_tokens);
1150
1151        let range = find_compact_range(&msgs, budget).expect("range");
1152        assert!(
1153            range.end <= msgs.len() - KEEP_RECENT_MESSAGES,
1154            "range was {range:?}"
1155        );
1156        assert!(msgs.len() - range.end >= KEEP_RECENT_MESSAGES);
1157        assert!(estimate_tokens_for_messages(&msgs[range.end..]) >= minimum_recent_tokens);
1158    }
1159
1160    #[test]
1161    fn find_compact_range_handles_four_to_twenty_one_message_histories() {
1162        for len in 4..=21 {
1163            let msgs = (0..len)
1164                .map(|_| user(&"x".repeat(5000)))
1165                .collect::<Vec<_>>();
1166            assert_eq!(
1167                find_compact_range(&msgs, 1).is_some(),
1168                len >= KEEP_RECENT_MESSAGES + 2,
1169                "len={len}"
1170            );
1171        }
1172    }
1173
1174    #[test]
1175    fn find_compact_range_recent_users_limit_mixed_history() {
1176        let msgs = vec![
1177            system("head"),
1178            assistant("a0"),
1179            user("u0"),
1180            tool_result("call-0", "r0", false),
1181            assistant("a1"),
1182            user("u1"),
1183            assistant("a2"),
1184            tool_result("call-1", "r1", false),
1185            user("u2"),
1186            assistant("a3"),
1187            system("note"),
1188            user("u3"),
1189            tool_result("call-2", "r2", false),
1190            assistant("a4"),
1191            user("u4"),
1192            assistant("a5"),
1193            tool_result("call-3", "r3", false),
1194            user("u5"),
1195            assistant("a6"),
1196            system("tail"),
1197            assistant("a7"),
1198        ];
1199
1200        let range = find_compact_range(&msgs, 1).expect("range");
1201        assert_eq!(range.end, 5);
1202        assert_eq!(msgs[range.end].role, MessageRole::User);
1203    }
1204
1205    #[test]
1206    fn find_compact_range_preserves_recent_user_turns_after_anchor() {
1207        let mut msgs = vec![system("head"), compaction_summary("summary")];
1208        msgs.extend((0..6).flat_map(|index| {
1209            [
1210                user(&format!("user {index}")),
1211                assistant("assistant"),
1212                assistant("tool fragment"),
1213            ]
1214        }));
1215        msgs.extend((0..12).map(|_| assistant("recent fragment")));
1216        let range = find_compact_range(&msgs, 1).expect("range");
1217        assert_eq!(range.start, 1);
1218        assert_eq!(range.end, 5);
1219        assert_eq!(msgs[range.end].role, MessageRole::User);
1220    }
1221
1222    #[test]
1223    fn find_compact_range_returns_none_when_end_cannot_cover_two_messages() {
1224        let msgs = vec![
1225            system("head"),
1226            compaction_summary("summary"),
1227            assistant("tail"),
1228            user("tail"),
1229        ];
1230        assert!(find_compact_range(&msgs, 1).is_none());
1231    }
1232
1233    #[test]
1234    fn extract_anchor_removes_leading_compact_summary() {
1235        let messages = vec![compaction_summary("anchor"), user("new")];
1236        let (anchor, remaining) = extract_anchor(&messages).expect("anchor");
1237        assert_eq!(anchor, "anchor");
1238        assert_eq!(remaining, &messages[1..]);
1239    }
1240
1241    #[test]
1242    fn extract_anchor_returns_none_without_leading_summary() {
1243        let messages = vec![user("new")];
1244        assert!(extract_anchor(&messages).is_none());
1245    }
1246
1247    #[test]
1248    fn compact_messages_on_handle_replaces_range_in_place() {
1249        let mut messages = vec![system("head")];
1250        messages.extend((0..9).map(|index| assistant(&format!("old {index}"))));
1251        messages.extend(
1252            (0..6).flat_map(|index| [user(&format!("user {index}")), assistant(&"x".repeat(4000))]),
1253        );
1254        messages.extend((0..10).map(|index| assistant(&format!("recent {index}"))));
1255        messages.push(user("tail"));
1256        let handle: std::sync::Arc<std::sync::Mutex<Vec<Message>>> =
1257            std::sync::Arc::new(std::sync::Mutex::new(messages));
1258        // The recent-message and recent-user limits meet at the fifth recent user.
1259        // The compacted prefix is replaced by one summary while the tail remains.
1260        let result = compact_messages_on_handle(&handle, "gist".into(), 100);
1261        let result = result.expect("should compact");
1262        assert!(result.after_tokens < result.before_tokens);
1263        let msgs = handle.lock().unwrap();
1264        assert_eq!(result.compacted_start, 0);
1265        assert_eq!(result.compacted_end, 13);
1266        assert!(is_compaction_summary(&msgs[0]));
1267        assert_eq!(msgs.last().unwrap().text_concat(), "tail");
1268    }
1269
1270    #[test]
1271    fn compact_messages_on_handle_none_when_under_budget() {
1272        let handle: std::sync::Arc<std::sync::Mutex<Vec<Message>>> =
1273            std::sync::Arc::new(std::sync::Mutex::new(vec![user("short")]));
1274        assert!(compact_messages_on_handle(&handle, "g".into(), 100_000).is_none());
1275    }
1276
1277    #[test]
1278    fn compact_messages_on_handle_none_when_summary_would_not_shrink() {
1279        // 4 messages, all tiny → find_compact_range returns a range but the
1280        // summary message itself is comparable in size, so after >= before.
1281        // Construct a case where find_compact_range returns Some but shrink
1282        // check rejects it: make the range cover near-empty messages so the
1283        // summary overhead exceeds the savings.
1284        let handle: std::sync::Arc<std::sync::Mutex<Vec<Message>>> =
1285            std::sync::Arc::new(std::sync::Mutex::new(vec![
1286                system("h"),
1287                user("."),
1288                assistant("."),
1289                user("."),
1290                assistant("."),
1291                user("t"),
1292            ]));
1293        // budget=1 forces a range, but messages are so small the summary won't help
1294        let result = compact_messages_on_handle(&handle, "x".into(), 1);
1295        // Either no range found (len 6 but tiny), or shrink rejected.
1296        // The key invariant: handle is unchanged if None.
1297        let before_len = handle.lock().unwrap().len();
1298        if result.is_none() {
1299            assert_eq!(handle.lock().unwrap().len(), before_len);
1300        }
1301    }
1302
1303    #[test]
1304    fn replace_range_puts_summary_system_message_in_place() {
1305        let msgs = vec![
1306            system("head"),
1307            user("m1"),
1308            assistant("m2"),
1309            user("m3"),
1310            assistant("m4"),
1311            user("tail"),
1312        ];
1313        let range = CompactRange {
1314            start: 1,
1315            end: 5,
1316            tokens_saved_estimate: 100,
1317        };
1318        let out = replace_range_with_summary(
1319            &msgs,
1320            &range,
1321            "gist: talked about m1..m4".into(),
1322            TurnId::now(),
1323        );
1324        assert_eq!(out.len(), 2, "summary + tail");
1325        assert_eq!(out[0].role, MessageRole::System);
1326        assert!(out[0].text_concat().contains("gist: talked about"));
1327        assert!(matches!(
1328            out[0].parts.as_slice(),
1329            [MessagePart::CompactSummary {
1330                seq_start: 1,
1331                seq_end: 4,
1332                count: 4,
1333                ..
1334            }]
1335        ));
1336        assert_eq!(out[1].role, MessageRole::User);
1337        assert_eq!(out[1].text_concat(), "tail");
1338    }
1339
1340    #[test]
1341    fn find_compact_range_anchors_on_latest_structured_summary() {
1342        let mut msgs = vec![
1343            system("head"),
1344            Message::system_compact_summary(TurnId::now(), "old", 0, 1, 2),
1345        ];
1346        msgs.extend(
1347            (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
1348        );
1349        msgs.extend((0..12).map(|index| assistant(&format!("recent {index}"))));
1350        let range = find_compact_range(&msgs, 1).expect("range");
1351        assert_eq!(range.start, 1);
1352        assert_eq!(range.end, 4);
1353    }
1354
1355    fn assistant_with_tool_use(text: &str, tool_name: &str, input: serde_json::Value) -> Message {
1356        Message {
1357            role: MessageRole::Assistant,
1358            parts: vec![
1359                MessagePart::Text { text: text.into() },
1360                MessagePart::ToolUse {
1361                    id: "call_test".into(),
1362                    name: tool_name.into(),
1363                    input,
1364                },
1365            ],
1366            turn_id: TurnId::now(),
1367            origin: MessageOrigin::User,
1368        }
1369    }
1370
1371    fn tool_result(id: &str, content: &str, is_error: bool) -> Message {
1372        Message {
1373            role: MessageRole::Tool,
1374            parts: vec![MessagePart::ToolResult {
1375                tool_use_id: id.into(),
1376                content: content.into(),
1377                is_error,
1378            }],
1379            turn_id: TurnId::now(),
1380            origin: MessageOrigin::User,
1381        }
1382    }
1383
1384    fn thinking(text: &str) -> Message {
1385        Message {
1386            role: MessageRole::Assistant,
1387            parts: vec![
1388                MessagePart::Thinking {
1389                    thinking: text.into(),
1390                    signature: None,
1391                },
1392                MessagePart::Text {
1393                    text: "after thinking".into(),
1394                },
1395            ],
1396            turn_id: TurnId::now(),
1397            origin: MessageOrigin::User,
1398        }
1399    }
1400
1401    #[test]
1402    fn format_slice_for_summary_includes_tool_use() {
1403        let slice = vec![
1404            user("read the file"),
1405            assistant_with_tool_use(
1406                "let me check",
1407                "fs.read",
1408                serde_json::json!({"path": "/tmp/foo.rs"}),
1409            ),
1410            tool_result("call_test", "fn main() {}", false),
1411        ];
1412        let out = format_slice_for_summary(&slice);
1413        assert!(out.contains("fs.read"), "missing tool name: {out}");
1414        assert!(out.contains("/tmp/foo.rs"), "missing tool input: {out}");
1415        assert!(
1416            out.contains("fn main()"),
1417            "missing tool_result content: {out}"
1418        );
1419        assert!(out.contains("tool_call"), "missing tool_call marker: {out}");
1420        assert!(
1421            out.contains("tool_result"),
1422            "missing tool_result marker: {out}"
1423        );
1424    }
1425
1426    #[test]
1427    fn format_slice_for_summary_includes_thinking() {
1428        let slice = vec![thinking("I should consider the edge case")];
1429        let out = format_slice_for_summary(&slice);
1430        assert!(out.contains("thinking"), "missing thinking marker: {out}");
1431        assert!(out.contains("edge case"), "missing thinking content: {out}");
1432    }
1433
1434    #[test]
1435    fn format_slice_for_summary_marks_error_tool_results() {
1436        let slice = vec![tool_result("call_1", "permission denied", true)];
1437        let out = format_slice_for_summary(&slice);
1438        assert!(out.contains("ERROR"), "missing ERROR marker: {out}");
1439    }
1440
1441    #[test]
1442    fn format_slice_for_summary_truncates_long_tool_input() {
1443        let long_input = serde_json::json!({"content": "x".repeat(5000)});
1444        let slice = vec![assistant_with_tool_use("check", "fs.write", long_input)];
1445        let out = format_slice_for_summary(&slice);
1446        let tool_call_line = out
1447            .lines()
1448            .find(|l| l.contains("tool_call"))
1449            .unwrap_or_else(|| panic!("no tool_call line in {out}"));
1450        assert!(
1451            tool_call_line.chars().count() < 2200,
1452            "tool_call line not truncated: {tool_call_line}"
1453        );
1454    }
1455
1456    fn compaction_summary(text: &str) -> Message {
1457        Message::system_compact_summary(TurnId::now(), text, 1, 5, 5)
1458    }
1459
1460    #[test]
1461    fn is_compaction_summary_detects_structured_variant() {
1462        assert!(is_compaction_summary(&compaction_summary("gist")));
1463        assert!(!is_compaction_summary(&system("plain system msg")));
1464        assert!(!is_compaction_summary(&user("user msg")));
1465    }
1466
1467    #[test]
1468    fn find_compact_range_spans_across_compaction_summaries() {
1469        let mut msgs = vec![
1470            system("head"),
1471            user(&"x".repeat(3000)),
1472            assistant(&"y".repeat(3000)),
1473            user(&"z".repeat(3000)),
1474            compaction_summary("first compaction summary"),
1475        ];
1476        msgs.extend(
1477            (0..6).flat_map(|index| [user(&format!("user {index}")), assistant(&"x".repeat(3000))]),
1478        );
1479        msgs.extend((0..10).map(|index| assistant(&format!("recent {index}"))));
1480        let range = find_compact_range(&msgs, 500).expect("expected range across summary");
1481        assert_eq!(
1482            range.start, 4,
1483            "range should anchor at the structured summary"
1484        );
1485        assert!(
1486            range.end > 4,
1487            "range should include later work, got {range:?}"
1488        );
1489        assert!(
1490            range.end - range.start >= 3,
1491            "range must cover >= 3 msgs, got {}",
1492            range.end - range.start
1493        );
1494    }
1495
1496    #[test]
1497    fn find_compact_starts_from_summary() {
1498        let mut msgs = vec![user("a"), assistant("b"), compaction_summary("summary 1")];
1499        msgs.extend(
1500            (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
1501        );
1502        msgs.extend((0..12).map(|index| assistant(&format!("recent {index}"))));
1503        let range = find_compact_range(&msgs, 10).expect("expected range");
1504        assert_eq!(
1505            range.start, 2,
1506            "should start from the compact summary anchor"
1507        );
1508        assert_eq!(range.end, 5, "the fifth recent user is retained");
1509    }
1510
1511    #[test]
1512    fn find_compact_range_includes_older_compaction_summaries() {
1513        let mut msgs = vec![compaction_summary("summary 0")];
1514        msgs.extend((0..6).flat_map(|index| {
1515            [
1516                user(&format!("old user {index}")),
1517                assistant(&"x".repeat(2000)),
1518            ]
1519        }));
1520        msgs.push(compaction_summary("summary 1"));
1521        msgs.extend((0..6).flat_map(|index| {
1522            [
1523                user(&format!("new user {index}")),
1524                assistant(&"z".repeat(2000)),
1525            ]
1526        }));
1527        msgs.extend((0..10).map(|index| assistant(&format!("tail {index}"))));
1528        let range = find_compact_range(&msgs, 500).expect("expected range");
1529        assert_eq!(range.start, 13, "should compact from the latest summary");
1530        assert!(
1531            range.end > range.start,
1532            "should include work after the latest summary"
1533        );
1534    }
1535
1536    #[test]
1537    fn compacted_message_tokens_detects_growth() {
1538        let msgs = vec![compaction_summary("summary 0"), user("a"), assistant("b")];
1539        let range = CompactRange {
1540            start: 1,
1541            end: 3,
1542            tokens_saved_estimate: 0,
1543        };
1544        let before = estimate_tokens_for_messages(&msgs);
1545        let after = estimate_compacted_message_tokens(
1546            &msgs,
1547            &range,
1548            "a very long summary that expands the transcript a lot",
1549        );
1550        assert!(after > before, "expected growth to be detectable");
1551    }
1552
1553    #[test]
1554    fn find_compact_starts_from_zero_without_summary() {
1555        let mut msgs = (0..26)
1556            .map(|index| assistant(&format!("old {index}")))
1557            .collect::<Vec<_>>();
1558        msgs.extend(
1559            (0..6).flat_map(|index| [user(&format!("user {index}")), assistant("assistant")]),
1560        );
1561        let range = find_compact_range(&msgs, 10).expect("expected range");
1562        assert_eq!(range.start, 0, "should start from 0 without summary");
1563        assert_eq!(range.end, 28, "the recent-message limit is retained");
1564    }
1565
1566    #[test]
1567    fn filter_orphan_tool_messages_removes_orphan_results() {
1568        use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
1569        let turn = TurnId::now();
1570        let msgs = vec![
1571            Message {
1572                role: MessageRole::Tool,
1573                parts: vec![MessagePart::ToolResult {
1574                    tool_use_id: "orphan".into(),
1575                    content: "no matching use".into(),
1576                    is_error: false,
1577                }],
1578                turn_id: turn.clone(),
1579                origin: MessageOrigin::User,
1580            },
1581            Message {
1582                role: MessageRole::Assistant,
1583                parts: vec![MessagePart::ToolUse {
1584                    id: "call_1".into(),
1585                    name: "fs.read".into(),
1586                    input: serde_json::json!({}),
1587                }],
1588                turn_id: turn.clone(),
1589                origin: MessageOrigin::User,
1590            },
1591            Message {
1592                role: MessageRole::Tool,
1593                parts: vec![MessagePart::ToolResult {
1594                    tool_use_id: "call_1".into(),
1595                    content: "ok".into(),
1596                    is_error: false,
1597                }],
1598                turn_id: turn,
1599                origin: MessageOrigin::User,
1600            },
1601        ];
1602        let mut filtered = msgs;
1603        filter_orphan_tool_messages(&mut filtered);
1604        assert_eq!(filtered.len(), 2, "orphan result should be removed");
1605    }
1606}