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
9fn 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
21fn 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
33fn 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#[derive(Debug, Clone)]
69pub enum CompressionPlanError {
70 ExposureGateNotMet {
72 usage_percent: f64,
73 trigger_percent: u8,
74 },
75 NoActiveMessages,
77 NotEnoughMessages { non_system_count: usize },
79 NothingToCompress {
81 anchor_index: usize,
82 non_system_count: usize,
83 },
84 ProtectedContentExceedsTarget {
87 projected_tokens: u32,
88 target_tokens: u32,
89 },
90 CandidateSetChanged,
92 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#[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#[derive(Debug, Clone)]
162pub struct CompressionPlan {
163 pub logical_pass_id: Option<String>,
166 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#[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
311pub 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
333pub 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
462pub 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
559pub fn estimate_context_compression_exposure(
562 session: &Session,
563 model_name: &str,
564 configured_budget: Option<&TokenBudget>,
565) -> ContextCompressionExposure {
566 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 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 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 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
642pub 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
660pub 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 let mut messages_to_summarize = non_system[..anchor_index].to_vec();
772
773 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 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 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 break;
823 };
824 let moved = messages_to_keep.remove(remove_index);
825 messages_to_summarize.push(moved);
826 }
827
828 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 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 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
912pub(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 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
950pub(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
987fn 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 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 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 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(§ions.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 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 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 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 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 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 session.add_message(Message::user("Old question about X"));
1648 session.add_message(Message::assistant("Old answer about X", None));
1649
1650 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}