返回 CodeWhale
compaction.rs
根目录 / crates / tui / src / core / engine / tests / compaction.rs
1 use super::*;
2
3 #[tokio::test]
4 async fn rejected_manual_compaction_route_closes_typed_lifecycle() {
5 let _env_lock = lock_test_env();
6 let _api_key = EnvVarGuard::remove("DEEPSEEK_API_KEY");
7 let route_config = Config {
8 provider: Some("deepseek".to_string()),
9 default_text_model: Some(crate::config::DEFAULT_TEXT_MODEL.to_string()),
10 ..Config::default()
11 }
12 .with_legacy_root(Some(String::new()), None);
13 let route = resolve_runtime_route(
14 &route_config,
15 ProviderKind::Deepseek,
16 Some(crate::config::DEFAULT_TEXT_MODEL),
17 )
18 .expect("structurally resolve route without credential");
19 assert!(
20 route.clone().validate().is_err(),
21 "fixture must fail at engine route installation"
22 );
23 let (mut engine, handle) = Engine::new(EngineConfig::default(), &route_config);
24
25 engine
26 .handle_manual_compaction_op(
27 "compact-route-invalid".to_string(),
28 route,
29 CompactionConfig::default(),
30 )
31 .await;
32
33 let mut started_id = None;
34 let mut failed_id = None;
35 let mut order = Vec::new();
36 let mut events = handle.rx_event.write().await;
37 while let Ok(event) = events.try_recv() {
38 match event {
39 Event::CompactionStarted { id, auto, .. } => {
40 assert!(!auto);
41 started_id = Some(id);
42 order.push("started");
43 }
44 Event::CompactionFailed { id, auto, message } => {
45 assert!(!auto);
46 assert!(message.contains("provider route is not ready"));
47 failed_id = Some(id);
48 order.push("failed");
49 }
50 Event::Error { .. } => order.push("error"),
51 _ => {}
52 }
53 }
54 assert_eq!(order, ["started", "failed", "error"]);
55 assert_eq!(started_id, failed_id);
56 }
57
58 #[tokio::test]
59 async fn queued_manual_compaction_cancellation_is_idempotent_and_skips_route_activation() {
60 let _env_lock = lock_test_env();
61 let _api_key = EnvVarGuard::remove("DEEPSEEK_API_KEY");
62 let route_config = Config {
63 provider: Some("deepseek".to_string()),
64 default_text_model: Some(crate::config::DEFAULT_TEXT_MODEL.to_string()),
65 ..Config::default()
66 }
67 .with_legacy_root(Some(String::new()), None);
68 let route = resolve_runtime_route(
69 &route_config,
70 ProviderKind::Deepseek,
71 Some(crate::config::DEFAULT_TEXT_MODEL),
72 )
73 .expect("structurally resolve route without credential");
74 let (mut engine, handle) = Engine::new(EngineConfig::default(), &route_config);
75 let id = "compact-cancel-before-start";
76
77 handle.cancel_compaction(id).expect("first cancel accepted");
78 handle
79 .cancel_compaction(id)
80 .expect("replayed cancel remains idempotent");
81 engine
82 .handle_manual_compaction_op(id.to_string(), route, CompactionConfig::default())
83 .await;
84
85 let mut events = handle.rx_event.write().await;
86 let drained = std::iter::from_fn(|| events.try_recv().ok()).collect::<Vec<_>>();
87 assert!(matches!(
88 drained.as_slice(),
89 [
90 Event::CompactionStarted { id: started, auto: false, .. },
91 Event::CompactionCancelled { id: cancelled, auto: false, .. },
92 Event::TurnComplete { status: TurnOutcomeStatus::Interrupted, .. }
93 ] if started == id && cancelled == id
94 ));
95 assert!(
96 !drained
97 .iter()
98 .any(|event| matches!(event, Event::Error { .. })),
99 "pre-start cancellation must not activate or validate the provider route"
100 );
101
102 let retry = engine
103 .claim_compaction(id)
104 .expect("the same stable id can be retried after terminal settlement");
105 assert!(!retry.is_cancelled());
106 handle
107 .cancel_compaction(id)
108 .expect("running cancel accepted");
109 assert!(
110 retry.is_cancelled(),
111 "running cancellation reaches its token"
112 );
113 engine.finish_compaction(id);
114 }
115
116 struct BlockingEmergencyCompactionModelClient {
117 entered: std::sync::Arc<tokio::sync::Notify>,
118 request_dropped: std::sync::Arc<std::sync::atomic::AtomicBool>,
119 }
120
121 #[tokio::test]
122 async fn failed_emergency_compaction_preserves_history_instead_of_trimming() {
123 use crate::llm_client::mock::MockLlmClient;
124 let _env_lock = lock_test_env();
125 let workspace = tempdir().unwrap();
126 // Emergency compaction persists a checkpoint before calling the model.
127 // Keep that prerequisite away from other tests' shared state fixtures so
128 // this exercises a summary failure, rather than an unrelated write failure.
129 let _home = EnvVarGuard::set("CODEWHALE_HOME", workspace.path());
130 let (mut engine, handle) = Engine::new(
131 deterministic_engine_config(workspace.path()),
132 &Config::default(),
133 );
134 engine.session.messages = (0..12)
135 .map(|i| Message {
136 role: if i % 2 == 0 {
137 Role::User
138 } else {
139 Role::Assistant
140 },
141 content: vec![ContentBlock::Text {
142 text: format!(
143 "must preserve instruction and evidence {i}: {}",
144 "x".repeat(20_000)
145 ),
146 cache_control: None,
147 }],
148 })
149 .collect::<Vec<_>>()
150 .into();
151 let before = engine.session.messages.clone();
152 let summary = engine.session.compaction_summary_prompt.clone();
153 let client = MockLlmClient::new(Vec::new());
154 let tools = vec![catalog_tool("read")];
155 let mut turn = TurnContext::new(1);
156 assert!(
157 !engine
158 .recover_context_overflow(
159 &client,
160 Some(&tools),
161 "provider rejection fixture",
162 &mut turn
163 )
164 .await
165 );
166 assert_eq!(engine.session.messages.as_slice(), before.as_slice());
167 assert_eq!(engine.session.compaction_summary_prompt, summary);
168 let mut events = handle.rx_event.write().await;
169 let drained = std::iter::from_fn(|| events.try_recv().ok()).collect::<Vec<_>>();
170 assert_eq!(
171 client.call_count(),
172 1,
173 "a deterministic summary failure is not retried unchanged: {drained:?}"
174 );
175 let requests = client.captured_requests();
176 assert_eq!(requests[0].tools.as_deref(), Some(tools.as_slice()));
177 assert_eq!(requests[0].tool_choice, Some(json!("none")));
178 assert_eq!(requests[0].system, engine.session.system_prompt);
179 }
180
181 #[tokio::test]
182 async fn manual_compaction_accounts_accepted_and_rejected_responses_once() {
183 use wiremock::matchers::{method, path};
184 use wiremock::{Mock, MockServer, ResponseTemplate};
185
186 let _env_lock = lock_test_env();
187 let _cost_scope = crate::cost_status::test_scope();
188 let workspace = tempdir().expect("isolated compaction workspace");
189 let _home = EnvVarGuard::set("CODEWHALE_HOME", workspace.path());
190 for (finish_reason, expected_status) in [
191 ("stop", TurnOutcomeStatus::Completed),
192 ("length", TurnOutcomeStatus::Failed),
193 ] {
194 let server = MockServer::start().await;
195 Mock::given(method("POST"))
196 .and(path("/v1/chat/completions"))
197 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
198 "id": format!("compaction-{finish_reason}"),
199 "object": "chat.completion",
200 "model": crate::config::DEFAULT_TEXT_MODEL,
201 "choices": [{
202 "index": 0,
203 "message": { "role": "assistant", "content": "Primary request: preserve the session migration. Completed: inspected the existing store. Constraints: keep every user message and failing test. Next: finish the transactional migration and rerun session_store::roundtrip." },
204 "finish_reason": finish_reason,
205 }],
206 "usage": { "prompt_tokens": 41, "completion_tokens": 7, "total_tokens": 48 },
207 })))
208 .expect(1)
209 .mount(&server)
210 .await;
211 let route_config = Config {
212 provider: Some("deepseek".to_string()),
213 ..Config::default()
214 }
215 .with_legacy_root(
216 Some("fixture-key".to_string()),
217 Some(format!("{}/v1", server.uri())),
218 );
219 let (mut engine, handle) = Engine::new(
220 EngineConfig {
221 workspace: workspace.path().to_path_buf(),
222 snapshots_enabled: false,
223 subagents_enabled: false,
224 ..EngineConfig::default()
225 },
226 &route_config,
227 );
228 engine.session.messages.push(Message {
229 role: Role::User,
230 content: vec![ContentBlock::Text {
231 text: "Preserve the transactional session migration.".to_string(),
232 cache_control: None,
233 }],
234 });
235 engine.config.goal_state.lock().unwrap().replace(
236 "Finish the session migration",
237 Some(1000),
238 None,
239 );
240 engine
241 .handle_manual_compaction("compact-accounting".to_string(), CancellationToken::new())
242 .await;
243
244 assert_eq!(engine.session.total_usage.input_tokens, 41);
245 assert_eq!(engine.session.total_usage.output_tokens, 7);
246 assert_eq!(
247 engine
248 .config
249 .goal_state
250 .lock()
251 .unwrap()
252 .snapshot()
253 .tokens_used,
254 48
255 );
256 let mut events = handle.rx_event.write().await;
257 let mut telemetry_count = 0;
258 let mut terminal_count = 0;
259 while let Ok(event) = events.try_recv() {
260 match event {
261 Event::RoutedTurnUsage { usage, .. } => {
262 telemetry_count += 1;
263 assert_eq!((usage.input_tokens, usage.output_tokens), (41, 7));
264 }
265 Event::TurnComplete {
266 usage,
267 parent_route_usage,
268 status,
269 ..
270 } => {
271 terminal_count += 1;
272 assert_eq!((usage.input_tokens, usage.output_tokens), (41, 7));
273 assert_eq!(parent_route_usage, Usage::default());
274 assert_eq!(status, expected_status);
275 }
276 _ => {}
277 }
278 }
279 assert_eq!((telemetry_count, terminal_count), (1, 1));
280 }
281 }
282
283 #[async_trait::async_trait]
284 impl crate::core::model_client::ModelClient for BlockingEmergencyCompactionModelClient {
285 fn provider_name(&self) -> &str {
286 "deepseek"
287 }
288
289 fn model(&self) -> &str {
290 crate::config::DEFAULT_TEXT_MODEL
291 }
292
293 async fn create_message(
294 &self,
295 _request: codewhale_models::MessageRequest,
296 ) -> anyhow::Result<codewhale_models::MessageResponse> {
297 let _drop_signal = DropSignal(std::sync::Arc::clone(&self.request_dropped));
298 self.entered.notify_one();
299 std::future::pending().await
300 }
301
302 async fn create_message_stream(
303 &self,
304 _request: codewhale_models::MessageRequest,
305 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
306 anyhow::bail!("emergency compaction uses the non-streaming model boundary")
307 }
308
309 async fn health_check(&self) -> anyhow::Result<bool> {
310 Ok(true)
311 }
312 }
313
314 #[tokio::test]
315 async fn emergency_compaction_cancellation_drops_provider_and_never_mutates_context() {
316 let route_config = Config {
317 provider: Some("deepseek".to_string()),
318 default_text_model: Some(crate::config::DEFAULT_TEXT_MODEL.to_string()),
319 ..Config::default()
320 };
321 let (mut engine, handle) = Engine::new(EngineConfig::default(), &route_config);
322 engine.session.messages = (0..8)
323 .map(|index| Message {
324 role: if index % 2 == 0 {
325 Role::User
326 } else {
327 Role::Assistant
328 },
329 content: vec![ContentBlock::Text {
330 text: format!("preserve emergency context item {index}"),
331 cache_control: None,
332 }],
333 })
334 .collect::<Vec<_>>()
335 .into();
336 let messages_before = engine.session.messages.clone();
337 let checkpoint_before = engine.session.compaction_summary_prompt.clone();
338 let entered = std::sync::Arc::new(tokio::sync::Notify::new());
339 let request_dropped = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
340 let client = std::sync::Arc::new(BlockingEmergencyCompactionModelClient {
341 entered: std::sync::Arc::clone(&entered),
342 request_dropped: std::sync::Arc::clone(&request_dropped),
343 });
344
345 let recovery = tokio::spawn(async move {
346 let mut turn = TurnContext::new(1);
347 let recovered = engine
348 .recover_context_overflow(client.as_ref(), None, "cancellation regression", &mut turn)
349 .await;
350 (engine, recovered)
351 });
352
353 let started_id = tokio::time::timeout(Duration::from_secs(1), async {
354 loop {
355 let event = handle
356 .rx_event
357 .write()
358 .await
359 .recv()
360 .await
361 .expect("emergency compaction start event");
362 if let Event::CompactionStarted { id, auto: true, .. } = event {
363 break id;
364 }
365 }
366 })
367 .await
368 .expect("emergency compaction publishes its stable id");
369 tokio::time::timeout(Duration::from_secs(1), entered.notified())
370 .await
371 .expect("emergency provider request starts");
372
373 handle
374 .cancel_compaction(started_id.clone())
375 .expect("exact emergency cancellation accepted");
376 let (engine, recovered) = tokio::time::timeout(Duration::from_secs(1), recovery)
377 .await
378 .expect("emergency cancellation settles promptly")
379 .expect("recovery task");
380
381 assert!(!recovered);
382 assert_eq!(&*engine.session.messages, &*messages_before);
383 assert_eq!(engine.session.compaction_summary_prompt, checkpoint_before);
384 assert!(
385 request_dropped.load(std::sync::atomic::Ordering::SeqCst),
386 "cancellation must drop the in-flight provider future"
387 );
388
389 let mut events = handle.rx_event.write().await;
390 let drained = std::iter::from_fn(|| events.try_recv().ok()).collect::<Vec<_>>();
391 assert!(matches!(
392 drained.as_slice(),
393 [Event::CompactionCancelled { id, auto: true, .. }] if id == &started_id
394 ));
395 assert!(
396 !drained.iter().any(|event| matches!(
397 event,
398 Event::CompactionCompleted { .. } | Event::CompactionFailed { .. }
399 )),
400 "a canceled emergency pass must have one canceled terminal event"
401 );
402 }
403
404 /// Experience mark 2: a one-message conversation has nothing to summarize.
405 /// Emergency recovery must not start a pass (no spinner, no model call)
406 /// before the failure the caller reports anyway.
407 #[tokio::test]
408 async fn emergency_recovery_skips_a_history_with_nothing_to_compact() {
409 use crate::llm_client::mock::MockLlmClient;
410 let _env_lock = lock_test_env();
411 let workspace = tempdir().unwrap();
412 let _home = EnvVarGuard::set("CODEWHALE_HOME", workspace.path());
413 let (mut engine, handle) = Engine::new(
414 deterministic_engine_config(workspace.path()),
415 &Config::default(),
416 );
417 engine.session.messages = vec![Message {
418 role: Role::User,
419 content: vec![ContentBlock::Text {
420 text: "hello".to_string(),
421 cache_control: None,
422 }],
423 }]
424 .into();
425 let client = MockLlmClient::new(Vec::new());
426 let mut turn = TurnContext::new(1);
427 assert!(
428 !engine
429 .recover_context_overflow(&client, None, "preflight token budget", &mut turn)
430 .await
431 );
432 assert_eq!(client.call_count(), 0, "no summary request for one message");
433 assert_eq!(turn.stop_diagnostics.emergency_compaction_attempts, 0);
434 let mut events = handle.rx_event.write().await;
435 let drained = std::iter::from_fn(|| events.try_recv().ok()).collect::<Vec<_>>();
436 assert!(
437 !drained.iter().any(|event| matches!(
438 event,
439 Event::CompactionStarted { .. } | Event::CompactionFailed { .. }
440 )),
441 "{drained:?}"
442 );
443 }
444
445 #[test]
446 fn recovery_failures_from_the_provider_are_told_apart_from_budget_failures() {
447 use super::super::compaction::is_provider_rejection;
448 use crate::llm_client::LlmError;
449 assert!(is_provider_rejection(&anyhow::Error::new(
450 LlmError::ModelError("\"nomic-embed-text:latest\" does not support chat".to_string())
451 )));
452 assert!(is_provider_rejection(&anyhow::anyhow!(
453 "connection refused while contacting http://localhost:11434"
454 )));
455 assert!(!is_provider_rejection(&anyhow::Error::new(
456 LlmError::ContextLengthError("prompt is too long".to_string())
457 )));
458 assert!(!is_provider_rejection(&anyhow::anyhow!(
459 "Making room did not shrink the context; the original conversation was preserved."
460 )));
461 }
462
463 #[test]
464 fn a_request_that_cannot_fit_names_the_cause_and_one_next_step() {
465 use super::super::context::context_does_not_fit_message;
466 let embed =
467 context_does_not_fit_message(true, true, "nomic-embed-text:latest", 5_200, 1_500, 5_100);
468 assert_eq!(
469 embed,
470 "nomic-embed-text:latest can't chat. Pick a chat model: /model."
471 );
472 let window = context_does_not_fit_message(true, true, "qwen3:4b", 5_300, 3_000, 5_100);
473 assert!(
474 window.contains("qwen3:4b's context window (~3000 tokens usable)"),
475 "{window}"
476 );
477 assert!(
478 window.contains("working instructions (~5100 tokens)"),
479 "{window}"
480 );
481 assert!(window.ends_with("raise num_ctx: /model."), "{window}");
482 assert!(!window.contains("compaction"), "{window}");
483 let message = context_does_not_fit_message(false, false, "small-model", 9_000, 6_000, 2_000);
484 assert!(
485 message.contains("there is not enough earlier conversation to summarize"),
486 "{message}"
487 );
488 assert!(message.ends_with("choose a larger model."), "{message}");
489 assert!(
490 !message.contains("/model"),
491 "headless has no command layer: {message}"
492 );
493 }
494
494 lines RUST