| 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 | api_key: Some(String::new()), |
| 10 | default_text_model: Some(crate::config::DEFAULT_TEXT_MODEL.to_string()), |
| 11 | ..Config::default() |
| 12 | }; |
| 13 | let route = resolve_runtime_route( |
| 14 | &route_config, |
| 15 | ApiProvider::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 | api_key: Some(String::new()), |
| 65 | default_text_model: Some(crate::config::DEFAULT_TEXT_MODEL.to_string()), |
| 66 | ..Config::default() |
| 67 | }; |
| 68 | let route = resolve_runtime_route( |
| 69 | &route_config, |
| 70 | ApiProvider::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 | api_key: Some("fixture-key".to_string()), |
| 214 | base_url: Some(format!("{}/v1", server.uri())), |
| 215 | ..Config::default() |
| 216 | }; |
| 217 | let (mut engine, handle) = Engine::new( |
| 218 | EngineConfig { |
| 219 | workspace: workspace.path().to_path_buf(), |
| 220 | snapshots_enabled: false, |
| 221 | subagents_enabled: false, |
| 222 | ..EngineConfig::default() |
| 223 | }, |
| 224 | &route_config, |
| 225 | ); |
| 226 | engine.session.messages.push(Message { |
| 227 | role: Role::User, |
| 228 | content: vec![ContentBlock::Text { |
| 229 | text: "Preserve the transactional session migration.".to_string(), |
| 230 | cache_control: None, |
| 231 | }], |
| 232 | }); |
| 233 | engine.config.goal_state.lock().unwrap().replace( |
| 234 | "Finish the session migration", |
| 235 | Some(1000), |
| 236 | None, |
| 237 | ); |
| 238 | engine |
| 239 | .handle_manual_compaction("compact-accounting".to_string(), CancellationToken::new()) |
| 240 | .await; |
| 241 | |
| 242 | assert_eq!(engine.session.total_usage.input_tokens, 41); |
| 243 | assert_eq!(engine.session.total_usage.output_tokens, 7); |
| 244 | assert_eq!( |
| 245 | engine |
| 246 | .config |
| 247 | .goal_state |
| 248 | .lock() |
| 249 | .unwrap() |
| 250 | .snapshot() |
| 251 | .tokens_used, |
| 252 | 48 |
| 253 | ); |
| 254 | let mut events = handle.rx_event.write().await; |
| 255 | let mut telemetry_count = 0; |
| 256 | let mut terminal_count = 0; |
| 257 | while let Ok(event) = events.try_recv() { |
| 258 | match event { |
| 259 | Event::RoutedTurnUsage { usage, .. } => { |
| 260 | telemetry_count += 1; |
| 261 | assert_eq!((usage.input_tokens, usage.output_tokens), (41, 7)); |
| 262 | } |
| 263 | Event::TurnComplete { |
| 264 | usage, |
| 265 | parent_route_usage, |
| 266 | status, |
| 267 | .. |
| 268 | } => { |
| 269 | terminal_count += 1; |
| 270 | assert_eq!((usage.input_tokens, usage.output_tokens), (41, 7)); |
| 271 | assert_eq!(parent_route_usage, Usage::default()); |
| 272 | assert_eq!(status, expected_status); |
| 273 | } |
| 274 | _ => {} |
| 275 | } |
| 276 | } |
| 277 | assert_eq!((telemetry_count, terminal_count), (1, 1)); |
| 278 | } |
| 279 | } |
| 280 | |
| 281 | #[async_trait::async_trait] |
| 282 | impl crate::core::model_client::ModelClient for BlockingEmergencyCompactionModelClient { |
| 283 | fn provider_name(&self) -> &str { |
| 284 | "deepseek" |
| 285 | } |
| 286 | |
| 287 | fn model(&self) -> &str { |
| 288 | crate::config::DEFAULT_TEXT_MODEL |
| 289 | } |
| 290 | |
| 291 | async fn create_message( |
| 292 | &self, |
| 293 | _request: codewhale_models::MessageRequest, |
| 294 | ) -> anyhow::Result<codewhale_models::MessageResponse> { |
| 295 | let _drop_signal = DropSignal(std::sync::Arc::clone(&self.request_dropped)); |
| 296 | self.entered.notify_one(); |
| 297 | std::future::pending().await |
| 298 | } |
| 299 | |
| 300 | async fn create_message_stream( |
| 301 | &self, |
| 302 | _request: codewhale_models::MessageRequest, |
| 303 | ) -> anyhow::Result<crate::llm_client::StreamEventBox> { |
| 304 | anyhow::bail!("emergency compaction uses the non-streaming model boundary") |
| 305 | } |
| 306 | |
| 307 | async fn health_check(&self) -> anyhow::Result<bool> { |
| 308 | Ok(true) |
| 309 | } |
| 310 | } |
| 311 | |
| 312 | #[tokio::test] |
| 313 | async fn emergency_compaction_cancellation_drops_provider_and_never_mutates_context() { |
| 314 | let route_config = Config { |
| 315 | provider: Some("deepseek".to_string()), |
| 316 | default_text_model: Some(crate::config::DEFAULT_TEXT_MODEL.to_string()), |
| 317 | ..Config::default() |
| 318 | }; |
| 319 | let (mut engine, handle) = Engine::new(EngineConfig::default(), &route_config); |
| 320 | engine.session.messages = (0..8) |
| 321 | .map(|index| Message { |
| 322 | role: if index % 2 == 0 { |
| 323 | Role::User |
| 324 | } else { |
| 325 | Role::Assistant |
| 326 | }, |
| 327 | content: vec![ContentBlock::Text { |
| 328 | text: format!("preserve emergency context item {index}"), |
| 329 | cache_control: None, |
| 330 | }], |
| 331 | }) |
| 332 | .collect::<Vec<_>>() |
| 333 | .into(); |
| 334 | let messages_before = engine.session.messages.clone(); |
| 335 | let checkpoint_before = engine.session.compaction_summary_prompt.clone(); |
| 336 | let entered = std::sync::Arc::new(tokio::sync::Notify::new()); |
| 337 | let request_dropped = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); |
| 338 | let client = std::sync::Arc::new(BlockingEmergencyCompactionModelClient { |
| 339 | entered: std::sync::Arc::clone(&entered), |
| 340 | request_dropped: std::sync::Arc::clone(&request_dropped), |
| 341 | }); |
| 342 | |
| 343 | let recovery = tokio::spawn(async move { |
| 344 | let mut turn = TurnContext::new(1); |
| 345 | let recovered = engine |
| 346 | .recover_context_overflow(client.as_ref(), None, "cancellation regression", &mut turn) |
| 347 | .await; |
| 348 | (engine, recovered) |
| 349 | }); |
| 350 | |
| 351 | let started_id = tokio::time::timeout(Duration::from_secs(1), async { |
| 352 | loop { |
| 353 | let event = handle |
| 354 | .rx_event |
| 355 | .write() |
| 356 | .await |
| 357 | .recv() |
| 358 | .await |
| 359 | .expect("emergency compaction start event"); |
| 360 | if let Event::CompactionStarted { id, auto: true, .. } = event { |
| 361 | break id; |
| 362 | } |
| 363 | } |
| 364 | }) |
| 365 | .await |
| 366 | .expect("emergency compaction publishes its stable id"); |
| 367 | tokio::time::timeout(Duration::from_secs(1), entered.notified()) |
| 368 | .await |
| 369 | .expect("emergency provider request starts"); |
| 370 | |
| 371 | handle |
| 372 | .cancel_compaction(started_id.clone()) |
| 373 | .expect("exact emergency cancellation accepted"); |
| 374 | let (engine, recovered) = tokio::time::timeout(Duration::from_secs(1), recovery) |
| 375 | .await |
| 376 | .expect("emergency cancellation settles promptly") |
| 377 | .expect("recovery task"); |
| 378 | |
| 379 | assert!(!recovered); |
| 380 | assert_eq!(&*engine.session.messages, &*messages_before); |
| 381 | assert_eq!(engine.session.compaction_summary_prompt, checkpoint_before); |
| 382 | assert!( |
| 383 | request_dropped.load(std::sync::atomic::Ordering::SeqCst), |
| 384 | "cancellation must drop the in-flight provider future" |
| 385 | ); |
| 386 | |
| 387 | let mut events = handle.rx_event.write().await; |
| 388 | let drained = std::iter::from_fn(|| events.try_recv().ok()).collect::<Vec<_>>(); |
| 389 | assert!(matches!( |
| 390 | drained.as_slice(), |
| 391 | [Event::CompactionCancelled { id, auto: true, .. }] if id == &started_id |
| 392 | )); |
| 393 | assert!( |
| 394 | !drained.iter().any(|event| matches!( |
| 395 | event, |
| 396 | Event::CompactionCompleted { .. } | Event::CompactionFailed { .. } |
| 397 | )), |
| 398 | "a canceled emergency pass must have one canceled terminal event" |
| 399 | ); |
| 400 | } |
| 401 |