| 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(§ion); |
| 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(¤t_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 |