返回 CodeWhale
compaction.rs
根目录 / crates / tui / src / core / engine / compaction.rs
1 //! Engine-owned compaction lifecycle, recovery and checkpoint installation.
2 //! Automatic, manual and emergency paths share the existing session and event authority.
3
4 use super::*;
5
6 pub(super) struct CompactionPass {
7 pub trigger: &'static str,
8 pub path: crate::compaction::CompactionPath,
9 pub tokens_before: usize,
10 pub threshold_tokens: usize,
11 pub usage: Usage,
12 }
13
14 /// Which branch the turn loop takes after [`Engine::run_auto_compaction`].
15 /// The phase only reports it; `run_turn` keeps the `continue` and `return`.
16 pub(super) enum AutoCompactionStep {
17 /// No pass was due, or a pass reached an outcome (compacted, skipped or
18 /// failed with the conversation unchanged): build this step's request.
19 Proceed,
20 /// The pass was stopped on its own while the turn stays live: start the
21 /// next loop iteration without sending a request.
22 Restart,
23 /// End the turn with this outcome: it was cancelled during the pass, or
24 /// the pass spent the turn's wall-clock budget.
25 EndTurn(TurnOutcomeStatus, Option<String>),
26 }
27
28 /// Engine-side sink for compaction downgrade notices.
29 ///
30 /// Delivered as `Event::Status`, the same channel other engine status lines
31 /// use, so a long recovery says what it is doing while it runs. A full event
32 /// channel drops the notice instead of stalling the pass; the same sentence
33 /// is already in the log via `logging::warn`.
34 #[derive(Debug)]
35 struct EngineCompactionNoticeSink {
36 tx: mpsc::Sender<Event>,
37 child: Option<Arc<crate::tools::subagent::engine::ChildAuthority>>,
38 }
39
40 impl crate::compaction::CompactionNoticeSink for EngineCompactionNoticeSink {
41 fn notice(&self, message: String) {
42 let _ = self.tx.try_send(Event::Status { message });
43 }
44 fn accounting_origin(&self) -> Option<(crate::cost_status::CostScopeToken, String, String)> {
45 self.child.as_ref().map(|child| child.accounting_origin())
46 }
47 fn settled_usage<'a>(
48 &'a self,
49 source: &'a str,
50 route: &'a crate::cost_status::EffectiveRouteEnvelope,
51 usage: &'a Usage,
52 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
53 Box::pin(async move {
54 if let Some(child) = self.child.as_ref() {
55 child.project_settled_response(source, route, usage).await;
56 }
57 })
58 }
59 }
60
61 impl Engine {
62 pub(super) async fn emit_compaction_started(
63 &mut self,
64 id: String,
65 auto: bool,
66 message: String,
67 ) {
68 let _ = self
69 .send_event(Event::CompactionStarted { id, auto, message })
70 .await;
71 }
72
73 pub(super) async fn emit_compaction_completed(
74 &mut self,
75 id: String,
76 auto: bool,
77 message: String,
78 messages_before: Option<usize>,
79 messages_after: Option<usize>,
80 pass: CompactionPass,
81 ) {
82 let summary_prompt = self.rendered_compaction_summary();
83 // Every call site runs after message replacement and checkpoint
84 // commit. Reuse the same complete estimate as context pressure.
85 let post_input_tokens = Some(self.estimated_input_tokens() as u64);
86 let reduction_ratio = messages_before
87 .zip(messages_after)
88 .filter(|(before, _)| *before > 0)
89 .map(|(before, after)| 1.0 - after as f64 / before as f64);
90 self.record_compaction_event(
91 "compaction.completed",
92 serde_json::json!({
93 "compaction_id": id,
94 "trigger": pass.trigger,
95 "path": match pass.path {
96 crate::compaction::CompactionPath::Summary => "summary",
97 crate::compaction::CompactionPath::PruneOnly => "pruning_only",
98 },
99 "messages_before": messages_before,
100 "messages_after": messages_after,
101 "estimated_tokens_before": pass.tokens_before,
102 "estimated_tokens_after": post_input_tokens,
103 "threshold_tokens": pass.threshold_tokens,
104 "summarizer_usage": pass.usage,
105 "reduction_ratio": reduction_ratio,
106 }),
107 )
108 .await;
109 let _ = self
110 .send_event(Event::CompactionCompleted {
111 id,
112 auto,
113 message,
114 messages_before,
115 messages_after,
116 summary_prompt,
117 post_input_tokens,
118 })
119 .await;
120 }
121
122 /// One audit producer shared by interactive, headless and Runtime hosts.
123 /// No transcript content or credentials enter this diagnostic record.
124 pub(super) async fn record_compaction_event(
125 &self,
126 event: &'static str,
127 mut details: serde_json::Value,
128 ) {
129 details["session_id"] = serde_json::json!(self.session.id);
130 details["thread_id"] = serde_json::json!(self.config.runtime_services.active_thread_id);
131 details["model"] = serde_json::json!(self.config.model);
132 #[cfg(test)]
133 let env_ticket = crate::test_support::env_scope_ticket();
134 if let Err(error) = tokio::task::spawn_blocking(move || {
135 #[cfg(test)]
136 let _membership = crate::test_support::join_env_scope(env_ticket);
137 crate::audit::log_sensitive_event(event, details);
138 })
139 .await
140 {
141 tracing::warn!(%error, "compaction audit writer failed");
142 }
143 }
144
145 pub(super) async fn emit_compaction_cancelled(
146 &mut self,
147 id: String,
148 auto: bool,
149 message: String,
150 ) {
151 let _ = self
152 .send_event(Event::CompactionCancelled { id, auto, message })
153 .await;
154 }
155
156 /// Render the accumulated compaction summary prompt to plain text so it
157 /// can travel in events and be persisted by host layers. All emit sites
158 /// run after `commit_compaction_checkpoint`, so this reflects the checkpoint
159 /// state the engine will use for subsequent requests.
160 pub(super) fn rendered_compaction_summary(&self) -> Option<String> {
161 self.session
162 .compaction_summary_prompt
163 .as_ref()
164 .map(|prompt| match prompt {
165 SystemPrompt::Text(text) => text.clone(),
166 SystemPrompt::Blocks(blocks) => blocks
167 .iter()
168 .map(|block| block.text.as_str())
169 .collect::<Vec<_>>()
170 .join("\n\n"),
171 })
172 .filter(|text| !text.trim().is_empty())
173 }
174
175 pub(super) async fn emit_compaction_failed(&mut self, id: String, auto: bool, message: String) {
176 let _ = self
177 .send_event(Event::CompactionFailed { id, auto, message })
178 .await;
179 }
180
181 pub(super) fn claim_compaction(&self, id: &str) -> Option<CancellationToken> {
182 self.compaction_cancellation
183 .lock()
184 .unwrap_or_else(std::sync::PoisonError::into_inner)
185 .claim(id)
186 }
187
188 pub(super) fn finish_compaction(&self, id: &str) {
189 self.compaction_cancellation
190 .lock()
191 .unwrap_or_else(std::sync::PoisonError::into_inner)
192 .finish(id);
193 }
194
195 /// Pressure and effective trigger in append-only turn metadata. Numeric
196 /// estimates never modify the session-pinned system/tool prefix.
197 pub(super) fn context_pressure_line(
198 &self,
199 current_text: &str,
200 prompt_context: &NextTurnPromptContext,
201 system_prompt: Option<&SystemPrompt>,
202 ) -> Option<String> {
203 // The engine owns automatic compaction. Asking the model to warn the
204 // user here created a competing save/compact ceremony before the
205 // automatic request-boundary guard could do its work (#5620).
206 if self.config.compaction.enabled {
207 return None;
208 }
209 let input_tokens = self.active_input_tokens_with_current_text(current_text, system_prompt);
210 let budget = route_context_budget_for_route(
211 prompt_context.provider,
212 &prompt_context.model,
213 prompt_context.route_limits,
214 input_tokens,
215 )?;
216 context_pressure_message(budget.usage_percent()).map(|warning| format!(
217 "{warning}. Estimated input: {input_tokens} tokens ({:.1}% of route budget). Making room automatically is off for this session. /compact saves the original conversation and its model-written handoff before replacing context.",
218 budget.usage_percent(),
219 ))
220 }
221
222 pub(super) fn prepare_compaction_envelope(
223 &self,
224 mut config: CompactionConfig,
225 ) -> PreparedCompactionEnvelope {
226 // Host-supplied configs may not carry the workspace; compaction needs
227 // it only to re-state the user's `/anchor` file after the summary.
228 config
229 .workspace
230 .get_or_insert_with(|| self.config.workspace.clone());
231 let mut prepared = PreparedCompactionEnvelope::new(config);
232 prepared.session_id = Some(self.session.id.clone());
233 prepared.notice_sink = Some(std::sync::Arc::new(EngineCompactionNoticeSink {
234 tx: self.tx_event.clone(),
235 child: self
236 .child_host
237 .as_ref()
238 .map(|child| child.authority.clone()),
239 }));
240 // The summary request must carry the reasoning tier the turn sends:
241 // reasoning routes render it at the head of the prompt, so omitting
242 // it forfeited the whole cached history prefix (#6540).
243 prepared.reasoning_effort = super::turn_loop::resolve_auto_effort(
244 self.session.reasoning_effort.as_deref(),
245 self.api_provider,
246 &self.api_config.active_route_base_url(),
247 &self.config.model,
248 );
249 prepared
250 }
251
252 pub(super) async fn handle_manual_compaction_op(
253 &mut self,
254 id: String,
255 route: ResolvedRuntimeRoute,
256 compaction: CompactionConfig,
257 ) {
258 self.emit_compaction_started(id.clone(), false, "Making room…".to_string())
259 .await;
260 let Some(cancel_token) = self.claim_compaction(&id) else {
261 let message = "Making room stopped before it started".to_string();
262 self.emit_compaction_cancelled(id, false, message).await;
263 let _ = self
264 .send_event(Event::TurnComplete {
265 usage: Usage::default(),
266 parent_route_usage: Usage::default(),
267 routed_usage_dropped_records: 0,
268 status: TurnOutcomeStatus::Interrupted,
269 error: None,
270 tool_catalog: None,
271 base_url: None,
272 })
273 .await;
274 return;
275 };
276 if let Err(err) = self.install_resolved_runtime_route(route) {
277 let message =
278 format!("Cannot compact context because its provider route is not ready: {err}");
279 self.finish_compaction(&id);
280 self.emit_compaction_failed(id, false, message.clone())
281 .await;
282 let _ = self
283 .send_event(Event::error(ErrorEnvelope::fatal_auth(message)))
284 .await;
285 return;
286 }
287 self.config.compaction = compaction;
288 self.handle_manual_compaction(id, cancel_token).await;
289 }
290
291 pub(super) async fn emit_compaction_usage(&self, usage: &Usage, elapsed: Duration) {
292 if *usage == Usage::default() {
293 return;
294 }
295 let _ = self
296 .send_event(Event::RoutedTurnUsage {
297 usage: usage.clone(),
298 duration_ms: u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX),
299 first_token_ms: None,
300 request_ms: None,
301 })
302 .await;
303 }
304
305 pub(super) async fn handle_manual_compaction(
306 &mut self,
307 id: String,
308 cancel_token: CancellationToken,
309 ) {
310 let zero_usage = Usage {
311 input_tokens: 0,
312 output_tokens: 0,
313 ..Usage::default()
314 };
315 let Some(client) = self.codewhale_client.clone() else {
316 let message = "Can't make room: no model is connected".to_string();
317 self.finish_compaction(&id);
318 self.emit_compaction_failed(id, false, message.clone())
319 .await;
320 let _ = self
321 .send_event(Event::error(ErrorEnvelope::fatal_auth(message.clone())))
322 .await;
323 let _ = self
324 .send_event(Event::TurnComplete {
325 usage: zero_usage,
326 parent_route_usage: Usage::default(),
327 routed_usage_dropped_records: 0,
328 status: TurnOutcomeStatus::Failed,
329 error: Some(message),
330 tool_catalog: None,
331 base_url: None,
332 })
333 .await;
334 return;
335 };
336
337 let messages_before = self.session.messages.len();
338 // Message counts alone do not show the win the user cares about: a
339 // compaction that drops few but enormous messages reads as a no-op.
340 // The emergency path already reports tokens; manual and auto now match.
341 let tokens_before = self.estimated_input_tokens();
342 let mut turn_status = TurnOutcomeStatus::Completed;
343 let mut turn_error = None;
344
345 let prepared = self.prepare_compaction_envelope(self.config.compaction.clone());
346
347 let started = Instant::now();
348 let mut compaction_usage = Usage::default();
349 let compaction_result = tokio::select! {
350 biased;
351 _ = cancel_token.cancelled() => None,
352 result = compact_messages_safe(
353 &client,
354 &self.session.messages,
355 self.session.system_prompt.as_ref(),
356 &prepared,
357 &mut compaction_usage,
358 ) => Some(result),
359 };
360 self.session.total_usage.add(&compaction_usage);
361 self.record_goal_usage_for_turn(&compaction_usage, started.elapsed());
362 self.emit_compaction_usage(&compaction_usage, started.elapsed())
363 .await;
364
365 let Some(compaction_result) = compaction_result else {
366 self.finish_compaction(&id);
367 self.emit_compaction_cancelled(
368 id,
369 false,
370 "Making room stopped; the conversation was not changed".to_string(),
371 )
372 .await;
373 let _ = self
374 .send_event(Event::TurnComplete {
375 usage: compaction_usage,
376 parent_route_usage: Usage::default(),
377 routed_usage_dropped_records: 0,
378 status: TurnOutcomeStatus::Interrupted,
379 error: None,
380 tool_catalog: None,
381 base_url: None,
382 })
383 .await;
384 return;
385 };
386
387 match compaction_result {
388 Ok(mut result) => {
389 if !result.messages.is_empty() || self.session.messages.is_empty() {
390 self.append_compaction_agent_topology(&mut result.messages)
391 .await;
392 if cancel_token.is_cancelled() {
393 self.finish_compaction(&id);
394 self.emit_compaction_cancelled(
395 id,
396 false,
397 "Making room stopped; the conversation was not changed".to_string(),
398 )
399 .await;
400 let _ = self
401 .send_event(Event::TurnComplete {
402 usage: compaction_usage,
403 parent_route_usage: Usage::default(),
404 routed_usage_dropped_records: 0,
405 status: TurnOutcomeStatus::Interrupted,
406 error: None,
407 tool_catalog: None,
408 base_url: None,
409 })
410 .await;
411 return;
412 }
413 let messages_after = result.messages.len();
414 let retries_used = result.retries_used;
415 let coverage_clause = result.coverage.receipt_clause();
416 let path = result.coverage.path;
417 if let Some(job) = self.child_job()
418 && let Err(error) = job
419 .before_replace(&self.session.messages, &result.messages)
420 .await
421 {
422 self.cancel_token.cancel();
423 tracing::error!(%error, "child compaction projection failed; original Session retained");
424 return;
425 }
426 self.session.replace_messages(result.messages);
427 if let Some(pm) = self.session.prefix_stability.as_mut() {
428 pm.note_history_reset("compaction");
429 }
430 self.commit_compaction_checkpoint(result.summary_prompt);
431 self.emit_session_updated().await;
432 let removed = messages_before.saturating_sub(messages_after);
433 let tokens_after = self.estimated_input_tokens();
434 let message = if retries_used > 0 {
435 format!(
436 "Made room: {messages_before} → {messages_after} messages ({removed} removed, {retries_used} retries), ~{tokens_before} → ~{tokens_after} tokens ({coverage_clause})"
437 )
438 } else {
439 format!(
440 "Made room: {messages_before} → {messages_after} messages ({removed} removed), ~{tokens_before} → ~{tokens_after} tokens ({coverage_clause})"
441 )
442 };
443 self.emit_compaction_completed(
444 id.clone(),
445 false,
446 message,
447 Some(messages_before),
448 Some(messages_after),
449 CompactionPass {
450 trigger: "manual",
451 path,
452 tokens_before,
453 threshold_tokens: prepared.config.token_threshold,
454 usage: compaction_usage.clone(),
455 },
456 )
457 .await;
458 } else {
459 let message = "Making room skipped: the summary came back empty".to_string();
460 self.emit_compaction_failed(id.clone(), false, message.clone())
461 .await;
462 turn_status = TurnOutcomeStatus::Failed;
463 turn_error = Some(message);
464 }
465 }
466 Err(err) => {
467 let message = crate::compaction::report_compaction_failure(
468 "Making room failed",
469 &id,
470 false,
471 &err,
472 );
473 self.emit_compaction_failed(id.clone(), false, message.clone())
474 .await;
475 let _ = self.send_event(Event::status(message.clone())).await;
476 turn_status = TurnOutcomeStatus::Failed;
477 turn_error = Some(message);
478 }
479 }
480
481 self.finish_compaction(&id);
482
483 let _ = self
484 .send_event(Event::TurnComplete {
485 usage: compaction_usage,
486 parent_route_usage: Usage::default(),
487 routed_usage_dropped_records: 0,
488 status: turn_status,
489 error: turn_error,
490 tool_catalog: None,
491 base_url: None,
492 })
493 .await;
494 }
495
496 /// The automatic compaction phase of `run_turn`, run once per loop
497 /// iteration before the request is built.
498 ///
499 /// When context pressure has reached the trigger, this summarizes history
500 /// with the tool prefix the next request will carry, installs the
501 /// checkpoint and reports the pass; a refusal is named once per turn.
502 /// `auto_compaction_suppressed` is the loop's turn-scoped latch: a failed,
503 /// cancelled or empty pass, or one that leaves pressure high, sets it so
504 /// the turn cannot become a paid summarization loop at every tool
505 /// boundary. The bounded hard-limit recovery
506 /// ([`Self::recover_context_overflow`]) stays available either way.
507 ///
508 /// The loop keeps its own control flow: the returned
509 /// [`AutoCompactionStep`] names the branch, and `run_turn` takes it.
510 pub(super) async fn run_auto_compaction(
511 &mut self,
512 client: &dyn crate::core::model_client::ModelClient,
513 active_tools: Option<&[Tool]>,
514 turn: &mut TurnContext,
515 auto_compaction_suppressed: &mut bool,
516 ) -> AutoCompactionStep {
517 let auto_compaction_config = self.config.compaction.clone();
518 // Billing usage accumulates every parent step and child-model
519 // call. Only the most recent parent-route request describes the
520 // live message list whose pressure we are checking here.
521 let billed_input_tokens = turn.live_input_tokens_for_compaction(
522 &self.session.messages,
523 self.session.system_prompt.as_ref(),
524 self.session.latest_parent_input_tokens,
525 );
526 let prepared = if !*auto_compaction_suppressed
527 && crate::compaction::compaction_pressure_reached_with_billed(
528 &self.session.messages,
529 self.session.system_prompt.as_ref(),
530 &auto_compaction_config,
531 billed_input_tokens,
532 ) {
533 let mut prepared = self.prepare_compaction_envelope(auto_compaction_config);
534 prepared.tools = active_tools.map(<[Tool]>::to_vec);
535 Some(prepared)
536 } else {
537 None
538 };
539
540 let compaction_go = match prepared.as_ref() {
541 None => false,
542 Some(prepared) => match crate::compaction::compaction_decision_with_billed(
543 &self.session.messages,
544 self.session.system_prompt.as_ref(),
545 prepared,
546 billed_input_tokens,
547 ) {
548 crate::compaction::CompactionDecision::Compact => true,
549 crate::compaction::CompactionDecision::NotNeeded => false,
550 crate::compaction::CompactionDecision::Refused(reason) => {
551 // A silent refusal looks like broken auto-compaction:
552 // the meter is full and nothing happens (#5577). Name
553 // the guard once per turn, in both the transcript
554 // status line and the trace.
555 if !turn.compaction_refusal_notified {
556 turn.compaction_refusal_notified = true;
557 let estimated_tokens_before = self.estimated_input_tokens();
558 self.record_compaction_event("compaction.refused", serde_json::json!({
559 "trigger": "auto",
560 "reason": match &reason {
561 crate::compaction::CompactionRefusal::TooFewMessages { .. } => "too_few_messages",
562 crate::compaction::CompactionRefusal::RetainedFloor { .. } => "retained_floor",
563 },
564 "messages_before": self.session.messages.len(),
565 "estimated_tokens_before": estimated_tokens_before,
566 "billed_input_tokens": billed_input_tokens,
567 "threshold_tokens": prepared.config.token_threshold,
568 })).await;
569 let message = match reason {
570 crate::compaction::CompactionRefusal::TooFewMessages { count } => {
571 format!(
572 "Context is filling up, but there is nothing to make room from yet: only {count} messages"
573 )
574 }
575 crate::compaction::CompactionRefusal::RetainedFloor {
576 floor,
577 threshold,
578 } => format!(
579 "Context is filling up, but making room would not help: retained context (~{}K tokens) cannot fall below the {}K trigger — /compact to force a pass, or trim pinned context",
580 floor / 1000,
581 threshold / 1000
582 ),
583 };
584 tracing::warn!(
585 target: "compaction",
586 ?reason,
587 billed = ?billed_input_tokens,
588 "auto-compaction refused under pressure"
589 );
590 let _ = self.send_event(Event::status(message)).await;
591 }
592 false
593 }
594 },
595 };
596 if let Some(prepared) = prepared
597 && compaction_go
598 {
599 let compaction_id = format!("compact_{}", &uuid::Uuid::new_v4().to_string()[..8]);
600 turn.stop_diagnostics.automatic_compaction_attempts = turn
601 .stop_diagnostics
602 .automatic_compaction_attempts
603 .saturating_add(1);
604 let compaction_cancel = self
605 .claim_compaction(&compaction_id)
606 .expect("a fresh automatic compaction id cannot be pre-canceled");
607 self.emit_compaction_started(compaction_id.clone(), true, "Making room…".to_string())
608 .await;
609 let auto_messages_before = self.session.messages.len();
610 let auto_tokens_before = self.estimated_input_tokens();
611 let turn_cancel = self.cancel_token.clone();
612 let started = Instant::now();
613 let mut compaction_usage = Usage::default();
614 // Parked: the compaction pass owns its own bound.
615 self.turn_heartbeat
616 .enter(super::turn_heartbeat::TurnPhase::Compacting, None, None);
617 let (compaction_result, turn_was_canceled) = tokio::select! {
618 biased;
619 _ = turn_cancel.cancelled() => (None, true),
620 _ = compaction_cancel.cancelled() => (None, false),
621 result = compact_messages_safe(
622 client,
623 &self.session.messages,
624 self.session.system_prompt.as_ref(),
625 &prepared,
626 &mut compaction_usage,
627 ) => (Some(result), false),
628 };
629 turn.add_usage(&compaction_usage);
630 self.emit_compaction_usage(&compaction_usage, started.elapsed())
631 .await;
632 let Some(compaction_result) = compaction_result else {
633 *auto_compaction_suppressed = true;
634 self.finish_compaction(&compaction_id);
635 let message = if turn_was_canceled {
636 "Making room stopped with the turn; the conversation was not changed"
637 } else {
638 "Making room stopped; the conversation was not changed"
639 }
640 .to_string();
641 self.emit_compaction_cancelled(compaction_id, true, message)
642 .await;
643 if turn_was_canceled {
644 return AutoCompactionStep::EndTurn(TurnOutcomeStatus::Interrupted, None);
645 }
646 return AutoCompactionStep::Restart;
647 };
648
649 match compaction_result {
650 Ok(mut result) => {
651 // Only update if we got valid messages (never corrupt state)
652 if !result.messages.is_empty() || self.session.messages.is_empty() {
653 self.append_compaction_agent_topology(&mut result.messages)
654 .await;
655 let turn_was_canceled = turn_cancel.is_cancelled();
656 if turn_was_canceled || compaction_cancel.is_cancelled() {
657 *auto_compaction_suppressed = true;
658 self.finish_compaction(&compaction_id);
659 let message = if turn_was_canceled {
660 "Making room stopped with the turn; the conversation was not changed"
661 } else {
662 "Making room stopped; the conversation was not changed"
663 }
664 .to_string();
665 self.emit_compaction_cancelled(compaction_id, true, message)
666 .await;
667 if turn_was_canceled {
668 return AutoCompactionStep::EndTurn(
669 TurnOutcomeStatus::Interrupted,
670 None,
671 );
672 }
673 return AutoCompactionStep::Restart;
674 }
675 let auto_messages_after = result.messages.len();
676 let retries_used = result.retries_used;
677 let coverage_clause = result.coverage.receipt_clause();
678 let path = result.coverage.path;
679 if let Some(job) = self.child_job()
680 && let Err(error) = job
681 .before_replace(&self.session.messages, &result.messages)
682 .await
683 {
684 return AutoCompactionStep::EndTurn(
685 TurnOutcomeStatus::Failed,
686 Some(format!("child compaction projection failed: {error:#}")),
687 );
688 }
689 self.session.replace_messages(result.messages);
690 turn.clear_parent_input_tokens();
691 if let Some(pm) = self.session.prefix_stability.as_mut() {
692 pm.note_history_reset("compaction");
693 }
694 self.commit_compaction_checkpoint(result.summary_prompt);
695 *auto_compaction_suppressed =
696 crate::compaction::compaction_pressure_reached(
697 &self.session.messages,
698 self.session.system_prompt.as_ref(),
699 &self.config.compaction,
700 );
701 self.emit_session_updated().await;
702 let removed = auto_messages_before.saturating_sub(auto_messages_after);
703 let auto_tokens_after = self.estimated_input_tokens();
704 let status = if retries_used > 0 {
705 format!(
706 "Made room: {auto_messages_before} → {auto_messages_after} messages ({removed} removed, {retries_used} retries), ~{auto_tokens_before} → ~{auto_tokens_after} tokens ({coverage_clause})"
707 )
708 } else {
709 format!(
710 "Made room: {auto_messages_before} → {auto_messages_after} messages ({removed} removed), ~{auto_tokens_before} → ~{auto_tokens_after} tokens ({coverage_clause})"
711 )
712 };
713 self.emit_compaction_completed(
714 compaction_id.clone(),
715 true,
716 status.clone(),
717 Some(auto_messages_before),
718 Some(auto_messages_after),
719 CompactionPass {
720 trigger: "auto",
721 path,
722 tokens_before: auto_tokens_before,
723 threshold_tokens: prepared.config.token_threshold,
724 usage: compaction_usage.clone(),
725 },
726 )
727 .await;
728 } else {
729 *auto_compaction_suppressed = true;
730 let message =
731 "Making room skipped: the summary came back empty".to_string();
732 self.emit_compaction_failed(compaction_id.clone(), true, message.clone())
733 .await;
734 let _ = self.send_event(Event::status(message)).await;
735 }
736 }
737 Err(err) => {
738 *auto_compaction_suppressed = true;
739 // Log error but continue with original messages (never corrupt)
740 let message = crate::compaction::report_compaction_failure(
741 "Making room failed",
742 &compaction_id,
743 true,
744 &err,
745 );
746 self.emit_compaction_failed(compaction_id.clone(), true, message.clone())
747 .await;
748 let _ = self.send_event(Event::status(message)).await;
749 }
750 }
751 self.finish_compaction(&compaction_id);
752 // C02-06: a compaction pass has its own bound, not the
753 // turn's. Recheck the wall clock before it can authorize the
754 // provider request that follows this phase.
755 if let Some(error) = self.turn_wall_clock_exhausted_error() {
756 let _ = self.send_event(Event::status(error.clone())).await;
757 return AutoCompactionStep::EndTurn(TurnOutcomeStatus::Failed, Some(error));
758 }
759 }
760 AutoCompactionStep::Proceed
761 }
762
763 pub(super) async fn recover_context_overflow(
764 &mut self,
765 client: &dyn crate::core::model_client::ModelClient,
766 tools: Option<&[Tool]>,
767 reason: &str,
768 turn: &mut TurnContext,
769 ) -> bool {
770 let Some(target_budget) = context_input_budget_for_route(
771 self.api_provider,
772 &self.session.model,
773 self.active_route_limits,
774 0,
775 ) else {
776 return false;
777 };
778 // Nothing to summarize or prune: a pass cannot help, so do not make
779 // the user wait on a model call before the failure the caller will
780 // report anyway.
781 if !crate::compaction::has_compactable_history(&self.session.messages) {
782 return false;
783 }
784
785 let id = format!("compact_{}", &uuid::Uuid::new_v4().to_string()[..8]);
786 turn.stop_diagnostics.emergency_compaction_attempts = turn
787 .stop_diagnostics
788 .emergency_compaction_attempts
789 .saturating_add(1);
790 let start_message = format!("Making room now ({reason})");
791 self.emit_compaction_started(id.clone(), true, start_message)
792 .await;
793 let Some(compaction_cancel) = self.claim_compaction(&id) else {
794 self.emit_compaction_cancelled(
795 id,
796 true,
797 "Making room stopped before it started; the conversation was not changed"
798 .to_string(),
799 )
800 .await;
801 return false;
802 };
803 let turn_cancel = self.cancel_token.clone();
804
805 // Measured with the estimator `after_tokens` and the preflight guard
806 // use, so `recovered` compares like with like and the receipt reads
807 // ~before → ~after on one scale.
808 let before_tokens = crate::compaction::estimate_input_tokens_for_pressure(
809 &self.session.messages,
810 self.session.system_prompt.as_ref(),
811 );
812 let before_count = self.session.messages.len();
813
814 let mut forced_config = self.config.compaction.clone();
815 forced_config.enabled = true;
816 forced_config.token_threshold = forced_config
817 .token_threshold
818 .min(target_budget.saturating_sub(1))
819 .max(1);
820 let mut prepared = self.prepare_compaction_envelope(forced_config);
821 prepared.tools = tools.map(<[Tool]>::to_vec);
822
823 let started = Instant::now();
824 let mut compaction_usage = Usage::default();
825 let (compaction_result, turn_was_canceled) = tokio::select! {
826 biased;
827 _ = turn_cancel.cancelled() => (None, true),
828 _ = compaction_cancel.cancelled() => (None, false),
829 result = compact_messages_safe(
830 client,
831 &self.session.messages,
832 self.session.system_prompt.as_ref(),
833 &prepared,
834 &mut compaction_usage,
835 ) => (Some(result), false),
836 };
837 turn.add_usage(&compaction_usage);
838 self.emit_compaction_usage(&compaction_usage, started.elapsed())
839 .await;
840 let Some(compaction_result) = compaction_result else {
841 self.finish_compaction(&id);
842 let message = if turn_was_canceled {
843 "Making room stopped with the turn; the conversation was not changed"
844 } else {
845 "Making room stopped; the conversation was not changed"
846 }
847 .to_string();
848 self.emit_compaction_cancelled(id, true, message).await;
849 return false;
850 };
851
852 let result = match compaction_result {
853 Ok(result) => result,
854 Err(err) => {
855 let message = if is_provider_rejection(&err) {
856 // The turn's error line carries the provider's answer;
857 // this receipt only closes the recovery attempt.
858 "Context recovery stopped: the provider rejected the request. Original conversation was preserved.".to_string()
859 } else {
860 let reason = format!("{err:#}");
861 let reason = reason.trim_end().trim_end_matches('.');
862 if reason
863 .to_ascii_lowercase()
864 .contains("conversation was preserved")
865 {
866 format!("Context recovery failed: {reason}.")
867 } else {
868 format!(
869 "Context recovery failed: {reason}. Original conversation was preserved."
870 )
871 }
872 };
873 if is_provider_rejection(&err) {
874 turn.context_recovery_rejection = Some(err);
875 }
876 self.emit_compaction_failed(id.clone(), true, message).await;
877 self.finish_compaction(&id);
878 return false;
879 }
880 };
881 let retries_used = result.retries_used;
882 let summary_prompt = result.summary_prompt;
883 let path = result.coverage.path;
884 let mut compacted_messages = result.messages;
885
886 let turn_was_canceled = turn_cancel.is_cancelled();
887 if turn_was_canceled || compaction_cancel.is_cancelled() {
888 self.finish_compaction(&id);
889 let message = if turn_was_canceled {
890 "Making room stopped with the turn; the conversation was not changed"
891 } else {
892 "Making room stopped; the conversation was not changed"
893 }
894 .to_string();
895 self.emit_compaction_cancelled(id, true, message).await;
896 return false;
897 }
898
899 if !compacted_messages.is_empty() || self.session.messages.is_empty() {
900 self.append_compaction_agent_topology(&mut compacted_messages)
901 .await;
902 let turn_was_canceled = turn_cancel.is_cancelled();
903 if turn_was_canceled || compaction_cancel.is_cancelled() {
904 self.finish_compaction(&id);
905 let message = if turn_was_canceled {
906 "Making room stopped with the turn; the conversation was not changed"
907 } else {
908 "Making room stopped; the conversation was not changed"
909 }
910 .to_string();
911 self.emit_compaction_cancelled(id, true, message).await;
912 return false;
913 }
914 }
915 // Validate the complete candidate before the only history swap. Bare
916 // front-trimming after a failed summary silently lost user state and
917 // could leave orphan tool results in an apparently recovered session.
918 let after_tokens = crate::compaction::estimate_input_tokens_for_pressure(
919 &compacted_messages,
920 self.session.system_prompt.as_ref(),
921 );
922 let after_count = compacted_messages.len();
923 let recovered = after_tokens <= target_budget && after_tokens < before_tokens;
924
925 if recovered {
926 if let Some(job) = self.child_job()
927 && let Err(error) = job
928 .before_replace(&self.session.messages, &compacted_messages)
929 .await
930 {
931 self.cancel_token.cancel();
932 tracing::error!(%error, "child recovery projection failed; original Session retained");
933 return false;
934 }
935 self.session.replace_messages(compacted_messages);
936 turn.clear_parent_input_tokens();
937 if let Some(pm) = self.session.prefix_stability.as_mut() {
938 pm.note_history_reset("compaction");
939 }
940 self.commit_compaction_checkpoint(summary_prompt);
941 self.emit_session_updated().await;
942 let removed = before_count.saturating_sub(after_count);
943 let mut details = format!(
944 "Made room: {before_count} → {after_count} messages ({removed} removed), ~{before_tokens} → ~{after_tokens} tokens"
945 );
946 if retries_used > 0 {
947 details.push_str(&format!(" ({retries_used} retries)"));
948 }
949 self.emit_compaction_completed(
950 id.clone(),
951 true,
952 details.clone(),
953 Some(before_count),
954 Some(after_count),
955 CompactionPass {
956 trigger: "emergency",
957 path,
958 tokens_before: before_tokens,
959 threshold_tokens: prepared.config.token_threshold,
960 usage: compaction_usage.clone(),
961 },
962 )
963 .await;
964 let _ = self.send_event(Event::status(details)).await;
965 self.finish_compaction(&id);
966 return true;
967 }
968
969 // Two distinct failures were previously conflated into one banner.
970 // When the provider rejected the request (its bill counts framing we
971 // cannot see), our estimate may already sit within the budget while
972 // the pass removed nothing — reporting that as "failed to reduce
973 // below model limit" with an estimate printed *under* the budget
974 // reads as self-contradictory. Name the actual outcome instead.
975 let message = if after_tokens > target_budget {
976 format!(
977 "Making room failed: the request is still over the model limit \
978 (estimate ~{after_tokens} tokens, budget ~{target_budget}). Original conversation was preserved."
979 )
980 } else {
981 format!(
982 "Making room made no progress (estimate ~{after_tokens} tokens \
983 is already within the ~{target_budget} budget; the provider may count the \
984 request differently). Original conversation was preserved."
985 )
986 };
987 self.emit_compaction_failed(id.clone(), true, message.clone())
988 .await;
989 let _ = self.send_event(Event::status(message)).await;
990 self.finish_compaction(&id);
991 false
992 }
993
994 /// Keep the rendered checkpoint for host persistence and repeat-compaction
995 /// metadata. The model sees the checkpoint exactly once through ordinary
996 /// conversation history; the stable system prefix never carries it.
997 pub(super) fn commit_compaction_checkpoint(&mut self, summary_prompt: Option<SystemPrompt>) {
998 let Some(summary_prompt) = summary_prompt else {
999 return;
1000 };
1001 self.session.compaction_summary_prompt = Some(summary_prompt);
1002 }
1003
1004 /// Capture the current session-owned Agent topology at the replacement
1005 /// history boundary. This is the Codewhale equivalent of Codex clearing
1006 /// its world-state reference after standalone compaction so the next turn
1007 /// receives fresh environment/subagent context instead of trusting the
1008 /// narrative summary as live process state.
1009 pub(super) async fn append_compaction_agent_topology(&self, messages: &mut Vec<Message>) {
1010 let snapshots = {
1011 let manager = self.subagent_manager.read().await;
1012 manager.list_for_session(&self.session.id)
1013 };
1014 crate::runtime_handoff::replace_agent_topology_checkpoint(messages, &snapshots);
1015 }
1016 }
1017
1018 /// A context-recovery failure that came from the provider refusing the
1019 /// request (capability, auth, reachability, quota) rather than from the
1020 /// summary itself. Context-length rejections are excluded: those really are
1021 /// the budget problem the caller already reports.
1022 pub(super) fn is_provider_rejection(err: &anyhow::Error) -> bool {
1023 use crate::error_taxonomy::{ErrorCategory, classify_error_message};
1024 let text = format!("{err:#}");
1025 if super::context::is_context_length_error_message(&text)
1026 || matches!(
1027 err.downcast_ref::<crate::llm_client::LlmError>(),
1028 Some(crate::llm_client::LlmError::ContextLengthError(_))
1029 )
1030 {
1031 return false;
1032 }
1033 err.downcast_ref::<crate::llm_client::LlmError>().is_some()
1034 || matches!(
1035 classify_error_message(&text),
1036 ErrorCategory::Authentication
1037 | ErrorCategory::Authorization
1038 | ErrorCategory::Network
1039 | ErrorCategory::RateLimit
1040 | ErrorCategory::Timeout
1041 )
1042 }
1043
1044 #[cfg(test)]
1045 mod tests {
1046 use super::*;
1047 use crate::compaction::CompactionNoticeSink as _;
1048
1049 /// The engine sink is the one link between a compaction downgrade and the
1050 /// person watching: the notice must land on the status line, not only in
1051 /// the log.
1052 #[tokio::test]
1053 async fn compaction_notice_sink_delivers_a_status_event() {
1054 let (tx, mut rx) = mpsc::channel(4);
1055 let sink = EngineCompactionNoticeSink { tx, child: None };
1056 sink.notice("Making room re-encoded 2 inline image(s)".to_string());
1057 match rx.recv().await {
1058 Some(Event::Status { message }) => {
1059 assert!(message.contains("re-encoded"), "{message}");
1060 }
1061 other => panic!("expected a Status event, got {other:?}"),
1062 }
1063 }
1064 }
1065
1065 lines RUST