返回 CodeWhale
compaction.rs
根目录 / crates / tui / src / compaction.rs
1 //! Context compaction for long conversations.
2
3 use anyhow::Result;
4 use std::collections::HashMap;
5 use std::fmt::Write;
6 use std::sync::Arc;
7 use std::time::Duration;
8
9 use crate::config::DEFAULT_TEXT_MODEL;
10 use crate::core::model_client::ModelClient;
11 use crate::logging;
12 use codewhale_models::Role;
13 use codewhale_models::{
14 CacheControl, ContentBlock, Message, MessageRequest, SystemBlock, SystemPrompt, Tool, Usage,
15 };
16
17 #[path = "compaction/last_round.rs"]
18 mod last_round;
19 #[cfg(test)]
20 #[path = "compaction/survival_contract.rs"]
21 mod survival_contract;
22 pub(crate) use last_round::last_round_start;
23 pub use last_round::{
24 CompactionCoverage, CompactionKeep, CompactionPath, LastCompactionSnapshot,
25 inspect_compaction_keep, last_round_kept_count, pinned_anchors_text,
26 };
27
28 /// Configuration for conversation compaction behavior.
29 ///
30 /// v0.8.11 simplified this from the prior token-OR-message-count trigger
31 /// to a token-only trigger. The
32 /// `message_threshold` field was removed: its only purpose was to fire
33 /// compaction on long sessions of small messages, which is exactly the
34 /// case where rewriting the prefix cache is least valuable. Token
35 /// budget is the right signal; message count was a 128K-era heuristic.
36 #[derive(Debug, Clone, PartialEq)]
37 pub struct CompactionConfig {
38 pub enabled: bool,
39 pub token_threshold: usize,
40 pub model: String,
41 /// Exact route image-input fact for the summarizer's outbound history.
42 pub image_input: crate::model_profile::SupportState,
43 /// Route-effective context window. `None` preserves compatibility for
44 /// callers that have not resolved a provider route yet.
45 pub effective_context_window: Option<u32>,
46 pub cache_summary: bool,
47 /// Optional user-supplied focus for a manual `/compact <focus>`: injected
48 /// into the summary request so the checkpoint weights what the user
49 /// said matters. `None` for automatic compaction.
50 pub focus: Option<String>,
51 /// Runtime turn that owns provider calls made by this compaction pass.
52 /// `None` for the foreground TUI. This is accounting provenance only and
53 /// is never included in a provider request.
54 pub runtime_cost_owner: Option<String>,
55 /// Workspace root, used only to re-state the user's `/anchor` file after
56 /// the summary. `None` skips anchors.
57 pub workspace: Option<std::path::PathBuf>,
58 /// Standing operator instructions from `[compaction] summary_instructions`
59 /// (#5956), appended to the summarizer prompt on every pass — manual and
60 /// automatic. `None` keeps the built-in prompt byte-identical. A manual
61 /// `/compact <focus>` still composes after this text.
62 pub summary_instructions: Option<String>,
63 /// Verbatim retention budget for recent plain user messages in the
64 /// replacement history (`[compaction] retained_user_message_tokens`,
65 /// #5956). Defaults to [`COMPACT_RETAINED_USER_MESSAGE_MAX_TOKENS`].
66 pub retained_user_message_tokens: usize,
67 }
68
69 /// Host callback for user-visible progress during a compaction pass.
70 ///
71 /// Compaction runs inside the engine's provider boundary, where the engine
72 /// owns the only channel that reaches the person. Injecting a sink lets a
73 /// downgrade (re-encoding images after an HTTP 413, say) say what it is doing
74 /// while it is doing it, instead of surfacing only when the pass ends.
75 pub trait CompactionNoticeSink: Send + Sync + std::fmt::Debug {
76 /// Deliver one already-rendered, user-visible sentence.
77 fn notice(&self, message: String);
78 /// Captured dispatch origin for a child; ordinary engines retain the existing billing path.
79 fn accounting_origin(&self) -> Option<(crate::cost_status::CostScopeToken, String, String)> {
80 None
81 }
82 /// Projection after the existing billing block has settled this exact response once.
83 fn settled_usage<'a>(
84 &'a self,
85 _source: &'a str,
86 _route: &'a crate::cost_status::EffectiveRouteEnvelope,
87 _usage: &'a Usage,
88 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
89 Box::pin(async {})
90 }
91 }
92
93 /// Host-prepared configuration carried from compaction eligibility through
94 /// the replacement-history commit.
95 #[derive(Clone)]
96 pub struct PreparedCompactionEnvelope {
97 pub config: CompactionConfig,
98 /// Durable handoff owner; set by the engine, never added to the stable prefix.
99 pub session_id: Option<String>,
100 /// Exact tool prefix of the interrupted request. Tool execution remains
101 /// disabled on the summary call; retaining schemas preserves cache reuse.
102 pub tools: Option<Vec<Tool>>,
103 /// Resolved reasoning tier the parent turn sends (#6540). Reasoning
104 /// routes render the effort into the head of the prompt, so a summary
105 /// request that omits it shares no cacheable prefix with the turn it
106 /// summarizes and re-bills the whole history uncached.
107 pub reasoning_effort: Option<String>,
108 /// User-visible progress sink supplied by the host that owns the run
109 /// (the interactive engine, typically). `None` keeps every notice in the
110 /// log only — the pass itself never depends on it.
111 pub notice_sink: Option<Arc<dyn CompactionNoticeSink>>,
112 }
113
114 impl std::fmt::Debug for PreparedCompactionEnvelope {
115 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
116 f.debug_struct("PreparedCompactionEnvelope")
117 .field("config", &self.config)
118 .field("session_id", &self.session_id)
119 .field("tools", &self.tools)
120 .field("reasoning_effort", &self.reasoning_effort)
121 .field("notice_sink", &self.notice_sink.as_ref().map(|_| "<sink>"))
122 .finish()
123 }
124 }
125
126 impl PartialEq for PreparedCompactionEnvelope {
127 /// The notice sink is host plumbing, not envelope content: two passes over
128 /// the same config are equal whether or not their host delivers notices.
129 fn eq(&self, other: &Self) -> bool {
130 self.config == other.config
131 && self.session_id == other.session_id
132 && self.tools == other.tools
133 && self.reasoning_effort == other.reasoning_effort
134 }
135 }
136
137 impl PreparedCompactionEnvelope {
138 #[must_use]
139 pub fn new(config: CompactionConfig) -> Self {
140 Self {
141 config,
142 session_id: None,
143 tools: None,
144 reasoning_effort: None,
145 notice_sink: None,
146 }
147 }
148 }
149
150 impl Default for CompactionConfig {
151 fn default() -> Self {
152 Self {
153 // ON BY DEFAULT since v0.8.6 (#402 P0 survivability). v0.8.64
154 // resolves the user-facing default through the active model's
155 // known context window, while explicit `auto_compact = false`
156 // remains the opt-out. This fallback covers code paths that build
157 // a `CompactionConfig` directly; real per-model values are still
158 // derived through the threshold helpers.
159 enabled: true,
160 // v0.8.11: 50K was a 128K-era leftover that biased every
161 // unconfigured caller toward "compact almost immediately on large-context routes."
162 // Bumped to 800K (80% of a 1M window) so the fallback
163 // default matches the hard automatic compaction guardrail. This
164 // keeps replacement compaction a late continuity guardrail.
165 // Real call sites override this via
166 // `compaction_threshold_for_model_and_effort`.
167 token_threshold: 800_000,
168 model: DEFAULT_TEXT_MODEL.to_string(),
169 image_input: crate::model_profile::SupportState::Unknown,
170 effective_context_window: None,
171 cache_summary: true,
172 focus: None,
173 runtime_cost_owner: None,
174 workspace: None,
175 summary_instructions: None,
176 retained_user_message_tokens: COMPACT_RETAINED_USER_MESSAGE_MAX_TOKENS,
177 }
178 }
179 }
180
181 /// A provider can return HTTP success with an empty, non-text, or known
182 /// placeholder response. Committing that response would discard the useful
183 /// history while leaving only a placeholder checkpoint. Keep this deliberately
184 /// conservative: it is a corruption guard, not a prose-length or language
185 /// scorer.
186 const COMPACTION_LANGUAGE_CONTRACT: &str = "Use the natural language of the most recent \
187 substantive user message for reasoning and user-facing prose. Keep code, identifiers, paths, \
188 commands, logs, tool payloads, quotations, and the English headings verbatim. English \
189 headings are not a request to switch languages.";
190
191 /// Failure kind for compaction LLM calls (deterministic vs transient).
192 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
193 pub enum CompactionFailureKind {
194 /// Same payload will fail again — do not sleep/retry unchanged.
195 Deterministic,
196 /// May resolve on retry (network, rate limit, timeout).
197 Transient,
198 /// Context overflow — drop the oldest history item and retry.
199 ContextOverflow,
200 }
201
202 impl CompactionFailureKind {
203 #[must_use]
204 pub fn is_transient(self) -> bool {
205 matches!(self, Self::Transient)
206 }
207 }
208
209 pub const KEEP_RECENT_MESSAGES: usize = 4;
210 const MIN_SUMMARIZE_MESSAGES: usize = 6;
211 const SUMMARY_TOOL_RESULT_SNIPPET_CHARS: usize = 240;
212 const TOOL_PRUNE_STOP_CHECK_BYTES: usize = 16 * 1024;
213 const RETAINED_TOOL_RESULT_MAX_CHARS: usize = 64 * 1024;
214 /// Token budget for the recent user messages retained verbatim in the
215 /// replacement history (Codex parity: COMPACT_USER_MESSAGE_MAX_TOKENS).
216 ///
217 /// This is now the *default* only: `[compaction] retained_user_message_tokens`
218 /// overrides it per session (#5956). Aliased to the config default so the two
219 /// names cannot drift apart.
220 pub(crate) const COMPACT_RETAINED_USER_MESSAGE_MAX_TOKENS: usize =
221 crate::config::DEFAULT_COMPACTION_RETAINED_USER_MESSAGE_TOKENS;
222 /// Handoff summarization prompt, appended to the live conversation as the
223 /// final user message. The headings are requested, not enforced:
224 /// `validate_compaction_summary` only rejects corrupt output, so a provider
225 /// that drifts from the layout still produces a usable note.
226 pub(crate) const COMPACT_PROMPT_OPENING: &str = "Write a handoff note so this session's work can \
227 continue after its earlier turns are condensed to make room.";
228
229 /// The note's layout, shared by the first request and the quality retry.
230 /// Each entry is a heading and what belongs under it.
231 const HANDOFF_SECTIONS: [(&str, &str); 10] = [
232 (
233 "Objective",
234 "what the user wants now, in their words where it matters; say if the goal changed and \
235 what replaced it.",
236 ),
237 (
238 "User direction",
239 "corrections, preferences and decisions, newest last; quote exactly anything the user \
240 corrected or insisted on.",
241 ),
242 (
243 "Permissions and limits",
244 "what the user explicitly allowed (pushing, deleting, spending, contacting someone) and \
245 what they ruled out. Record only what the user said; never infer a permission.",
246 ),
247 (
248 "Done",
249 "finished work with evidence: commands and results, commits, test counts. Separate \
250 verified from assumed.",
251 ),
252 (
253 "Changed files",
254 "each path and its state: edited, created or deleted; committed or not; which branch.",
255 ),
256 (
257 "Still running",
258 "shell commands, servers and ports started in this session that may still be live, each \
259 with the command or ID that checks or stops it. Leave agents out; their status is reported \
260 separately.",
261 ),
262 (
263 "Verification left",
264 "exact commands or checks still needed, and the last failure verbatim.",
265 ),
266 (
267 "Open questions",
268 "what waits on the user or is undecided, and what it blocks.",
269 ),
270 (
271 "Next action",
272 "the one next step, concrete enough to start without rereading history.",
273 ),
274 (
275 "Reference",
276 "exact paths, identifiers, URLs, values and short snippets the work depends on.",
277 ),
278 ];
279
280 const HANDOFF_FOLD_IN_RULE: &str = "If the history already holds an earlier handoff note, fold \
281 it in: carry forward what is still true, update what later work changed, and drop what is \
282 finished or superseded. Keep quoted user wording exact rather than paraphrasing it again. Do not \
283 copy the note's opening or closing lines or the user-pinned anchors; Codewhale adds those itself.";
284
285 const HANDOFF_CLOSING_RULE: &str = "Write about the task, not about condensing context. Be \
286 specific and brief. Do not call tools.";
287
288 const HANDOFF_FOCUS_LINE: &str = "The user asked this handoff to focus on:";
289
290 fn handoff_sections_text() -> String {
291 let mut text = String::from(
292 "Use these headings in this order, and write None under any heading with nothing to \
293 report:",
294 );
295 for (heading, guidance) in HANDOFF_SECTIONS {
296 let _ = write!(text, "\n## {heading} - {guidance}");
297 }
298 text
299 }
300
301 /// The full handoff request before the language contract, operator
302 /// instructions and focus line are appended.
303 fn compact_prompt_body() -> String {
304 format!(
305 "{COMPACT_PROMPT_OPENING} The reader sees only the note, the most recent messages and \
306 live tool state; anything left out is gone.\n\n{}\n\n{HANDOFF_FOLD_IN_RULE}\n\n{HANDOFF_CLOSING_RULE}",
307 handoff_sections_text()
308 )
309 }
310
311 /// First line of every handoff note Codewhale writes. It must describe what
312 /// `last_round::replacement_messages` actually keeps (see the
313 /// `summary_header_matches_what_replacement_history_keeps` test). It opens the
314 /// checkpoint message, so together with [`COMPACTION_CHECKPOINT_PROVENANCE`]
315 /// it is how a new-format checkpoint is recognised.
316 ///
317 /// A history rebuilt from turn records has to recognise the checkpoint
318 /// messages a document carries, so the header is crate-visible
319 /// (`runtime_threads`' recovery projection tests, #6664).
320 pub(crate) const SUMMARY_HEADER: &str = "Codewhale handoff note. Earlier turns of this session were \
321 condensed to make room. Kept above are the most recent user messages and the last steps of the \
322 current round. Long tool output there is shortened, with a marker where it was cut, and the \
323 oldest kept message may be shortened too. Everything earlier, including earlier steps of this \
324 round, exists only in this note; rerun a command or reread a file when its full output matters. \
325 Build on this note instead of redoing finished work, and check live state (files, git, running \
326 commands) before relying on anything it reports. Agent status, when there is any, comes from the separate agent status message, not from this note.";
327
328 const SUMMARY_CLOSING: &str = "Continue the user's task from here. Permissions and limits in \
329 this note restate what the user already decided; the note grants nothing new, and anything \
330 unclear gets checked before acting. Take the next action without asking the user to restate the \
331 task, save, or approve continuing because earlier turns were condensed.";
332
333 /// Detection marker for a new-format checkpoint: the opening words of
334 /// [`SUMMARY_HEADER`]. A new checkpoint is recognised only when its first
335 /// block starts with this marker and its second block is the provenance
336 /// block, so a user message that merely quotes the phrase stays a user
337 /// message.
338 pub const COMPACTION_SUMMARY_MARKER: &str = "Codewhale handoff note";
339 /// Legacy detection marker only. Checkpoints written before the handoff note
340 /// opened with this sentence. Saved sessions that still carry it must be
341 /// recognised so the next pass replaces that checkpoint instead of stacking a
342 /// second one.
343 pub const LEGACY_V2_COMPACTION_SUMMARY_MARKER: &str =
344 "Another language model started to solve this problem";
345 /// Legacy detection marker only, written by pre-v0.9.6 compaction; sessions
346 /// saved under that format must still be recognized so their summary is
347 /// replaced, not stacked.
348 pub const LEGACY_COMPACTION_SUMMARY_MARKER: &str = "Conversation Summary (Auto-Generated)";
349 /// Markers that identify a checkpoint by substring. Only the legacy markers
350 /// qualify: older checkpoints may lack the provenance block, while the new
351 /// marker is a plain phrase a user or a project instruction can quote.
352 const LEGACY_COMPACTION_SUMMARY_MARKERS: [&str; 2] = [
353 LEGACY_V2_COMPACTION_SUMMARY_MARKER,
354 LEGACY_COMPACTION_SUMMARY_MARKER,
355 ];
356 const COMPACTION_CHECKPOINT_PROVENANCE: &str = "<!-- codewhale.compaction-checkpoint.v1 -->";
357 const COMPACTION_SUMMARY_BEGIN: &str = "<!-- compaction-summary:begin -->";
358 const COMPACTION_SUMMARY_END: &str = "<!-- compaction-summary:end -->";
359
360 /// Heading the pre-v0.9.6 builder wrote, optionally after a
361 /// `## Pinned Facts (User Anchors)` section.
362 const LEGACY_SUMMARY_HEADING: &str = "## 📋 Conversation Summary (Auto-Generated)";
363 const LEGACY_ANCHORS_HEADING: &str = "## Pinned Facts (User Anchors)";
364
365 /// Whether text opens the way a legacy checkpoint opened. A message that only
366 /// quotes a marker later in its text is the user's, not a checkpoint (#6680).
367 /// Pre-v0.9.6 checkpoints opened with [`LEGACY_SUMMARY_HEADING`], or with the
368 /// pinned-anchors section followed by that heading on its own line.
369 fn is_legacy_compaction_summary_text(text: &str) -> bool {
370 let text = text.trim_start();
371 LEGACY_COMPACTION_SUMMARY_MARKERS
372 .iter()
373 .any(|marker| text.starts_with(marker))
374 || text.starts_with(LEGACY_SUMMARY_HEADING)
375 || (text.starts_with(LEGACY_ANCHORS_HEADING)
376 && text
377 .lines()
378 .any(|line| line.trim_end() == LEGACY_SUMMARY_HEADING))
379 }
380
381 /// Byte offset of the earliest legacy marker in `text`. New-format carriers
382 /// are always wrapped in the begin/end delimiters, so a bare new marker in a
383 /// system prompt is host text (for example a project instruction quoting it)
384 /// and must not truncate the prompt.
385 fn legacy_summary_start(text: &str) -> Option<usize> {
386 LEGACY_COMPACTION_SUMMARY_MARKERS
387 .iter()
388 .filter_map(|marker| text.find(marker))
389 .min()
390 }
391
392 fn summary_section(text: &str) -> Option<&str> {
393 let begin = text.find(COMPACTION_SUMMARY_BEGIN)? + COMPACTION_SUMMARY_BEGIN.len();
394 let remainder = &text[begin..];
395 let end = remainder.find(COMPACTION_SUMMARY_END)?;
396 let summary = remainder[..end].trim();
397 (!summary.is_empty()).then_some(summary)
398 }
399
400 fn strip_summary_text(mut text: String) -> Option<String> {
401 while let Some(begin) = text.find(COMPACTION_SUMMARY_BEGIN) {
402 let after_begin = begin + COMPACTION_SUMMARY_BEGIN.len();
403 let end = text[after_begin..]
404 .find(COMPACTION_SUMMARY_END)
405 .map_or(text.len(), |offset| {
406 after_begin + offset + COMPACTION_SUMMARY_END.len()
407 });
408 text.replace_range(begin..end, "");
409 }
410 if let Some(marker) = legacy_summary_start(&text) {
411 text.truncate(marker);
412 }
413 let text = text.trim().to_string();
414 (!text.is_empty()).then_some(text)
415 }
416
417 /// Extract the persisted checkpoint payload from a legacy system-prompt
418 /// carrier. Runtime-thread storage used that carrier before checkpoints moved
419 /// into conversation history; the engine strips it before provider dispatch.
420 #[must_use]
421 pub fn extract_compaction_summary(prompt: Option<&SystemPrompt>) -> Option<SystemPrompt> {
422 match prompt? {
423 SystemPrompt::Text(text) => summary_section(text)
424 .map(str::to_string)
425 .or_else(|| legacy_summary_start(text).map(|start| text[start..].trim().to_string()))
426 .map(SystemPrompt::Text),
427 SystemPrompt::Blocks(blocks) => {
428 let blocks = blocks
429 .iter()
430 .filter_map(|block| {
431 let text = summary_section(&block.text)
432 .map(str::to_string)
433 .or_else(|| {
434 legacy_summary_start(&block.text)
435 .map(|start| block.text[start..].trim().to_string())
436 })?;
437 let mut summary = block.clone();
438 summary.text = text;
439 Some(summary)
440 })
441 .collect::<Vec<_>>();
442 (!blocks.is_empty()).then_some(SystemPrompt::Blocks(blocks))
443 }
444 }
445 }
446
447 /// Remove every committed compaction-summary block from a system prompt.
448 ///
449 /// Compaction commits exactly one live summary: the newest one replaces its
450 /// predecessors. Before this existed, each compaction appended another
451 /// summary block to the successor system prompt, so the stable prefix grew by
452 /// up to a full summary per pass — which re-latched compaction pressure and
453 /// retriggered compaction on the next turn, forever.
454 #[must_use]
455 pub fn strip_compaction_summaries(prompt: Option<&SystemPrompt>) -> Option<SystemPrompt> {
456 match prompt.cloned()? {
457 SystemPrompt::Text(text) => strip_summary_text(text).map(SystemPrompt::Text),
458 SystemPrompt::Blocks(blocks) => {
459 let blocks = blocks
460 .into_iter()
461 .filter_map(|mut block| {
462 block.text = strip_summary_text(block.text)?;
463 Some(block)
464 })
465 .collect::<Vec<_>>();
466 (!blocks.is_empty()).then_some(SystemPrompt::Blocks(blocks))
467 }
468 }
469 }
470
471 /// Flatten a committed summary prompt to the text stored in history.
472 #[must_use]
473 pub fn summary_prompt_text(prompt: &SystemPrompt) -> String {
474 match prompt {
475 SystemPrompt::Text(text) => text.clone(),
476 SystemPrompt::Blocks(blocks) => blocks
477 .iter()
478 .map(|block| block.text.as_str())
479 .collect::<Vec<_>>()
480 .join("\n\n"),
481 }
482 }
483
484 #[must_use]
485 pub(crate) fn compaction_checkpoint_message(prompt: &SystemPrompt) -> Message {
486 Message {
487 role: Role::User,
488 content: vec![
489 ContentBlock::Text {
490 text: summary_prompt_text(prompt),
491 cache_control: None,
492 },
493 ContentBlock::Text {
494 text: COMPACTION_CHECKPOINT_PROVENANCE.to_string(),
495 cache_control: None,
496 },
497 ],
498 }
499 }
500
501 /// Whether a history message is a compaction checkpoint. New-format
502 /// checkpoints are recognised only structurally (see
503 /// [`is_wire_compaction_checkpoint_message`]); substring matching is kept for
504 /// the two legacy markers, whose older checkpoints may lack provenance.
505 #[must_use]
506 pub(crate) fn is_compaction_checkpoint_message(message: &Message) -> bool {
507 is_wire_compaction_checkpoint_message(message)
508 || user_text_of(message).is_some_and(|text| is_legacy_compaction_summary_text(&text))
509 }
510
511 /// Request-time recognition: exactly the header text block plus the
512 /// provenance block. User text merely quoting a marker keeps its original
513 /// wire position. Checkpoints saved before the handoff note still pass, so
514 /// their position survives restore and the next pass replaces them.
515 pub(crate) fn is_wire_compaction_checkpoint_message(message: &Message) -> bool {
516 let [
517 ContentBlock::Text {
518 text,
519 cache_control: None,
520 },
521 ContentBlock::Text {
522 text: provenance,
523 cache_control: None,
524 },
525 ] = message.content.as_slice()
526 else {
527 return false;
528 };
529 message.role == Role::User
530 && (text.starts_with(COMPACTION_SUMMARY_MARKER)
531 || text.starts_with(LEGACY_V2_COMPACTION_SUMMARY_MARKER))
532 && provenance == COMPACTION_CHECKPOINT_PROVENANCE
533 }
534
535 /// Keep the checkpoint at its original historical boundary on session load.
536 /// Later user turns must remain after the saved compaction boundary.
537 pub(crate) fn restore_compaction_checkpoint(
538 mut messages: Vec<Message>,
539 checkpoint: Option<&SystemPrompt>,
540 ) -> Vec<Message> {
541 let typed_position = messages
542 .iter()
543 .position(is_wire_compaction_checkpoint_message);
544 let checkpoint_index = if let Some(index) = typed_position {
545 messages.retain(|message| !is_wire_compaction_checkpoint_message(message));
546 index
547 } else {
548 // Legacy sessions have no independent provenance. Only replace
549 // header-prefixed messages when a saved checkpoint exists; identical
550 // user text is still ambiguous in that legacy format.
551 if checkpoint.is_some() {
552 messages.retain(|message| !is_compaction_checkpoint_message(message));
553 }
554 messages.len()
555 };
556 if let Some(checkpoint) = checkpoint {
557 messages.insert(
558 checkpoint_index.min(messages.len()),
559 compaction_checkpoint_message(checkpoint),
560 );
561 }
562 messages
563 }
564
565 pub(crate) fn estimate_tokens_for_message(message: &Message, include_thinking: bool) -> usize {
566 message
567 .content
568 .iter()
569 .map(|c| match c {
570 ContentBlock::Text { text, .. } => text.len() / 4,
571 // Replay-capable routes retain reasoning even on text-only
572 // assistant messages and across later user turns.
573 ContentBlock::Thinking { thinking, .. } if include_thinking => thinking.len() / 4,
574 ContentBlock::Thinking { .. } => 0,
575 ContentBlock::ToolUse { input, .. } => serde_json::to_string(input)
576 .map(|s| s.len() / 4)
577 .unwrap_or(100),
578 ContentBlock::ToolResult {
579 content,
580 content_blocks,
581 ..
582 } => {
583 let images = content_blocks.as_ref().map_or(0, |blocks| {
584 blocks
585 .iter()
586 .filter(|block| {
587 block.get("type").and_then(serde_json::Value::as_str) == Some("image")
588 })
589 .count()
590 });
591 content.len() / 4 + images * IMAGE_TOKEN_ESTIMATE
592 }
593 // An inline image is real input the model pays for; estimating it
594 // at 0 undercounts the budget and risks overflow in image-heavy
595 // sessions. Use a conservative flat per-image estimate (vision
596 // tiles are typically ~1k tokens); erring high compacts slightly
597 // early rather than overflowing.
598 ContentBlock::ImageUrl { .. } => IMAGE_TOKEN_ESTIMATE,
599 ContentBlock::ServerToolUse { input, .. } => input.to_string().len() / 4,
600 ContentBlock::ToolSearchToolResult { content, .. }
601 | ContentBlock::CodeExecutionToolResult { content, .. } => {
602 content.to_string().len() / 4
603 }
604 })
605 .sum::<usize>()
606 }
607
608 /// Conservative flat token estimate for an inline image (`ContentBlock::ImageUrl`).
609 /// Vision models bill images by resized tile count; ~1k tokens is a safe
610 /// mid-range estimate that keeps the compaction trigger from under-reading an
611 /// image-heavy session.
612 const IMAGE_TOKEN_ESTIMATE: usize = 1000;
613
614 pub fn estimate_tokens(messages: &[Message]) -> usize {
615 // Rough estimate: ~4 bytes per token. Count every retained reasoning
616 // block: DeepSeek/Kimi replay text-only assistant reasoning too. This
617 // route-neutral estimate cannot assume a transport will omit it.
618 messages
619 .iter()
620 .map(|message| estimate_tokens_for_message(message, true))
621 .sum()
622 }
623
624 pub(crate) fn message_has_tool_use(message: &Message) -> bool {
625 message
626 .content
627 .iter()
628 .any(|block| matches!(block, ContentBlock::ToolUse { .. }))
629 }
630
631 /// Conservative text estimate: three characters per token, but never below
632 /// the bytes/4 estimate the message path uses. Characters alone read
633 /// multibyte (CJK) text at ~0.33 tokens each, *under* the byte estimate's
634 /// ~0.75, which made the "conservative" figure the smaller one exactly where
635 /// it matters.
636 ///
637 /// Known limitation: both are heuristics, not a tokenizer. CJK costs roughly
638 /// 0.6-1.5 tokens per character depending on the provider's tokenizer, so a
639 /// CJK-heavy prompt can still be underestimated; the compaction trigger
640 /// bounds that by taking the larger of this estimate and the provider-billed
641 /// prompt size.
642 pub(crate) fn estimate_text_tokens_conservative(text: &str) -> usize {
643 text.chars().count().div_ceil(3).max(text.len().div_ceil(4))
644 }
645
646 fn estimate_system_tokens_conservative(system: Option<&SystemPrompt>) -> usize {
647 match system {
648 Some(SystemPrompt::Text(text)) => estimate_text_tokens_conservative(text),
649 Some(SystemPrompt::Blocks(blocks)) => blocks
650 .iter()
651 .map(|block| estimate_text_tokens_conservative(&block.text))
652 .sum(),
653 None => 0,
654 }
655 }
656
657 /// Conservative estimate for full request input tokens (messages + system + framing).
658 #[must_use]
659 pub fn estimate_input_tokens_conservative(
660 messages: &[Message],
661 system: Option<&SystemPrompt>,
662 ) -> usize {
663 let message_tokens = estimate_tokens(messages).saturating_mul(3).div_ceil(2);
664 let system_tokens = estimate_system_tokens_conservative(system);
665 let framing_overhead = messages.len().saturating_mul(12).saturating_add(48);
666 message_tokens
667 .saturating_add(system_tokens)
668 .saturating_add(framing_overhead)
669 }
670
671 /// Best-effort estimate of real request input tokens, without the 1.5×
672 /// safety inflation used by overflow math.
673 ///
674 /// Compaction *pressure* compares against a threshold whose percentage means
675 /// "fraction of the context window" on the user-facing meter. Feeding the
676 /// inflated overflow estimate into that comparison made an 80% setting fire
677 /// at roughly half the real usage. Overflow protection keeps its inflated
678 /// estimator; the pressure trigger uses this one, preferring provider-billed
679 /// prompt tokens when the caller has them.
680 #[must_use]
681 pub fn estimate_input_tokens_for_pressure(
682 messages: &[Message],
683 system: Option<&SystemPrompt>,
684 ) -> usize {
685 let message_tokens = estimate_tokens(messages);
686 let system_tokens = estimate_system_tokens_conservative(system);
687 let framing_overhead = messages.len().saturating_mul(12).saturating_add(48);
688 message_tokens
689 .saturating_add(system_tokens)
690 .saturating_add(framing_overhead)
691 }
692
693 fn estimate_retained_floor_conservative(
694 messages: &[Message],
695 system_prompt: Option<&SystemPrompt>,
696 prepared: &PreparedCompactionEnvelope,
697 ) -> usize {
698 let config = &prepared.config;
699 let retained = last_round::replacement_messages(messages, config.retained_user_message_tokens);
700 let retained_tokens = estimate_tokens(&retained).saturating_mul(3).div_ceil(2);
701 let framing = retained.len().saturating_mul(12).saturating_add(48);
702 let anchors = user_anchors_section(config.workspace.as_deref());
703 let summary_scaffolding_tokens =
704 estimate_text_tokens_conservative(&build_compaction_summary_block_text("", &anchors));
705
706 // Post-compaction the committed summary is REPLACED, not stacked, so prior
707 // summary blocks must not inflate the floor. Count only the exact installed
708 // scaffolding here; the model owns the concise summary length, just as it
709 // owns the answer length on an ordinary turn.
710 let retained_system_prompt = strip_compaction_summaries(system_prompt);
711 retained_tokens
712 .saturating_add(estimate_system_tokens_conservative(
713 retained_system_prompt.as_ref(),
714 ))
715 .saturating_add(framing)
716 .saturating_add(summary_scaffolding_tokens)
717 }
718
719 /// Whether the current canonical request has reached the configured automatic
720 /// compaction pressure. This deliberately excludes eligibility/reclaimability:
721 /// local tool-result pruning uses it to decide when pressure has actually
722 /// cleared, even if the remaining transcript cannot support an LLM summary.
723 #[must_use]
724 pub fn compaction_pressure_reached(
725 messages: &[Message],
726 system_prompt: Option<&SystemPrompt>,
727 config: &CompactionConfig,
728 ) -> bool {
729 compaction_pressure_reached_with_billed(messages, system_prompt, config, None)
730 }
731
732 /// Pressure check that additionally honors provider-billed prompt tokens.
733 ///
734 /// Billed usage is the ground truth for how large the context actually is;
735 /// the estimator undercounts non-ASCII text and cannot see server-side
736 /// framing. Whichever signal is higher decides, so an undercounting estimate
737 /// cannot hide pressure the provider already billed for. Callers must only
738 /// pass a billed count that describes the message list being checked —
739 /// post-prune re-checks pass `None` and fall back to the estimate.
740 #[must_use]
741 pub fn compaction_pressure_reached_with_billed(
742 messages: &[Message],
743 system_prompt: Option<&SystemPrompt>,
744 config: &CompactionConfig,
745 billed_input_tokens: Option<u64>,
746 ) -> bool {
747 if !config.enabled {
748 return false;
749 }
750 let billed = billed_input_tokens
751 .and_then(|tokens| usize::try_from(tokens).ok())
752 .unwrap_or(0);
753 // Billing alone proving pressure short-circuits the walk (#perf-r5):
754 // `estimated.max(billed) >= threshold` is unconditionally true when
755 // `billed >= threshold`, so estimating cannot change the answer and the
756 // O(transcript) pass is skipped. Over-pressure sessions pay this check
757 // multiple times per step (pressure gate + decision re-check).
758 if billed >= config.token_threshold {
759 return true;
760 }
761 let estimated = estimate_input_tokens_for_pressure(messages, system_prompt);
762 estimated.max(billed) >= config.token_threshold
763 }
764
765 /// Estimate-only eligibility check ([`should_compact_with_billed`] with no
766 /// billed tokens): used by the request preview, where no provider bill exists.
767 pub fn should_compact(
768 messages: &[Message],
769 system_prompt: Option<&SystemPrompt>,
770 prepared: &PreparedCompactionEnvelope,
771 ) -> bool {
772 should_compact_with_billed(messages, system_prompt, prepared, None)
773 }
774
775 /// Why an over-pressure context still did not start an automatic pass.
776 ///
777 /// A refusal is not a bug by itself — each guard exists for a reason — but a
778 /// silent refusal is: the user watches a full context meter while
779 /// auto-compaction appears broken (#5577). Callers surface these.
780 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
781 pub enum CompactionRefusal {
782 /// Too few messages for a summary pass to mean anything.
783 TooFewMessages { count: usize },
784 /// The conservative retained floor (system prompt + kept messages +
785 /// summary allowance) cannot get below the trigger, so a pass would
786 /// recur on every step without relieving pressure.
787 RetainedFloor { floor: usize, threshold: usize },
788 }
789
790 /// Outcome of the automatic-compaction eligibility check.
791 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
792 pub enum CompactionDecision {
793 /// Disabled, or pressure not reached: nothing to do, nothing to explain.
794 NotNeeded,
795 /// Start a pass.
796 Compact,
797 /// Pressure is real but a guard declined; the reason names the guard.
798 Refused(CompactionRefusal),
799 }
800
801 /// Eligibility check that honors provider-billed prompt tokens for the
802 /// pressure gate, mirroring [`compaction_pressure_reached_with_billed`].
803 pub fn should_compact_with_billed(
804 messages: &[Message],
805 system_prompt: Option<&SystemPrompt>,
806 prepared: &PreparedCompactionEnvelope,
807 billed_input_tokens: Option<u64>,
808 ) -> bool {
809 matches!(
810 compaction_decision_with_billed(messages, system_prompt, prepared, billed_input_tokens),
811 CompactionDecision::Compact
812 )
813 }
814
815 /// Full eligibility decision, including *why* an over-pressure context was
816 /// refused, so hosts can tell the user instead of silently holding.
817 #[must_use]
818 pub fn compaction_decision_with_billed(
819 messages: &[Message],
820 system_prompt: Option<&SystemPrompt>,
821 prepared: &PreparedCompactionEnvelope,
822 billed_input_tokens: Option<u64>,
823 ) -> CompactionDecision {
824 let config = &prepared.config;
825 if !config.enabled {
826 return CompactionDecision::NotNeeded;
827 }
828 // Pressure gate + prune projection share one estimate (#perf-r5): both
829 // consume `estimate_input_tokens_for_pressure` over the same
830 // `(messages, system_prompt)`, a pure function, so it is computed at
831 // most once. `billed >= threshold` proves pressure without estimating
832 // (max is unconditionally >= threshold then); the estimate is deferred
833 // until something actually needs it — the prune projection below — so
834 // the billed-corner still reaches the TooFew and RetainedFloor guards
835 // unchanged, and skips the walk entirely when no prune candidates exist.
836 let billed = billed_input_tokens
837 .and_then(|tokens| usize::try_from(tokens).ok())
838 .unwrap_or(0);
839 let estimated: Option<usize> = if billed < config.token_threshold {
840 let estimate = estimate_input_tokens_for_pressure(messages, system_prompt);
841 if estimate.max(billed) < config.token_threshold {
842 return CompactionDecision::NotNeeded;
843 }
844 Some(estimate)
845 } else {
846 None
847 };
848
849 // The execution path mechanically prunes old verbose tool results before
850 // asking the model for a summary. Local pruning alone may be enough to
851 // clear pressure even when the transcript is too small for an LLM pass.
852 // Project that outcome from the measured plan — per-block deltas use the
853 // estimator's own arithmetic, so this equals re-estimating a pruned copy
854 // without cloning a multi-megabyte transcript on every step.
855 let prune_plan = plan_tool_result_prunes(messages, KEEP_RECENT_MESSAGES);
856 if !prune_plan.is_empty() {
857 let estimate = match estimated {
858 Some(value) => value,
859 None => estimate_input_tokens_for_pressure(messages, system_prompt),
860 };
861 let reclaimed_tokens: usize = prune_plan.iter().map(PlannedPrune::tokens_reclaimed).sum();
862 let projected = estimate.saturating_sub(reclaimed_tokens);
863 if projected < config.token_threshold.saturating_mul(4) / 5 {
864 return CompactionDecision::Compact;
865 }
866 }
867
868 if messages.len() < MIN_SUMMARIZE_MESSAGES {
869 return CompactionDecision::Refused(CompactionRefusal::TooFewMessages {
870 count: messages.len(),
871 });
872 }
873
874 // Reclaimability guard: do not start a pass whose replacement request
875 // (system prompt + retained user messages + summary allowance)
876 // cannot get below the trigger, or a large stable prefix would cause
877 // auto-compaction on every tool step.
878 let floor = estimate_retained_floor_conservative(messages, system_prompt, prepared);
879 if floor >= config.token_threshold {
880 return CompactionDecision::Refused(CompactionRefusal::RetainedFloor {
881 floor,
882 threshold: config.token_threshold,
883 });
884 }
885 CompactionDecision::Compact
886 }
887
888 /// Whether a compaction pass could shrink this history at all: enough
889 /// messages to summarize, or old tool output to prune. A one- or two-message
890 /// conversation that is over budget is over budget because of its fixed
891 /// prefix or its newest message, and summarizing it only spends a model call
892 /// before the same failure (experience mark 2).
893 #[must_use]
894 pub fn has_compactable_history(messages: &[Message]) -> bool {
895 messages.len() >= MIN_SUMMARIZE_MESSAGES
896 || !plan_tool_result_prunes(messages, KEEP_RECENT_MESSAGES).is_empty()
897 }
898
899 fn truncate_chars(text: &str, max_chars: usize) -> &str {
900 if max_chars == 0 {
901 return "";
902 }
903 match text.char_indices().nth(max_chars) {
904 Some((idx, _)) => &text[..idx],
905 None => text,
906 }
907 }
908
909 fn tail_chars(text: &str, max_chars: usize) -> String {
910 if max_chars == 0 {
911 return String::new();
912 }
913 let total_chars = text.chars().count();
914 if total_chars <= max_chars {
915 return text.to_string();
916 }
917 let start_char = total_chars.saturating_sub(max_chars);
918 let start_idx = text
919 .char_indices()
920 .nth(start_char)
921 .map_or(0, |(idx, _)| idx);
922 text[start_idx..].to_string()
923 }
924
925 #[derive(Debug, Clone)]
926 struct ToolUseInfo {
927 name: String,
928 key: String,
929 args_preview: String,
930 }
931
932 fn tool_use_key(name: &str, input: &serde_json::Value) -> String {
933 format!(
934 "{name}:{}",
935 serde_json::to_string(input).unwrap_or_else(|_| input.to_string())
936 )
937 }
938
939 fn tool_args_preview(input: &serde_json::Value) -> String {
940 let redacted = codewhale_config::persistence::redact_json_secrets(input);
941 let raw = serde_json::to_string(&redacted).unwrap_or_else(|_| redacted.to_string());
942 truncate_chars(&raw, 120).to_string()
943 }
944
945 fn collect_tool_uses(
946 messages: &[Message],
947 ) -> HashMap<codewhale_models::ToolCallKey<'_>, Option<(&str, ToolUseInfo)>> {
948 let mut tool_uses = HashMap::new();
949 for message in messages {
950 for block in &message.content {
951 if let ContentBlock::ToolUse {
952 id, name, input, ..
953 } = block
954 && let Some(key) = block.tool_call_key()
955 && !key.as_str().trim().is_empty()
956 {
957 // Ambiguous old or malformed new history is left intact; a
958 // later call must not supply the earlier result's metadata.
959 tool_uses
960 .entry(key)
961 .and_modify(|entry| *entry = None)
962 .or_insert_with(|| {
963 Some((
964 id.as_str(),
965 ToolUseInfo {
966 name: name.clone(),
967 key: tool_use_key(name, input),
968 args_preview: tool_args_preview(input),
969 },
970 ))
971 });
972 }
973 }
974 }
975 tool_uses
976 }
977
978 struct ToolResultPruneCandidate {
979 message_idx: usize,
980 block_idx: usize,
981 key: String,
982 tool_name: String,
983 args_preview: String,
984 original_len: usize,
985 }
986
987 fn tool_result_content_blocks_len(content_blocks: Option<&[serde_json::Value]>) -> usize {
988 content_blocks
989 .and_then(|blocks| serde_json::to_vec(blocks).ok())
990 .map_or(0, |bytes| bytes.len())
991 }
992
993 #[cfg(test)]
994 fn prune_tool_results(messages: &mut [Message], protected_window: usize) -> usize {
995 prune_tool_results_until(messages, protected_window, |_, _| false)
996 }
997
998 /// Mechanically prune old verbose tool results before paying for an LLM summary.
999 ///
1000 /// The most recent `protected_window` messages stay byte-for-byte intact. Older
1001 /// duplicate tool results keep the freshest full body and replace earlier
1002 /// copies with one-line summaries; non-duplicate old results are summarized only
1003 /// when they exceed the normal summary snippet size.
1004 fn prune_tool_results_until<F>(
1005 messages: &mut [Message],
1006 protected_window: usize,
1007 mut should_stop: F,
1008 ) -> usize
1009 where
1010 F: FnMut(&[Message], usize) -> bool,
1011 {
1012 let plan = plan_tool_result_prunes(messages, protected_window);
1013 let mut bytes_saved = 0usize;
1014 for planned in plan {
1015 if let ContentBlock::ToolResult {
1016 content,
1017 content_blocks,
1018 ..
1019 } = &mut messages[planned.message_idx].content[planned.block_idx]
1020 {
1021 bytes_saved = bytes_saved.saturating_add(planned.bytes_reclaimed());
1022 *content = planned.summary;
1023 *content_blocks = None;
1024
1025 if should_stop(messages, bytes_saved) {
1026 break;
1027 }
1028 }
1029 }
1030 bytes_saved
1031 }
1032
1033 /// One tool-result replacement the pruner has decided on, measured up front
1034 /// so eligibility checks can project the outcome without cloning a
1035 /// multi-megabyte transcript ([`compaction_decision_with_billed`] used to
1036 /// copy the entire message list every over-pressure step just to ask "would
1037 /// pruning be enough?").
1038 struct PlannedPrune {
1039 message_idx: usize,
1040 block_idx: usize,
1041 summary: String,
1042 content_len: usize,
1043 blocks_len: usize,
1044 image_count: usize,
1045 }
1046
1047 impl PlannedPrune {
1048 /// Byte reduction this replacement realizes, matching the pruner's
1049 /// accounting exactly.
1050 fn bytes_reclaimed(&self) -> usize {
1051 self.content_len
1052 .saturating_sub(self.summary.len())
1053 .saturating_add(self.blocks_len)
1054 }
1055
1056 /// Estimator-token reduction, using the same per-block arithmetic as
1057 /// [`estimate_tokens_for_message`] so a projection built from these
1058 /// deltas equals re-estimating the pruned transcript.
1059 fn tokens_reclaimed(&self) -> usize {
1060 let before = self.content_len / 4 + self.image_count * IMAGE_TOKEN_ESTIMATE;
1061 let after = self.summary.len() / 4;
1062 before.saturating_sub(after)
1063 }
1064 }
1065
1066 /// Decide, without mutating anything, which old tool results pruning would
1067 /// replace. The most recent `protected_window` messages stay untouched; older
1068 /// duplicate results keep the freshest full body; non-duplicates are replaced
1069 /// only when they exceed the summary snippet size.
1070 fn plan_tool_result_prunes(messages: &[Message], protected_window: usize) -> Vec<PlannedPrune> {
1071 let cutoff = messages.len().saturating_sub(protected_window);
1072 if cutoff == 0 {
1073 return Vec::new();
1074 }
1075
1076 let tool_uses = collect_tool_uses(messages);
1077 let mut candidates = Vec::new();
1078 let mut latest_by_key: HashMap<String, usize> = HashMap::new();
1079 let mut count_by_key: HashMap<String, usize> = HashMap::new();
1080
1081 for (message_idx, message) in messages.iter().take(cutoff).enumerate() {
1082 for (block_idx, block) in message.content.iter().enumerate() {
1083 let ContentBlock::ToolResult {
1084 tool_use_id,
1085 content,
1086 content_blocks,
1087 ..
1088 } = block
1089 else {
1090 continue;
1091 };
1092 let Some((provider_id, info)) = block
1093 .tool_call_key()
1094 .and_then(|key| tool_uses.get(&key))
1095 .and_then(Option::as_ref)
1096 else {
1097 continue;
1098 };
1099 if provider_id != &tool_use_id.as_str() {
1100 continue;
1101 }
1102 latest_by_key.insert(info.key.clone(), message_idx);
1103 *count_by_key.entry(info.key.clone()).or_insert(0) += 1;
1104 candidates.push(ToolResultPruneCandidate {
1105 message_idx,
1106 block_idx,
1107 key: info.key.clone(),
1108 tool_name: info.name.clone(),
1109 args_preview: info.args_preview.clone(),
1110 original_len: content
1111 .len()
1112 .saturating_add(tool_result_content_blocks_len(content_blocks.as_deref())),
1113 });
1114 }
1115 }
1116
1117 // The maps above are fully populated before planning completes, so the order
1118 // below only changes which message bytes are rewritten first. Planning from
1119 // newest to oldest lets the pruner stop as soon as enough bytes were saved,
1120 // preserving the earlier JSON request prefix for byte-level KV caches.
1121 candidates.reverse();
1122
1123 let mut plan = Vec::new();
1124 for candidate in candidates {
1125 let duplicate_count = count_by_key.get(&candidate.key).copied().unwrap_or(0);
1126 let is_latest_duplicate = duplicate_count > 1
1127 && latest_by_key.get(&candidate.key) == Some(&candidate.message_idx);
1128 if is_latest_duplicate {
1129 continue;
1130 }
1131 if duplicate_count <= 1 && candidate.original_len <= SUMMARY_TOOL_RESULT_SNIPPET_CHARS {
1132 continue;
1133 }
1134
1135 let summary = format!(
1136 "[{}] tool result pruned ({} bytes; args: {})",
1137 candidate.tool_name, candidate.original_len, candidate.args_preview
1138 );
1139 if summary.len() >= candidate.original_len {
1140 continue;
1141 }
1142
1143 let ContentBlock::ToolResult {
1144 content,
1145 content_blocks,
1146 ..
1147 } = &messages[candidate.message_idx].content[candidate.block_idx]
1148 else {
1149 continue;
1150 };
1151 plan.push(PlannedPrune {
1152 message_idx: candidate.message_idx,
1153 block_idx: candidate.block_idx,
1154 summary,
1155 content_len: content.len(),
1156 blocks_len: tool_result_content_blocks_len(content_blocks.as_deref()),
1157 image_count: content_blocks.as_ref().map_or(0, |blocks| {
1158 blocks
1159 .iter()
1160 .filter(|block| {
1161 block.get("type").and_then(serde_json::Value::as_str) == Some("image")
1162 })
1163 .count()
1164 }),
1165 });
1166 }
1167 plan
1168 }
1169
1170 fn truncate_retained_block(label: &str, content: &mut String, max_chars: usize) -> bool {
1171 let char_count = content.chars().count();
1172 if char_count <= max_chars {
1173 return false;
1174 }
1175
1176 let snippet_budget = max_chars.saturating_sub(256).max(1024);
1177 let head_chars = snippet_budget / 2;
1178 let tail_chars_budget = snippet_budget.saturating_sub(head_chars);
1179 let head = truncate_chars(content, head_chars).to_string();
1180 let tail = tail_chars(content, tail_chars_budget);
1181 *content =
1182 format!("[{label} retained-history truncated from {char_count} chars]\n{head}\n…\n{tail}");
1183 true
1184 }
1185
1186 // Retained reasoning is replay protocol even without a signature (DeepSeek
1187 // tool turns). Summarize older exchanges as units; do not rewrite their peers.
1188 fn sanitize_retained_messages(mut messages: Vec<Message>) -> Vec<Message> {
1189 for message in &mut messages {
1190 for block in &mut message.content {
1191 if let ContentBlock::ToolResult {
1192 content,
1193 content_blocks,
1194 ..
1195 } = block
1196 && truncate_retained_block("tool result", content, RETAINED_TOOL_RESULT_MAX_CHARS)
1197 {
1198 *content_blocks = None;
1199 }
1200 }
1201 }
1202 messages
1203 }
1204
1205 /// Result of a compaction operation with metadata.
1206 #[derive(Debug)]
1207 pub struct CompactionResult {
1208 /// Compacted messages
1209 pub messages: Vec<Message>,
1210 /// Host-persistence copy of the history checkpoint.
1211 pub summary_prompt: Option<SystemPrompt>,
1212 /// Number of retries used before success
1213 pub retries_used: u32,
1214 /// Last-round coverage for inspector receipts.
1215 pub coverage: CompactionCoverage,
1216 }
1217
1218 /// Classify a compaction LLM failure for the retry / input-ladder policy.
1219 fn classify_compaction_failure(e: &anyhow::Error) -> CompactionFailureKind {
1220 if let Some(error) = llm_error_in_chain(e) {
1221 return match error {
1222 crate::llm_client::LlmError::ContextLengthError(_) => {
1223 CompactionFailureKind::ContextOverflow
1224 }
1225 crate::llm_client::LlmError::QuotaExhausted(_) => CompactionFailureKind::Deterministic,
1226 error if error.is_retryable() => CompactionFailureKind::Transient,
1227 _ => CompactionFailureKind::Deterministic,
1228 };
1229 }
1230
1231 let text = e.to_string();
1232 if is_context_window_error_message(&text) {
1233 return CompactionFailureKind::ContextOverflow;
1234 }
1235 let category = crate::error_taxonomy::classify_error_message(&text);
1236 match category {
1237 crate::error_taxonomy::ErrorCategory::Network
1238 | crate::error_taxonomy::ErrorCategory::RateLimit
1239 | crate::error_taxonomy::ErrorCategory::Timeout => CompactionFailureKind::Transient,
1240 _ => CompactionFailureKind::Deterministic,
1241 }
1242 }
1243
1244 fn llm_error_in_chain(error: &anyhow::Error) -> Option<&crate::llm_client::LlmError> {
1245 error
1246 .chain()
1247 .find_map(|cause| cause.downcast_ref::<crate::llm_client::LlmError>())
1248 }
1249
1250 /// Record and render a compaction failure as actionable, credential-safe text.
1251 ///
1252 /// This classifies only the error supplied by the failed request; it never
1253 /// infers a cause from later provider failures. Unknown diagnostics stay
1254 /// visible after central secret/path redaction, and the same safe detail is
1255 /// written to the runtime log so a transient status message remains auditable.
1256 #[must_use]
1257 pub fn report_compaction_failure(
1258 prefix: &str,
1259 id: &str,
1260 auto: bool,
1261 error: &anyhow::Error,
1262 ) -> String {
1263 let raw = error.to_string();
1264 let safe_raw = crate::safe_label::safe_error_text(&raw);
1265 tracing::warn!(
1266 compaction_id = %id,
1267 auto,
1268 error = %safe_raw,
1269 "context compaction failed"
1270 );
1271 let detail = match llm_error_in_chain(error) {
1272 Some(crate::llm_client::LlmError::QuotaExhausted(_)) => {
1273 "provider plan quota exhausted — switch provider/model or renew the provider plan"
1274 .to_string()
1275 }
1276 Some(crate::llm_client::LlmError::RateLimited { .. }) => {
1277 "provider rate limit blocked making room — retry after the limit resets or switch provider/model"
1278 .to_string()
1279 }
1280 Some(crate::llm_client::LlmError::AuthenticationError(_)) => {
1281 "provider authentication failed — sign in or replace the credential, then retry"
1282 .to_string()
1283 }
1284 Some(crate::llm_client::LlmError::AuthorizationError(_)) => {
1285 "provider authorization rejected making room — verify account access or switch provider/model"
1286 .to_string()
1287 }
1288 _ => match crate::error_taxonomy::classify_error_message(&raw) {
1289 crate::error_taxonomy::ErrorCategory::RateLimit => {
1290 "provider rate limit blocked making room — retry after the limit resets or switch provider/model"
1291 .to_string()
1292 }
1293 crate::error_taxonomy::ErrorCategory::Authentication => {
1294 "provider authentication failed — sign in or replace the credential, then retry"
1295 .to_string()
1296 }
1297 crate::error_taxonomy::ErrorCategory::Authorization => {
1298 "provider authorization rejected making room — verify account access or switch provider/model"
1299 .to_string()
1300 }
1301 _ => safe_raw,
1302 },
1303 };
1304
1305 format!("{prefix}: {detail}")
1306 }
1307
1308 /// Check if an error is transient and worth retrying. Categories that map to
1309 /// transient retry: Network, RateLimit, Timeout. Context overflow is *not*
1310 /// transient — it needs a smaller input (ladder), not the same payload.
1311 fn is_transient_error(e: &anyhow::Error) -> bool {
1312 classify_compaction_failure(e).is_transient()
1313 }
1314
1315 fn is_context_window_error_message(text: &str) -> bool {
1316 let lower = text.to_lowercase();
1317 lower.contains("too long for this model")
1318 || lower.contains("prompt is too long")
1319 || lower.contains("maximum prompt length")
1320 || lower.contains("maximum context length")
1321 || lower.contains("context_length_exceeded")
1322 || lower.contains("context window")
1323 || (lower.contains("context")
1324 && (lower.contains("token") || lower.contains("too long") || lower.contains("maximum")))
1325 }
1326
1327 /// Compact messages with retry and backoff for transient errors.
1328 ///
1329 /// This function wraps `compact_messages` with retry logic to handle
1330 /// transient network errors and rate limits. It uses exponential backoff
1331 /// with delays of 1s, 2s, 4s between retries.
1332 ///
1333 /// # Safety
1334 /// - Never panics
1335 /// - Never corrupts the original messages (returns error instead)
1336 /// - Only retries on transient errors (network, rate limit, etc.)
1337 ///
1338 /// `invocation_usage` retains every decoded response, including rejected
1339 /// summaries, across retries and cancellation of this future.
1340 pub async fn compact_messages_safe(
1341 client: &dyn ModelClient,
1342 messages: &[Message],
1343 system_prompt: Option<&SystemPrompt>,
1344 prepared: &PreparedCompactionEnvelope,
1345 invocation_usage: &mut Usage,
1346 ) -> Result<CompactionResult> {
1347 const MAX_RETRIES: u32 = 3;
1348 const BASE_DELAY_MS: u64 = 1000;
1349
1350 // Persist the complete pre-compaction history before any pruning or provider
1351 // call. Failure leaves the original context intact. The model-authored
1352 // handoff is saved separately before replacement is returned to the engine.
1353 let checkpoint_id = uuid::Uuid::new_v4().to_string();
1354 if let Some(session_id) = prepared.session_id.as_deref() {
1355 let bytes = serde_json::to_vec(&codewhale_config::persistence::redact_json_secrets(
1356 &serde_json::to_value(messages)?,
1357 ))?;
1358 crate::artifacts::write_session_relative_immutable(
1359 session_id,
1360 &std::path::PathBuf::from("artifacts")
1361 .join(format!("context-transfer-{checkpoint_id}.json")),
1362 &bytes,
1363 )?;
1364 }
1365
1366 let config = &prepared.config;
1367 let was_over_threshold = compaction_pressure_reached(messages, system_prompt, config);
1368 // Leave room for useful work after a local prune. Clearing the trigger
1369 // by a few tokens caused another prefix rewrite on the next tool result.
1370 let prune_target = config.token_threshold.saturating_mul(4) / 5;
1371 let mut pruned_messages = messages.to_vec();
1372 let mut now_under_threshold = false;
1373 let mut next_stop_check_bytes = 0usize;
1374 let pruned_bytes = prune_tool_results_until(
1375 &mut pruned_messages,
1376 KEEP_RECENT_MESSAGES,
1377 |candidate_messages, bytes_saved| {
1378 if !was_over_threshold || bytes_saved < next_stop_check_bytes {
1379 return false;
1380 }
1381
1382 // Stop at the first suffix-side prune check that clears the target.
1383 // The check itself is a full compaction-plan pass, so bound it by saved
1384 // bytes instead of running it after every candidate in huge sessions.
1385 next_stop_check_bytes = bytes_saved.saturating_add(TOOL_PRUNE_STOP_CHECK_BYTES);
1386 now_under_threshold =
1387 estimate_input_tokens_for_pressure(candidate_messages, system_prompt)
1388 < prune_target;
1389 now_under_threshold
1390 },
1391 );
1392 if was_over_threshold && pruned_bytes > 0 && !now_under_threshold {
1393 // The throttled in-loop check may skip the exact candidate that clears the
1394 // budget. Do one final pass so a successful local prune still avoids LLM compaction.
1395 now_under_threshold =
1396 estimate_input_tokens_for_pressure(&pruned_messages, system_prompt) < prune_target;
1397 }
1398
1399 if pruned_bytes > 0 {
1400 logging::info(format!(
1401 "Local tool-result prune saved {pruned_bytes} bytes before LLM compaction"
1402 ));
1403 if was_over_threshold && now_under_threshold {
1404 let kept = sanitize_retained_messages(pruned_messages);
1405 last_round::validate_last_round_coverage(messages, &kept)?;
1406 let coverage = last_round::measure_coverage(
1407 messages,
1408 &kept,
1409 CompactionPath::PruneOnly,
1410 pinned_anchors_text(config.workspace.as_deref())
1411 .map(|text| text.chars().count())
1412 .unwrap_or(0),
1413 );
1414 return Ok(CompactionResult {
1415 messages: kept,
1416 summary_prompt: None,
1417 retries_used: 0,
1418 coverage,
1419 });
1420 }
1421 }
1422
1423 let mut last_error: Option<anyhow::Error> = None;
1424 let mut quality_retries = 0u32;
1425
1426 for attempt in 0..MAX_RETRIES {
1427 if attempt > 0 {
1428 // Exponential backoff: 1s, 2s, 4s
1429 let delay = Duration::from_millis(BASE_DELAY_MS * (1 << (attempt - 1)));
1430 tokio::time::sleep(delay).await;
1431 }
1432
1433 match compact_messages_with_metadata(
1434 client,
1435 // If a local prune cannot clear pressure, summarize the original
1436 // evidence. Pruning first both erased facts the handoff needs and
1437 // invalidated the cached history prefix for the summary request.
1438 messages,
1439 config,
1440 system_prompt,
1441 prepared.tools.as_deref(),
1442 prepared.reasoning_effort.as_deref(),
1443 prepared.notice_sink.as_deref(),
1444 &mut quality_retries,
1445 invocation_usage,
1446 )
1447 .await
1448 {
1449 Ok((msgs, prompt, mut coverage)) => {
1450 let kept = sanitize_retained_messages(msgs);
1451 last_round::validate_last_round_coverage(messages, &kept)?;
1452 if config.enabled
1453 && compaction_pressure_reached(messages, system_prompt, config)
1454 && estimate_input_tokens_for_pressure(&kept, system_prompt)
1455 >= estimate_input_tokens_for_pressure(messages, system_prompt)
1456 {
1457 anyhow::bail!(
1458 "Making room did not shrink the context; the original conversation was preserved."
1459 );
1460 }
1461 let keep: CompactionKeep = inspect_compaction_keep(&kept);
1462 coverage.last_round_messages = keep.last_round_messages;
1463 coverage.last_round_tool_results = keep.last_round_tool_results;
1464 coverage.last_round_assistant = keep.last_round_assistant;
1465 if let (Some(session_id), Some(summary)) =
1466 (prepared.session_id.as_deref(), prompt.as_ref())
1467 {
1468 let text = summary_prompt_text(summary);
1469 let redacted = codewhale_config::persistence::redact_json_secrets(
1470 &serde_json::Value::String(text),
1471 );
1472 crate::artifacts::write_session_relative_immutable(
1473 session_id,
1474 &std::path::PathBuf::from("artifacts")
1475 .join(format!("context-transfer-{checkpoint_id}.md")),
1476 redacted.as_str().unwrap_or_default().as_bytes(),
1477 )?;
1478 }
1479 return Ok(CompactionResult {
1480 messages: kept,
1481 summary_prompt: prompt,
1482 retries_used: attempt.saturating_add(quality_retries),
1483 coverage,
1484 });
1485 }
1486 Err(e) => {
1487 // Only retry on transient errors
1488 if !is_transient_error(&e) {
1489 return Err(e);
1490 }
1491 last_error = Some(e);
1492 }
1493 }
1494 }
1495
1496 Err(last_error
1497 .unwrap_or_else(|| anyhow::anyhow!("Making room failed after {MAX_RETRIES} retries")))
1498 }
1499
1500 pub(crate) fn build_compaction_summary_block_text(summary: &str, anchors: &str) -> String {
1501 let summary = summary.trim();
1502 let summary = if summary.is_empty() {
1503 "(no summary available)"
1504 } else {
1505 summary
1506 };
1507 let mut text = format!("{SUMMARY_HEADER}\n\n{summary}");
1508 text.push_str(anchors);
1509 text.push_str("\n\n");
1510 text.push_str(SUMMARY_CLOSING);
1511 text
1512 }
1513
1514 /// Codex-parity replacement history: the most recent user-role messages,
1515 /// selected newest-first within a fixed token budget and restored to
1516 /// transcript order. Content boundaries carry runtime provenance and image
1517 /// turns, so structured messages are retained whole or dropped whole. Only a
1518 /// single text block can be truncated to fit the remaining budget.
1519 /// Result blocks answer a tool call that lives in an earlier message. A
1520 /// retained older turn has already lost that call to the summary, so a kept
1521 /// result block becomes an orphan providers reject outright (#6119).
1522 fn is_orphaned_result_block(block: &ContentBlock) -> bool {
1523 matches!(
1524 block,
1525 ContentBlock::ToolResult { .. }
1526 | ContentBlock::ToolSearchToolResult { .. }
1527 | ContentBlock::CodeExecutionToolResult { .. }
1528 )
1529 }
1530
1531 pub(crate) fn retained_user_messages(messages: &[Message], max_tokens: usize) -> Vec<Message> {
1532 let mut selected: Vec<Message> = Vec::new();
1533 let mut remaining = max_tokens;
1534 for msg in messages.iter().rev() {
1535 if remaining == 0 {
1536 break;
1537 }
1538 if msg.role != Role::User
1539 || crate::runtime_handoff::is_runtime_owned_user_message(msg)
1540 || (user_text_of(msg).is_none()
1541 && !msg
1542 .content
1543 .iter()
1544 .any(|block| matches!(block, ContentBlock::ImageUrl { .. })))
1545 {
1546 continue;
1547 }
1548 if is_compaction_checkpoint_message(msg) {
1549 continue;
1550 }
1551 let tokens: usize = msg
1552 .content
1553 .iter()
1554 .map(|block| match block {
1555 ContentBlock::Text { text, .. } => estimate_text_tokens_conservative(text),
1556 ContentBlock::ImageUrl { .. } => IMAGE_TOKEN_ESTIMATE,
1557 _ => 0,
1558 })
1559 .sum();
1560 let mut retained = msg.clone();
1561 if tokens <= remaining {
1562 // Keep the text and images; never a result block whose call was
1563 // summarized away with the region around it (#6119).
1564 retained
1565 .content
1566 .retain(|block| !is_orphaned_result_block(block));
1567 remaining -= tokens;
1568 } else {
1569 let [ContentBlock::Text { text, .. }] = retained.content.as_mut_slice() else {
1570 // Never flatten or partially retain an engine metadata block:
1571 // that would turn runtime-owned traffic into a user prompt.
1572 break;
1573 };
1574 *text = truncate_chars(text, remaining.saturating_mul(3)).to_string();
1575 remaining = 0;
1576 }
1577 selected.push(retained);
1578 }
1579 selected.reverse();
1580 selected
1581 }
1582
1583 /// User-pinned facts from `/anchor` (`.codewhale/anchors.md`). These are the
1584 /// user's own words, re-stated after the summary because the command promises
1585 /// they survive compaction.
1586 fn user_anchors_section(workspace: Option<&std::path::Path>) -> String {
1587 match pinned_anchors_text(workspace) {
1588 Some(contents) => format!("\n\nUser-pinned anchors (verbatim):\n{contents}"),
1589 None => String::new(),
1590 }
1591 }
1592
1593 #[cfg(test)]
1594 async fn compact_messages(
1595 client: &dyn ModelClient,
1596 messages: &[Message],
1597 config: &CompactionConfig,
1598 ) -> Result<(Vec<Message>, Option<SystemPrompt>, Vec<Message>)> {
1599 let mut quality_retries = 0;
1600 let mut invocation_usage = Usage::default();
1601 let (messages, summary_prompt, _coverage) = compact_messages_with_metadata(
1602 client,
1603 messages,
1604 config,
1605 None,
1606 None,
1607 None,
1608 None,
1609 &mut quality_retries,
1610 &mut invocation_usage,
1611 )
1612 .await?;
1613 Ok((messages, summary_prompt, Vec::new()))
1614 }
1615
1616 async fn compact_messages_with_metadata(
1617 client: &dyn ModelClient,
1618 messages: &[Message],
1619 config: &CompactionConfig,
1620 system_prompt: Option<&SystemPrompt>,
1621 tools: Option<&[Tool]>,
1622 reasoning_effort: Option<&str>,
1623 notice_sink: Option<&dyn CompactionNoticeSink>,
1624 quality_retries: &mut u32,
1625 invocation_usage: &mut Usage,
1626 ) -> Result<(Vec<Message>, Option<SystemPrompt>, CompactionCoverage)> {
1627 if messages.is_empty() {
1628 return Ok((Vec::new(), None, CompactionCoverage::default()));
1629 }
1630
1631 let summary = create_summary(
1632 client,
1633 messages,
1634 config,
1635 system_prompt,
1636 tools,
1637 reasoning_effort,
1638 notice_sink,
1639 quality_retries,
1640 invocation_usage,
1641 )
1642 .await?;
1643 let anchors = user_anchors_section(config.workspace.as_deref());
1644 let checkpoint_text = build_compaction_summary_block_text(&summary, &anchors);
1645 let summary_block = SystemBlock {
1646 block_type: "text".to_string(),
1647 text: checkpoint_text.clone(),
1648 cache_control: config.cache_summary.then(|| CacheControl {
1649 cache_type: "ephemeral".to_string(),
1650 }),
1651 };
1652
1653 let retained = last_round::build_replacement_history(
1654 messages,
1655 &checkpoint_text,
1656 pinned_anchors_text(config.workspace.as_deref()).as_deref(),
1657 config.retained_user_message_tokens,
1658 )?;
1659 let mut coverage = last_round::measure_coverage(
1660 messages,
1661 &retained,
1662 CompactionPath::Summary,
1663 pinned_anchors_text(config.workspace.as_deref())
1664 .map(|text| text.chars().count())
1665 .unwrap_or(0),
1666 );
1667 // Report the tuning actually in force so the receipt shows the operator
1668 // their knobs took effect (#5956).
1669 coverage.retained_user_message_tokens = config.retained_user_message_tokens;
1670 coverage.operator_instructions_applied =
1671 operator_instructions_section(config.summary_instructions.as_deref()).is_some();
1672 Ok((
1673 retained,
1674 Some(SystemPrompt::Blocks(vec![summary_block])),
1675 coverage,
1676 ))
1677 }
1678
1679 /// Delimiters around the operator's standing summarizer instructions. The
1680 /// summarizer sees a plain user message, so the section must announce itself:
1681 /// unfenced free text reads as more conversation to summarize.
1682 const OPERATOR_INSTRUCTIONS_HEADER: &str = "--- Additional instructions from the operator ---";
1683 const OPERATOR_INSTRUCTIONS_FOOTER: &str = "--- End of additional instructions ---";
1684
1685 /// Render `[compaction] summary_instructions` as a delimited prompt suffix.
1686 ///
1687 /// `None` (unset, or whitespace-only) produces no section at all, which is
1688 /// what keeps the default prompt byte-identical to the pre-#5956 constant.
1689 /// The cap is enforced here rather than at config load so the warning fires
1690 /// once per compaction pass instead of once per turn.
1691 fn operator_instructions_section(instructions: Option<&str>) -> Option<String> {
1692 let text = instructions
1693 .map(str::trim)
1694 .filter(|text| !text.is_empty())?;
1695 let max_chars = crate::config::COMPACTION_SUMMARY_INSTRUCTIONS_MAX_CHARS;
1696 let text = if text.chars().count() > max_chars {
1697 logging::warn(format!(
1698 "[compaction] summary_instructions is longer than {max_chars} characters; \
1699 the summarizer prompt suffix was truncated"
1700 ));
1701 truncate_chars(text, max_chars)
1702 } else {
1703 text
1704 };
1705 Some(format!(
1706 "\n\n{OPERATOR_INSTRUCTIONS_HEADER}\n{text}\n{OPERATOR_INSTRUCTIONS_FOOTER}"
1707 ))
1708 }
1709
1710 fn compact_prompt(focus: Option<&str>, instructions: Option<&str>) -> String {
1711 with_instructions_and_focus(
1712 format!("{} {COMPACTION_LANGUAGE_CONTRACT}", compact_prompt_body()),
1713 focus,
1714 instructions,
1715 )
1716 }
1717
1718 fn compact_quality_retry_prompt(focus: Option<&str>, instructions: Option<&str>) -> String {
1719 with_instructions_and_focus(
1720 format!(
1721 "The previous reply was empty or a placeholder, so there is still no handoff note. \
1722 Write it now, with real content from this session under each heading.\n\n{}\n\n\
1723 {HANDOFF_FOLD_IN_RULE}\n\n{HANDOFF_CLOSING_RULE} Do not refuse or return a placeholder. \
1724 {COMPACTION_LANGUAGE_CONTRACT}",
1725 handoff_sections_text()
1726 ),
1727 focus,
1728 instructions,
1729 )
1730 }
1731
1732 /// Standing operator instructions first, then the one-off `/compact <focus>`.
1733 fn with_instructions_and_focus(
1734 mut prompt: String,
1735 focus: Option<&str>,
1736 instructions: Option<&str>,
1737 ) -> String {
1738 if let Some(section) = operator_instructions_section(instructions) {
1739 prompt.push_str(&section);
1740 }
1741 if let Some(focus) = focus.map(str::trim).filter(|focus| !focus.is_empty()) {
1742 let _ = write!(prompt, "\n\n{HANDOFF_FOCUS_LINE} {focus}");
1743 }
1744 prompt
1745 }
1746
1747 fn validate_compaction_summary(summary: &str) -> Result<()> {
1748 let trimmed = summary.trim();
1749 if trimmed.is_empty() {
1750 anyhow::bail!("The summary for making room was unusable: no text was returned.");
1751 }
1752
1753 // Strip every non-word edge, not just ASCII punctuation. Providers can
1754 // return visually non-empty Unicode punctuation or emoji-only payloads;
1755 // neither is a usable continuation checkpoint. `is_alphanumeric` keeps
1756 // this language-neutral for CJK and other scripts without imposing a
1757 // prose-length heuristic.
1758 let normalized = trimmed
1759 .trim_matches(|ch: char| !ch.is_alphanumeric())
1760 .to_ascii_lowercase();
1761 if normalized.is_empty() {
1762 anyhow::bail!(
1763 "The summary for making room was unusable: only whitespace or punctuation was returned."
1764 );
1765 }
1766 if matches!(
1767 normalized.as_str(),
1768 "no summary available"
1769 | "summary unavailable"
1770 | "no summary"
1771 | "n/a"
1772 | "na"
1773 | "not available"
1774 | "i cannot provide a summary"
1775 | "i can't provide a summary"
1776 | "unable to provide a summary"
1777 ) {
1778 anyhow::bail!("The summary for making room was unusable: a placeholder was returned.");
1779 }
1780 Ok(())
1781 }
1782
1783 /// Drop the oldest history message before retrying an over-window summary
1784 /// request, plus any tool results the removal orphans (strict providers
1785 /// reject unpaired results). `messages` ends with the handoff instruction.
1786 ///
1787 /// The newest checkpoint is never dropped: it carries the previous handoff
1788 /// note, and without it the fold-in has nothing to carry forward, so user
1789 /// corrections and limits would vanish without notice. Returns `false`, having
1790 /// changed nothing, when no other history message can go and at least one
1791 /// would remain; the caller then fails the pass instead of summarizing
1792 /// without the note.
1793 fn drop_oldest_history_messages(messages: &mut Vec<Message>) -> bool {
1794 let instruction = messages.len().saturating_sub(1);
1795 let history = &messages[..instruction];
1796 let mut keep = history
1797 .iter()
1798 .rposition(is_wire_compaction_checkpoint_message);
1799 let Some(index) = (0..instruction).find(|&index| Some(index) != keep) else {
1800 return false;
1801 };
1802 if history.len() <= 1 {
1803 return false;
1804 }
1805 messages.remove(index);
1806 if keep.is_some_and(|keep_at| keep_at > index) {
1807 keep = keep.map(|keep_at| keep_at - 1);
1808 }
1809 while index + 1 < messages.len()
1810 && Some(index) != keep
1811 && messages[index].content.iter().any(is_orphaned_result_block)
1812 {
1813 messages.remove(index);
1814 if keep.is_some_and(|keep_at| keep_at > index) {
1815 keep = keep.map(|keep_at| keep_at - 1);
1816 }
1817 }
1818 true
1819 }
1820
1821 /// The summary request for one compaction pass: the parent turn's exact
1822 /// request inputs — model, system prompt, tools, reasoning tier, and stored
1823 /// history in order — plus one trailing user instruction. Only per-request
1824 /// controls differ (output cap, streaming, tool choice), so the provider's
1825 /// prefix cache covers the history the turn already paid for (#6540).
1826 ///
1827 /// Known limit: `tool_choice: "none"` keeps the summary from executing tools.
1828 /// Anthropic Messages documents a `tool_choice` change as invalidating the
1829 /// cached *message* blocks (system and tools stay cached), so on those routes
1830 /// the history is still re-read uncached. The Responses (Codex) builder does
1831 /// not send this field; its cache effect on Chat Completions routes has not
1832 /// been measured live.
1833 pub(crate) fn compaction_summary_request(
1834 history: Vec<Message>,
1835 config: &CompactionConfig,
1836 system_prompt: Option<&SystemPrompt>,
1837 tools: Option<&[Tool]>,
1838 reasoning_effort: Option<&str>,
1839 max_tokens: u32,
1840 ) -> MessageRequest {
1841 MessageRequest {
1842 model: config.model.clone(),
1843 messages: history,
1844 max_tokens,
1845 system: system_prompt.cloned(),
1846 tools: tools.map(<[Tool]>::to_vec),
1847 // Tool schemas stay for prefix parity; execution stays off.
1848 tool_choice: tools
1849 .filter(|tools| !tools.is_empty())
1850 .map(|_| serde_json::json!("none")),
1851 metadata: None,
1852 thinking: None,
1853 reasoning_effort: reasoning_effort.map(str::to_string),
1854 stream: Some(false),
1855 // Route parity with ordinary turns: turns send no sampling
1856 // params, so every provider's own normalization/defaults apply.
1857 // A hard-coded 0.3 leaked to the wire on routes that pass
1858 // temperature through (e.g. Kimi Code membership), where the
1859 // fixed-sampling contract rejects it and the whole compaction
1860 // pass fails.
1861 temperature: None,
1862 top_p: None,
1863 }
1864 }
1865
1866 async fn create_summary(
1867 client: &dyn ModelClient,
1868 messages: &[Message],
1869 config: &CompactionConfig,
1870 system_prompt: Option<&SystemPrompt>,
1871 tools: Option<&[Tool]>,
1872 reasoning_effort: Option<&str>,
1873 notice_sink: Option<&dyn CompactionNoticeSink>,
1874 quality_retries: &mut u32,
1875 invocation_usage: &mut Usage,
1876 ) -> Result<String> {
1877 // The summarization request IS the live conversation plus one final user
1878 // message asking for the handoff summary, so the provider's prefix cache
1879 // covers everything already sent this session.
1880 let mut request_messages = messages.to_vec();
1881 let stripped_images = crate::image_attach::strip_images_when_unsupported(
1882 &mut request_messages,
1883 config.image_input,
1884 &config.model,
1885 );
1886 if stripped_images > 0 {
1887 logging::warn(format!(
1888 "Compaction omitted {stripped_images} image block(s) unsupported by its route"
1889 ));
1890 }
1891 request_messages.push(Message {
1892 role: Role::User,
1893 content: vec![ContentBlock::Text {
1894 text: compact_prompt(
1895 config.focus.as_deref(),
1896 config.summary_instructions.as_deref(),
1897 ),
1898 cache_control: None,
1899 }],
1900 });
1901
1902 let mut quality_retry_used = false;
1903 // Request-size ladder: an HTTP 413 caps the request *body*, which the
1904 // token-side budget cannot predict (a flat per-image token estimate can
1905 // hide megabytes of base64). A refused summary call gets two byte-side
1906 // downgrades before it fails — re-encode the inline images smaller, then
1907 // replace them with text notes — each followed by exactly one retry.
1908 let mut size_ladder = RequestSizeLadder::Start;
1909 loop {
1910 // Codex compaction is a normal model generation over the existing
1911 // cached prefix. Do the same here: the resolved route decides how
1912 // much output the model may need instead of imposing a smaller,
1913 // compaction-only ceiling that can be consumed by hidden reasoning.
1914 let cost_route = client.effective_route_envelope(&config.model, chrono::Utc::now());
1915 let request = compaction_summary_request(
1916 request_messages.clone(),
1917 config,
1918 system_prompt,
1919 tools,
1920 reasoning_effort,
1921 client.effective_max_output_tokens(&cost_route.model),
1922 );
1923
1924 // Capture the session scope before awaiting so a late response cannot
1925 // accrue into a subsequently loaded/new session.
1926 let accounting_origin = notice_sink.and_then(CompactionNoticeSink::accounting_origin);
1927 let cost_scope = accounting_origin
1928 .as_ref()
1929 .map_or_else(crate::cost_status::scope_token, |origin| origin.0);
1930 let response = match client.create_message(request).await {
1931 Ok(response) => response,
1932 // A byte-side size rejection can also read like a length problem
1933 // (gateway HTML pages carry their own wording); the size ladder
1934 // owns it, not the drop-oldest ladder.
1935 Err(err) if !is_request_too_large_error(&err) && is_context_window_error(&err) => {
1936 if !drop_oldest_history_messages(&mut request_messages) {
1937 logging::warn(
1938 "Compaction summary input is over the context window with nothing \
1939 left to drop but the previous handoff note; the pass fails instead \
1940 of losing that note",
1941 );
1942 return Err(err);
1943 }
1944 logging::warn(format!(
1945 "Compaction summary input over the context window ({err}); \
1946 dropped the oldest history item and retrying"
1947 ));
1948 continue;
1949 }
1950 Err(err) if is_request_too_large_error(&err) => match size_ladder {
1951 RequestSizeLadder::Start => {
1952 // Decoding, resizing and re-encoding megabytes of inline
1953 // images is CPU-bound work; run it off the async worker so
1954 // the engine keeps servicing events while it happens.
1955 let mut outbound = std::mem::take(&mut request_messages);
1956 let joined = tokio::task::spawn_blocking(move || {
1957 let shrunk = crate::image_attach::shrink_images_for_request(&mut outbound);
1958 (outbound, shrunk)
1959 })
1960 .await;
1961 let (outbound, shrunk) = match joined {
1962 Ok(joined) => joined,
1963 Err(join_error) => {
1964 return Err(err.context(format!(
1965 "The summary request exceeded the provider's request-body limit and re-encoding its inline images failed: {join_error}"
1966 )));
1967 }
1968 };
1969 request_messages = outbound;
1970 if shrunk.images_seen == 0 {
1971 return Err(err.context(
1972 "The summary request exceeded the provider's request-body limit and the history carries no inline images to re-encode",
1973 ));
1974 }
1975 size_ladder = RequestSizeLadder::ImagesShrunk;
1976 if shrunk.images > 0 {
1977 let message = format!(
1978 "Making room exceeded the provider's request-body limit (HTTP 413); re-encoded {} inline image(s) smaller ({} to {}) and is retrying the summary.",
1979 shrunk.images,
1980 crate::image_attach::human_bytes(shrunk.bytes_before),
1981 crate::image_attach::human_bytes(shrunk.bytes_after),
1982 );
1983 logging::warn(&message);
1984 deliver_compaction_notice(notice_sink, message);
1985 continue;
1986 }
1987 // Images are present but every one already fits its share of
1988 // the byte budget, and the body was still refused: the cap
1989 // sits below the budget. A retry would send identical bytes,
1990 // so replace the images now instead of failing the pass.
1991 if let Err(replace_error) =
1992 replace_inline_images_for_retry(&mut request_messages, notice_sink)
1993 {
1994 return Err(err.context(format!(
1995 "The summary request exceeded the provider's request-body limit; {replace_error}"
1996 )));
1997 }
1998 size_ladder = RequestSizeLadder::ImagesReplaced;
1999 continue;
2000 }
2001 RequestSizeLadder::ImagesShrunk => {
2002 if let Err(replace_error) =
2003 replace_inline_images_for_retry(&mut request_messages, notice_sink)
2004 {
2005 return Err(err.context(format!(
2006 "The summary request still exceeded the provider's request-body limit and {replace_error}"
2007 )));
2008 }
2009 size_ladder = RequestSizeLadder::ImagesReplaced;
2010 continue;
2011 }
2012 RequestSizeLadder::ImagesReplaced => {
2013 return Err(err.context(
2014 "The summary request exceeded the provider's request-body limit even after re-encoding and then replacing every inline image",
2015 ));
2016 }
2017 },
2018 Err(err) => return Err(err),
2019 };
2020
2021 // Keep the caller's total before any validation or subsequent await.
2022 // A rejected summary or canceled retry still consumed these tokens.
2023 crate::core::turn::add_usage_to(invocation_usage, &response.usage);
2024
2025 let source_id = format!(
2026 "compaction:dispatch:{}:response:{}",
2027 cost_route
2028 .dispatched_at
2029 .timestamp_nanos_opt()
2030 .unwrap_or_default(),
2031 response.id
2032 );
2033 if let Some((scope, session, agent)) = accounting_origin.as_ref() {
2034 if crate::core::engine::turn_loop::usage_has_reported_data(&response.usage) {
2035 if let Some(owner) = config.runtime_cost_owner.as_deref() {
2036 crate::cost_status::report_effective_route_for_runtime(
2037 *scope,
2038 Some(owner),
2039 &source_id,
2040 &cost_route,
2041 &response.usage,
2042 );
2043 } else {
2044 crate::cost_status::report_effective_route_for_interactive_origin(
2045 *scope,
2046 session,
2047 agent,
2048 &source_id,
2049 &cost_route,
2050 &response.usage,
2051 );
2052 }
2053 } else if let Some(owner) = config.runtime_cost_owner.as_deref() {
2054 crate::cost_status::report_unreceipted_provider_success(
2055 *scope,
2056 Some(owner),
2057 &source_id,
2058 &cost_route,
2059 );
2060 } else {
2061 crate::cost_status::report_unreceipted_for_interactive_origin(
2062 *scope,
2063 session,
2064 agent,
2065 &source_id,
2066 &cost_route,
2067 );
2068 }
2069 } else {
2070 crate::cost_status::report_effective_route_for_runtime(
2071 cost_scope,
2072 config.runtime_cost_owner.as_deref(),
2073 &source_id,
2074 &cost_route,
2075 &response.usage,
2076 );
2077 }
2078 if let Some(sink) = notice_sink {
2079 sink.settled_usage(&source_id, &cost_route, &response.usage)
2080 .await;
2081 }
2082
2083 // Usage above is already billed; a provider-declared incomplete
2084 // summary must still fail rather than replace the session history
2085 // with a fragment.
2086 if codewhale_models::is_incomplete_stop_reason(response.stop_reason.as_deref()) {
2087 anyhow::bail!(
2088 "The summary for making room was incomplete: provider stop reason `{}`; the partial summary was not accepted.",
2089 codewhale_models::stop_reason_detail(response.stop_reason.as_deref())
2090 );
2091 }
2092 if response
2093 .content
2094 .iter()
2095 .any(|block| matches!(block, ContentBlock::ToolUse { .. }))
2096 {
2097 anyhow::bail!(
2098 "Making room returned a tool call instead of a summary; the original conversation was preserved."
2099 );
2100 }
2101
2102 let summary = response
2103 .content
2104 .iter()
2105 .filter_map(|block| match block {
2106 ContentBlock::Text { text, .. } => Some(text.clone()),
2107 _ => None,
2108 })
2109 .collect::<Vec<_>>()
2110 .join("\n");
2111
2112 if let Err(error) = validate_compaction_summary(&summary) {
2113 if quality_retry_used {
2114 return Err(error.context(
2115 "Compaction summary remained unusable after one conservative retry; \
2116 no replacement checkpoint was committed",
2117 ));
2118 }
2119
2120 quality_retry_used = true;
2121 *quality_retries = (*quality_retries).saturating_add(1);
2122 logging::warn(
2123 "Compaction provider returned an unusable successful response; retrying once with the conservative handoff prompt",
2124 );
2125 let Some(instruction) = request_messages.last_mut() else {
2126 return Err(error.context(
2127 "Compaction summary validation failed and the retry instruction was missing",
2128 ));
2129 };
2130 instruction.content = vec![ContentBlock::Text {
2131 text: compact_quality_retry_prompt(
2132 config.focus.as_deref(),
2133 config.summary_instructions.as_deref(),
2134 ),
2135 cache_control: None,
2136 }];
2137 continue;
2138 }
2139
2140 return Ok(summary);
2141 }
2142 }
2143
2144 /// How far the request-size ladder for one summary call has descended.
2145 ///
2146 /// The ladder exists because HTTP 413 rejects the request *body* by bytes,
2147 /// which the token-side context budget that governs compaction cannot see.
2148 /// Each rung is one retry: re-encode the inline images under a byte budget,
2149 /// then replace them with text notes.
2150 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
2151 enum RequestSizeLadder {
2152 /// No size downgrade applied yet.
2153 Start,
2154 /// Inline images were re-encoded smaller.
2155 ImagesShrunk,
2156 /// Inline images were replaced with text notes.
2157 ImagesReplaced,
2158 }
2159
2160 /// Whether the provider refused the request body for size (HTTP 413 and its
2161 /// common wordings). A smaller payload can succeed where this one did not,
2162 /// which is exactly what the request-size ladder trades on.
2163 ///
2164 /// Walks the whole error chain: the rejection may be stated by a fronting
2165 /// gateway (an HTML page from openresty reading "413 Request Entity Too
2166 /// Large") and then wrapped again by the client's own error text, so the
2167 /// top-level message alone is not reliable.
2168 fn is_request_too_large_error(e: &anyhow::Error) -> bool {
2169 e.chain().any(|cause| {
2170 let lower = cause.to_string().to_lowercase();
2171 // The client always renders the status code into the message
2172 // ("HTTP 413"), so that token is the primary signal; the phrase list
2173 // below catches gateways that state the condition without it.
2174 lower.contains("http 413")
2175 || lower.contains("payload too large")
2176 || lower.contains("request entity too large")
2177 || lower.contains("request body too large")
2178 || lower.contains("length limit exceeded")
2179 })
2180 }
2181
2182 /// Forward one user-visible progress sentence, when the host supplied a sink.
2183 fn deliver_compaction_notice(sink: Option<&dyn CompactionNoticeSink>, message: String) {
2184 if let Some(sink) = sink {
2185 sink.notice(message);
2186 }
2187 }
2188
2189 /// Replace every inline image for one retry and announce it.
2190 ///
2191 /// Shared by the two rungs that can reach the replace step: after a shrink
2192 /// pass that changed something, and directly when the images already fit the
2193 /// byte budget but the body was refused anyway.
2194 fn replace_inline_images_for_retry(
2195 messages: &mut [Message],
2196 notice_sink: Option<&dyn CompactionNoticeSink>,
2197 ) -> Result<usize> {
2198 let replaced = crate::image_attach::replace_images_with_placeholders(
2199 messages,
2200 "the summary request exceeded the provider's request-body limit (HTTP 413)",
2201 );
2202 if replaced == 0 {
2203 anyhow::bail!("no inline images were left to replace");
2204 }
2205 let message = format!(
2206 "Making room still exceeded the provider's request-body limit; replaced {replaced} inline image(s) with text notes for this summary pass and is retrying."
2207 );
2208 logging::warn(&message);
2209 deliver_compaction_notice(notice_sink, message);
2210 Ok(replaced)
2211 }
2212
2213 fn is_context_window_error(e: &anyhow::Error) -> bool {
2214 let text = e.to_string();
2215 if crate::error_taxonomy::classify_error_message(&text)
2216 != crate::error_taxonomy::ErrorCategory::InvalidInput
2217 {
2218 return false;
2219 }
2220
2221 let lower = text.to_lowercase();
2222 lower.contains("context")
2223 || lower.contains("token")
2224 || lower.contains("prompt is too long")
2225 || lower.contains("requested")
2226 || lower.contains("maximum")
2227 }
2228
2229 /// Collect text from a user message without treating tool-result payloads
2230 /// as new user instructions.
2231 fn user_text_of(msg: &Message) -> Option<String> {
2232 if msg.role != "user" {
2233 return None;
2234 }
2235 let text = msg
2236 .content
2237 .iter()
2238 .filter_map(|block| match block {
2239 ContentBlock::Text { text, .. } => Some(text.as_str()),
2240 _ => None,
2241 })
2242 .collect::<Vec<_>>()
2243 .join("\n");
2244 let text = text.trim();
2245 (!text.is_empty()).then(|| text.to_string())
2246 }
2247
2248 #[cfg(test)]
2249 #[path = "compaction/tests.rs"]
2250 mod quota_tests;
2251
2252 #[cfg(test)]
2253 mod tests {
2254 use codewhale_models::{ImageUrlContent, Message};
2255
2256 #[test]
2257 fn legacy_checkpoint_headings_are_recognised_but_quotes_are_not() {
2258 for text in [
2259 "## 📋 Conversation Summary (Auto-Generated)\n\nkey facts.",
2260 "## Pinned Facts (User Anchors)\n\n- keep tabs\n\n---\n\n## 📋 Conversation Summary (Auto-Generated)\n\nfacts",
2261 "Conversation Summary (Auto-Generated)\nold",
2262 "Another language model started to solve this problem\nold",
2263 ] {
2264 assert!(is_legacy_compaction_summary_text(text), "{text}");
2265 }
2266 for text in [
2267 "Explain: ## 📋 Conversation Summary (Auto-Generated)",
2268 "## Pinned Facts (User Anchors)\nsee Conversation Summary (Auto-Generated) above",
2269 "Codewhale handoff note\nnot structural",
2270 ] {
2271 assert!(!is_legacy_compaction_summary_text(text), "{text}");
2272 }
2273 }
2274
2275 #[test]
2276 fn restore_without_typed_checkpoint_preserves_user_marker_quotes() {
2277 // The current marker is recognised only structurally (header block +
2278 // provenance block), so only the legacy substring markers apply here.
2279 for marker in [
2280 LEGACY_V2_COMPACTION_SUMMARY_MARKER,
2281 LEGACY_COMPACTION_SUMMARY_MARKER,
2282 ] {
2283 let quoted = Message {
2284 role: Role::User,
2285 content: vec![ContentBlock::Text {
2286 text: format!("Explain this marker: {marker}"),
2287 cache_control: None,
2288 }],
2289 };
2290 let legacy = Message {
2291 role: Role::User,
2292 content: vec![ContentBlock::Text {
2293 text: format!("{marker}\nold summary"),
2294 cache_control: None,
2295 }],
2296 };
2297 let messages = vec![quoted.clone(), legacy.clone()];
2298 assert_eq!(
2299 restore_compaction_checkpoint(messages.clone(), None),
2300 messages
2301 );
2302
2303 let summary = SystemPrompt::Text(build_compaction_summary_block_text("Summary", ""));
2304 let restored =
2305 restore_compaction_checkpoint(vec![quoted.clone(), legacy], Some(&summary));
2306 assert_eq!(restored.len(), 2);
2307 assert_eq!(restored[0], quoted);
2308 assert!(is_wire_compaction_checkpoint_message(&restored[1]));
2309 }
2310 }
2311
2312 #[test]
2313 fn restore_replaces_duplicate_generated_checkpoints_without_deleting_user_quote() {
2314 let summary = SystemPrompt::Text(build_compaction_summary_block_text("Summary", ""));
2315 let generated = compaction_checkpoint_message(&summary);
2316 let user_quote = Message {
2317 role: Role::User,
2318 content: vec![ContentBlock::Text {
2319 text: summary_prompt_text(&summary),
2320 cache_control: None,
2321 }],
2322 };
2323 assert!(!is_wire_compaction_checkpoint_message(&user_quote));
2324 let restored = restore_compaction_checkpoint(
2325 vec![generated.clone(), user_quote.clone(), generated],
2326 Some(&summary),
2327 );
2328 assert_eq!(restored.len(), 2);
2329 assert!(is_wire_compaction_checkpoint_message(&restored[0]));
2330 assert_eq!(restored[1], user_quote);
2331
2332 // With a saved summary, legacy header-prefixed copies are replaced.
2333 let legacy = Message {
2334 role: Role::User,
2335 content: vec![ContentBlock::Text {
2336 text: format!("{LEGACY_V2_COMPACTION_SUMMARY_MARKER}\nold summary"),
2337 cache_control: None,
2338 }],
2339 };
2340 let legacy_restored =
2341 restore_compaction_checkpoint(vec![legacy.clone(), legacy], Some(&summary));
2342 assert_eq!(legacy_restored.len(), 1);
2343 assert!(is_wire_compaction_checkpoint_message(&legacy_restored[0]));
2344 }
2345
2346 #[test]
2347 fn inline_image_estimates_nonzero_tokens() {
2348 let msg = Message {
2349 role: Role::User,
2350 content: vec![ContentBlock::ImageUrl {
2351 image_url: ImageUrlContent {
2352 url: "data:image/png;base64,AAAA".to_string(),
2353 },
2354 }],
2355 };
2356 assert!(
2357 estimate_tokens_for_message(&msg, false) >= IMAGE_TOKEN_ESTIMATE,
2358 "an inline image must not estimate to 0 tokens"
2359 );
2360 }
2361
2362 use super::*;
2363 use serde_json::json;
2364
2365 fn msg(role: &str, text: &str) -> Message {
2366 Message {
2367 role: Role::from(role),
2368 content: vec![ContentBlock::Text {
2369 text: text.to_string(),
2370 cache_control: None,
2371 }],
2372 }
2373 }
2374
2375 fn prepared(config: &CompactionConfig) -> PreparedCompactionEnvelope {
2376 PreparedCompactionEnvelope::new(config.clone())
2377 }
2378
2379 fn tool_use(id: &str, name: &str, input: serde_json::Value) -> Message {
2380 Message {
2381 role: Role::Assistant,
2382 content: vec![ContentBlock::ToolUse {
2383 execution_id: None,
2384 id: id.to_string(),
2385 name: name.to_string(),
2386 input,
2387 caller: None,
2388 thought_signature: None,
2389 }],
2390 }
2391 }
2392
2393 fn tool_result(id: &str, content: &str) -> Message {
2394 Message {
2395 role: Role::User,
2396 content: vec![ContentBlock::ToolResult {
2397 execution_id: None,
2398 tool_use_id: id.to_string(),
2399 content: content.to_string(),
2400 is_error: None,
2401 content_blocks: None,
2402 }],
2403 }
2404 }
2405
2406 #[test]
2407 fn truncate_chars_respects_unicode_boundaries() {
2408 let text = "abc😀é";
2409 assert_eq!(truncate_chars(text, 0), "");
2410 assert_eq!(truncate_chars(text, 1), "a");
2411 assert_eq!(truncate_chars(text, 3), "abc");
2412 assert_eq!(truncate_chars(text, 4), "abc😀");
2413 assert_eq!(truncate_chars(text, 5), "abc😀é");
2414 }
2415
2416 #[test]
2417 fn prune_tool_results_summarizes_old_verbose_outputs() {
2418 let verbose = "x".repeat(SUMMARY_TOOL_RESULT_SNIPPET_CHARS + 80);
2419 let mut messages = vec![
2420 tool_use("call-1", "read_file", json!({"path": "Cargo.toml"})),
2421 tool_result("call-1", &verbose),
2422 msg("user", "recent question"),
2423 msg("assistant", "recent answer"),
2424 ];
2425
2426 let saved = prune_tool_results(&mut messages, 2);
2427
2428 assert!(saved > 0);
2429 let ContentBlock::ToolResult { content, .. } = &messages[1].content[0] else {
2430 panic!("expected tool result");
2431 };
2432 assert!(content.contains("[read_file] tool result pruned"));
2433 assert!(content.contains("Cargo.toml"));
2434 assert!(content.len() < verbose.len());
2435 }
2436
2437 #[test]
2438 fn prune_tool_results_preserves_protected_tail() {
2439 let verbose = "x".repeat(SUMMARY_TOOL_RESULT_SNIPPET_CHARS + 80);
2440 let mut messages = vec![
2441 msg("user", "older context"),
2442 tool_use("call-1", "read_file", json!({"path": "Cargo.toml"})),
2443 tool_result("call-1", &verbose),
2444 ];
2445
2446 let saved = prune_tool_results(&mut messages, 2);
2447
2448 assert_eq!(saved, 0);
2449 let ContentBlock::ToolResult { content, .. } = &messages[2].content[0] else {
2450 panic!("expected tool result");
2451 };
2452 assert_eq!(content, &verbose);
2453 }
2454
2455 #[test]
2456 fn prune_tool_results_preserves_prefix_bytes_when_reverse_prune_is_enough() {
2457 let older_verbose = "old ".repeat(SUMMARY_TOOL_RESULT_SNIPPET_CHARS + 40);
2458 let newer_verbose = "new ".repeat(SUMMARY_TOOL_RESULT_SNIPPET_CHARS + 40);
2459 let mut messages = vec![
2460 tool_use("call-old", "read_file", json!({"path": "old.txt"})),
2461 tool_result("call-old", &older_verbose),
2462 tool_use("call-new", "read_file", json!({"path": "new.txt"})),
2463 tool_result("call-new", &newer_verbose),
2464 msg("user", "protected tail"),
2465 ];
2466 let original = messages.clone();
2467
2468 // Simulate the caller clearing its token budget after one suffix prune.
2469 let saved = prune_tool_results_until(&mut messages, 1, |_, saved| saved > 0);
2470
2471 assert!(saved > 0);
2472 assert_eq!(&messages[..3], &original[..3]);
2473 assert_eq!(&messages[4..], &original[4..]);
2474 let ContentBlock::ToolResult { content, .. } = &messages[3].content[0] else {
2475 panic!("expected pruned tool result");
2476 };
2477 assert!(content.contains("[read_file] tool result pruned"));
2478 assert!(content.contains("new.txt"));
2479 assert!(content.len() < newer_verbose.len());
2480 }
2481
2482 #[test]
2483 fn prune_tool_results_stops_after_newest_duplicate_prune() {
2484 let oldest = "oldest ".repeat(80);
2485 let middle = "middle ".repeat(80);
2486 let latest = "latest ".repeat(80);
2487 let mut messages = vec![
2488 tool_use("call-1", "read_file", json!({"path": "Cargo.toml"})),
2489 tool_result("call-1", &oldest),
2490 tool_use("call-2", "read_file", json!({"path": "Cargo.toml"})),
2491 tool_result("call-2", &middle),
2492 tool_use("call-3", "read_file", json!({"path": "Cargo.toml"})),
2493 tool_result("call-3", &latest),
2494 msg("user", "protected tail"),
2495 ];
2496 let original = messages.clone();
2497
2498 let saved = prune_tool_results_until(&mut messages, 1, |_, saved| saved > 0);
2499
2500 assert!(saved > 0);
2501 assert_eq!(&messages[..3], &original[..3]);
2502 assert_eq!(&messages[4..], &original[4..]);
2503 let ContentBlock::ToolResult { content, .. } = &messages[3].content[0] else {
2504 panic!("expected middle duplicate to be pruned");
2505 };
2506 assert!(content.contains("[read_file] tool result pruned"));
2507 }
2508
2509 #[test]
2510 fn prune_tool_results_dedupes_identical_reads_but_keeps_latest_full_body() {
2511 let first = "first ".repeat(80);
2512 let second = "second ".repeat(80);
2513 let mut messages = vec![
2514 tool_use("call-1", "read_file", json!({"path": "Cargo.toml"})),
2515 tool_result("call-1", &first),
2516 tool_use("call-2", "read_file", json!({"path": "Cargo.toml"})),
2517 tool_result("call-2", &second),
2518 msg("user", "tail"),
2519 ];
2520
2521 let saved = prune_tool_results(&mut messages, 1);
2522
2523 assert!(saved > 0);
2524 let ContentBlock::ToolResult { content: older, .. } = &messages[1].content[0] else {
2525 panic!("expected older tool result");
2526 };
2527 assert!(older.contains("tool result pruned"));
2528 let ContentBlock::ToolResult {
2529 content: latest, ..
2530 } = &messages[3].content[0]
2531 else {
2532 panic!("expected latest tool result");
2533 };
2534 assert_eq!(latest, &second);
2535 }
2536
2537 #[test]
2538 fn context_window_errors_are_detected_for_summary_fallback() {
2539 for msg in [
2540 "HTTP 400 Bad Request: maximum context length is 1000000 tokens",
2541 "invalid_request_error: prompt is too long for the current model",
2542 "You requested 1000001 tokens but the maximum is 1000000",
2543 "request exceeds context window",
2544 ] {
2545 assert!(
2546 is_context_window_error(&anyhow::anyhow!(msg)),
2547 "expected context-window detection for `{msg}`",
2548 );
2549 }
2550
2551 assert!(!is_context_window_error(&anyhow::anyhow!(
2552 "Invalid request: missing required field"
2553 )));
2554 assert!(!is_context_window_error(&anyhow::anyhow!(
2555 "503 Service Unavailable"
2556 )));
2557 }
2558
2559 #[test]
2560 fn tool_args_preview_redacts_sensitive_first_without_dropping_siblings() {
2561 let input: serde_json::Value = serde_json::from_str(
2562 r#"{"api_key":"sk-tool-secret-value","command":"cargo test -p auth"}"#,
2563 )
2564 .unwrap();
2565
2566 let preview: serde_json::Value = serde_json::from_str(&tool_args_preview(&input)).unwrap();
2567
2568 assert_eq!(preview["api_key"], codewhale_config::persistence::REDACTED);
2569 assert_eq!(preview["command"], "cargo test -p auth");
2570 }
2571
2572 #[test]
2573 fn tool_args_preview_redacts_sensitive_later_without_touching_earlier_fields() {
2574 let input: serde_json::Value =
2575 serde_json::from_str(r#"{"command":"cargo test","api_key":"plain-secret-value"}"#)
2576 .unwrap();
2577
2578 let preview: serde_json::Value = serde_json::from_str(&tool_args_preview(&input)).unwrap();
2579
2580 assert_eq!(preview["command"], "cargo test");
2581 assert_eq!(preview["api_key"], codewhale_config::persistence::REDACTED);
2582 }
2583
2584 #[test]
2585 fn tool_args_preview_redacts_nested_sensitive_values_recursively() {
2586 let input: serde_json::Value = serde_json::from_str(
2587 r#"{"meta":{"token":"nested-secret","keep":"yes"},"steps":[{"password":"pw","name":"a"}]}"#,
2588 )
2589 .unwrap();
2590
2591 let preview: serde_json::Value = serde_json::from_str(&tool_args_preview(&input)).unwrap();
2592
2593 assert_eq!(
2594 preview["meta"]["token"],
2595 codewhale_config::persistence::REDACTED
2596 );
2597 assert_eq!(preview["meta"]["keep"], "yes");
2598 assert_eq!(
2599 preview["steps"][0]["password"],
2600 codewhale_config::persistence::REDACTED
2601 );
2602 assert_eq!(preview["steps"][0]["name"], "a");
2603 }
2604
2605 #[test]
2606 fn tool_args_preview_redacts_complete_multi_word_secret_value() {
2607 let input: serde_json::Value =
2608 serde_json::from_str(r#"{"command":"run this","password":"hunter two words"}"#)
2609 .unwrap();
2610
2611 let serialized = tool_args_preview(&input);
2612 let preview: serde_json::Value = serde_json::from_str(&serialized).unwrap();
2613
2614 assert_eq!(preview["command"], "run this");
2615 assert_eq!(preview["password"], codewhale_config::persistence::REDACTED);
2616 assert!(!serialized.contains("hunter"));
2617 assert!(!serialized.contains("two words"));
2618 }
2619
2620 struct FixedSummaryClient {
2621 request: std::sync::Mutex<Option<MessageRequest>>,
2622 provider: &'static str,
2623 model: &'static str,
2624 }
2625
2626 impl Default for FixedSummaryClient {
2627 fn default() -> Self {
2628 Self {
2629 request: std::sync::Mutex::new(None),
2630 provider: "test",
2631 model: "test-model",
2632 }
2633 }
2634 }
2635
2636 impl FixedSummaryClient {
2637 fn for_route(provider: &'static str, model: &'static str) -> Self {
2638 Self {
2639 request: std::sync::Mutex::new(None),
2640 provider,
2641 model,
2642 }
2643 }
2644 }
2645
2646 const FIXED_SUMMARY: &str = "1. Primary request and intent — migrate the session store. \
2647 2. Key technical concepts — sqlite. 7. Pending tasks — finish the fixed clock. \
2648 8. Current work — rerunning the session tests.";
2649
2650 /// A real PNG whose bytes are worth shrinking. Noise defeats compression,
2651 /// which is the point: the shrink ladder must actually re-encode.
2652 fn noisy_png_bytes(width: u32, height: u32) -> Vec<u8> {
2653 use image::ImageEncoder as _;
2654 let mut pixels = image::RgbImage::new(width, height);
2655 for (x, y, pixel) in pixels.enumerate_pixels_mut() {
2656 *pixel = image::Rgb([
2657 (x.wrapping_mul(31) ^ y.wrapping_mul(17)) as u8,
2658 (x.wrapping_mul(7) ^ y.wrapping_mul(29)) as u8,
2659 (x.wrapping_add(y).wrapping_mul(13)) as u8,
2660 ]);
2661 }
2662 let mut bytes = Vec::new();
2663 image::codecs::png::PngEncoder::new_with_quality(
2664 &mut bytes,
2665 image::codecs::png::CompressionType::Fast,
2666 image::codecs::png::FilterType::NoFilter,
2667 )
2668 .write_image(
2669 pixels.as_raw(),
2670 width,
2671 height,
2672 image::ExtendedColorType::Rgb8,
2673 )
2674 .expect("encode fixture png");
2675 bytes
2676 }
2677
2678 fn png_data_url(width: u32, height: u32) -> String {
2679 use base64::Engine as _;
2680 format!(
2681 "data:image/png;base64,{}",
2682 base64::engine::general_purpose::STANDARD.encode(noisy_png_bytes(width, height))
2683 )
2684 }
2685
2686 fn user_image_message(data_url: &str) -> Message {
2687 Message {
2688 role: Role::User,
2689 content: vec![
2690 ContentBlock::Text {
2691 text: "look at this screenshot".to_string(),
2692 cache_control: None,
2693 },
2694 ContentBlock::ImageUrl {
2695 image_url: ImageUrlContent {
2696 url: data_url.to_string(),
2697 },
2698 },
2699 ],
2700 }
2701 }
2702
2703 /// Total inline-image URL bytes a request carries, the byte side an HTTP
2704 /// 413 boundary actually measures.
2705 fn inline_image_bytes(request: &MessageRequest) -> usize {
2706 request
2707 .messages
2708 .iter()
2709 .flat_map(|message| message.content.iter())
2710 .map(|block| match block {
2711 ContentBlock::ImageUrl { image_url } => image_url.url.len(),
2712 ContentBlock::ToolResult { content_blocks, .. } => {
2713 content_blocks.as_ref().map_or(0, |blocks| {
2714 blocks.iter().fold(0, |sum, block| {
2715 let nested = block
2716 .get("data")
2717 .and_then(serde_json::Value::as_str)
2718 .map_or(0, str::len);
2719 sum + nested
2720 })
2721 })
2722 }
2723 _ => 0,
2724 })
2725 .sum()
2726 }
2727
2728 fn summary_content() -> Vec<ContentBlock> {
2729 vec![ContentBlock::Text {
2730 text: "1. Primary request: keep working on the session store. \
2731 7. Pending: rerun the tests."
2732 .to_string(),
2733 cache_control: None,
2734 }]
2735 }
2736
2737 fn request_body_413() -> anyhow::Error {
2738 anyhow::anyhow!(
2739 "LLM error: HTTP 413: Failed to buffer the request body: length limit exceeded"
2740 )
2741 }
2742
2743 /// The same boundary stated by a fronting gateway (openresty) instead of
2744 /// the API's own body reader. Reported on 2026-09-26; the ladder must
2745 /// treat it as the same byte-side rejection.
2746 fn request_body_413_html_gateway() -> anyhow::Error {
2747 anyhow::anyhow!(
2748 "LLM error: HTTP 413: DeepSeek API returned an HTML error page (HTTP 413): \
2749 413 Request Entity Too Large 413 Request Entity Too Large openresty"
2750 )
2751 }
2752
2753 #[derive(Debug, Default)]
2754 struct RecordingNoticeSink {
2755 messages: std::sync::Mutex<Vec<String>>,
2756 }
2757
2758 impl RecordingNoticeSink {
2759 fn messages(&self) -> Vec<String> {
2760 self.messages.lock().expect("notice sink").clone()
2761 }
2762 }
2763
2764 impl CompactionNoticeSink for RecordingNoticeSink {
2765 fn notice(&self, message: String) {
2766 self.messages.lock().expect("notice sink").push(message);
2767 }
2768 }
2769
2770 /// A 900x900 noise PNG exceeds the 2 MiB inline-image budget, so the
2771 /// ladder must actually rewrite it.
2772 fn body_413_retry_envelope(
2773 sink: std::sync::Arc<RecordingNoticeSink>,
2774 ) -> PreparedCompactionEnvelope {
2775 let mut envelope = prepared(&CompactionConfig {
2776 enabled: false,
2777 ..CompactionConfig::default()
2778 });
2779 envelope.notice_sink = Some(sink);
2780 envelope
2781 }
2782
2783 #[tokio::test]
2784 async fn request_body_413_reencodes_images_smaller_and_retries() {
2785 let _environment = crate::test_support::lock_test_env();
2786 let messages = vec![
2787 msg("user", "context before the screenshots"),
2788 user_image_message(&png_data_url(900, 900)),
2789 ];
2790 let client = ScriptedSummaryClient::with_outcomes(vec![
2791 Err(request_body_413()),
2792 Ok(summary_content()),
2793 ]);
2794 let sink = std::sync::Arc::new(RecordingNoticeSink::default());
2795 let envelope = body_413_retry_envelope(sink.clone());
2796 let mut usage = Usage::default();
2797
2798 let result = compact_messages_safe(&client, &messages, None, &envelope, &mut usage).await;
2799 assert!(result.is_ok(), "the retry after shrinking must succeed");
2800
2801 let requests = client.requests.lock().expect("requests").clone();
2802 assert_eq!(requests.len(), 2, "one rejection, one retry");
2803 let first = inline_image_bytes(&requests[0]);
2804 let second = inline_image_bytes(&requests[1]);
2805 assert!(first > 0, "the fixture must carry the image");
2806 assert!(
2807 second > 0 && second < first,
2808 "the retry must carry smaller image bytes ({second} < {first})"
2809 );
2810 let notices = sink.messages();
2811 assert_eq!(notices.len(), 1, "one notice per downgrade: {notices:?}");
2812 assert!(
2813 notices[0].contains("re-encoded") && notices[0].contains("413"),
2814 "the notice must name the downgrade: {notices:?}"
2815 );
2816 }
2817
2818 #[tokio::test]
2819 async fn request_body_413_after_shrinking_replaces_images_with_notes() {
2820 let _environment = crate::test_support::lock_test_env();
2821 let messages = vec![
2822 msg("user", "context before the screenshots"),
2823 user_image_message(&png_data_url(900, 900)),
2824 ];
2825 let client = ScriptedSummaryClient::with_outcomes(vec![
2826 Err(request_body_413()),
2827 Err(request_body_413()),
2828 Ok(summary_content()),
2829 ]);
2830 let sink = std::sync::Arc::new(RecordingNoticeSink::default());
2831 let envelope = body_413_retry_envelope(sink.clone());
2832 let mut usage = Usage::default();
2833
2834 let result = compact_messages_safe(&client, &messages, None, &envelope, &mut usage).await;
2835 assert!(
2836 result.is_ok(),
2837 "the retry after replacing images must succeed"
2838 );
2839
2840 let requests = client.requests.lock().expect("requests").clone();
2841 assert_eq!(requests.len(), 3, "reject, shrink retry, replace retry");
2842 let first = inline_image_bytes(&requests[0]);
2843 let second = inline_image_bytes(&requests[1]);
2844 let third = inline_image_bytes(&requests[2]);
2845 assert!(second > 0 && second < first);
2846 assert_eq!(third, 0, "the last rung carries no image bytes");
2847 let note_present = requests[2].messages.iter().any(|message| {
2848 message.content.iter().any(|block| match block {
2849 ContentBlock::Text { text, .. } => text.contains("omitted from this summary pass"),
2850 _ => false,
2851 })
2852 });
2853 assert!(note_present, "the summarizer is told what was there");
2854 let notices = sink.messages();
2855 assert_eq!(notices.len(), 2, "one notice per rung: {notices:?}");
2856 assert!(notices[1].contains("replaced"), "{notices:?}");
2857 }
2858
2859 #[tokio::test]
2860 async fn request_body_413_from_a_gateway_html_page_enters_the_same_ladder() {
2861 let _environment = crate::test_support::lock_test_env();
2862 let messages = vec![
2863 msg("user", "context before the screenshots"),
2864 user_image_message(&png_data_url(900, 900)),
2865 ];
2866 let client = ScriptedSummaryClient::with_outcomes(vec![
2867 Err(request_body_413_html_gateway()),
2868 Ok(summary_content()),
2869 ]);
2870 let sink = std::sync::Arc::new(RecordingNoticeSink::default());
2871 let envelope = body_413_retry_envelope(sink.clone());
2872 let mut usage = Usage::default();
2873
2874 let result = compact_messages_safe(&client, &messages, None, &envelope, &mut usage).await;
2875 assert!(
2876 result.is_ok(),
2877 "the gateway-page rejection must enter the ladder"
2878 );
2879 let requests = client.requests.lock().expect("requests").clone();
2880 assert_eq!(requests.len(), 2, "one rejection, one retry");
2881 assert!(
2882 inline_image_bytes(&requests[1]) < inline_image_bytes(&requests[0]),
2883 "the retry must carry smaller image bytes"
2884 );
2885 let notices = sink.messages();
2886 assert_eq!(
2887 notices.len(),
2888 1,
2889 "the user hears about the downgrade: {notices:?}"
2890 );
2891 assert!(notices[0].contains("413"), "{notices:?}");
2892 }
2893
2894 #[tokio::test]
2895 async fn request_body_413_with_in_budget_images_skips_the_noop_retry_and_replaces() {
2896 // The images fit their share of the 2 MiB budget, yet the endpoint
2897 // refused the body: the cap sits below the budget. The ladder must not
2898 // report "no images" (which would fail the pass outright) and must not
2899 // resend identical bytes — it goes straight to the replace rung.
2900 let _environment = crate::test_support::lock_test_env();
2901 let messages = vec![
2902 msg("user", "context before the screenshot"),
2903 user_image_message(&png_data_url(64, 64)),
2904 ];
2905 let client = ScriptedSummaryClient::with_outcomes(vec![
2906 Err(request_body_413()),
2907 Ok(summary_content()),
2908 ]);
2909 let sink = std::sync::Arc::new(RecordingNoticeSink::default());
2910 let envelope = body_413_retry_envelope(sink.clone());
2911 let mut usage = Usage::default();
2912
2913 let result = compact_messages_safe(&client, &messages, None, &envelope, &mut usage).await;
2914 assert!(
2915 result.is_ok(),
2916 "in-budget images must still let the ladder finish"
2917 );
2918 let requests = client.requests.lock().expect("requests").clone();
2919 assert_eq!(requests.len(), 2, "replace directly, no identical retry");
2920 assert!(
2921 inline_image_bytes(&requests[0]) > 0,
2922 "the fixture carried the image"
2923 );
2924 assert_eq!(
2925 inline_image_bytes(&requests[1]),
2926 0,
2927 "the retry carries notes, not bytes"
2928 );
2929 let notices = sink.messages();
2930 assert_eq!(
2931 notices.len(),
2932 1,
2933 "one notice for the replace rung: {notices:?}"
2934 );
2935 assert!(notices[0].contains("replaced"), "{notices:?}");
2936 }
2937
2938 #[tokio::test]
2939 async fn request_body_413_without_inline_images_fails_with_context() {
2940 let _environment = crate::test_support::lock_test_env();
2941 let messages = vec![msg("user", "no images anywhere in this history")];
2942 let client = ScriptedSummaryClient::with_outcomes(vec![Err(request_body_413())]);
2943 let sink = std::sync::Arc::new(RecordingNoticeSink::default());
2944 let envelope = body_413_retry_envelope(sink.clone());
2945 let mut usage = Usage::default();
2946
2947 let error = compact_messages_safe(&client, &messages, None, &envelope, &mut usage)
2948 .await
2949 .expect_err("nothing left to shrink, the pass must fail");
2950 let text = format!("{error:#}");
2951 assert!(
2952 text.contains("request-body limit"),
2953 "the failure must name the boundary: {text}"
2954 );
2955 let requests = client.requests.lock().expect("requests").clone();
2956 assert_eq!(requests.len(), 1, "no pointless identical retry");
2957 assert!(sink.messages().is_empty(), "nothing was downgraded");
2958 }
2959
2960 #[tokio::test]
2961 async fn request_body_413_after_replacements_fails_with_the_full_ladder() {
2962 let _environment = crate::test_support::lock_test_env();
2963 let messages = vec![
2964 msg("user", "context before the screenshots"),
2965 user_image_message(&png_data_url(900, 900)),
2966 ];
2967 let client = ScriptedSummaryClient::with_outcomes(vec![
2968 Err(request_body_413()),
2969 Err(request_body_413()),
2970 Err(request_body_413()),
2971 ]);
2972 let sink = std::sync::Arc::new(RecordingNoticeSink::default());
2973 let envelope = body_413_retry_envelope(sink.clone());
2974 let mut usage = Usage::default();
2975
2976 let error = compact_messages_safe(&client, &messages, None, &envelope, &mut usage)
2977 .await
2978 .expect_err("every rung refused, the pass must fail");
2979 let text = format!("{error:#}");
2980 assert!(
2981 text.contains("even after re-encoding and then replacing"),
2982 "the failure must report the full ladder: {text}"
2983 );
2984 let requests = client.requests.lock().expect("requests").clone();
2985 assert_eq!(requests.len(), 3, "two retries, then stop");
2986 assert_eq!(sink.messages().len(), 2, "both rungs announced");
2987 }
2988
2989 struct ScriptedSummaryClient {
2990 responses: std::sync::Mutex<std::collections::VecDeque<anyhow::Result<Vec<ContentBlock>>>>,
2991 requests: std::sync::Mutex<Vec<MessageRequest>>,
2992 retry_started: Option<std::sync::Arc<tokio::sync::Notify>>,
2993 }
2994
2995 impl ScriptedSummaryClient {
2996 fn new(responses: Vec<Vec<ContentBlock>>) -> Self {
2997 Self::with_outcomes(responses.into_iter().map(Ok).collect())
2998 }
2999
3000 fn with_outcomes(responses: Vec<anyhow::Result<Vec<ContentBlock>>>) -> Self {
3001 Self {
3002 responses: std::sync::Mutex::new(responses.into()),
3003 requests: std::sync::Mutex::new(Vec::new()),
3004 retry_started: None,
3005 }
3006 }
3007 }
3008
3009 #[async_trait::async_trait]
3010 impl crate::core::model_client::ModelClient for ScriptedSummaryClient {
3011 fn provider_name(&self) -> &str {
3012 "test"
3013 }
3014
3015 fn model(&self) -> &str {
3016 "test-model"
3017 }
3018
3019 async fn create_message(
3020 &self,
3021 request: MessageRequest,
3022 ) -> anyhow::Result<codewhale_models::MessageResponse> {
3023 self.requests
3024 .lock()
3025 .expect("capture scripted summary request")
3026 .push(request);
3027 let outcome = self
3028 .responses
3029 .lock()
3030 .expect("read scripted summary response")
3031 .pop_front();
3032 if outcome.is_none()
3033 && let Some(retry_started) = &self.retry_started
3034 {
3035 retry_started.notify_one();
3036 return std::future::pending().await;
3037 }
3038 let content = outcome
3039 .ok_or_else(|| anyhow::anyhow!("scripted summary responses exhausted"))??;
3040 Ok(codewhale_models::MessageResponse {
3041 id: "summary-scripted".to_string(),
3042 r#type: "message".to_string(),
3043 role: "assistant".to_string(),
3044 content,
3045 model: "test-model".to_string(),
3046 stop_reason: None,
3047 stop_sequence: None,
3048 container: None,
3049 usage: Usage {
3050 input_tokens: 17,
3051 output_tokens: 3,
3052 prompt_cache_hit_tokens: Some(5),
3053 reasoning_tokens: Some(2),
3054 ..Usage::default()
3055 },
3056 })
3057 }
3058
3059 async fn create_message_stream(
3060 &self,
3061 _request: MessageRequest,
3062 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
3063 anyhow::bail!("streaming is unused by compaction")
3064 }
3065
3066 async fn health_check(&self) -> anyhow::Result<bool> {
3067 Ok(true)
3068 }
3069 }
3070
3071 #[async_trait::async_trait]
3072 impl crate::core::model_client::ModelClient for FixedSummaryClient {
3073 fn provider_name(&self) -> &str {
3074 self.provider
3075 }
3076
3077 fn model(&self) -> &str {
3078 self.model
3079 }
3080
3081 async fn create_message(
3082 &self,
3083 request: MessageRequest,
3084 ) -> anyhow::Result<codewhale_models::MessageResponse> {
3085 *self.request.lock().expect("capture summary request") = Some(request);
3086 Ok(codewhale_models::MessageResponse {
3087 id: "summary-fixture".to_string(),
3088 r#type: "message".to_string(),
3089 role: "assistant".to_string(),
3090 content: vec![ContentBlock::Text {
3091 text: FIXED_SUMMARY.to_string(),
3092 cache_control: None,
3093 }],
3094 model: self.model.to_string(),
3095 stop_reason: None,
3096 stop_sequence: None,
3097 container: None,
3098 usage: codewhale_models::Usage::default(),
3099 })
3100 }
3101
3102 async fn create_message_stream(
3103 &self,
3104 _request: MessageRequest,
3105 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
3106 anyhow::bail!("streaming is unused by compaction")
3107 }
3108
3109 async fn health_check(&self) -> anyhow::Result<bool> {
3110 Ok(true)
3111 }
3112 }
3113
3114 #[tokio::test]
3115 async fn compaction_persists_original_and_model_handoff_before_returning_replacement() {
3116 let _environment = crate::test_support::lock_test_env();
3117 let root = tempfile::tempdir().unwrap();
3118 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
3119 let original = (0..12)
3120 .map(|i| {
3121 msg(
3122 if i % 2 == 0 { "user" } else { "assistant" },
3123 &format!("Work item {i}"),
3124 )
3125 })
3126 .collect::<Vec<_>>();
3127 let mut envelope = prepared(&CompactionConfig::default());
3128 envelope.session_id = Some("handoff-test".into());
3129 let client = FixedSummaryClient::default();
3130 let mut usage = Usage::default();
3131 let result = compact_messages_safe(&client, &original, None, &envelope, &mut usage)
3132 .await
3133 .unwrap();
3134 assert!(result.summary_prompt.is_some());
3135 let files = std::fs::read_dir(root.path().join("sessions/handoff-test/artifacts"))
3136 .unwrap()
3137 .map(|entry| entry.unwrap().path())
3138 .collect::<Vec<_>>();
3139 let json = files
3140 .iter()
3141 .find(|path| path.extension().is_some_and(|ext| ext == "json"))
3142 .unwrap();
3143 let restored: Vec<Message> = serde_json::from_slice(&std::fs::read(json).unwrap()).unwrap();
3144 assert_eq!(restored, original);
3145 let markdown = files
3146 .iter()
3147 .find(|path| path.extension().is_some_and(|ext| ext == "md"))
3148 .unwrap();
3149 assert!(
3150 std::fs::read_to_string(markdown)
3151 .unwrap()
3152 .contains("migrate the session store")
3153 );
3154 // An unwritable artifact destination must abort before a provider call.
3155 envelope.session_id = Some("blocked-handoff".into());
3156 std::fs::write(
3157 root.path().join("sessions/blocked-handoff"),
3158 b"not a directory",
3159 )
3160 .unwrap();
3161 let blocked = FixedSummaryClient::default();
3162 assert!(
3163 compact_messages_safe(&blocked, &original, None, &envelope, &mut usage)
3164 .await
3165 .is_err()
3166 );
3167 assert!(blocked.request.lock().unwrap().is_none());
3168 }
3169
3170 #[tokio::test]
3171 async fn compaction_commits_summary_and_retains_recent_user_messages() {
3172 let messages = vec![
3173 msg(
3174 "user",
3175 "Objective: migrate the session store to sqlite without breaking existing logins",
3176 ),
3177 msg("assistant", "Working on it."),
3178 tool_use(
3179 "t1",
3180 "Bash",
3181 json!({"command": "cargo test -p session-store"}),
3182 ),
3183 tool_result("t1", "test session_store::roundtrip ... ok\nexit code 0"),
3184 msg("user", "Sounds good, do it"),
3185 msg("assistant", "Nearly done, rerunning the suite."),
3186 ];
3187 let config = CompactionConfig {
3188 model: "test-model".to_string(),
3189 cache_summary: false,
3190 ..Default::default()
3191 };
3192 let client = FixedSummaryClient::default();
3193
3194 let (retained, summary_prompt, _) =
3195 compact_messages(&client, &messages, &config).await.unwrap();
3196
3197 let request = client
3198 .request
3199 .lock()
3200 .expect("read summary request")
3201 .clone()
3202 .expect("summary request was captured");
3203 assert_eq!(&request.messages[..messages.len()], messages.as_slice());
3204 assert_eq!(request.messages.len(), messages.len() + 1);
3205 let ContentBlock::Text { text, .. } = &request.messages.last().unwrap().content[0] else {
3206 panic!("final compaction instruction must be text");
3207 };
3208 assert!(!text.contains(COMPACTION_SUMMARY_MARKER));
3209 assert_eq!(request.temperature, None);
3210 assert_eq!(request.top_p, None);
3211 assert_eq!(
3212 request.max_tokens,
3213 crate::route_budget::effective_max_output_tokens_for_route(
3214 crate::config::ProviderKind::Custom,
3215 "test-model",
3216 None,
3217 )
3218 );
3219
3220 let Some(SystemPrompt::Blocks(blocks)) = summary_prompt else {
3221 panic!("compaction must produce a summary system block");
3222 };
3223 let text = &blocks[0].text;
3224 assert!(text.contains(FIXED_SUMMARY));
3225 assert!(text.starts_with(COMPACTION_SUMMARY_MARKER));
3226
3227 // Replacement history keeps older user turns, then the open round
3228 // verbatim (user + assistant + tools), then one checkpoint.
3229 assert!(retained.iter().any(|message| {
3230 user_text_of(message).is_some_and(|text| text.contains("Objective: migrate"))
3231 }));
3232 assert!(
3233 retained
3234 .iter()
3235 .any(|message| { user_text_of(message).as_deref() == Some("Sounds good, do it") })
3236 );
3237 assert!(retained.iter().any(|message| {
3238 message.role.is_assistant_like()
3239 && message.content.iter().any(|block| {
3240 matches!(
3241 block,
3242 ContentBlock::Text { text, .. }
3243 if text.contains("Nearly done, rerunning the suite.")
3244 )
3245 })
3246 }));
3247 assert!(retained.iter().any(|message| {
3248 message.content.iter().any(|block| {
3249 matches!(
3250 block,
3251 ContentBlock::ToolResult { content, .. }
3252 if content.contains("session_store::roundtrip")
3253 )
3254 })
3255 }));
3256 assert!(is_compaction_checkpoint_message(retained.last().unwrap()));
3257 assert!(is_wire_compaction_checkpoint_message(
3258 retained.last().unwrap()
3259 ));
3260 assert!(matches!(
3261 &retained.last().unwrap().content[0],
3262 ContentBlock::Text { text: checkpoint, .. } if checkpoint == text
3263 ));
3264 last_round::validate_last_round_coverage(&messages, &retained[..retained.len() - 1])
3265 .unwrap();
3266 }
3267
3268 #[tokio::test]
3269 async fn uninterrupted_task_compacts_repeatedly_with_original_prefix_and_recent_tool_pairs() {
3270 let system = SystemPrompt::Text("stable project instructions and permissions".into());
3271 let mut prepared = PreparedCompactionEnvelope::new(CompactionConfig {
3272 token_threshold: 40_000,
3273 model: "test-model".into(),
3274 ..Default::default()
3275 });
3276 prepared.tools = Some(vec![
3277 serde_json::from_value(json!({
3278 "name": "File", "description": "Read a file", "input_schema": {"type": "object"}
3279 }))
3280 .unwrap(),
3281 ]);
3282 let client = FixedSummaryClient::default();
3283 let mut messages = vec![msg(
3284 "user",
3285 "Finish the migration. Preserve logins; do not publish.",
3286 )];
3287 for epoch in 0..3 {
3288 for step in 0..20 {
3289 let id = format!("{epoch}-{step}");
3290 let mut call = tool_use(&id, "File", json!({"path":"session.rs"}));
3291 call.content.insert(
3292 0,
3293 ContentBlock::Text {
3294 text: format!("Evidence {id}: {}", "x".repeat(12_000)),
3295 cache_control: None,
3296 },
3297 );
3298 if step == 19 {
3299 call.content.insert(
3300 0,
3301 ContentBlock::Thinking {
3302 thinking: "retained reasoning ".repeat(2000),
3303 signature: None,
3304 state: None,
3305 },
3306 );
3307 }
3308 messages.push(call);
3309 messages.push(tool_result(
3310 &id,
3311 &format!("Observed {id}: {}", "e".repeat(1000)),
3312 ));
3313 }
3314 let original = messages.clone();
3315 let mut usage = Usage::default();
3316 let result =
3317 compact_messages_safe(&client, &messages, Some(&system), &prepared, &mut usage)
3318 .await
3319 .unwrap();
3320 let request = client.request.lock().unwrap().clone().unwrap();
3321 assert_eq!(request.system.as_ref(), Some(&system));
3322 assert_eq!(request.tools, prepared.tools);
3323 assert_eq!(request.tool_choice, Some(json!("none")));
3324 assert_eq!(
3325 &request.messages[..original.len()],
3326 original.as_slice(),
3327 "summary must see the original evidence and reusable history prefix"
3328 );
3329 assert!(estimate_tokens(&result.messages) < estimate_tokens(&original) / 2);
3330 assert_eq!(
3331 result
3332 .messages
3333 .iter()
3334 .filter(|m| is_compaction_checkpoint_message(m))
3335 .count(),
3336 1
3337 );
3338 for step in [18, 19] {
3339 let id = format!("{epoch}-{step}");
3340 let expected = original.iter().find(|m| m.content.iter().any(|b| matches!(b, ContentBlock::ToolUse { id: found, .. } if found == &id))).unwrap();
3341 assert!(
3342 result.messages.contains(expected),
3343 "retained assistant text, calls and reasoning must survive unchanged"
3344 );
3345 assert!(result.messages.iter().any(|m| m.content.iter().any(|b| matches!(b, ContentBlock::ToolResult { tool_use_id, .. } if tool_use_id == &id))));
3346 }
3347 assert_eq!(result.messages[0], original[0]);
3348 messages = result.messages;
3349 }
3350 }
3351
3352 #[test]
3353 fn coverage_floor_rejects_a_replacement_that_drops_last_round_assistant() {
3354 let original = vec![
3355 msg("user", "What failed?"),
3356 msg("assistant", "session_store::roundtrip panics on reload."),
3357 ];
3358 let gutting = vec![msg("user", "What failed?")];
3359 let error = last_round::validate_last_round_coverage(&original, &gutting)
3360 .expect_err("dropping last-round assistant text must fail closed");
3361 assert!(error.to_string().contains("assistant"), "{error}");
3362 }
3363
3364 #[tokio::test]
3365 async fn compaction_preserves_runtime_provenance_and_saved_user_title() {
3366 let runtime = crate::runtime_handoff::operate_contract_runtime_message();
3367 let mut user = msg("user", "Build a focus timer");
3368 user.content.push(ContentBlock::Text {
3369 text: "<turn_meta>\nInput provenance: external_user\nInput authority: external_current_turn\n</turn_meta>".to_string(),
3370 cache_control: None,
3371 });
3372 let messages = vec![
3373 runtime.clone(),
3374 msg("assistant", "Ready."),
3375 user.clone(),
3376 msg("assistant", "Building the timer."),
3377 msg("user", "Keep the controls simple"),
3378 msg("assistant", "Adding start and stop."),
3379 ];
3380 let (retained, summary, _) = compact_messages(
3381 &FixedSummaryClient::default(),
3382 &messages,
3383 &CompactionConfig::default(),
3384 )
3385 .await
3386 .unwrap();
3387 assert_eq!(retained.first(), Some(&runtime));
3388 assert!(retained.contains(&user));
3389 assert!(crate::runtime_handoff::is_operate_contract_message(
3390 &retained[0]
3391 ));
3392 let saved = crate::session_manager::create_saved_session_with_mode(
3393 &retained,
3394 "test-model",
3395 std::path::Path::new("."),
3396 0,
3397 summary.as_ref(),
3398 None,
3399 );
3400 let restored: crate::session_manager::SavedSession =
3401 serde_json::from_slice(&serde_json::to_vec(&saved).unwrap()).unwrap();
3402 assert_eq!(restored.metadata.title, "Build a focus timer");
3403 assert_eq!(restored.messages, retained);
3404 }
3405
3406 #[test]
3407 fn retained_literal_runtime_xml_remains_user_authored() {
3408 let runtime = crate::runtime_handoff::operate_contract_runtime_message();
3409 // A user may paste the exact bytes, including a metadata example, in
3410 // one ordinary text block. Compaction must not split it into authority.
3411 let literal = user_text_of(&runtime).unwrap();
3412 let user = msg("user", &literal);
3413 let retained = retained_user_messages(std::slice::from_ref(&user), usize::MAX);
3414 assert_eq!(retained, vec![user]);
3415 assert_eq!(
3416 crate::runtime_handoff::classify_user_turn_prompt(&retained[0]),
3417 crate::runtime_handoff::UserTurnPromptKind::Editable,
3418 );
3419 assert_eq!(
3420 crate::session_manager::conversation_title_prompt(&retained),
3421 Some(literal.as_str()),
3422 );
3423 assert!(!crate::runtime_handoff::is_operate_contract_message(
3424 &retained[0]
3425 ));
3426 }
3427
3428 #[test]
3429 fn retained_image_turn_and_structured_content_keep_their_boundaries() {
3430 let image = Message {
3431 role: Role::User,
3432 content: vec![ContentBlock::ImageUrl {
3433 image_url: ImageUrlContent {
3434 url: "data:image/png;base64,AAAA".to_string(),
3435 },
3436 }],
3437 };
3438 let mut mixed = msg("user", "Compare these views");
3439 mixed.content.extend(image.content.clone());
3440 mixed.content.push(ContentBlock::Text {
3441 text: "Keep the original colors".to_string(),
3442 cache_control: Some(CacheControl {
3443 cache_type: "ephemeral".to_string(),
3444 }),
3445 });
3446 let messages = vec![image, mixed];
3447 let retained = retained_user_messages(&messages, 3_000);
3448 assert_eq!(retained, messages);
3449 let saved = crate::session_manager::create_saved_session_with_mode(
3450 &retained,
3451 "test-model",
3452 std::path::Path::new("."),
3453 0,
3454 None,
3455 None,
3456 );
3457 assert_eq!(
3458 saved.metadata.title,
3459 crate::session_manager::DEFAULT_SESSION_TITLE
3460 );
3461 assert!(retained_user_messages(&messages[..1], IMAGE_TOKEN_ESTIMATE - 1).is_empty());
3462 }
3463
3464 #[test]
3465 fn retained_budget_never_partially_promotes_structural_metadata() {
3466 let runtime = crate::runtime_handoff::operate_contract_runtime_message();
3467 let text = msg("user", "αβγδεζηθικ");
3468 let retained = retained_user_messages(&[runtime.clone(), text.clone()], 5);
3469 assert_eq!(retained, vec![text]);
3470 assert!(retained_user_messages(&[runtime], 1).is_empty());
3471 assert_eq!(
3472 user_text_of(&retained_user_messages(&[msg("user", "αβγδεζηθικ")], 2)[0]).as_deref(),
3473 Some("αβγδεζ"),
3474 );
3475 }
3476
3477 #[test]
3478 fn retained_older_turn_drops_result_blocks_whose_call_was_summarized() {
3479 // #6119: a host-supplied user message can mix text with a tool
3480 // result; the tool_use it answers lives in the summarized region, so
3481 // the retained copy must keep the text and drop the orphaned result.
3482 let mut mixed = msg("user", "Please keep this context.");
3483 mixed.content.push(ContentBlock::ToolResult {
3484 execution_id: None,
3485 tool_use_id: "toolu_orphan_1".to_string(),
3486 content: "{\"ok\":true}".to_string(),
3487 is_error: None,
3488 content_blocks: None,
3489 });
3490 let retained = retained_user_messages(std::slice::from_ref(&mixed), usize::MAX);
3491 assert_eq!(retained.len(), 1);
3492 assert!(
3493 !retained[0]
3494 .content
3495 .iter()
3496 .any(|block| matches!(block, ContentBlock::ToolResult { .. })),
3497 "the retained copy must not keep an orphaned tool_result"
3498 );
3499 assert_eq!(
3500 user_text_of(&retained[0]).as_deref(),
3501 Some("Please keep this context.")
3502 );
3503 // Insufficient budget still refuses to partially retain a multi-block
3504 // turn; the structural-metadata guard is unchanged.
3505 assert!(retained_user_messages(std::slice::from_ref(&mixed), 1).is_empty());
3506 }
3507
3508 #[test]
3509 fn summary_quality_gate_rejects_empty_and_known_placeholder_text() {
3510 for summary in [
3511 "",
3512 " \n\t ",
3513 "...",
3514 "。。。",
3515 "🫧",
3516 "N/A",
3517 "(no summary available)",
3518 "I cannot provide a summary.",
3519 ] {
3520 let error = validate_compaction_summary(summary)
3521 .expect_err("degenerate summary must fail closed");
3522 assert!(error.to_string().contains("unusable"), "{error}");
3523 }
3524 validate_compaction_summary(FIXED_SUMMARY)
3525 .expect("a substantive continuation handoff must be accepted");
3526 validate_compaction_summary(
3527 "目的: #4394の空要約を防止。完了: 検証と実装。制約: 履歴を変更しない。次: テスト実行。",
3528 )
3529 .expect("a concise multilingual handoff must not be rejected by prose length");
3530 }
3531
3532 #[tokio::test]
3533 async fn empty_successful_summary_retries_once_without_replacing_history() {
3534 let original = vec![
3535 msg(
3536 "user",
3537 "Keep the migration transactional and preserve existing sessions.",
3538 ),
3539 msg(
3540 "assistant",
3541 "I am updating the session store and its fixtures.",
3542 ),
3543 ];
3544 let client = ScriptedSummaryClient::new(vec![
3545 vec![ContentBlock::Text {
3546 text: " \n\t ".to_string(),
3547 cache_control: None,
3548 }],
3549 vec![ContentBlock::Text {
3550 text: FIXED_SUMMARY.to_string(),
3551 cache_control: None,
3552 }],
3553 ]);
3554 let config = CompactionConfig {
3555 model: "test-model".to_string(),
3556 cache_summary: false,
3557 ..Default::default()
3558 };
3559
3560 let mut invocation_usage = Usage::default();
3561 let result = compact_messages_safe(
3562 &client,
3563 &original,
3564 None,
3565 &prepared(&config),
3566 &mut invocation_usage,
3567 )
3568 .await
3569 .expect("the conservative retry should recover a usable summary");
3570 assert_eq!(invocation_usage.input_tokens, 34);
3571 assert_eq!(invocation_usage.output_tokens, 6);
3572 assert_eq!(invocation_usage.prompt_cache_hit_tokens, Some(10));
3573 assert_eq!(invocation_usage.reasoning_tokens, Some(4));
3574
3575 let requests = client
3576 .requests
3577 .lock()
3578 .expect("read scripted summary requests");
3579 assert_eq!(requests.len(), 2, "quality failure retries exactly once");
3580 let ContentBlock::Text { text, .. } = &requests[1]
3581 .messages
3582 .last()
3583 .expect("retry instruction")
3584 .content[0]
3585 else {
3586 panic!("retry instruction must be text");
3587 };
3588 assert!(text.contains("previous reply was empty or a placeholder"));
3589 drop(requests);
3590
3591 assert_eq!(
3592 result.retries_used, 1,
3593 "quality retry must reach diagnostics"
3594 );
3595 assert_eq!(original[0].role, "user", "source history remains untouched");
3596 assert!(result.messages.iter().any(is_compaction_checkpoint_message));
3597 let Some(SystemPrompt::Blocks(blocks)) = result.summary_prompt else {
3598 panic!("recovered summary must be committed");
3599 };
3600 assert!(blocks[0].text.contains(FIXED_SUMMARY));
3601 assert!(!blocks[0].text.contains("(no summary available)"));
3602 }
3603
3604 #[tokio::test]
3605 async fn quality_retry_count_survives_a_later_transient_failure() {
3606 let client = ScriptedSummaryClient::with_outcomes(vec![
3607 Ok(vec![ContentBlock::Text {
3608 text: "...".to_string(),
3609 cache_control: None,
3610 }]),
3611 Err(anyhow::anyhow!("request timed out")),
3612 Ok(vec![ContentBlock::Text {
3613 text: FIXED_SUMMARY.to_string(),
3614 cache_control: None,
3615 }]),
3616 ]);
3617 let config = CompactionConfig {
3618 model: "test-model".to_string(),
3619 cache_summary: false,
3620 ..Default::default()
3621 };
3622
3623 let mut invocation_usage = Usage::default();
3624 let result = compact_messages_safe(
3625 &client,
3626 &[msg("user", "Preserve the current migration state.")],
3627 None,
3628 &prepared(&config),
3629 &mut invocation_usage,
3630 )
3631 .await
3632 .expect("the outer retry should recover after the transient failure");
3633 assert_eq!(invocation_usage.input_tokens, 34);
3634 assert_eq!(invocation_usage.output_tokens, 6);
3635 assert_eq!(invocation_usage.prompt_cache_hit_tokens, Some(10));
3636 assert_eq!(invocation_usage.reasoning_tokens, Some(4));
3637
3638 assert_eq!(
3639 result.retries_used, 2,
3640 "one quality retry plus one outer transient retry must be reported"
3641 );
3642 assert_eq!(
3643 client
3644 .requests
3645 .lock()
3646 .expect("read scripted summary requests")
3647 .len(),
3648 3,
3649 "the diagnostic count must match the two calls after the initial request"
3650 );
3651 }
3652
3653 #[tokio::test]
3654 async fn compaction_usage_survives_cancellation_during_quality_retry() {
3655 let _cost_scope = crate::cost_status::test_scope();
3656 let retry_started = std::sync::Arc::new(tokio::sync::Notify::new());
3657 let mut client = ScriptedSummaryClient::new(vec![vec![ContentBlock::Text {
3658 text: "...".to_string(),
3659 cache_control: None,
3660 }]]);
3661 client.retry_started = Some(std::sync::Arc::clone(&retry_started));
3662 let messages = vec![msg("user", "Preserve the migration state.")];
3663 let prepared = prepared(&CompactionConfig::default());
3664 let mut invocation_usage = Usage::default();
3665 {
3666 let compaction =
3667 compact_messages_safe(&client, &messages, None, &prepared, &mut invocation_usage);
3668 tokio::pin!(compaction);
3669 tokio::select! {
3670 result = &mut compaction => panic!("retry must remain pending: {result:?}"),
3671 _ = retry_started.notified() => {},
3672 _ = tokio::time::sleep(Duration::from_secs(10)) => panic!("quality retry did not start"),
3673 }
3674 }
3675 assert_eq!(invocation_usage.input_tokens, 17);
3676 assert_eq!(invocation_usage.output_tokens, 3);
3677 assert_eq!(invocation_usage.prompt_cache_hit_tokens, Some(5));
3678 assert_eq!(invocation_usage.reasoning_tokens, Some(2));
3679 assert_eq!(client.requests.lock().unwrap().len(), 2);
3680 }
3681
3682 #[tokio::test]
3683 async fn non_text_summary_failure_preserves_history_after_one_retry() {
3684 let original = vec![
3685 msg(
3686 "user",
3687 "Do not lose the current branch or the failing test name.",
3688 ),
3689 msg("assistant", "The failing test is session_store::roundtrip."),
3690 ];
3691 let client = ScriptedSummaryClient::new(vec![
3692 vec![ContentBlock::thinking("internal-only response")],
3693 vec![ContentBlock::thinking("still no user-visible handoff")],
3694 ]);
3695 let config = CompactionConfig {
3696 model: "test-model".to_string(),
3697 cache_summary: false,
3698 ..Default::default()
3699 };
3700
3701 let mut invocation_usage = Usage::default();
3702 let error = compact_messages_safe(
3703 &client,
3704 &original,
3705 None,
3706 &prepared(&config),
3707 &mut invocation_usage,
3708 )
3709 .await
3710 .expect_err("two non-text responses must not replace history");
3711 assert_eq!(invocation_usage.input_tokens, 34);
3712 assert_eq!(invocation_usage.output_tokens, 6);
3713 assert_eq!(invocation_usage.prompt_cache_hit_tokens, Some(10));
3714 assert_eq!(invocation_usage.reasoning_tokens, Some(4));
3715
3716 assert!(
3717 error
3718 .to_string()
3719 .contains("remained unusable after one conservative retry"),
3720 "{error}"
3721 );
3722 assert_eq!(
3723 client
3724 .requests
3725 .lock()
3726 .expect("read scripted summary requests")
3727 .len(),
3728 2,
3729 "quality failure gets one retry, not the transient retry ladder"
3730 );
3731 assert_eq!(
3732 original,
3733 vec![
3734 msg(
3735 "user",
3736 "Do not lose the current branch or the failing test name."
3737 ),
3738 msg("assistant", "The failing test is session_store::roundtrip."),
3739 ],
3740 "borrowed source history must remain byte-for-byte unchanged"
3741 );
3742 }
3743 #[tokio::test]
3744 async fn compaction_uses_the_resolved_route_output_allowance() {
3745 for (route_label, provider, model) in [
3746 (
3747 "thinking-default route",
3748 crate::config::ProviderKind::Deepseek,
3749 "deepseek-v4-flash",
3750 ),
3751 (
3752 "fixed-sampling route",
3753 crate::config::ProviderKind::Moonshot,
3754 "k3",
3755 ),
3756 ] {
3757 let client = FixedSummaryClient::for_route(provider.as_str(), model);
3758 let config = CompactionConfig {
3759 model: model.to_string(),
3760 cache_summary: false,
3761 ..Default::default()
3762 };
3763 compact_messages(&client, &[msg("user", "summarize this task")], &config)
3764 .await
3765 .expect("route compaction should complete");
3766
3767 let request = client
3768 .request
3769 .lock()
3770 .expect("read summary request")
3771 .clone()
3772 .expect("summary request was captured");
3773 assert_eq!(
3774 request.max_tokens,
3775 crate::route_budget::effective_max_output_tokens_for_route(provider, model, None),
3776 "{route_label} must use the ordinary route output policy"
3777 );
3778 assert_eq!(request.temperature, None);
3779 assert_eq!(request.top_p, None);
3780 }
3781 }
3782
3783 struct TruncatedSummaryClient;
3784
3785 #[async_trait::async_trait]
3786 impl crate::core::model_client::ModelClient for TruncatedSummaryClient {
3787 fn provider_name(&self) -> &str {
3788 "test"
3789 }
3790
3791 fn model(&self) -> &str {
3792 "test-model"
3793 }
3794
3795 async fn create_message(
3796 &self,
3797 _request: MessageRequest,
3798 ) -> anyhow::Result<codewhale_models::MessageResponse> {
3799 Ok(codewhale_models::MessageResponse {
3800 id: "summary-truncated".to_string(),
3801 r#type: "message".to_string(),
3802 role: "assistant".to_string(),
3803 content: vec![ContentBlock::Text {
3804 text: "1. Primary request and intent — mig".to_string(),
3805 cache_control: None,
3806 }],
3807 model: "test-model".to_string(),
3808 stop_reason: Some("max_tokens".to_string()),
3809 stop_sequence: None,
3810 container: None,
3811 usage: codewhale_models::Usage::default(),
3812 })
3813 }
3814
3815 async fn create_message_stream(
3816 &self,
3817 _request: MessageRequest,
3818 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
3819 anyhow::bail!("streaming is unused by compaction")
3820 }
3821
3822 async fn health_check(&self) -> anyhow::Result<bool> {
3823 Ok(true)
3824 }
3825 }
3826
3827 /// A provider-truncated summary must fail compaction instead of replacing
3828 /// session history with a fragment.
3829 #[tokio::test]
3830 async fn truncated_summary_response_fails_compaction() {
3831 let messages: Vec<Message> = (0..40)
3832 .map(|index| {
3833 msg(
3834 if index % 2 == 0 { "user" } else { "assistant" },
3835 &format!("padding message {index} with enough text to compact"),
3836 )
3837 })
3838 .collect();
3839 let config = CompactionConfig {
3840 model: "test-model".to_string(),
3841 cache_summary: false,
3842 ..Default::default()
3843 };
3844
3845 let error = compact_messages(&TruncatedSummaryClient, &messages, &config)
3846 .await
3847 .expect_err("a truncated summary must not be committed");
3848 let text = error.to_string();
3849 assert!(text.contains("incomplete"), "{text}");
3850 assert!(text.contains("max_tokens"), "{text}");
3851 }
3852
3853 #[test]
3854 fn estimate_tokens_empty_messages() {
3855 let messages: Vec<Message> = vec![];
3856 assert_eq!(estimate_tokens(&messages), 0);
3857 }
3858
3859 #[test]
3860 fn estimate_tokens_with_text() {
3861 let messages = vec![Message {
3862 role: Role::User,
3863 content: vec![ContentBlock::Text {
3864 text: "Hello, world!".to_string(), // 13 chars = ~3 tokens
3865 cache_control: None,
3866 }],
3867 }];
3868 let tokens = estimate_tokens(&messages);
3869 assert!(tokens > 0 && tokens < 10);
3870 }
3871
3872 #[test]
3873 fn conservative_estimate_never_reads_cjk_below_the_byte_estimate() {
3874 let cjk = "中".repeat(1200); // 3,600 UTF-8 bytes
3875 assert!(estimate_text_tokens_conservative(&cjk) >= cjk.len() / 4);
3876 // ASCII keeps its three-characters-per-token reading.
3877 assert_eq!(estimate_text_tokens_conservative(&"a".repeat(300)), 100);
3878 }
3879
3880 #[test]
3881 fn pressure_counts_text_only_reasoning_and_server_tool_payloads() {
3882 let payload = "retained evidence ".repeat(1000);
3883 let blocks = vec![
3884 ContentBlock::thinking(payload.clone()),
3885 ContentBlock::ServerToolUse {
3886 id: "server-call".into(),
3887 name: "code_execution".into(),
3888 input: json!({"code": payload}),
3889 },
3890 ContentBlock::CodeExecutionToolResult {
3891 tool_use_id: "server-call".into(),
3892 content: json!({"stdout": payload}),
3893 },
3894 ContentBlock::ToolSearchToolResult {
3895 tool_use_id: "search-call".into(),
3896 content: json!({"description": payload}),
3897 },
3898 ];
3899 for block in blocks {
3900 let messages = vec![Message {
3901 role: Role::Assistant,
3902 content: vec![block],
3903 }];
3904 assert!(estimate_tokens(&messages) >= payload.len() / 4);
3905 assert!(estimate_input_tokens_for_pressure(&messages, None) >= payload.len() / 4);
3906 }
3907 }
3908
3909 #[test]
3910 fn estimate_tokens_counts_tool_round_thinking_across_turns() {
3911 // Per DeepSeek thinking-mode rules, any assistant message that
3912 // performed a tool call keeps its reasoning_content in the request
3913 // forever, including across new user turns. Token estimates must
3914 // count those bytes.
3915 let thinking = "reasoning ".repeat(800);
3916 let current_messages = vec![
3917 Message {
3918 role: Role::User,
3919 content: vec![ContentBlock::Text {
3920 text: "Use a tool".to_string(),
3921 cache_control: None,
3922 }],
3923 },
3924 Message {
3925 role: Role::Assistant,
3926 content: vec![
3927 ContentBlock::Thinking {
3928 signature: None,
3929 state: None,
3930 thinking: thinking.clone(),
3931 },
3932 ContentBlock::ToolUse {
3933 execution_id: None,
3934 id: "tool-1".to_string(),
3935 name: "read_file".to_string(),
3936 input: serde_json::json!({"path": "Cargo.toml"}),
3937 caller: None,
3938 thought_signature: None,
3939 },
3940 ],
3941 },
3942 Message {
3943 role: Role::User,
3944 content: vec![ContentBlock::ToolResult {
3945 execution_id: None,
3946 tool_use_id: "tool-1".to_string(),
3947 content: "manifest".to_string(),
3948 is_error: None,
3949 content_blocks: None,
3950 }],
3951 },
3952 ];
3953 let historical_messages = {
3954 let mut messages = current_messages.clone();
3955 messages.push(Message {
3956 role: Role::Assistant,
3957 content: vec![ContentBlock::Text {
3958 text: "Done.".to_string(),
3959 cache_control: None,
3960 }],
3961 });
3962 messages.push(Message {
3963 role: Role::User,
3964 content: vec![ContentBlock::Text {
3965 text: "Next question.".to_string(),
3966 cache_control: None,
3967 }],
3968 });
3969 messages
3970 };
3971 let completed_messages = {
3972 let mut messages = current_messages.clone();
3973 messages.push(Message {
3974 role: Role::Assistant,
3975 content: vec![ContentBlock::Text {
3976 text: "Done.".to_string(),
3977 cache_control: None,
3978 }],
3979 });
3980 messages
3981 };
3982
3983 let lower_bound = thinking.len() / 5;
3984 assert!(estimate_tokens(&current_messages) > lower_bound);
3985 assert!(estimate_tokens(&completed_messages) > lower_bound);
3986 assert!(estimate_tokens(&historical_messages) > lower_bound);
3987 }
3988
3989 #[test]
3990 fn should_compact_respects_enabled_flag() {
3991 let config = CompactionConfig {
3992 enabled: false,
3993 ..Default::default()
3994 };
3995 // Even with many messages, disabled compaction should return false
3996 let messages: Vec<Message> = (0..100)
3997 .map(|_| Message {
3998 role: Role::User,
3999 content: vec![ContentBlock::Text {
4000 text: "test".to_string(),
4001 cache_control: None,
4002 }],
4003 })
4004 .collect();
4005 assert!(!should_compact(&messages, None, &prepared(&config)));
4006 }
4007
4008 /// The #5577 acceptance case: a session whose provider bills 842K prompt
4009 /// tokens on a 1M window (threshold 800K) MUST compact even when the
4010 /// local estimate is far lower — the bounded working list undercounts
4011 /// what the provider actually saw, and billed truth wins.
4012 #[test]
4013 fn billed_842k_on_a_1m_window_compacts_despite_a_small_estimate() {
4014 let config = CompactionConfig {
4015 enabled: true,
4016 token_threshold: 800_000,
4017 ..Default::default()
4018 };
4019 let messages: Vec<Message> = (0..40)
4020 .map(|i| Message {
4021 role: if i % 2 == 0 {
4022 Role::User
4023 } else {
4024 Role::Assistant
4025 },
4026 content: vec![ContentBlock::Text {
4027 text: format!("short message {i}"),
4028 cache_control: None,
4029 }],
4030 })
4031 .collect();
4032 // Estimate alone stays far under the trigger…
4033 assert!(!should_compact(&messages, None, &prepared(&config)));
4034 // …but the provider's billed prompt total decides.
4035 assert_eq!(
4036 compaction_decision_with_billed(&messages, None, &prepared(&config), Some(842_000)),
4037 CompactionDecision::Compact
4038 );
4039 }
4040
4041 /// A refusal under real pressure must name its guard so the host can
4042 /// tell the user, instead of the silent hold that reads as a broken
4043 /// auto-compactor (#5577).
4044 #[test]
4045 fn refusals_under_pressure_name_their_guard() {
4046 let config = CompactionConfig {
4047 enabled: true,
4048 token_threshold: 100,
4049 ..Default::default()
4050 };
4051 // Too few messages to summarize: over-pressure, short transcript.
4052 let few: Vec<Message> = (0..3)
4053 .map(|i| Message {
4054 role: Role::User,
4055 content: vec![ContentBlock::Text {
4056 text: format!("message {i} {}", "x".repeat(300)),
4057 cache_control: None,
4058 }],
4059 })
4060 .collect();
4061 assert_eq!(
4062 compaction_decision_with_billed(&few, None, &prepared(&config), None),
4063 CompactionDecision::Refused(CompactionRefusal::TooFewMessages { count: few.len() })
4064 );
4065
4066 // Retained floor above the trigger: a giant system prompt no pass
4067 // can reclaim. The refusal carries the numbers the user needs.
4068 let many: Vec<Message> = (0..12)
4069 .map(|i| Message {
4070 role: if i % 2 == 0 {
4071 Role::User
4072 } else {
4073 Role::Assistant
4074 },
4075 content: vec![ContentBlock::Text {
4076 text: format!("message {i} {}", "y".repeat(200)),
4077 cache_control: None,
4078 }],
4079 })
4080 .collect();
4081 let system = SystemPrompt::Text("s".repeat(4_000));
4082 match compaction_decision_with_billed(&many, Some(&system), &prepared(&config), None) {
4083 CompactionDecision::Refused(CompactionRefusal::RetainedFloor { floor, threshold }) => {
4084 assert_eq!(threshold, 100);
4085 assert!(floor >= threshold, "floor {floor} must be over {threshold}");
4086 }
4087 other => panic!("expected a retained-floor refusal, got {other:?}"),
4088 }
4089 }
4090
4091 /// v0.8.11: message-count is no longer a compaction trigger. Long
4092 /// chats of small messages stay uncompacted because rewriting the
4093 /// prefix cache for a tiny budget reclaim is net-negative. Only token
4094 /// pressure (and the explicit `/compact` slash command) trigger
4095 /// compaction.
4096 #[test]
4097 fn message_count_no_longer_triggers_compaction() {
4098 let config = CompactionConfig {
4099 enabled: true,
4100 token_threshold: 1_000_000,
4101 ..Default::default()
4102 };
4103
4104 // 200 tiny messages, well above the prior message threshold.
4105 let many_messages: Vec<Message> = (0..200)
4106 .map(|_| Message {
4107 role: Role::User,
4108 content: vec![ContentBlock::Text {
4109 text: "x".to_string(),
4110 cache_control: None,
4111 }],
4112 })
4113 .collect();
4114 // Token total stays minuscule so the token threshold is not hit;
4115 // without the prior message-count trigger, no compaction.
4116 assert!(!should_compact(&many_messages, None, &prepared(&config)));
4117 }
4118
4119 // ========================================================================
4120 // Additional Compaction Trigger Tests
4121 // ========================================================================
4122
4123 #[test]
4124 fn full_request_pressure_crosses_token_threshold() {
4125 let config = CompactionConfig {
4126 enabled: true,
4127 token_threshold: 20_000,
4128 ..Default::default()
4129 };
4130
4131 // Create messages that exceed token threshold
4132 let messages: Vec<Message> = (0..20).map(|_| msg("user", &"x".repeat(5_000))).collect();
4133
4134 assert!(compaction_pressure_reached(&messages, None, &config));
4135 }
4136
4137 #[test]
4138 fn auto_compaction_uses_full_request_pressure_across_context_sizes() {
4139 for (window, output_reserve) in [
4140 (128_000_u64, 4_096_u64),
4141 (272_000, 4_096),
4142 // Large windows use the same ordinary request reservation; there
4143 // is no second, non-wire reasoning allowance.
4144 (1_000_000, 65_536),
4145 ] {
4146 let budget = crate::context_budget::ContextBudget::new(window, 0, output_reserve);
4147 let threshold = usize::try_from(budget.compaction_trigger_for_percent(80.0))
4148 .expect("test threshold fits usize");
4149 let raw_target = threshold.saturating_mul(7) / 10;
4150 let chars_per_message = raw_target.saturating_mul(4) / 14;
4151 let messages: Vec<Message> = (0..14)
4152 .map(|index| {
4153 msg(
4154 if index % 2 == 0 { "user" } else { "assistant" },
4155 &"x".repeat(chars_per_message),
4156 )
4157 })
4158 .collect();
4159 let raw = estimate_tokens(&messages);
4160 let full = estimate_input_tokens_for_pressure(&messages, None);
4161 let config = CompactionConfig {
4162 enabled: true,
4163 token_threshold: threshold,
4164 ..Default::default()
4165 };
4166
4167 assert!(
4168 raw < threshold,
4169 "raw message estimator alone must not cross {window}"
4170 );
4171 // The pressure estimate adds per-message framing on top of the
4172 // raw message tokens; billed usage from the provider can also
4173 // cross the trigger on its own.
4174 assert!(
4175 full < threshold,
4176 "70%-filled fixture must stay under the {window} trigger: {full} >= {threshold}"
4177 );
4178 assert!(
4179 crate::compaction::compaction_pressure_reached_with_billed(
4180 &messages,
4181 None,
4182 &config,
4183 Some(threshold as u64),
4184 ),
4185 "billed prompt tokens at the trigger must reach pressure for {window}"
4186 );
4187 assert!(
4188 crate::compaction::should_compact_with_billed(
4189 &messages,
4190 None,
4191 &prepared(&config),
4192 Some(threshold as u64),
4193 ),
4194 "billed pressure must trigger eligibility for a {window}-token route"
4195 );
4196 }
4197 }
4198
4199 #[test]
4200 fn auto_compaction_skips_pressure_that_cannot_be_reclaimed_below_trigger() {
4201 let messages: Vec<Message> = (0..20)
4202 .map(|index| {
4203 msg(
4204 if index % 2 == 0 { "user" } else { "assistant" },
4205 &"x".repeat(500),
4206 )
4207 })
4208 .collect();
4209 let system = SystemPrompt::Text("s".repeat(24_000));
4210 let config = CompactionConfig {
4211 enabled: true,
4212 token_threshold: 10_000,
4213 ..Default::default()
4214 };
4215
4216 assert!(
4217 estimate_input_tokens_conservative(&messages, Some(&system)) >= config.token_threshold,
4218 "fixture must be under full-request pressure"
4219 );
4220 assert!(
4221 !should_compact(&messages, Some(&system), &prepared(&config)),
4222 "a pinned/system floor above the trigger would loop every tool step"
4223 );
4224 }
4225
4226 #[test]
4227 fn full_request_threshold_is_inclusive() {
4228 let messages: Vec<Message> = (0..10)
4229 .map(|index| msg(if index % 2 == 0 { "user" } else { "assistant" }, "payload"))
4230 .collect();
4231 let threshold = estimate_input_tokens_for_pressure(&messages, None);
4232 let config = CompactionConfig {
4233 enabled: true,
4234 token_threshold: threshold,
4235 ..Default::default()
4236 };
4237
4238 assert!(compaction_pressure_reached(&messages, None, &config));
4239 }
4240
4241 #[test]
4242 fn test_should_compact_below_token_threshold() {
4243 let config = CompactionConfig {
4244 enabled: true,
4245 token_threshold: 1000,
4246 ..Default::default()
4247 };
4248
4249 // Create short messages
4250 let messages: Vec<Message> = (0..5).map(|_| msg("user", "short")).collect();
4251
4252 assert!(!should_compact(&messages, None, &prepared(&config)));
4253 }
4254
4255 #[test]
4256 fn auto_compaction_uses_token_threshold_without_fixed_floor() {
4257 let config = CompactionConfig {
4258 enabled: true,
4259 token_threshold: 20_000,
4260 ..Default::default()
4261 };
4262
4263 // Long sessions are dominated by assistant/tool output; the retained
4264 // user tail stays small, so the pass is reclaimable.
4265 let messages: Vec<Message> = (0..20)
4266 .map(|index| {
4267 if index % 2 == 0 {
4268 msg("user", &"x".repeat(100))
4269 } else {
4270 msg("assistant", &"x".repeat(10_000))
4271 }
4272 })
4273 .collect();
4274 assert!(should_compact(&messages, None, &prepared(&config)));
4275 }
4276
4277 #[test]
4278 fn test_compaction_result_retries_used() {
4279 // This test verifies the CompactionResult structure
4280 let result = CompactionResult {
4281 messages: vec![],
4282 summary_prompt: None,
4283 retries_used: 2,
4284 coverage: CompactionCoverage::default(),
4285 };
4286
4287 assert_eq!(result.retries_used, 2);
4288 assert!(result.messages.is_empty());
4289 }
4290 }
4291
4291 lines RUST