Skip to main content

bamboo_compression/
compression_tooling.rs

1use crate::counter::{TiktokenTokenCounter, TokenCounter};
2use crate::limits::{create_budget_for_model, ModelLimitsRegistry};
3use crate::{BudgetStrategy, TokenBudget};
4use bamboo_domain::MessagePhase;
5use bamboo_domain::{
6    CompressionEvent, CompressionTriggerType, ConversationSummary, Message, Session,
7};
8
9/// Checks if a message is part of a skill tool chain (load_skill / read_skill_resource).
10fn is_skill_tool_chain_message(message: &Message) -> bool {
11    message.tool_calls.as_ref().is_some_and(|calls| {
12        calls.iter().any(|call| {
13            matches!(
14                call.function.name.as_str(),
15                "load_skill" | "read_skill_resource"
16            )
17        })
18    })
19}
20
21/// The `tool_call_id`s a message participates in: the ids of the tool calls an
22/// assistant message initiates, plus the id of the call a `tool` result answers.
23/// Two messages belong to the same tool chain iff these overlap.
24fn tool_chain_call_ids(message: &Message) -> impl Iterator<Item = String> + '_ {
25    message
26        .tool_calls
27        .iter()
28        .flatten()
29        .map(|call| call.id.clone())
30        .chain(message.tool_call_id.clone())
31}
32
33/// Close the compressed (`messages_to_summarize`) set over tool chains: repeatedly
34/// move any message still in `messages_to_keep` that shares a `tool_call_id` with
35/// an already-compressed message into the summarize set. This keeps an assistant
36/// `tool_calls` message and its matching `tool` result(s) on the same side — a
37/// split leaves an orphan `tool_result` (or a `tool_use` with no result) in the
38/// active set, which providers reject with a 400 that then poisons every
39/// subsequent request in the session (#340). Protected messages are left in place
40/// (skill chains are already fully protected upstream, so are never partially
41/// compressed; protected user messages carry no `tool_call_id`).
42fn close_compressed_set_over_tool_chains(
43    messages_to_keep: &mut Vec<Message>,
44    messages_to_summarize: &mut Vec<Message>,
45    protected_user_ids: &HashSet<String>,
46    never_compress_ids: &[String],
47) {
48    loop {
49        let compressed_call_ids: HashSet<String> = messages_to_summarize
50            .iter()
51            .flat_map(tool_chain_call_ids)
52            .collect();
53        let split_index = messages_to_keep.iter().position(|message| {
54            !protected_user_ids.contains(message.id.as_str())
55                && !never_compress_ids.contains(&message.id)
56                && tool_chain_call_ids(message).any(|id| compressed_call_ids.contains(&id))
57        });
58        match split_index {
59            Some(index) => messages_to_summarize.push(messages_to_keep.remove(index)),
60            None => break,
61        }
62    }
63}
64use chrono::Utc;
65use std::collections::HashSet;
66
67/// Structured reason why a compression plan could not be built.
68#[derive(Debug, Clone)]
69pub enum CompressionPlanError {
70    /// The exposure gate (threshold not reached) prevented building.
71    ExposureGateNotMet {
72        usage_percent: f64,
73        trigger_percent: u8,
74    },
75    /// No active messages in the session.
76    NoActiveMessages,
77    /// Not enough non-system messages to compress (need >=3).
78    NotEnoughMessages { non_system_count: usize },
79    /// Nothing to compress after anchor/keep splitting.
80    NothingToCompress {
81        anchor_index: usize,
82        non_system_count: usize,
83    },
84    /// Eligible history was exhausted while protected/recent content still kept
85    /// the active prompt above the configured post-compression target.
86    ProtectedContentExceedsTarget {
87        projected_tokens: u32,
88        target_tokens: u32,
89    },
90    /// The session changed after candidate selection and before finalization.
91    CandidateSetChanged,
92    /// The real generated summary was larger than the source-derived reserve,
93    /// so applying the candidate set would miss the post-compression target.
94    SummaryExceedsTarget {
95        projected_tokens: u32,
96        target_tokens: u32,
97    },
98}
99
100impl std::fmt::Display for CompressionPlanError {
101    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
102        match self {
103            Self::ExposureGateNotMet {
104                usage_percent,
105                trigger_percent,
106            } => write!(
107                f,
108                "compression threshold not reached (usage={:.1}%, trigger={}%)",
109                usage_percent, trigger_percent
110            ),
111            Self::NoActiveMessages => write!(f, "no active messages to compress"),
112            Self::NotEnoughMessages { non_system_count } => write!(
113                f,
114                "not enough non-system messages to compress ({}, need >=3)",
115                non_system_count
116            ),
117            Self::NothingToCompress {
118                anchor_index,
119                non_system_count,
120            } => write!(
121                f,
122                "nothing to compress after anchor/keep splitting (anchor_index={}, non_system={})",
123                anchor_index, non_system_count
124            ),
125            Self::ProtectedContentExceedsTarget {
126                projected_tokens,
127                target_tokens,
128            } => write!(
129                f,
130                "protected active content prevents compression target (projected={}, target={})",
131                projected_tokens, target_tokens
132            ),
133            Self::CandidateSetChanged => {
134                write!(f, "compression candidate set changed before finalization")
135            }
136            Self::SummaryExceedsTarget {
137                projected_tokens,
138                target_tokens,
139            } => write!(
140                f,
141                "actual summary misses compression target (projected={}, target={})",
142                projected_tokens, target_tokens
143            ),
144        }
145    }
146}
147
148/// Metadata about current context pressure, used to decide when compression
149/// should be requested by host-side control flow.
150#[derive(Debug, Clone)]
151pub struct ContextCompressionExposure {
152    pub budget: TokenBudget,
153    pub active_tokens: u32,
154    pub active_usage_percent: f64,
155    pub active_usage_percent_rounded: u8,
156    pub should_expose_tool: bool,
157}
158
159/// A compression plan describing which active historical messages should be
160/// archived and summarized.
161#[derive(Debug, Clone)]
162pub struct CompressionPlan {
163    /// Stable identifier for all requests and the persisted event belonging to
164    /// one logical compression pass.
165    pub logical_pass_id: Option<String>,
166    /// Tokens from active prompt blocks rendered outside `Session.messages`.
167    pub fixed_prompt_tokens: u32,
168    pub compressed_message_ids: Vec<String>,
169    pub messages_to_summarize: Vec<Message>,
170    pub summary_tokens: u32,
171    pub summary_content: String,
172    pub active_usage_before_percent: f64,
173    pub active_usage_after_percent: f64,
174    pub trigger_percent: u8,
175    pub target_percent: u8,
176    pub segments_removed: usize,
177    pub trigger_type: CompressionTriggerType,
178    pub compression_ratio: f64,
179    pub model_used: Option<String>,
180    pub latency_ms: u64,
181    pub source_tokens: u32,
182    pub represented_source_tokens: u32,
183    pub target_summary_tokens: u32,
184    pub actual_summary_content_tokens: u32,
185    pub summary_target_ratio: f64,
186    pub summary_budget_clamped: bool,
187    pub summary_budget_clamp_reason: Option<String>,
188    pub summarization_map_calls: u32,
189    pub summarization_reduce_calls: u32,
190    pub summarization_fallback_used: bool,
191}
192
193/// Immutable archive selection produced before any summarization request is
194/// sent. The final summary must represent exactly this set.
195#[derive(Debug, Clone)]
196pub struct CompressionCandidatePlan {
197    pub compressed_message_ids: Vec<String>,
198    pub messages_to_summarize: Vec<Message>,
199    pub source_tokens: u32,
200    pub previous_represented_source_tokens: u32,
201    pub represented_source_tokens: u32,
202    pub target_summary_tokens: u32,
203    pub summary_target_ratio: f64,
204    pub active_usage_before_percent: f64,
205    pub projected_usage_after_percent: f64,
206    pub trigger_percent: u8,
207    pub target_percent: u8,
208    pub segments_removed: usize,
209    pub trigger_type: CompressionTriggerType,
210    additional_fixed_tokens: u32,
211    context_window: u32,
212    target_limit: u32,
213}
214
215pub const DEFAULT_SUMMARY_TARGET_RATIO: f64 = 0.20;
216
217pub fn normalized_summary_target_ratio(value: f64) -> f64 {
218    if value.is_finite() && value > 0.0 {
219        value.clamp(0.01, 0.50)
220    } else {
221        DEFAULT_SUMMARY_TARGET_RATIO
222    }
223}
224
225fn target_summary_content_tokens(
226    session: &Session,
227    counter: &impl TokenCounter,
228    newly_represented_tokens: u32,
229    target_ratio: f64,
230) -> (u32, u32) {
231    let target_ratio = normalized_summary_target_ratio(target_ratio);
232    match session.conversation_summary.as_ref() {
233        Some(summary) if summary.represented_source_tokens > 0 => {
234            let represented = summary
235                .represented_source_tokens
236                .saturating_add(newly_represented_tokens);
237            (
238                ((represented as f64) * target_ratio).ceil() as u32,
239                summary.represented_source_tokens,
240            )
241        }
242        Some(summary) => {
243            let existing_tokens = counter.count_text(&summary.content);
244            let inferred_previous_source = ((existing_tokens as f64) / target_ratio).ceil() as u32;
245            (
246                existing_tokens.saturating_add(
247                    ((newly_represented_tokens as f64) * target_ratio).ceil() as u32,
248                ),
249                inferred_previous_source,
250            )
251        }
252        None => (
253            ((newly_represented_tokens as f64) * target_ratio).ceil() as u32,
254            0,
255        ),
256    }
257}
258
259fn protected_compression_message_ids(non_system: &[Message]) -> HashSet<String> {
260    let user_indexes = non_system
261        .iter()
262        .enumerate()
263        .filter_map(|(index, message)| {
264            matches!(message.role, bamboo_domain::Role::User).then_some(index)
265        })
266        .collect::<Vec<_>>();
267    let keep_user_count = user_indexes.len().min(3);
268    let mut protected = user_indexes[user_indexes.len().saturating_sub(keep_user_count)..]
269        .iter()
270        .filter_map(|index| non_system.get(*index))
271        .map(|message| message.id.clone())
272        .collect::<HashSet<_>>();
273
274    let skill_call_ids = non_system
275        .iter()
276        .filter(|message| is_skill_tool_chain_message(message))
277        .flat_map(|message| {
278            message
279                .tool_calls
280                .iter()
281                .flatten()
282                .map(|call| call.id.clone())
283        })
284        .collect::<HashSet<_>>();
285
286    for message in non_system {
287        if message.never_compress
288            || is_skill_tool_chain_message(message)
289            || message
290                .tool_call_id
291                .as_ref()
292                .is_some_and(|id| skill_call_ids.contains(id))
293        {
294            protected.insert(message.id.clone());
295        }
296    }
297    protected
298}
299
300fn post_compaction_recovery_tokens(
301    compressed_messages: &[Message],
302    session: &Session,
303    counter: &impl TokenCounter,
304) -> u32 {
305    build_post_compaction_recovery_message(compressed_messages, session)
306        .as_ref()
307        .map(|message| counter.count_message(message))
308        .unwrap_or(0)
309}
310
311/// Select the exact archive candidates before invoking the summarization model.
312///
313/// Candidate segments are considered oldest-first. Generic tool chains remain
314/// atomic because selection operates on [`crate::MessageSegmenter`] output;
315/// newest user turns, `never_compress`, and skill chains remain active.
316pub fn build_forced_compression_candidate_plan(
317    session: &Session,
318    model_name: &str,
319    configured_budget: Option<&TokenBudget>,
320    summary_target_ratio: f64,
321    trigger_type: CompressionTriggerType,
322) -> Result<CompressionCandidatePlan, CompressionPlanError> {
323    build_forced_compression_candidate_plan_with_fixed_tokens(
324        session,
325        model_name,
326        configured_budget,
327        summary_target_ratio,
328        trigger_type,
329        0,
330    )
331}
332
333/// Variant of [`build_forced_compression_candidate_plan`] that accounts for
334/// fixed prompt blocks rendered outside `Session.messages` (for example task,
335/// plan, workflow, project-resource, and external-memory context blocks).
336pub fn build_forced_compression_candidate_plan_with_fixed_tokens(
337    session: &Session,
338    model_name: &str,
339    configured_budget: Option<&TokenBudget>,
340    summary_target_ratio: f64,
341    trigger_type: CompressionTriggerType,
342    additional_fixed_tokens: u32,
343) -> Result<CompressionCandidatePlan, CompressionPlanError> {
344    let exposure = estimate_context_compression_exposure(session, model_name, configured_budget);
345    let budget = &exposure.budget;
346    let counter = TiktokenTokenCounter::default();
347    let active_messages = active_messages_for_budget(session);
348    if active_messages.is_empty() {
349        return Err(CompressionPlanError::NoActiveMessages);
350    }
351
352    let system_messages = active_messages
353        .iter()
354        .filter(|message| matches!(message.role, bamboo_domain::Role::System))
355        .cloned()
356        .collect::<Vec<_>>();
357    let non_system = active_messages
358        .into_iter()
359        .filter(|message| !matches!(message.role, bamboo_domain::Role::System))
360        .collect::<Vec<_>>();
361    if non_system.len() < 3 {
362        return Err(CompressionPlanError::NotEnoughMessages {
363            non_system_count: non_system.len(),
364        });
365    }
366
367    let context_window = budget.max_context_tokens;
368    let target_limit = budget.compression_target_context_tokens();
369    let system_tokens = counter.count_messages(&system_messages);
370    let protected_ids = protected_compression_message_ids(&non_system);
371    let segments = crate::segmenter::MessageSegmenter::new().segment(non_system.clone());
372    let mut remaining_tokens = counter.count_messages(&non_system);
373    let mut source_tokens = 0u32;
374    let mut selected_messages = Vec::new();
375    let mut selected_segment_count = 0usize;
376    let ratio = normalized_summary_target_ratio(summary_target_ratio);
377    let summary_envelope_tokens = counter.count_messages(&[compression_summary_message("")]);
378    let mut projected_tokens = system_tokens
379        .saturating_add(remaining_tokens)
380        .saturating_add(additional_fixed_tokens);
381    let mut target_summary_tokens = 0u32;
382    let mut previous_represented_source_tokens = session
383        .conversation_summary
384        .as_ref()
385        .map(|summary| summary.represented_source_tokens)
386        .unwrap_or(0);
387
388    for segment in segments {
389        if segment
390            .messages
391            .iter()
392            .any(|message| protected_ids.contains(&message.id))
393        {
394            continue;
395        }
396
397        let segment_tokens = counter.count_messages(&segment.messages);
398        source_tokens = source_tokens.saturating_add(segment_tokens);
399        remaining_tokens = remaining_tokens.saturating_sub(segment_tokens);
400        selected_messages.extend(segment.messages);
401        selected_segment_count += 1;
402
403        let (desired, previous_represented) =
404            target_summary_content_tokens(session, &counter, source_tokens, ratio);
405        target_summary_tokens = desired;
406        previous_represented_source_tokens = previous_represented;
407        let recovery_tokens =
408            post_compaction_recovery_tokens(&selected_messages, session, &counter);
409        projected_tokens = system_tokens
410            .saturating_add(remaining_tokens)
411            .saturating_add(additional_fixed_tokens)
412            .saturating_add(summary_envelope_tokens)
413            .saturating_add(target_summary_tokens)
414            .saturating_add(recovery_tokens);
415        if projected_tokens <= target_limit {
416            break;
417        }
418    }
419
420    if selected_messages.is_empty() {
421        return Err(CompressionPlanError::NothingToCompress {
422            anchor_index: 0,
423            non_system_count: non_system.len(),
424        });
425    }
426    if projected_tokens > target_limit {
427        return Err(CompressionPlanError::ProtectedContentExceedsTarget {
428            projected_tokens,
429            target_tokens: target_limit,
430        });
431    }
432
433    let represented_source_tokens =
434        previous_represented_source_tokens.saturating_add(source_tokens);
435    let projected_usage_after_percent =
436        context_window_usage_percent(projected_tokens, context_window);
437    let compressed_message_ids = selected_messages
438        .iter()
439        .map(|message| message.id.clone())
440        .collect();
441
442    Ok(CompressionCandidatePlan {
443        compressed_message_ids,
444        messages_to_summarize: selected_messages,
445        source_tokens,
446        previous_represented_source_tokens,
447        represented_source_tokens,
448        target_summary_tokens,
449        summary_target_ratio: ratio,
450        active_usage_before_percent: exposure.active_usage_percent,
451        projected_usage_after_percent,
452        trigger_percent: budget.compression_trigger_percent,
453        target_percent: budget.compression_target_percent,
454        segments_removed: selected_segment_count,
455        trigger_type,
456        additional_fixed_tokens,
457        context_window,
458        target_limit,
459    })
460}
461
462/// Bind a completed summary to a previously selected candidate set and verify
463/// the real post-compression prompt before any session mutation occurs.
464pub fn finalize_compression_candidate_plan(
465    session: &Session,
466    candidate: CompressionCandidatePlan,
467    summary_content: String,
468) -> Result<CompressionPlan, CompressionPlanError> {
469    let active_ids = session
470        .messages
471        .iter()
472        .filter(|message| !message.compressed)
473        .map(|message| message.id.as_str())
474        .collect::<HashSet<_>>();
475    if candidate
476        .compressed_message_ids
477        .iter()
478        .any(|id| !active_ids.contains(id.as_str()))
479    {
480        return Err(CompressionPlanError::CandidateSetChanged);
481    }
482
483    let candidate_ids = candidate
484        .compressed_message_ids
485        .iter()
486        .map(String::as_str)
487        .collect::<HashSet<_>>();
488    let counter = TiktokenTokenCounter::default();
489    let summary_tokens = counter.count_messages(&[compression_summary_message(&summary_content)]);
490    let actual_summary_content_tokens = counter.count_text(&summary_content);
491    let recovery_tokens =
492        post_compaction_recovery_tokens(&candidate.messages_to_summarize, session, &counter);
493    let remaining_tokens = session
494        .messages
495        .iter()
496        .filter(|message| !message.compressed)
497        .filter(|message| !candidate_ids.contains(message.id.as_str()))
498        .cloned()
499        .collect::<Vec<_>>();
500    let projected_tokens = counter
501        .count_messages(&remaining_tokens)
502        .saturating_add(candidate.additional_fixed_tokens)
503        .saturating_add(summary_tokens)
504        .saturating_add(recovery_tokens);
505    if projected_tokens > candidate.target_limit {
506        return Err(CompressionPlanError::SummaryExceedsTarget {
507            projected_tokens,
508            target_tokens: candidate.target_limit,
509        });
510    }
511
512    Ok(CompressionPlan {
513        logical_pass_id: None,
514        fixed_prompt_tokens: candidate.additional_fixed_tokens,
515        compressed_message_ids: candidate.compressed_message_ids,
516        messages_to_summarize: candidate.messages_to_summarize,
517        summary_tokens,
518        summary_content,
519        active_usage_before_percent: candidate.active_usage_before_percent,
520        active_usage_after_percent: context_window_usage_percent(
521            projected_tokens,
522            candidate.context_window,
523        ),
524        trigger_percent: candidate.trigger_percent,
525        target_percent: candidate.target_percent,
526        segments_removed: candidate.segments_removed,
527        trigger_type: candidate.trigger_type,
528        compression_ratio: 0.0,
529        model_used: None,
530        latency_ms: 0,
531        source_tokens: candidate.source_tokens,
532        represented_source_tokens: candidate.represented_source_tokens,
533        target_summary_tokens: candidate.target_summary_tokens,
534        actual_summary_content_tokens,
535        summary_target_ratio: candidate.summary_target_ratio,
536        summary_budget_clamped: false,
537        summary_budget_clamp_reason: None,
538        summarization_map_calls: 0,
539        summarization_reduce_calls: 0,
540        summarization_fallback_used: false,
541    })
542}
543
544pub fn context_window_usage_percent(total_tokens: u32, context_window_tokens: u32) -> f64 {
545    if context_window_tokens == 0 {
546        return 0.0;
547    }
548    (total_tokens as f64 / context_window_tokens as f64) * 100.0
549}
550
551pub fn normalized_trigger_percent(trigger_percent: u8) -> f64 {
552    match trigger_percent {
553        0 => 100.0,
554        1..=100 => trigger_percent as f64,
555        _ => 100.0,
556    }
557}
558
559/// Estimate whether context pressure has crossed the configured threshold for
560/// compression eligibility.
561pub fn estimate_context_compression_exposure(
562    session: &Session,
563    model_name: &str,
564    configured_budget: Option<&TokenBudget>,
565) -> ContextCompressionExposure {
566    // When a budget was already resolved upstream (the production path — see
567    // `resolve_token_budget`, which publishes the current-round snapshot in
568    // `session.resolved_token_budget` (#180),
569    // issue #20 bug 1), use it directly. Only when none is available do we fall
570    // back to a model-derived budget. No `model_limits.json` registry is in
571    // scope synchronously here, so this fallback resolves to the global default
572    // rather than silently fabricating an empty override registry (#20 bug 2).
573    let budget = configured_budget.cloned().unwrap_or_else(|| {
574        create_budget_for_model(
575            model_name,
576            BudgetStrategy::default(),
577            &ModelLimitsRegistry::new(),
578        )
579    });
580    let counter = TiktokenTokenCounter::default();
581    let active_messages = active_messages_for_budget(session);
582    let active_message_tokens = counter.count_messages(&active_messages);
583    let summary_tokens = session
584        .conversation_summary
585        .as_ref()
586        .map(|summary| counter.count_messages(&[compression_summary_message(&summary.content)]))
587        .unwrap_or(0);
588    let active_tokens = active_message_tokens.saturating_add(summary_tokens);
589    // Use context window as the denominator for a single, provider-aligned
590    // pressure scale across backend and frontend.
591    let context_window = budget.max_context_tokens;
592    let estimated_usage = context_window_usage_percent(active_tokens, context_window);
593    let usage = session
594        .token_usage
595        .as_ref()
596        .and_then(|token_usage| {
597            let denominator = if token_usage.max_context_tokens > 0 {
598                token_usage.max_context_tokens
599            } else if token_usage.budget_limit > 0 {
600                // Legacy payload compatibility.
601                token_usage.budget_limit
602            } else {
603                context_window
604            };
605            (denominator > 0).then_some(context_window_usage_percent(
606                token_usage.total_tokens,
607                denominator,
608            ))
609        })
610        .map(|persisted_usage| persisted_usage.max(estimated_usage))
611        .unwrap_or(estimated_usage);
612
613    let rounded = usage.clamp(0.0, 100.0).round() as u8;
614    let trigger_tokens = budget.compression_trigger_context_tokens();
615    let trigger_percent = if budget.max_context_tokens > 0 {
616        (trigger_tokens as f64 / budget.max_context_tokens as f64) * 100.0
617    } else {
618        0.0
619    };
620    let threshold_reached = usage >= trigger_percent;
621
622    // Check non-system message count to stay consistent with the plan
623    // building requirement of >=3 non-system messages.  Using
624    // active_messages.len() would include system messages and expose the
625    // tool even when plan building would immediately fail.
626    let non_system_count = active_messages
627        .iter()
628        .filter(|m| !matches!(m.role, bamboo_domain::Role::System))
629        .count();
630
631    let should_expose_tool = threshold_reached && non_system_count >= 3;
632
633    ContextCompressionExposure {
634        budget,
635        active_tokens,
636        active_usage_percent: usage,
637        active_usage_percent_rounded: rounded,
638        should_expose_tool,
639    }
640}
641
642/// Build a compression plan that archives older active messages and replaces
643/// them with a caller-provided summary.
644pub fn build_compression_plan_with_summary(
645    session: &Session,
646    model_name: &str,
647    configured_budget: Option<&TokenBudget>,
648    summary_content: String,
649) -> Result<CompressionPlan, CompressionPlanError> {
650    build_compression_plan_with_summary_internal(
651        session,
652        model_name,
653        configured_budget,
654        summary_content,
655        true,
656        CompressionTriggerType::Auto,
657    )
658}
659
660/// Build a compression plan while bypassing "tool exposure" gating.
661///
662/// This is intended for host-enforced fallback paths when context pressure is
663/// critically high and compression must be attempted regardless of the normal
664/// trigger gate.
665pub fn build_forced_compression_plan_with_summary(
666    session: &Session,
667    model_name: &str,
668    configured_budget: Option<&TokenBudget>,
669    summary_content: String,
670    trigger_type: CompressionTriggerType,
671) -> Result<CompressionPlan, CompressionPlanError> {
672    build_compression_plan_with_summary_internal(
673        session,
674        model_name,
675        configured_budget,
676        summary_content,
677        false,
678        trigger_type,
679    )
680}
681
682fn build_compression_plan_with_summary_internal(
683    session: &Session,
684    model_name: &str,
685    configured_budget: Option<&TokenBudget>,
686    summary_content: String,
687    require_exposure_gate: bool,
688    trigger_type: CompressionTriggerType,
689) -> Result<CompressionPlan, CompressionPlanError> {
690    let exposure = estimate_context_compression_exposure(session, model_name, configured_budget);
691    if require_exposure_gate && !exposure.should_expose_tool {
692        return Err(CompressionPlanError::ExposureGateNotMet {
693            usage_percent: exposure.active_usage_percent,
694            trigger_percent: exposure.budget.compression_trigger_percent,
695        });
696    }
697
698    let budget = &exposure.budget;
699    let counter = TiktokenTokenCounter::default();
700    let summary_message = compression_summary_message(&summary_content);
701    let summary_tokens = counter.count_messages(&[summary_message]);
702
703    let context_window = budget.max_context_tokens;
704    let target_limit = budget.compression_target_context_tokens();
705
706    let mut active_messages = active_messages_for_budget(session);
707    if active_messages.is_empty() {
708        tracing::debug!("compression plan: no active messages, cannot build plan");
709        return Err(CompressionPlanError::NoActiveMessages);
710    }
711
712    let system_messages: Vec<Message> = active_messages
713        .iter()
714        .filter(|m| matches!(m.role, bamboo_domain::Role::System))
715        .cloned()
716        .collect();
717    let system_tokens = counter.count_messages(&system_messages);
718    let reserved_non_window_tokens = system_tokens.saturating_add(summary_tokens);
719    let window_limit = target_limit.saturating_sub(reserved_non_window_tokens);
720
721    let non_system: Vec<Message> = active_messages
722        .drain(..)
723        .filter(|m| !matches!(m.role, bamboo_domain::Role::System))
724        .collect();
725
726    if non_system.len() < 3 {
727        tracing::debug!(
728            "compression plan: not enough non-system messages ({}), need at least 3",
729            non_system.len()
730        );
731        return Err(CompressionPlanError::NotEnoughMessages {
732            non_system_count: non_system.len(),
733        });
734    }
735
736    let user_indexes = non_system
737        .iter()
738        .enumerate()
739        .filter_map(|(index, message)| {
740            matches!(message.role, bamboo_domain::Role::User).then_some(index)
741        })
742        .collect::<Vec<_>>();
743    let keep_user_count = user_indexes.len().min(3);
744    let anchor_index = if keep_user_count > 0 {
745        user_indexes[user_indexes.len() - keep_user_count]
746    } else {
747        non_system
748            .iter()
749            .rposition(|m| matches!(m.role, bamboo_domain::Role::User))
750            .unwrap_or(non_system.len().saturating_sub(1))
751    };
752    let protected_user_ids: HashSet<String> = if keep_user_count > 0 {
753        user_indexes[user_indexes.len() - keep_user_count..]
754            .iter()
755            .filter_map(|idx| non_system.get(*idx))
756            .map(|message| message.id.clone())
757            .collect()
758    } else {
759        HashSet::new()
760    };
761
762    tracing::debug!(
763        "compression plan: context_window={}, target_limit={}, system_tokens={}, summary_tokens={}, window_limit={}, non_system_messages={}, keep_user_count={}, keep_from_index={}",
764        context_window, target_limit, system_tokens, summary_tokens, window_limit, non_system.len(), keep_user_count, anchor_index
765    );
766
767    // Keep the newest 3 user turns (or fewer if there are not enough user
768    // turns) as active context and summarize older history before that
769    // boundary. If budget is still too high, continue moving the oldest
770    // non-protected messages into the summarize set.
771    let mut messages_to_summarize = non_system[..anchor_index].to_vec();
772
773    // Protected messages must never be summarized — move them to the keep set.
774    let mut never_compress_ids: Vec<String> = messages_to_summarize
775        .iter()
776        .filter(|m| m.never_compress || is_skill_tool_chain_message(m))
777        .map(|m| m.id.clone())
778        .collect();
779
780    // Also protect tool result messages that correspond to skill tool calls.
781    let skill_call_ids: Vec<String> = messages_to_summarize
782        .iter()
783        .filter(|m| is_skill_tool_chain_message(m))
784        .flat_map(|m| m.tool_calls.iter().flatten().map(|c| c.id.clone()))
785        .collect();
786    if !skill_call_ids.is_empty() {
787        for m in &*messages_to_summarize {
788            if let Some(ref call_id) = m.tool_call_id {
789                if skill_call_ids.contains(call_id) && !never_compress_ids.contains(&m.id) {
790                    never_compress_ids.push(m.id.clone());
791                }
792            }
793        }
794    }
795
796    if !never_compress_ids.is_empty() {
797        messages_to_summarize.retain(|m| !never_compress_ids.contains(&m.id));
798    }
799
800    let non_system_count = non_system.len();
801    let mut messages_to_keep = non_system[anchor_index..].to_vec();
802    // Add never_compress / skill messages to the keep set.
803    for id in &never_compress_ids {
804        if let Some(msg) = non_system.iter().find(|m| &m.id == id) {
805            if !messages_to_keep.iter().any(|m| m.id == *id) {
806                messages_to_keep.push(msg.clone());
807            }
808        }
809    }
810
811    while !messages_to_keep.is_empty() {
812        let keep_tokens = counter.count_messages(&messages_to_keep);
813        if keep_tokens <= window_limit {
814            break;
815        }
816
817        let Some(remove_index) = messages_to_keep.iter().position(|message| {
818            !protected_user_ids.contains(message.id.as_str())
819                && !never_compress_ids.contains(&message.id)
820        }) else {
821            // Remaining messages are all protected; stop shrinking.
822            break;
823        };
824        let moved = messages_to_keep.remove(remove_index);
825        messages_to_summarize.push(moved);
826    }
827
828    // Tool-chain atomicity. The one-at-a-time eviction above can split a generic
829    // (non-skill) tool chain — compressing an assistant `tool_calls` message
830    // while keeping its `tool` result active, or vice versa. That leaves an
831    // orphan `tool_result` (or a `tool_use` with no result) in the active set,
832    // which providers reject with a 400 that then poisons EVERY subsequent
833    // request in the session. Close the compressed set over tool chains: keep
834    // moving any kept message that shares a `tool_call_id` with an
835    // already-compressed message into the summarize set until none remain.
836    // (Skill chains are fully protected above, so are never partially
837    // compressed; protected user messages carry no `tool_call_id`.)
838    close_compressed_set_over_tool_chains(
839        &mut messages_to_keep,
840        &mut messages_to_summarize,
841        &protected_user_ids,
842        &never_compress_ids,
843    );
844
845    if messages_to_summarize.is_empty() {
846        tracing::debug!(
847            "compression plan: messages_to_summarize is empty after anchor/keep splitting"
848        );
849        return Err(CompressionPlanError::NothingToCompress {
850            anchor_index,
851            non_system_count,
852        });
853    }
854
855    let compressed_message_ids = messages_to_summarize
856        .iter()
857        .map(|message| message.id.clone())
858        .collect::<Vec<_>>();
859
860    let keep_tokens = counter.count_messages(&messages_to_keep);
861    let active_before = exposure.active_usage_percent;
862    // Use context_window as denominator, consistent with
863    // estimate_context_compression_exposure().
864    let active_after = if context_window == 0 {
865        0.0
866    } else {
867        let after_total = reserved_non_window_tokens.saturating_add(keep_tokens);
868        (after_total as f64 / context_window as f64) * 100.0
869    };
870
871    // Count actual segments being compressed using the same segmenter that
872    // prepare_hybrid_context uses, so the segment count is accurate.
873    let segmenter = crate::segmenter::MessageSegmenter::new();
874    let segments_removed = segmenter.segment(messages_to_summarize.clone()).len();
875    let source_tokens = counter.count_messages(&messages_to_summarize);
876    let actual_summary_content_tokens = counter.count_text(&summary_content);
877    let previous_represented_source_tokens = session
878        .conversation_summary
879        .as_ref()
880        .map(|summary| summary.represented_source_tokens)
881        .unwrap_or(0);
882
883    Ok(CompressionPlan {
884        logical_pass_id: None,
885        fixed_prompt_tokens: 0,
886        compressed_message_ids,
887        messages_to_summarize,
888        summary_tokens,
889        summary_content,
890        active_usage_before_percent: active_before,
891        active_usage_after_percent: active_after,
892        trigger_percent: budget.compression_trigger_percent,
893        target_percent: budget.compression_target_percent,
894        segments_removed,
895        trigger_type,
896        compression_ratio: 0.0,
897        model_used: None,
898        latency_ms: 0,
899        source_tokens,
900        represented_source_tokens: previous_represented_source_tokens.saturating_add(source_tokens),
901        target_summary_tokens: actual_summary_content_tokens,
902        actual_summary_content_tokens,
903        summary_target_ratio: 0.0,
904        summary_budget_clamped: false,
905        summary_budget_clamp_reason: None,
906        summarization_map_calls: 0,
907        summarization_reduce_calls: 0,
908        summarization_fallback_used: false,
909    })
910}
911
912/// Apply a previously computed compression plan to the session.
913/// Extract recently modified files from tool calls in the given messages.
914pub(super) fn extract_recently_modified_files(messages: &[Message]) -> Vec<(String, String)> {
915    let mut files = Vec::new();
916    for message in messages {
917        if let Some(ref tool_calls) = message.tool_calls {
918            for call in tool_calls {
919                let tool_name = call.function.name.as_str();
920                if !matches!(tool_name, "Write" | "Edit" | "Bash") {
921                    continue;
922                }
923                let args = &call.function.arguments;
924                if let Ok(parsed) = serde_json::from_str::<serde_json::Value>(args) {
925                    if let Some(path) = parsed.get("file_path").and_then(|v| v.as_str()) {
926                        files.push((path.to_string(), tool_name.to_string()));
927                    } else if let Some(cmd) = parsed.get("command").and_then(|v| v.as_str()) {
928                        // Extract file paths from shell commands heuristically
929                        for part in cmd.split_whitespace() {
930                            if part.contains('/')
931                                && (part.ends_with(".rs")
932                                    || part.ends_with(".ts")
933                                    || part.ends_with(".js")
934                                    || part.ends_with(".toml")
935                                    || part.ends_with(".json")
936                                    || part.ends_with(".md"))
937                            {
938                                files.push((part.to_string(), "Bash".to_string()));
939                            }
940                        }
941                    }
942                }
943            }
944        }
945    }
946    files.truncate(10);
947    files
948}
949
950/// Extract key decision snippets from assistant messages.
951pub(super) fn extract_key_decisions(messages: &[Message], limit: usize) -> Vec<String> {
952    let decision_keywords = [
953        "decided to",
954        "approach is",
955        "use ",
956        "using ",
957        "we'll go with",
958        "the plan is",
959        "strategy:",
960        "solution:",
961        "chose to",
962        "switched to",
963        "refactored to",
964        "migrated to",
965        "replaced with",
966    ];
967    let mut decisions = Vec::new();
968    for message in messages {
969        if !matches!(message.role, bamboo_domain::Role::Assistant) {
970            continue;
971        }
972        let content = &message.content;
973        for line in content.lines() {
974            let line_lower = line.to_lowercase();
975            if decision_keywords.iter().any(|kw| line_lower.contains(kw)) {
976                let truncated: String = line.chars().take(200).collect();
977                decisions.push(truncated);
978                if decisions.len() >= limit {
979                    return decisions;
980                }
981            }
982        }
983    }
984    decisions
985}
986
987/// Build a post-compaction recovery message that preserves critical context
988/// from the compressed messages so the LLM can continue work without losing
989/// track of active files, tasks, and decisions.
990fn build_post_compaction_recovery_message(
991    compressed_messages: &[Message],
992    session: &Session,
993) -> Option<Message> {
994    if compressed_messages.is_empty() {
995        return None;
996    }
997
998    let mut sections = Vec::new();
999
1000    // 1. Recently modified files
1001    let files = extract_recently_modified_files(compressed_messages);
1002    if !files.is_empty() {
1003        let mut section = String::from("## Recently Modified Files\n");
1004        for (path, tool) in &files {
1005            section.push_str(&format!("- {} ({})\n", path, tool));
1006        }
1007        sections.push(section);
1008    }
1009
1010    // 2. Active tasks from task list
1011    if let Some(ref task_list) = session.task_list {
1012        let active_items: Vec<_> = task_list
1013            .items
1014            .iter()
1015            .filter(|item| !matches!(item.status, bamboo_domain::TaskItemStatus::Completed))
1016            .collect();
1017        if !active_items.is_empty() {
1018            let mut section = String::from("## Active Tasks\n");
1019            for item in active_items.iter().take(10) {
1020                section.push_str(&format!("- [{:?}] {}\n", item.status, item.description));
1021            }
1022            sections.push(section);
1023        }
1024    }
1025
1026    // 3. Key decisions
1027    let decisions = extract_key_decisions(compressed_messages, 5);
1028    if !decisions.is_empty() {
1029        let mut section = String::from("## Key Decisions\n");
1030        for decision in &decisions {
1031            section.push_str(&format!("- {}\n", decision));
1032        }
1033        sections.push(section);
1034    }
1035
1036    if sections.is_empty() {
1037        return None;
1038    }
1039
1040    let mut content = String::from("[post-compaction-recovery]\nContext extracted from compressed messages for continued work.\n\n");
1041    content.push_str(&sections.join("\n"));
1042
1043    let mut message = Message::assistant(content, None);
1044    message.never_compress = true;
1045    Some(message)
1046}
1047
1048struct SummaryQualityMetrics {
1049    file_coverage: f64,
1050    decision_coverage: f64,
1051}
1052
1053fn validate_summary_quality(summary: &str, messages: &[Message]) -> SummaryQualityMetrics {
1054    let files = extract_recently_modified_files(messages);
1055    let decisions = extract_key_decisions(messages, 10);
1056
1057    let files_mentioned = files
1058        .iter()
1059        .filter(|(path, _)| summary.contains(path.as_str()))
1060        .count();
1061    let file_coverage = if files.is_empty() {
1062        1.0
1063    } else {
1064        files_mentioned as f64 / files.len() as f64
1065    };
1066
1067    let decisions_mentioned = decisions
1068        .iter()
1069        .filter(|d| {
1070            let check_str: String = d.chars().take(50).collect();
1071            summary.contains(&check_str)
1072        })
1073        .count();
1074    let decision_coverage = if decisions.is_empty() {
1075        1.0
1076    } else {
1077        decisions_mentioned as f64 / decisions.len() as f64
1078    };
1079
1080    SummaryQualityMetrics {
1081        file_coverage,
1082        decision_coverage,
1083    }
1084}
1085
1086pub fn apply_compression_plan(session: &mut Session, plan: CompressionPlan) -> usize {
1087    let compressed_ids: HashSet<&str> = plan
1088        .compressed_message_ids
1089        .iter()
1090        .map(String::as_str)
1091        .collect();
1092
1093    let mut changed_indexes = Vec::new();
1094    for (index, message) in session.messages.iter_mut().enumerate() {
1095        if message.compressed || !compressed_ids.contains(message.id.as_str()) {
1096            continue;
1097        }
1098        message.compressed = true;
1099        changed_indexes.push(index);
1100    }
1101
1102    if changed_indexes.is_empty() {
1103        return 0;
1104    }
1105
1106    let mut event = CompressionEvent::new(
1107        changed_indexes.len(),
1108        plan.segments_removed,
1109        plan.active_usage_before_percent,
1110        plan.active_usage_after_percent,
1111        plan.summary_tokens,
1112        plan.trigger_type,
1113        plan.compression_ratio,
1114        plan.model_used.clone(),
1115        plan.latency_ms,
1116    );
1117    if let Some(logical_pass_id) = plan.logical_pass_id.as_ref() {
1118        event.id.clone_from(logical_pass_id);
1119    }
1120    event.source_tokens = plan.source_tokens;
1121    event.fixed_prompt_tokens = plan.fixed_prompt_tokens;
1122    event.actual_summary_tokens = plan.actual_summary_content_tokens;
1123    event.target_summary_tokens = plan.target_summary_tokens;
1124    event.summary_target_ratio = plan.summary_target_ratio;
1125    event.actual_summary_ratio = if plan.represented_source_tokens == 0 {
1126        0.0
1127    } else {
1128        plan.actual_summary_content_tokens as f64 / plan.represented_source_tokens as f64
1129    };
1130    event.summary_budget_clamped = plan.summary_budget_clamped;
1131    event.summary_budget_clamp_reason = plan.summary_budget_clamp_reason.clone();
1132    event.summarization_map_calls = plan.summarization_map_calls;
1133    event.summarization_reduce_calls = plan.summarization_reduce_calls;
1134    event.summarization_fallback_used = plan.summarization_fallback_used;
1135    let event_id = event.id.clone();
1136    for index in changed_indexes {
1137        session.messages[index].compressed_by_event_id = Some(event_id.clone());
1138    }
1139    session.compression_events.push(event);
1140    session.conversation_summary = Some(
1141        ConversationSummary::new(
1142            &plan.summary_content,
1143            plan.compressed_message_ids.len(),
1144            plan.summary_tokens,
1145        )
1146        .with_compression_metrics(
1147            plan.represented_source_tokens,
1148            plan.target_summary_tokens,
1149            plan.summary_target_ratio,
1150            plan.summary_budget_clamped,
1151            plan.summary_budget_clamp_reason.clone(),
1152        ),
1153    );
1154
1155    // Inject a post-compaction recovery message to preserve critical context
1156    // from the compressed messages (files, tasks, decisions).
1157    let compressed_messages: Vec<Message> = session
1158        .messages
1159        .iter()
1160        .filter(|m| compressed_ids.contains(m.id.as_str()))
1161        .cloned()
1162        .collect();
1163    if let Some(recovery) = build_post_compaction_recovery_message(&compressed_messages, session) {
1164        // Insert just before the last user message, or at the end
1165        let insert_pos = session
1166            .messages
1167            .iter()
1168            .rposition(|m| matches!(m.role, bamboo_domain::Role::User) && !m.compressed)
1169            .map(|pos| pos + 1)
1170            .unwrap_or(session.messages.len());
1171        session.messages.insert(insert_pos, recovery);
1172    }
1173
1174    let quality = validate_summary_quality(&plan.summary_content, &compressed_messages);
1175    if quality.file_coverage < 0.5 || quality.decision_coverage < 0.3 {
1176        tracing::warn!(
1177            "[{}] Summary quality: file_coverage={:.0}%, decision_coverage={:.0}%",
1178            session.id,
1179            quality.file_coverage * 100.0,
1180            quality.decision_coverage * 100.0
1181        );
1182    }
1183
1184    // Instead of clearing token_usage entirely (which forces the next round
1185    // to rely on heuristic estimates that don't account for tool schema
1186    // tokens), recompute an approximate post-compression snapshot.  We
1187    // preserve both the total context-window denominator and the request input
1188    // limit from the previous usage snapshot so their meanings stay distinct
1189    // across rounds.
1190    let counter = TiktokenTokenCounter::default();
1191    let remaining_active: Vec<_> = session
1192        .messages
1193        .iter()
1194        .filter(|m| !m.compressed)
1195        .cloned()
1196        .collect();
1197    let system_msgs: Vec<_> = remaining_active
1198        .iter()
1199        .filter(|m| matches!(m.role, bamboo_domain::Role::System))
1200        .cloned()
1201        .collect();
1202    let window_msgs: Vec<_> = remaining_active
1203        .iter()
1204        .filter(|m| !matches!(m.role, bamboo_domain::Role::System))
1205        .cloned()
1206        .collect();
1207    let system_tokens = counter
1208        .count_messages(&system_msgs)
1209        .saturating_add(plan.fixed_prompt_tokens);
1210    let new_summary_tokens = plan.summary_tokens;
1211    let window_tokens = counter.count_messages(&window_msgs);
1212    let total_tokens = system_tokens
1213        .saturating_add(new_summary_tokens)
1214        .saturating_add(window_tokens);
1215    let previous_usage = session.token_usage.take();
1216    let budget_limit = previous_usage
1217        .as_ref()
1218        .map(|u| {
1219            if u.budget_limit > 0 {
1220                u.budget_limit
1221            } else {
1222                // Legacy snapshots used the total context window as the only
1223                // denominator.
1224                u.max_context_tokens
1225            }
1226        })
1227        .unwrap_or(0);
1228    let max_context_tokens = previous_usage
1229        .as_ref()
1230        .map(|u| u.max_context_tokens)
1231        .unwrap_or(0);
1232    session.token_usage = Some(bamboo_domain::TokenBudgetUsage {
1233        system_tokens,
1234        summary_tokens: new_summary_tokens,
1235        window_tokens,
1236        total_tokens,
1237        max_context_tokens,
1238        budget_limit,
1239        truncation_occurred: false,
1240        segments_removed: 0,
1241        prompt_cached_tool_outputs: 0,
1242        prompt_cached_tool_tokens_saved: 0,
1243        thinking_tokens: 0,
1244        cache_read_input_tokens: 0,
1245    });
1246
1247    session.updated_at = Utc::now();
1248    plan.compressed_message_ids.len()
1249}
1250
1251pub fn compression_summary_message(summary_content: &str) -> Message {
1252    Message::system(format!(
1253        "<!-- CONVERSATION_SUMMARY_START -->\n\
1254         ## Previous Conversation Summary\n\
1255         The following is compressed historical context for continuity only.\n\
1256         It is background memory, not a new user request. Follow the current task list and recent messages over this summary when they conflict.\n\n\
1257         {}\n\
1258         <!-- CONVERSATION_SUMMARY_END -->",
1259        summary_content
1260    ))
1261}
1262
1263pub fn active_messages_for_budget(session: &Session) -> Vec<Message> {
1264    session
1265        .messages
1266        .iter()
1267        .filter(|message| !message.compressed)
1268        .cloned()
1269        .collect()
1270}
1271
1272pub fn summary_source_messages(session: &Session) -> Vec<Message> {
1273    session
1274        .messages
1275        .iter()
1276        .filter(|message| !message.compressed)
1277        .filter(|message| !matches!(message.role, bamboo_domain::Role::System))
1278        .cloned()
1279        .collect()
1280}
1281
1282pub fn build_summary_prompt(
1283    session: &Session,
1284    messages: &[Message],
1285    existing_summary: Option<&str>,
1286) -> String {
1287    let mut content = String::new();
1288    content.push_str(
1289        "You are compressing conversation history for continued work. Produce a compact but reliable working-memory summary.\n\n",
1290    );
1291    content.push_str(
1292        "Critical requirements:\n- First capture the in-flight work right before compression (what was being done, where, and with which tool/file)\n- Distinguish clearly between ACTIVE work, COMPLETED work, and OBSOLETE or superseded work\n- Do not restate old tasks as active unless they are still unresolved\n- The current task list is the source of truth for what is actively being worked on\n- Preserve constraints, decisions, file paths, code changes, errors, tool findings, blockers, and the next step\n- If earlier plans conflict with the current task list or newer messages, treat the earlier plans as obsolete or completed\n- Explicitly evaluate each clear user requirement (e.g. requirement 1, requirement 2) with a status and evidence\n- Return only summary text in the same language as the conversation\n\n",
1293    );
1294
1295    if let Some(existing) = existing_summary.map(str::trim).filter(|s| !s.is_empty()) {
1296        content.push_str("## Existing Summary\n");
1297        content.push_str(existing);
1298        content.push_str("\n\n");
1299    }
1300
1301    let task_list_prompt = session.format_task_list_for_prompt();
1302    if !task_list_prompt.trim().is_empty() {
1303        content.push_str("## Current Task List\n");
1304        content.push_str(task_list_prompt.trim());
1305        content.push_str("\n\n");
1306    }
1307
1308    content.push_str(
1309        "## Required Output Sections\n1. Pre-compression in-flight work (what was being done immediately before compression)\n2. Current active objective\n3. Requirement checklist (Requirement | Status: completed/in_progress/pending/blocked/obsolete | Evidence)\n4. Active tasks\n5. Completed tasks\n6. Obsolete or superseded tasks\n7. Important context and constraints\n8. Files, code, and tool findings\n9. Open issues and next step\n\n",
1310    );
1311
1312    content.push_str("## Messages To Compress\n\n");
1313    for message in messages {
1314        let role = match message.role {
1315            bamboo_domain::Role::System => continue,
1316            bamboo_domain::Role::User => "User",
1317            bamboo_domain::Role::Assistant => match message.phase {
1318                Some(MessagePhase::Commentary) => "Assistant Commentary",
1319                Some(MessagePhase::FinalAnswer) => "Assistant Final",
1320                None => "Assistant",
1321            },
1322            bamboo_domain::Role::Tool => "Tool Result",
1323        };
1324
1325        content.push_str("### ");
1326        content.push_str(role);
1327        content.push('\n');
1328        if let Some(tool_calls) = &message.tool_calls {
1329            if !tool_calls.is_empty() {
1330                let names = tool_calls
1331                    .iter()
1332                    .map(|call| call.function.name.as_str())
1333                    .collect::<Vec<_>>()
1334                    .join(", ");
1335                content.push_str("Called tools: ");
1336                content.push_str(&names);
1337                content.push('\n');
1338            }
1339        }
1340        if let Some(tool_call_id) = &message.tool_call_id {
1341            content.push_str("Tool call id: ");
1342            content.push_str(tool_call_id);
1343            content.push('\n');
1344        }
1345        // Keep this compatibility renderer lossless. Production callers bound
1346        // the fully rendered request through the hierarchical summarizer;
1347        // clipping each message would discard tails without bounding the total.
1348        content.push_str(&message.content);
1349        content.push_str("\n\n");
1350    }
1351
1352    content.push_str(
1353        "Return only the summary text. Be explicit about what is active now versus what is already done or no longer relevant.",
1354    );
1355    content
1356}
1357
1358#[cfg(test)]
1359mod tests {
1360    use super::*;
1361    use bamboo_domain::TokenBudgetUsage;
1362    use bamboo_domain::{FunctionCall, TaskItem, TaskItemStatus, TaskList, ToolCall};
1363    use chrono::Utc;
1364
1365    fn make_budget() -> TokenBudget {
1366        TokenBudget {
1367            max_context_tokens: 1000,
1368            max_output_tokens: 100,
1369            strategy: BudgetStrategy::Hybrid {
1370                window_size: 20,
1371                enable_summarization: true,
1372            },
1373            safety_margin: 0,
1374            compression_trigger_percent: 50,
1375            compression_target_percent: 20,
1376            working_reserve_tokens: 0,
1377            fallback_trigger_percent: 75,
1378            prompt_cache_min_tool_output_chars: 1_200,
1379            prompt_cache_head_chars: 280,
1380            prompt_cache_tail_chars: 180,
1381            prompt_cache_recent_user_turns: 2,
1382            prompt_cache_recent_tool_chains: 2,
1383            max_tool_output_tokens: 0,
1384        }
1385    }
1386
1387    fn make_session_with_pressure() -> Session {
1388        let mut session = Session::new("compression-hysteresis", "gpt-4o-mini");
1389        session.token_budget = Some(make_budget());
1390        session.add_message(Message::system("system"));
1391        for i in 0..3 {
1392            session.add_message(Message::user(format!(
1393                "User message {i}: {}",
1394                "alpha beta gamma delta epsilon ".repeat(2)
1395            )));
1396            session.add_message(Message::assistant(
1397                format!(
1398                    "Assistant message {i}: {}",
1399                    "work log decisions next steps ".repeat(2)
1400                ),
1401                None,
1402            ));
1403        }
1404        session
1405    }
1406
1407    #[test]
1408    fn context_window_usage_percent_uses_context_window_denominator() {
1409        assert_eq!(context_window_usage_percent(0, 0), 0.0);
1410        assert_eq!(context_window_usage_percent(500, 1000), 50.0);
1411    }
1412
1413    #[test]
1414    fn estimate_context_compression_exposure_crosses_trigger_when_usage_is_high_enough() {
1415        let mut session = make_session_with_pressure();
1416        if let Some(budget) = session.token_budget.as_mut() {
1417            budget.compression_trigger_percent = 10;
1418        }
1419        let exposure = estimate_context_compression_exposure(
1420            &session,
1421            "gpt-4o-mini",
1422            session.token_budget.as_ref(),
1423        );
1424        assert!(exposure.active_usage_percent >= 10.0);
1425        assert!(exposure.should_expose_tool);
1426    }
1427
1428    #[test]
1429    fn estimate_context_compression_exposure_stays_below_trigger_when_usage_is_low() {
1430        let mut session = make_session_with_pressure();
1431        if let Some(budget) = session.token_budget.as_mut() {
1432            budget.compression_trigger_percent = 99;
1433        }
1434
1435        let exposure = estimate_context_compression_exposure(
1436            &session,
1437            "gpt-4o-mini",
1438            session.token_budget.as_ref(),
1439        );
1440
1441        assert!(exposure.active_usage_percent < 99.0);
1442        assert!(!exposure.should_expose_tool);
1443    }
1444
1445    #[test]
1446    fn build_summary_prompt_includes_task_list_and_state_sections() {
1447        let mut session = Session::new("summary-prompt", "gpt-4o-mini");
1448        session.set_task_list(TaskList {
1449            session_id: session.id.clone(),
1450            title: "Task List".to_string(),
1451            items: vec![
1452                TaskItem {
1453                    id: "task_1".to_string(),
1454                    description: "检查 51% 又回落到 50% 的触发逻辑".to_string(),
1455                    status: TaskItemStatus::InProgress,
1456                    depends_on: Vec::new(),
1457                    notes: "避免刚压缩完又立刻再次压缩".to_string(),
1458                    ..TaskItem::default()
1459                },
1460                TaskItem {
1461                    id: "task_2".to_string(),
1462                    description: "重写 summarizer prompt 并纳入 task list".to_string(),
1463                    status: TaskItemStatus::Pending,
1464                    depends_on: Vec::new(),
1465                    notes: String::new(),
1466                    ..TaskItem::default()
1467                },
1468            ],
1469            created_at: Utc::now(),
1470            updated_at: Utc::now(),
1471        });
1472        let prompt = build_summary_prompt(
1473            &session,
1474            &[
1475                Message::user("继续修复 context compression"),
1476                Message::assistant("先分析 trigger / target / summary", None),
1477            ],
1478            Some("old summary"),
1479        );
1480
1481        assert!(prompt.contains("## Current Task List"));
1482        assert!(prompt.contains("Current active objective"));
1483        assert!(prompt.contains("Requirement checklist"));
1484        assert!(prompt.contains("Active tasks"));
1485        assert!(prompt.contains("Completed tasks"));
1486        assert!(prompt.contains("Obsolete or superseded tasks"));
1487        assert!(prompt.contains("检查 51% 又回落到 50% 的触发逻辑"));
1488        assert!(prompt.contains("old summary"));
1489    }
1490
1491    #[test]
1492    fn forced_plan_keeps_last_three_user_messages_active() {
1493        let budget = TokenBudget {
1494            max_context_tokens: 1200,
1495            max_output_tokens: 100,
1496            strategy: BudgetStrategy::Hybrid {
1497                window_size: 20,
1498                enable_summarization: true,
1499            },
1500            safety_margin: 0,
1501            compression_trigger_percent: 80,
1502            compression_target_percent: 20,
1503            working_reserve_tokens: 0,
1504            fallback_trigger_percent: 75,
1505            prompt_cache_min_tool_output_chars: 1_200,
1506            prompt_cache_head_chars: 280,
1507            prompt_cache_tail_chars: 180,
1508            prompt_cache_recent_user_turns: 2,
1509            prompt_cache_recent_tool_chains: 2,
1510            max_tool_output_tokens: 0,
1511        };
1512        let mut session = Session::new("keep-last-three-user-turns", "gpt-4o-mini");
1513        session.token_budget = Some(budget.clone());
1514        session.add_message(Message::system("system"));
1515        for i in 0..6 {
1516            session.add_message(Message::user(format!(
1517                "U{i}: {}",
1518                "alpha beta gamma ".repeat(8)
1519            )));
1520            session.add_message(Message::assistant(
1521                format!("A{i}: {}", "analysis plan steps ".repeat(8)),
1522                None,
1523            ));
1524        }
1525
1526        let plan = build_forced_compression_plan_with_summary(
1527            &session,
1528            "gpt-4o-mini",
1529            Some(&budget),
1530            "summary".to_string(),
1531            CompressionTriggerType::CriticalOverflow,
1532        )
1533        .expect("forced plan should build");
1534
1535        let compressed_ids = plan
1536            .compressed_message_ids
1537            .iter()
1538            .map(String::as_str)
1539            .collect::<HashSet<_>>();
1540        let kept_user_contents = session
1541            .messages
1542            .iter()
1543            .filter(|message| !matches!(message.role, bamboo_domain::Role::System))
1544            .filter(|message| !compressed_ids.contains(message.id.as_str()))
1545            .filter(|message| matches!(message.role, bamboo_domain::Role::User))
1546            .map(|message| message.content.clone())
1547            .collect::<Vec<_>>();
1548
1549        assert!(
1550            kept_user_contents.len() >= 3,
1551            "expected to keep at least 3 user messages, got {}",
1552            kept_user_contents.len()
1553        );
1554        assert!(kept_user_contents
1555            .iter()
1556            .any(|content| content.starts_with("U3:")));
1557        assert!(kept_user_contents
1558            .iter()
1559            .any(|content| content.starts_with("U4:")));
1560        assert!(kept_user_contents
1561            .iter()
1562            .any(|content| content.starts_with("U5:")));
1563    }
1564
1565    #[test]
1566    fn estimate_exposure_prefers_persisted_budget_usage_when_higher() {
1567        let mut session = Session::new("persisted-usage", "gpt-4o-mini");
1568        session.token_budget = Some(TokenBudget {
1569            max_context_tokens: 100_000,
1570            max_output_tokens: 1_000,
1571            strategy: BudgetStrategy::Hybrid {
1572                window_size: 20,
1573                enable_summarization: true,
1574            },
1575            safety_margin: 0,
1576            compression_trigger_percent: 80,
1577            compression_target_percent: 50,
1578            working_reserve_tokens: 0,
1579            fallback_trigger_percent: 75,
1580            prompt_cache_min_tool_output_chars: 1_200,
1581            prompt_cache_head_chars: 280,
1582            prompt_cache_tail_chars: 180,
1583            prompt_cache_recent_user_turns: 2,
1584            prompt_cache_recent_tool_chains: 2,
1585            max_tool_output_tokens: 0,
1586        });
1587        session.add_message(Message::system("system"));
1588        session.add_message(Message::user("short"));
1589        session.add_message(Message::assistant("short", None));
1590        session.add_message(Message::user("follow-up"));
1591        session.add_message(Message::assistant("reply", None));
1592        session.token_usage = Some(TokenBudgetUsage {
1593            system_tokens: 100,
1594            summary_tokens: 0,
1595            window_tokens: 95_900,
1596            total_tokens: 96_000,
1597            max_context_tokens: 100_000,
1598            budget_limit: 10_000,
1599            truncation_occurred: true,
1600            segments_removed: 12,
1601            prompt_cached_tool_outputs: 0,
1602            prompt_cached_tool_tokens_saved: 0,
1603            thinking_tokens: 0,
1604            cache_read_input_tokens: 0,
1605        });
1606
1607        let exposure = estimate_context_compression_exposure(
1608            &session,
1609            "gpt-4o-mini",
1610            session.token_budget.as_ref(),
1611        );
1612
1613        assert!(
1614            exposure.active_usage_percent >= 96.0,
1615            "expected persisted context-window usage to drive exposure, got {}",
1616            exposure.active_usage_percent
1617        );
1618        assert!(exposure.should_expose_tool);
1619    }
1620
1621    #[test]
1622    fn never_compress_messages_are_excluded_from_summarize_set() {
1623        let budget = TokenBudget {
1624            max_context_tokens: 1200,
1625            max_output_tokens: 100,
1626            strategy: BudgetStrategy::Hybrid {
1627                window_size: 20,
1628                enable_summarization: true,
1629            },
1630            safety_margin: 0,
1631            compression_trigger_percent: 80,
1632            compression_target_percent: 20,
1633            working_reserve_tokens: 0,
1634            fallback_trigger_percent: 75,
1635            prompt_cache_min_tool_output_chars: 1_200,
1636            prompt_cache_head_chars: 280,
1637            prompt_cache_tail_chars: 180,
1638            prompt_cache_recent_user_turns: 2,
1639            prompt_cache_recent_tool_chains: 2,
1640            max_tool_output_tokens: 0,
1641        };
1642        let mut session = Session::new("never-compress-test", "gpt-4o-mini");
1643        session.token_budget = Some(budget.clone());
1644        session.add_message(Message::system("system"));
1645
1646        // Old user message that should be summarized
1647        session.add_message(Message::user("Old question about X"));
1648        session.add_message(Message::assistant("Old answer about X", None));
1649
1650        // Protected user message (never_compress = true)
1651        let mut protected = Message::user("Critical context that must survive");
1652        protected.never_compress = true;
1653        session.add_message(protected);
1654        session.add_message(Message::assistant("Response to critical", None));
1655
1656        // Recent user messages that anchor the keep window
1657        for i in 0..4 {
1658            session.add_message(Message::user(format!(
1659                "Recent U{i}: {}",
1660                "padding text to fill budget ".repeat(6)
1661            )));
1662            session.add_message(Message::assistant(
1663                format!("Recent A{i}: {}", "reply padding text ".repeat(6)),
1664                None,
1665            ));
1666        }
1667
1668        let plan = build_forced_compression_plan_with_summary(
1669            &session,
1670            "gpt-4o-mini",
1671            Some(&budget),
1672            "summary".to_string(),
1673            CompressionTriggerType::Auto,
1674        )
1675        .expect("plan should build");
1676
1677        let compressed_ids: HashSet<&str> = plan
1678            .compressed_message_ids
1679            .iter()
1680            .map(String::as_str)
1681            .collect();
1682
1683        // Find the never_compress message
1684        let protected_msg = session
1685            .messages
1686            .iter()
1687            .find(|m| m.never_compress)
1688            .expect("should find the protected message");
1689
1690        assert!(
1691            !compressed_ids.contains(protected_msg.id.as_str()),
1692            "never_compress message should NOT be in the compressed set"
1693        );
1694    }
1695
1696    #[test]
1697    fn skill_tool_chain_messages_are_protected_from_compression() {
1698        let budget = TokenBudget {
1699            max_context_tokens: 1200,
1700            max_output_tokens: 100,
1701            strategy: BudgetStrategy::Hybrid {
1702                window_size: 20,
1703                enable_summarization: true,
1704            },
1705            safety_margin: 0,
1706            compression_trigger_percent: 80,
1707            compression_target_percent: 20,
1708            working_reserve_tokens: 0,
1709            fallback_trigger_percent: 75,
1710            prompt_cache_min_tool_output_chars: 1_200,
1711            prompt_cache_head_chars: 280,
1712            prompt_cache_tail_chars: 180,
1713            prompt_cache_recent_user_turns: 2,
1714            prompt_cache_recent_tool_chains: 2,
1715            max_tool_output_tokens: 0,
1716        };
1717        let mut session = Session::new("skill-chain-test", "gpt-4o-mini");
1718        session.token_budget = Some(budget.clone());
1719        session.add_message(Message::system("system"));
1720
1721        // Skill tool chain (load_skill + read_skill_resource)
1722        let mut skill_call = Message::assistant(String::new(), None);
1723        skill_call.tool_calls = Some(vec![ToolCall {
1724            id: "tc-skill".to_string(),
1725            tool_type: "function".to_string(),
1726            function: FunctionCall {
1727                name: "load_skill".to_string(),
1728                arguments: r#"{"skill_id":"my-skill"}"#.to_string(),
1729            },
1730        }]);
1731        session.add_message(skill_call);
1732
1733        let mut skill_result = Message::tool_result("tc-skill", "skill loaded");
1734        skill_result.tool_success = Some(true);
1735        session.add_message(skill_result);
1736
1737        // Regular messages to fill budget
1738        for i in 0..6 {
1739            session.add_message(Message::user(format!(
1740                "U{i}: {}",
1741                "alpha beta gamma delta ".repeat(8)
1742            )));
1743            session.add_message(Message::assistant(
1744                format!("A{i}: {}", "analysis steps plan ".repeat(8)),
1745                None,
1746            ));
1747        }
1748
1749        let plan = build_forced_compression_plan_with_summary(
1750            &session,
1751            "gpt-4o-mini",
1752            Some(&budget),
1753            "summary".to_string(),
1754            CompressionTriggerType::Auto,
1755        )
1756        .expect("plan should build");
1757
1758        let compressed_ids: HashSet<&str> = plan
1759            .compressed_message_ids
1760            .iter()
1761            .map(String::as_str)
1762            .collect();
1763
1764        // Skill tool chain messages should not be compressed
1765        let skill_messages: Vec<&Message> = session
1766            .messages
1767            .iter()
1768            .filter(|m| {
1769                m.tool_calls
1770                    .as_ref()
1771                    .is_some_and(|calls| calls.iter().any(|c| c.function.name == "load_skill"))
1772                    || m.tool_call_id.as_deref() == Some("tc-skill")
1773            })
1774            .collect();
1775
1776        for msg in &skill_messages {
1777            assert!(
1778                !compressed_ids.contains(msg.id.as_str()),
1779                "skill tool chain message {} should NOT be compressed",
1780                msg.id
1781            );
1782        }
1783    }
1784
1785    #[test]
1786    fn generic_tool_chain_is_never_split_by_forced_compression() {
1787        // A generic (non-skill) tool_use and its tool_result must never be split
1788        // across the compression boundary — a split orphans one of them in the
1789        // active set and the provider 400s, poisoning the session. #340.
1790        let budget = TokenBudget {
1791            max_context_tokens: 1200,
1792            max_output_tokens: 100,
1793            strategy: BudgetStrategy::Hybrid {
1794                window_size: 20,
1795                enable_summarization: true,
1796            },
1797            safety_margin: 0,
1798            compression_trigger_percent: 80,
1799            compression_target_percent: 20,
1800            working_reserve_tokens: 0,
1801            fallback_trigger_percent: 75,
1802            prompt_cache_min_tool_output_chars: 1_200,
1803            prompt_cache_head_chars: 280,
1804            prompt_cache_tail_chars: 180,
1805            prompt_cache_recent_user_turns: 2,
1806            prompt_cache_recent_tool_chains: 2,
1807            max_tool_output_tokens: 0,
1808        };
1809        let mut session = Session::new("generic-chain-test", "gpt-4o-mini");
1810        session.token_budget = Some(budget.clone());
1811        session.add_message(Message::system("system"));
1812
1813        // A generic tool chain (search + its result) placed early so it is an
1814        // eviction candidate; the large result makes it a prime compression target.
1815        let mut call = Message::assistant(String::new(), None);
1816        call.tool_calls = Some(vec![ToolCall {
1817            id: "tc-gen".to_string(),
1818            tool_type: "function".to_string(),
1819            function: FunctionCall {
1820                name: "search".to_string(),
1821                arguments: r#"{"q":"rust"}"#.to_string(),
1822            },
1823        }]);
1824        session.add_message(call);
1825        let mut result = Message::tool_result("tc-gen", &"search result payload ".repeat(20));
1826        result.tool_success = Some(true);
1827        session.add_message(result);
1828
1829        // Filler to push usage over budget and force eviction.
1830        for i in 0..8 {
1831            session.add_message(Message::user(format!(
1832                "U{i}: {}",
1833                "alpha beta gamma delta ".repeat(8)
1834            )));
1835            session.add_message(Message::assistant(
1836                format!("A{i}: {}", "analysis steps plan ".repeat(8)),
1837                None,
1838            ));
1839        }
1840
1841        let plan = build_forced_compression_plan_with_summary(
1842            &session,
1843            "gpt-4o-mini",
1844            Some(&budget),
1845            "summary".to_string(),
1846            CompressionTriggerType::Auto,
1847        )
1848        .expect("plan should build");
1849
1850        let compressed: HashSet<&str> = plan
1851            .compressed_message_ids
1852            .iter()
1853            .map(String::as_str)
1854            .collect();
1855
1856        let assistant_id = session
1857            .messages
1858            .iter()
1859            .find(|m| {
1860                m.tool_calls
1861                    .as_ref()
1862                    .is_some_and(|c| c.iter().any(|tc| tc.id == "tc-gen"))
1863            })
1864            .map(|m| m.id.as_str())
1865            .expect("assistant tool_use message present");
1866        let result_id = session
1867            .messages
1868            .iter()
1869            .find(|m| m.tool_call_id.as_deref() == Some("tc-gen"))
1870            .map(|m| m.id.as_str())
1871            .expect("tool_result message present");
1872
1873        // The tool_use and its tool_result must land on the SAME side.
1874        assert_eq!(
1875            compressed.contains(assistant_id),
1876            compressed.contains(result_id),
1877            "generic tool_use ({assistant_id}) and its tool_result ({result_id}) must not be split"
1878        );
1879    }
1880
1881    fn tool_use_message(call_id: &str) -> Message {
1882        let mut assistant = Message::assistant(String::new(), None);
1883        assistant.tool_calls = Some(vec![ToolCall {
1884            id: call_id.to_string(),
1885            tool_type: "function".to_string(),
1886            function: FunctionCall {
1887                name: "search".to_string(),
1888                arguments: "{}".to_string(),
1889            },
1890        }]);
1891        assistant
1892    }
1893
1894    #[test]
1895    fn close_compressed_set_over_tool_chains_reunites_a_split_chain() {
1896        // Pre-split state: the assistant tool_use is compressed while its
1897        // tool_result was left in the keep (active) set — an orphan. #340.
1898        let assistant = tool_use_message("tc-1");
1899        let result = Message::tool_result("tc-1", "result payload");
1900        let assistant_id = assistant.id.clone();
1901        let result_id = result.id.clone();
1902
1903        let mut messages_to_keep = vec![result];
1904        let mut messages_to_summarize = vec![assistant];
1905
1906        close_compressed_set_over_tool_chains(
1907            &mut messages_to_keep,
1908            &mut messages_to_summarize,
1909            &HashSet::new(),
1910            &[],
1911        );
1912
1913        // The orphan tool_result must be pulled into the compressed set.
1914        assert!(
1915            messages_to_keep.is_empty(),
1916            "orphaned tool_result must be moved into the compressed set"
1917        );
1918        let summarized: HashSet<&str> = messages_to_summarize
1919            .iter()
1920            .map(|m| m.id.as_str())
1921            .collect();
1922        assert!(summarized.contains(assistant_id.as_str()));
1923        assert!(summarized.contains(result_id.as_str()));
1924    }
1925
1926    #[test]
1927    fn close_compressed_set_over_tool_chains_respects_protected_messages() {
1928        // A protected (never-compress) chain member must NOT be force-compressed.
1929        let assistant = tool_use_message("tc-2");
1930        let result = Message::tool_result("tc-2", "result");
1931        let result_id = result.id.clone();
1932
1933        let mut messages_to_keep = vec![result];
1934        let mut messages_to_summarize = vec![assistant];
1935        let never_compress_ids = vec![result_id.clone()];
1936
1937        close_compressed_set_over_tool_chains(
1938            &mut messages_to_keep,
1939            &mut messages_to_summarize,
1940            &HashSet::new(),
1941            &never_compress_ids,
1942        );
1943
1944        assert_eq!(
1945            messages_to_keep.len(),
1946            1,
1947            "protected result must stay in keep"
1948        );
1949        assert_eq!(messages_to_keep[0].id, result_id);
1950    }
1951
1952    #[test]
1953    fn recovery_message_returns_none_for_empty_messages() {
1954        let session = Session::new("recovery-empty", "model");
1955        let result = build_post_compaction_recovery_message(&[], &session);
1956        assert!(result.is_none());
1957    }
1958
1959    #[test]
1960    fn recovery_message_has_never_compress_flag() {
1961        let mut session = Session::new("recovery-flag", "model");
1962        let messages = vec![Message::assistant("no decisions here", None)];
1963        session.set_task_list(TaskList {
1964            session_id: session.id.clone(),
1965            title: "Tasks".to_string(),
1966            items: vec![TaskItem {
1967                id: "t1".to_string(),
1968                description: "Active task".to_string(),
1969                status: TaskItemStatus::InProgress,
1970                ..TaskItem::default()
1971            }],
1972            created_at: Utc::now(),
1973            updated_at: Utc::now(),
1974        });
1975        let recovery = build_post_compaction_recovery_message(&messages, &session)
1976            .expect("should return recovery message");
1977        assert!(recovery.never_compress);
1978        assert!(recovery.content.contains("[post-compaction-recovery]"));
1979    }
1980
1981    #[test]
1982    fn recovery_message_extracts_file_paths_from_tool_calls() {
1983        let session = Session::new("recovery-files", "model");
1984        let mut write_call = Message::assistant("writing file", None);
1985        write_call.tool_calls = Some(vec![ToolCall {
1986            id: "tc1".to_string(),
1987            tool_type: "function".to_string(),
1988            function: FunctionCall {
1989                name: "Write".to_string(),
1990                arguments: r#"{"file_path":"/src/main.rs","content":"fn main() {}"}"#.to_string(),
1991            },
1992        }]);
1993        let mut edit_call = Message::assistant("editing file", None);
1994        edit_call.tool_calls = Some(vec![ToolCall {
1995            id: "tc2".to_string(),
1996            tool_type: "function".to_string(),
1997            function: FunctionCall {
1998                name: "Edit".to_string(),
1999                arguments: r#"{"file_path":"/lib/utils.rs","old":"x","new":"y"}"#.to_string(),
2000            },
2001        }]);
2002        let messages = vec![write_call, edit_call];
2003
2004        let recovery = build_post_compaction_recovery_message(&messages, &session)
2005            .expect("should return recovery");
2006        assert!(recovery.content.contains("/src/main.rs"));
2007        assert!(recovery.content.contains("/lib/utils.rs"));
2008        assert!(recovery.content.contains("Recently Modified Files"));
2009    }
2010
2011    #[test]
2012    fn recovery_message_includes_active_tasks() {
2013        let mut session = Session::new("recovery-tasks", "model");
2014        session.set_task_list(TaskList {
2015            session_id: session.id.clone(),
2016            title: "Tasks".to_string(),
2017            items: vec![
2018                TaskItem {
2019                    id: "t1".to_string(),
2020                    description: "Fix auth middleware".to_string(),
2021                    status: TaskItemStatus::InProgress,
2022                    ..TaskItem::default()
2023                },
2024                TaskItem {
2025                    id: "t2".to_string(),
2026                    description: "Add tests".to_string(),
2027                    status: TaskItemStatus::Pending,
2028                    ..TaskItem::default()
2029                },
2030                TaskItem {
2031                    id: "t3".to_string(),
2032                    description: "Done task".to_string(),
2033                    status: TaskItemStatus::Completed,
2034                    ..TaskItem::default()
2035                },
2036            ],
2037            created_at: Utc::now(),
2038            updated_at: Utc::now(),
2039        });
2040        let messages = vec![Message::assistant("some work", None)];
2041
2042        let recovery = build_post_compaction_recovery_message(&messages, &session)
2043            .expect("should return recovery");
2044        assert!(recovery.content.contains("Active Tasks"));
2045        assert!(recovery.content.contains("Fix auth middleware"));
2046        assert!(recovery.content.contains("Add tests"));
2047        // Completed tasks should NOT appear in active tasks
2048        assert!(!recovery.content.contains("Done task"));
2049    }
2050
2051    #[test]
2052    fn apply_compression_plan_injects_recovery_message() {
2053        let budget = TokenBudget {
2054            max_context_tokens: 1200,
2055            max_output_tokens: 100,
2056            strategy: BudgetStrategy::Hybrid {
2057                window_size: 20,
2058                enable_summarization: true,
2059            },
2060            safety_margin: 0,
2061            compression_trigger_percent: 80,
2062            compression_target_percent: 20,
2063            working_reserve_tokens: 0,
2064            fallback_trigger_percent: 75,
2065            prompt_cache_min_tool_output_chars: 1_200,
2066            prompt_cache_head_chars: 280,
2067            prompt_cache_tail_chars: 180,
2068            prompt_cache_recent_user_turns: 2,
2069            prompt_cache_recent_tool_chains: 2,
2070            max_tool_output_tokens: 0,
2071        };
2072        let mut session = Session::new("recovery-inject", "gpt-4o-mini");
2073        session.token_budget = Some(budget.clone());
2074        session.add_message(Message::system("system"));
2075
2076        // Old messages with tool calls containing file paths
2077        let mut write_msg = Message::assistant("writing", None);
2078        write_msg.tool_calls = Some(vec![ToolCall {
2079            id: "tc-w".to_string(),
2080            tool_type: "function".to_string(),
2081            function: FunctionCall {
2082                name: "Write".to_string(),
2083                arguments: r#"{"file_path":"/src/lib.rs","content":"pub fn hello() {}"}"#
2084                    .to_string(),
2085            },
2086        }]);
2087        session.add_message(Message::user("Write the file"));
2088        session.add_message(write_msg);
2089
2090        // Fill with enough messages to force compression
2091        for i in 0..6 {
2092            session.add_message(Message::user(format!(
2093                "U{i}: {}",
2094                "alpha beta gamma delta ".repeat(8)
2095            )));
2096            session.add_message(Message::assistant(
2097                format!("A{i}: {}", "analysis plan ".repeat(8)),
2098                None,
2099            ));
2100        }
2101
2102        let plan = build_forced_compression_plan_with_summary(
2103            &session,
2104            "gpt-4o-mini",
2105            Some(&budget),
2106            "summary text".to_string(),
2107            CompressionTriggerType::Auto,
2108        )
2109        .expect("plan should build");
2110
2111        assert!(!plan.compressed_message_ids.is_empty());
2112
2113        let compressed_count = apply_compression_plan(&mut session, plan);
2114        assert!(compressed_count > 0);
2115
2116        // Verify recovery message was injected
2117        let has_recovery = session.messages.iter().any(|m| {
2118            m.never_compress
2119                && m.content.contains("[post-compaction-recovery]")
2120                && m.content.contains("/src/lib.rs")
2121        });
2122        assert!(
2123            has_recovery,
2124            "session should contain a post-compaction recovery message with the file path"
2125        );
2126    }
2127
2128    #[test]
2129    fn summary_quality_full_coverage_when_all_files_mentioned() {
2130        let messages = vec![{
2131            let mut m = Message::assistant("writing", None);
2132            m.tool_calls = Some(vec![ToolCall {
2133                id: "tc1".to_string(),
2134                tool_type: "function".to_string(),
2135                function: FunctionCall {
2136                    name: "Write".to_string(),
2137                    arguments: r#"{"file_path":"/src/main.rs","content":"fn main() {}"}"#
2138                        .to_string(),
2139                },
2140            }]);
2141            m
2142        }];
2143        let summary = "Modified /src/main.rs to add main function";
2144        let quality = validate_summary_quality(summary, &messages);
2145        assert!(
2146            quality.file_coverage >= 0.99,
2147            "file_coverage should be ~1.0, got {:.2}",
2148            quality.file_coverage
2149        );
2150    }
2151
2152    #[test]
2153    fn summary_quality_zero_coverage_when_no_files_mentioned() {
2154        let messages = vec![{
2155            let mut m = Message::assistant("writing", None);
2156            m.tool_calls = Some(vec![ToolCall {
2157                id: "tc1".to_string(),
2158                tool_type: "function".to_string(),
2159                function: FunctionCall {
2160                    name: "Write".to_string(),
2161                    arguments: r#"{"file_path":"/src/main.rs","content":"fn main() {}"}"#
2162                        .to_string(),
2163                },
2164            }]);
2165            m
2166        }];
2167        let summary = "Summary that mentions nothing about files";
2168        let quality = validate_summary_quality(summary, &messages);
2169        assert!(
2170            quality.file_coverage < 0.01,
2171            "file_coverage should be ~0.0, got {:.2}",
2172            quality.file_coverage
2173        );
2174    }
2175
2176    #[test]
2177    fn summary_quality_handles_empty_messages() {
2178        let quality = validate_summary_quality("some summary", &[]);
2179        assert_eq!(quality.file_coverage, 1.0);
2180        assert_eq!(quality.decision_coverage, 1.0);
2181    }
2182
2183    fn candidate_budget(max_context_tokens: u32, target_percent: u8) -> TokenBudget {
2184        TokenBudget {
2185            max_context_tokens,
2186            max_output_tokens: max_context_tokens / 4,
2187            strategy: BudgetStrategy::Hybrid {
2188                window_size: 20,
2189                enable_summarization: true,
2190            },
2191            safety_margin: 0,
2192            compression_trigger_percent: 80,
2193            compression_target_percent: target_percent,
2194            working_reserve_tokens: 0,
2195            fallback_trigger_percent: 75,
2196            prompt_cache_min_tool_output_chars: 1_200,
2197            prompt_cache_head_chars: 280,
2198            prompt_cache_tail_chars: 180,
2199            prompt_cache_recent_user_turns: 2,
2200            prompt_cache_recent_tool_chains: 2,
2201            max_tool_output_tokens: 0,
2202        }
2203    }
2204
2205    #[test]
2206    fn candidate_plan_is_selected_before_summary_and_reserves_twenty_percent_of_source() {
2207        let budget = candidate_budget(6_000, 40);
2208        let mut session = Session::new("candidate-first", "main-model");
2209        session.add_message(Message::system("system"));
2210
2211        let mut tool_call = Message::assistant("searching", None);
2212        tool_call.tool_calls = Some(vec![ToolCall {
2213            id: "candidate-chain".to_string(),
2214            tool_type: "function".to_string(),
2215            function: FunctionCall {
2216                name: "search".to_string(),
2217                arguments: r#"{"q":"compression"}"#.to_string(),
2218            },
2219        }]);
2220        let tool_call_id = tool_call.id.clone();
2221        session.add_message(tool_call);
2222        let tool_result = Message::tool_result(
2223            "candidate-chain",
2224            "result payload with decisions and paths ".repeat(80),
2225        );
2226        let tool_result_id = tool_result.id.clone();
2227        session.add_message(tool_result);
2228
2229        let mut never = Message::assistant("protected runtime state ".repeat(80), None);
2230        never.never_compress = true;
2231        let never_id = never.id.clone();
2232        session.add_message(never);
2233
2234        for index in 0..12 {
2235            session.add_message(Message::user(format!(
2236                "U{index}: {}",
2237                "requirement context decision ".repeat(40)
2238            )));
2239            session.add_message(Message::assistant(
2240                format!("A{index}: {}", "implementation evidence result ".repeat(40)),
2241                None,
2242            ));
2243        }
2244        let newest_user_ids = session
2245            .messages
2246            .iter()
2247            .filter(|message| matches!(message.role, bamboo_domain::Role::User))
2248            .rev()
2249            .take(3)
2250            .map(|message| message.id.clone())
2251            .collect::<HashSet<_>>();
2252
2253        let candidate = build_forced_compression_candidate_plan(
2254            &session,
2255            "main-model",
2256            Some(&budget),
2257            0.20,
2258            CompressionTriggerType::Auto,
2259        )
2260        .expect("candidate plan should reach target");
2261        let selected = candidate
2262            .compressed_message_ids
2263            .iter()
2264            .cloned()
2265            .collect::<HashSet<_>>();
2266
2267        assert_eq!(
2268            selected,
2269            candidate
2270                .messages_to_summarize
2271                .iter()
2272                .map(|message| message.id.clone())
2273                .collect()
2274        );
2275        assert_eq!(
2276            candidate.target_summary_tokens,
2277            ((candidate.source_tokens as f64) * 0.20).ceil() as u32
2278        );
2279        assert!(!selected.contains(&never_id));
2280        assert!(newest_user_ids.is_disjoint(&selected));
2281        assert_eq!(
2282            selected.contains(&tool_call_id),
2283            selected.contains(&tool_result_id),
2284            "generic tool chain must be selected atomically"
2285        );
2286        assert!(
2287            candidate.projected_usage_after_percent <= budget.compression_target_percent as f64
2288        );
2289    }
2290
2291    #[test]
2292    fn cumulative_twenty_percent_uses_represented_raw_tokens_not_previous_summary_length() {
2293        let budget = candidate_budget(30_000, 50);
2294        let mut session = Session::new("cumulative-ratio", "main-model");
2295        session.add_message(Message::system("system"));
2296        session.conversation_summary = Some(
2297            ConversationSummary::new("existing detailed summary ".repeat(200), 40, 2_000)
2298                .with_compression_metrics(10_000, 2_000, 0.20, false, None),
2299        );
2300        for index in 0..80 {
2301            session.add_message(Message::user(format!(
2302                "U{index}: {}",
2303                "raw source requirement and evidence ".repeat(30)
2304            )));
2305            session.add_message(Message::assistant(
2306                format!(
2307                    "A{index}: {}",
2308                    "implementation result and next step ".repeat(30)
2309                ),
2310                None,
2311            ));
2312        }
2313
2314        let candidate = build_forced_compression_candidate_plan(
2315            &session,
2316            "main-model",
2317            Some(&budget),
2318            0.20,
2319            CompressionTriggerType::Auto,
2320        )
2321        .expect("cumulative candidate plan");
2322        assert_eq!(
2323            candidate.target_summary_tokens,
2324            ((10_000u32.saturating_add(candidate.source_tokens) as f64) * 0.20).ceil() as u32
2325        );
2326        assert_eq!(
2327            candidate.represented_source_tokens,
2328            10_000u32.saturating_add(candidate.source_tokens)
2329        );
2330    }
2331
2332    #[test]
2333    fn legacy_summary_uses_conservative_growth_and_migrates_represented_source_metadata() {
2334        let budget = candidate_budget(20_000, 50);
2335        let counter = TiktokenTokenCounter::default();
2336        let existing_content = "legacy summary fact and decision ".repeat(180);
2337        let existing_tokens = counter.count_text(&existing_content);
2338        let mut session = Session::new("legacy-summary-ratio", "main-model");
2339        session.add_message(Message::system("system"));
2340        session.conversation_summary = Some(ConversationSummary::new(
2341            existing_content,
2342            40,
2343            existing_tokens,
2344        ));
2345        for index in 0..40 {
2346            session.add_message(Message::user(format!(
2347                "U{index}: {}",
2348                "raw requirement evidence ".repeat(40)
2349            )));
2350            session.add_message(Message::assistant(
2351                format!("A{index}: {}", "result next step ".repeat(40)),
2352                None,
2353            ));
2354        }
2355
2356        let candidate = build_forced_compression_candidate_plan(
2357            &session,
2358            "main-model",
2359            Some(&budget),
2360            0.20,
2361            CompressionTriggerType::Auto,
2362        )
2363        .expect("legacy candidate plan");
2364        assert_eq!(
2365            candidate.target_summary_tokens,
2366            existing_tokens.saturating_add(((candidate.source_tokens as f64) * 0.20).ceil() as u32)
2367        );
2368        assert_eq!(
2369            candidate.previous_represented_source_tokens,
2370            ((existing_tokens as f64) / 0.20).ceil() as u32
2371        );
2372        let expected_represented = candidate.represented_source_tokens;
2373        let expected_target = candidate.target_summary_tokens;
2374        let plan = finalize_compression_candidate_plan(
2375            &session,
2376            candidate,
2377            "migrated legacy summary with new evidence".to_string(),
2378        )
2379        .expect("short real summary should satisfy target");
2380        assert!(apply_compression_plan(&mut session, plan) > 0);
2381        let migrated = session
2382            .conversation_summary
2383            .as_ref()
2384            .expect("migrated summary");
2385        assert_eq!(migrated.represented_source_tokens, expected_represented);
2386        assert_eq!(migrated.target_token_count, expected_target);
2387        assert_eq!(migrated.target_ratio, 0.20);
2388    }
2389
2390    #[test]
2391    fn final_postcondition_counts_the_recovery_message_inserted_during_apply() {
2392        let budget = candidate_budget(8_000, 40);
2393        let mut session = Session::new("recovery-postcondition", "main-model");
2394        session.add_message(Message::system("system"));
2395        for index in 0..20 {
2396            session.add_message(Message::user(format!(
2397                "U{index}: {}",
2398                "requirement source content ".repeat(45)
2399            )));
2400            let mut assistant = Message::assistant(
2401                format!("A{index}: {}", "implementation evidence ".repeat(45)),
2402                None,
2403            );
2404            if index == 0 {
2405                assistant.tool_calls = Some(vec![ToolCall {
2406                    id: "write-recovery-763".to_string(),
2407                    tool_type: "function".to_string(),
2408                    function: FunctionCall {
2409                        name: "Write".to_string(),
2410                        arguments:
2411                            r#"{"file_path":"/workspace/src/recovery_763.rs","content":"fixed"}"#
2412                                .to_string(),
2413                    },
2414                }]);
2415            }
2416            session.add_message(assistant);
2417            if index == 0 {
2418                session.add_message(Message::tool_result(
2419                    "write-recovery-763",
2420                    "write completed",
2421                ));
2422            }
2423        }
2424
2425        let candidate = build_forced_compression_candidate_plan(
2426            &session,
2427            "main-model",
2428            Some(&budget),
2429            0.20,
2430            CompressionTriggerType::CriticalOverflow,
2431        )
2432        .expect("candidate should reserve recovery-message tokens");
2433        let plan = finalize_compression_candidate_plan(
2434            &session,
2435            candidate,
2436            "summary with /workspace/src/recovery_763.rs".to_string(),
2437        )
2438        .expect("final plan");
2439        assert!(apply_compression_plan(&mut session, plan) > 0);
2440        assert!(session
2441            .messages
2442            .iter()
2443            .any(|message| message.content.contains("[post-compaction-recovery]")));
2444        assert!(
2445            session.token_usage.as_ref().is_some_and(
2446                |usage| usage.total_tokens <= budget.compression_target_context_tokens()
2447            ),
2448            "the real active context, including recovery, must remain at or below target"
2449        );
2450    }
2451
2452    #[test]
2453    fn finalization_revalidates_actual_summary_without_mutating_session() {
2454        let budget = candidate_budget(6_000, 40);
2455        let mut session = Session::new("atomic-finalize", "main-model");
2456        session.add_message(Message::system("system"));
2457        for index in 0..20 {
2458            session.add_message(Message::user(format!(
2459                "U{index}: {}",
2460                "source content ".repeat(50)
2461            )));
2462            session.add_message(Message::assistant(
2463                format!("A{index}: {}", "response content ".repeat(50)),
2464                None,
2465            ));
2466        }
2467        let candidate = build_forced_compression_candidate_plan(
2468            &session,
2469            "main-model",
2470            Some(&budget),
2471            0.20,
2472            CompressionTriggerType::CriticalOverflow,
2473        )
2474        .expect("candidate plan");
2475        let before_flags = session
2476            .messages
2477            .iter()
2478            .map(|message| (message.id.clone(), message.compressed))
2479            .collect::<Vec<_>>();
2480
2481        let result = finalize_compression_candidate_plan(
2482            &session,
2483            candidate,
2484            "oversized summary ".repeat(20_000),
2485        );
2486        assert!(matches!(
2487            result,
2488            Err(CompressionPlanError::SummaryExceedsTarget { .. })
2489        ));
2490        assert_eq!(
2491            before_flags,
2492            session
2493                .messages
2494                .iter()
2495                .map(|message| (message.id.clone(), message.compressed))
2496                .collect::<Vec<_>>()
2497        );
2498        assert!(session.conversation_summary.is_none());
2499        assert!(session.compression_events.is_empty());
2500    }
2501
2502    #[test]
2503    fn candidate_plan_reports_when_protected_content_makes_target_impossible() {
2504        let budget = candidate_budget(2_000, 20);
2505        let mut session = Session::new("protected-capacity", "main-model");
2506        session.add_message(Message::system("system"));
2507        session.add_message(Message::assistant(
2508            "one eligible old message ".repeat(200),
2509            None,
2510        ));
2511        for index in 0..3 {
2512            session.add_message(Message::user(format!(
2513                "protected recent user {index} {}",
2514                "large active content ".repeat(300)
2515            )));
2516        }
2517
2518        let result = build_forced_compression_candidate_plan(
2519            &session,
2520            "main-model",
2521            Some(&budget),
2522            0.20,
2523            CompressionTriggerType::CriticalOverflow,
2524        );
2525        assert!(matches!(
2526            result,
2527            Err(CompressionPlanError::ProtectedContentExceedsTarget { .. })
2528        ));
2529    }
2530}