| 1 | const WORKING_SET_SUMMARY_MARKER: &str = "## Repo Working Set"; |
| 2 | |
| 3 | #[tokio::test] |
| 4 | async fn event_capacity_cancel_before_admission_has_no_started_turn_and_keeps_classifier_cost() { |
| 5 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 6 | let _cost = crate::cost_status::test_scope(); |
| 7 | let workspace = tempdir().unwrap(); |
| 8 | let config = Config::default(); |
| 9 | let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn( |
| 10 | "must not run", |
| 11 | )])); |
| 12 | let mut engine_config = deterministic_engine_config(workspace.path()); |
| 13 | engine_config.features.disable(Feature::Mcp); |
| 14 | let (mut engine, handle) = Engine::new_with_model_client(engine_config, &config, mock.clone()); |
| 15 | let (tx, rx) = mpsc::channel(1); |
| 16 | engine.tx_event = tx; |
| 17 | let mut handle = handle; |
| 18 | handle.rx_event = Arc::new(RwLock::new(rx)); |
| 19 | engine |
| 20 | .tx_event |
| 21 | .try_send(Event::status("existing idle receipt")) |
| 22 | .unwrap(); |
| 23 | let mut op = external_user_message_op("not admitted", AppMode::Agent, &config); |
| 24 | let Op::SendMessage(spec) = &mut op else { |
| 25 | unreachable!() |
| 26 | }; |
| 27 | spec.initial_routed_usage |
| 28 | .records |
| 29 | .push(crate::cost_status::RuntimeUsageRecord { |
| 30 | source_id: "event-capacity:classifier".into(), |
| 31 | usage: crate::cost_status::EffectiveRouteUsage { |
| 32 | route: crate::cost_status::EffectiveRouteEnvelope::capture( |
| 33 | None, |
| 34 | ProviderKind::Openai, |
| 35 | "openai", |
| 36 | "classifier", |
| 37 | None, |
| 38 | chrono::Utc::now(), |
| 39 | ), |
| 40 | usage: Usage { |
| 41 | input_tokens: 7, |
| 42 | output_tokens: 3, |
| 43 | ..Usage::default() |
| 44 | }, |
| 45 | }, |
| 46 | }); |
| 47 | let controls = Arc::clone(&engine.turn_controls); |
| 48 | handle.send(op).await.unwrap(); |
| 49 | let task = tokio::spawn(engine.run()); |
| 50 | tokio::time::timeout(Duration::from_secs(2), async { |
| 51 | loop { |
| 52 | if controls.lock().unwrap().active.is_some() { |
| 53 | break; |
| 54 | } |
| 55 | tokio::task::yield_now().await; |
| 56 | } |
| 57 | }) |
| 58 | .await |
| 59 | .expect("queued op enters the existing control scope"); |
| 60 | handle.cancel(); |
| 61 | // The oneshot snapshot can settle while the event receiver remains full. |
| 62 | let snapshot = tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot()) |
| 63 | .await |
| 64 | .expect("cancel releases admission reservation") |
| 65 | .unwrap(); |
| 66 | assert!( |
| 67 | snapshot.messages.is_empty(), |
| 68 | "no session mutation before admission" |
| 69 | ); |
| 70 | assert_eq!(mock.call_count(), 0); |
| 71 | assert!(controls.lock().unwrap().active.is_none()); |
| 72 | let cost = crate::cost_status::drain(); |
| 73 | assert!( |
| 74 | cost.usage_source_fingerprints |
| 75 | .contains(&crate::cost_status::usage_source_fingerprint( |
| 76 | "event-capacity:classifier" |
| 77 | ),), |
| 78 | "already billed classifier work remains accounted" |
| 79 | ); |
| 80 | let mut events = handle.rx_event.write().await; |
| 81 | assert!(matches!(events.try_recv(), Ok(Event::Status { .. }))); |
| 82 | assert!( |
| 83 | events.try_recv().is_err(), |
| 84 | "unadmitted work has no fabricated terminal event" |
| 85 | ); |
| 86 | drop(events); |
| 87 | handle.send(Op::Shutdown).await.unwrap(); |
| 88 | tokio::time::timeout(Duration::from_secs(2), task) |
| 89 | .await |
| 90 | .unwrap() |
| 91 | .unwrap(); |
| 92 | } |
| 93 | |
| 94 | #[tokio::test] |
| 95 | async fn event_capacity_admitted_full_channel_settles_once_with_partial_usage_without_drain() { |
| 96 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 97 | let _cost = crate::cost_status::test_scope(); |
| 98 | let workspace = tempdir().unwrap(); |
| 99 | let config = Config::default(); |
| 100 | let mock = Arc::new(MockLlmClient::new(Vec::new())); |
| 101 | let mut engine_config = deterministic_engine_config(workspace.path()); |
| 102 | engine_config.features.disable(Feature::Mcp); |
| 103 | let (engine, handle) = Engine::new_with_model_client(engine_config, &config, mock.clone()); |
| 104 | let tx = engine.tx_event.clone(); |
| 105 | let entered = Arc::new(tokio::sync::Notify::new()); |
| 106 | let filled = Arc::clone(&entered); |
| 107 | mock.push_factory(move |_| { |
| 108 | while tx.try_send(Event::status("backpressure receipt")).is_ok() {} |
| 109 | filled.notify_one(); |
| 110 | vec![ |
| 111 | canned::message_start("billed-before-cancel"), |
| 112 | canned::message_delta( |
| 113 | "end_turn", |
| 114 | Some(Usage { |
| 115 | input_tokens: 11, |
| 116 | output_tokens: 5, |
| 117 | ..Usage::default() |
| 118 | }), |
| 119 | ), |
| 120 | canned::text_block_start(0), |
| 121 | canned::text_delta(0, "must not render after cancellation"), |
| 122 | canned::message_stop(), |
| 123 | ] |
| 124 | }); |
| 125 | let controls = Arc::clone(&engine.turn_controls); |
| 126 | let task = tokio::spawn(engine.run()); |
| 127 | handle |
| 128 | .send(external_user_message_op( |
| 129 | "admitted", |
| 130 | AppMode::Agent, |
| 131 | &config, |
| 132 | )) |
| 133 | .await |
| 134 | .unwrap(); |
| 135 | tokio::time::timeout(Duration::from_secs(2), entered.notified()) |
| 136 | .await |
| 137 | .unwrap(); |
| 138 | handle.cancel(); |
| 139 | let snapshot = tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot()) |
| 140 | .await |
| 141 | .expect("admitted cancellation settles before any event drain") |
| 142 | .unwrap(); |
| 143 | assert_eq!(mock.call_count(), 1); |
| 144 | assert_eq!( |
| 145 | snapshot.total_tokens, 16, |
| 146 | "known partial provider usage survives cancellation" |
| 147 | ); |
| 148 | assert!(controls.lock().unwrap().active.is_none()); |
| 149 | let mut events = handle.rx_event.write().await; |
| 150 | let mut started = 0; |
| 151 | let mut completed = 0; |
| 152 | let mut terminal_was_last = false; |
| 153 | while let Ok(event) = events.try_recv() { |
| 154 | assert!( |
| 155 | !terminal_was_last, |
| 156 | "completion stays after prior accepted observations" |
| 157 | ); |
| 158 | match event { |
| 159 | Event::TurnStarted { .. } => started += 1, |
| 160 | Event::TurnComplete { |
| 161 | status, |
| 162 | usage, |
| 163 | parent_route_usage, |
| 164 | error, |
| 165 | .. |
| 166 | } => { |
| 167 | completed += 1; |
| 168 | terminal_was_last = true; |
| 169 | assert_eq!(status, TurnOutcomeStatus::Interrupted); |
| 170 | assert!( |
| 171 | error.is_none(), |
| 172 | "cancellation is not a provider failure: {error:?}" |
| 173 | ); |
| 174 | assert_eq!((usage.input_tokens, usage.output_tokens), (11, 5)); |
| 175 | assert_eq!(usage, parent_route_usage); |
| 176 | } |
| 177 | Event::MessageDelta { content, .. } => assert!(!content.contains("must not render")), |
| 178 | _ => {} |
| 179 | } |
| 180 | } |
| 181 | assert_eq!((started, completed), (1, 1)); |
| 182 | assert!(terminal_was_last); |
| 183 | drop(events); |
| 184 | handle.send(Op::Shutdown).await.unwrap(); |
| 185 | tokio::time::timeout(Duration::from_secs(2), task) |
| 186 | .await |
| 187 | .unwrap() |
| 188 | .unwrap(); |
| 189 | } |
| 190 | |
| 191 | #[tokio::test] |
| 192 | async fn event_capacity_invalid_images_release_queued_control_and_keep_classifier_receipt() { |
| 193 | use crate::llm_client::mock::MockLlmClient; |
| 194 | for full in [false, true] { |
| 195 | let _cost = crate::cost_status::test_scope(); |
| 196 | let workspace = tempdir().unwrap(); |
| 197 | let config = Config::default(); |
| 198 | let mock = Arc::new(MockLlmClient::new(Vec::new())); |
| 199 | let mut engine_config = deterministic_engine_config(workspace.path()); |
| 200 | engine_config.features.disable(Feature::Mcp); |
| 201 | let (engine, handle) = Engine::new_with_model_client(engine_config, &config, mock.clone()); |
| 202 | if full { |
| 203 | while engine.tx_event.try_send(Event::status("occupied")).is_ok() {} |
| 204 | } |
| 205 | let controls = Arc::clone(&engine.turn_controls); |
| 206 | let mut op = external_user_message_op("invalid attachment", AppMode::Agent, &config); |
| 207 | let Op::SendMessage(spec) = &mut op else { |
| 208 | unreachable!() |
| 209 | }; |
| 210 | spec.images |
| 211 | .push(codewhale_protocol::runtime::RuntimeImageInput { |
| 212 | mime: "image/png".into(), |
| 213 | data_base64: "invalid@base64".into(), |
| 214 | }); |
| 215 | spec.initial_routed_usage |
| 216 | .records |
| 217 | .push(crate::cost_status::RuntimeUsageRecord { |
| 218 | source_id: "event-capacity:invalid-image-classifier".into(), |
| 219 | usage: crate::cost_status::EffectiveRouteUsage { |
| 220 | route: crate::cost_status::EffectiveRouteEnvelope::capture( |
| 221 | None, |
| 222 | ProviderKind::Openai, |
| 223 | "openai", |
| 224 | "classifier", |
| 225 | None, |
| 226 | chrono::Utc::now(), |
| 227 | ), |
| 228 | usage: Usage { |
| 229 | input_tokens: 7, |
| 230 | output_tokens: 3, |
| 231 | ..Usage::default() |
| 232 | }, |
| 233 | }, |
| 234 | }); |
| 235 | handle.send(op).await.unwrap(); |
| 236 | let task = tokio::spawn(engine.run()); |
| 237 | if full { |
| 238 | tokio::time::timeout(Duration::from_secs(2), async { |
| 239 | loop { |
| 240 | if controls.lock().unwrap().active.is_some() { |
| 241 | break; |
| 242 | } |
| 243 | tokio::task::yield_now().await; |
| 244 | } |
| 245 | }) |
| 246 | .await |
| 247 | .unwrap(); |
| 248 | handle.cancel(); |
| 249 | } |
| 250 | let snapshot = tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot()) |
| 251 | .await |
| 252 | .expect("rejected input releases the current control") |
| 253 | .unwrap(); |
| 254 | assert!(snapshot.messages.is_empty()); |
| 255 | assert_eq!(mock.call_count(), 0); |
| 256 | assert!(controls.lock().unwrap().active.is_none()); |
| 257 | let cost = crate::cost_status::drain(); |
| 258 | assert!(cost.usage_source_fingerprints.contains( |
| 259 | &crate::cost_status::usage_source_fingerprint( |
| 260 | "event-capacity:invalid-image-classifier" |
| 261 | ), |
| 262 | )); |
| 263 | let mut events = handle.rx_event.write().await; |
| 264 | let mut invalid_errors = 0; |
| 265 | while let Ok(event) = events.try_recv() { |
| 266 | assert!(!matches!( |
| 267 | event, |
| 268 | Event::TurnStarted { .. } | Event::TurnComplete { .. } |
| 269 | )); |
| 270 | if let Event::Error { envelope, .. } = event { |
| 271 | assert_eq!(envelope.code, "image_input_invalid"); |
| 272 | invalid_errors += 1; |
| 273 | } |
| 274 | } |
| 275 | assert_eq!( |
| 276 | invalid_errors, |
| 277 | usize::from(!full), |
| 278 | "a cancelled full-channel rejection is not fabricated as delivered" |
| 279 | ); |
| 280 | drop(events); |
| 281 | handle.send(Op::Shutdown).await.unwrap(); |
| 282 | tokio::time::timeout(Duration::from_secs(2), task) |
| 283 | .await |
| 284 | .unwrap() |
| 285 | .unwrap(); |
| 286 | } |
| 287 | } |
| 288 | |
| 289 | #[tokio::test] |
| 290 | async fn event_capacity_releases_all_admitted_senders_and_keeps_idle_receipts_lossless() { |
| 291 | let workspace = tempdir().unwrap(); |
| 292 | let (mut engine, _handle) = Engine::new( |
| 293 | deterministic_engine_config(workspace.path()), |
| 294 | &Config::default(), |
| 295 | ); |
| 296 | let (tx, mut rx) = mpsc::channel(1); |
| 297 | engine.tx_event = tx; |
| 298 | engine.tx_event.try_send(Event::status("occupied")).unwrap(); |
| 299 | let turn = engine.begin_turn_control(); |
| 300 | let mut senders = |
| 301 | Box::pin(futures_util::future::join_all((0..8).map(|_| { |
| 302 | engine.send_event(Event::status("admitted observation")) |
| 303 | }))); |
| 304 | assert!( |
| 305 | tokio::time::timeout(Duration::from_millis(20), &mut senders) |
| 306 | .await |
| 307 | .is_err() |
| 308 | ); |
| 309 | engine.cancel_token.cancel(); |
| 310 | let outcomes = tokio::time::timeout(Duration::from_secs(1), &mut senders) |
| 311 | .await |
| 312 | .unwrap(); |
| 313 | assert_eq!(outcomes, vec![Err(streaming::EventSendError::Cancelled); 8]); |
| 314 | drop(senders); |
| 315 | assert!(matches!(rx.try_recv(), Ok(Event::Status { .. }))); |
| 316 | // A cancelled turn still delivers its usage/status receipts when they |
| 317 | // fit immediately. The stream suffix remains strictly cancelled. |
| 318 | engine |
| 319 | .send_event(Event::status( |
| 320 | "cancelled turn receipt with available capacity", |
| 321 | )) |
| 322 | .await |
| 323 | .unwrap(); |
| 324 | assert!( |
| 325 | matches!(rx.try_recv(), Ok(Event::Status { message, .. }) if message == "cancelled turn receipt with available capacity") |
| 326 | ); |
| 327 | assert!( |
| 328 | !engine |
| 329 | .send_stream_event(Event::status("forbidden stream suffix")) |
| 330 | .await |
| 331 | ); |
| 332 | assert!(rx.try_recv().is_err()); |
| 333 | engine |
| 334 | .tx_event |
| 335 | .try_send(Event::status("occupied before idle receipt")) |
| 336 | .unwrap(); |
| 337 | drop(turn); |
| 338 | let mut idle = Box::pin(engine.send_event(Event::status("idle receipt after cancelled turn"))); |
| 339 | assert!( |
| 340 | tokio::time::timeout(Duration::from_millis(20), &mut idle) |
| 341 | .await |
| 342 | .is_err() |
| 343 | ); |
| 344 | assert!(matches!(rx.try_recv(), Ok(Event::Status { .. }))); |
| 345 | tokio::time::timeout(Duration::from_secs(1), &mut idle) |
| 346 | .await |
| 347 | .unwrap() |
| 348 | .unwrap(); |
| 349 | drop(idle); |
| 350 | assert!( |
| 351 | matches!(rx.try_recv(), Ok(Event::Status { message, .. }) if message == "idle receipt after cancelled turn") |
| 352 | ); |
| 353 | assert!(rx.try_recv().is_err()); |
| 354 | } |
| 355 | |
| 356 | #[tokio::test] |
| 357 | async fn event_capacity_nested_vm_fanout_is_bounded_per_program_not_per_turn() { |
| 358 | struct Counter(std::sync::atomic::AtomicUsize); |
| 359 | #[async_trait::async_trait] |
| 360 | impl codewhale_workflow_js::ToolInvoker for Counter { |
| 361 | async fn invoke( |
| 362 | &self, |
| 363 | _: codewhale_workflow_js::ToolCallRequest, |
| 364 | ) -> Result<codewhale_workflow_js::ToolCallResponse, codewhale_workflow_js::DriverError> |
| 365 | { |
| 366 | self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst); |
| 367 | tokio::task::yield_now().await; |
| 368 | Ok(codewhale_workflow_js::ToolCallResponse { |
| 369 | ok: true, |
| 370 | result: json!(null), |
| 371 | }) |
| 372 | } |
| 373 | } |
| 374 | let invoker = Arc::new(Counter(std::sync::atomic::AtomicUsize::new(0))); |
| 375 | for program in 1..=2 { |
| 376 | let result = codewhale_workflow_js::WorkflowVm::new().run_tools_script( |
| 377 | r#"const calls = await Promise.allSettled(Array.from({length:256}, () => tools.call('read', {}))); |
| 378 | return {ok:calls.filter(x=>x.status==='fulfilled').length, |
| 379 | rejected:calls.filter(x=>x.status==='rejected' && x.reason.kind==='admission').length};"#, |
| 380 | json!(null), |
| 381 | Arc::new(codewhale_workflow_js::testing::FakeDriver::new()), |
| 382 | invoker.clone(), codewhale_workflow_js::WorkflowRunCancel::new(), |
| 383 | ).await.unwrap(); |
| 384 | assert_eq!(result, json!({"ok":50,"rejected":206})); |
| 385 | assert_eq!( |
| 386 | invoker.0.load(std::sync::atomic::Ordering::SeqCst), |
| 387 | 50 * program |
| 388 | ); |
| 389 | } |
| 390 | } |
| 391 | |
| 392 | #[tokio::test] |
| 393 | async fn event_capacity_admitted_user_shell_cancels_before_side_effect_and_settles_without_drain() { |
| 394 | let workspace = tempdir().unwrap(); |
| 395 | let marker = workspace.path().join("shell-must-not-start.txt"); |
| 396 | let mut engine_config = deterministic_engine_config(workspace.path()); |
| 397 | engine_config.features.disable(Feature::Mcp); |
| 398 | let (engine, handle) = Engine::new(engine_config, &Config::default()); |
| 399 | // Leave precisely the two lifecycle slots free. TurnStarted fills one; |
| 400 | // the reserved terminal owns the other. The next tool observation must |
| 401 | // wait before the human-provenance command can reach its executor. |
| 402 | for _ in 0..engine.tx_event.max_capacity() - 2 { |
| 403 | engine |
| 404 | .tx_event |
| 405 | .try_send(Event::status("prior receipt")) |
| 406 | .unwrap(); |
| 407 | } |
| 408 | let controls = Arc::clone(&engine.turn_controls); |
| 409 | handle |
| 410 | .send(Op::RunShellCommand { |
| 411 | command: format!("echo must-not-run > \"{}\"", marker.display()), |
| 412 | mode: AppMode::Agent, |
| 413 | allow_shell: true, |
| 414 | trust_mode: true, |
| 415 | auto_approve: true, |
| 416 | approval_mode: ApprovalMode::Bypass, |
| 417 | }) |
| 418 | .await |
| 419 | .unwrap(); |
| 420 | let task = tokio::spawn(engine.run()); |
| 421 | tokio::time::timeout(Duration::from_secs(2), async { |
| 422 | loop { |
| 423 | if controls.lock().unwrap().active.is_some() { |
| 424 | break; |
| 425 | } |
| 426 | tokio::task::yield_now().await; |
| 427 | } |
| 428 | }) |
| 429 | .await |
| 430 | .unwrap(); |
| 431 | handle.cancel(); |
| 432 | tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot()) |
| 433 | .await |
| 434 | .expect("shell cancellation settles before any event drain") |
| 435 | .unwrap(); |
| 436 | assert!( |
| 437 | !marker.exists(), |
| 438 | "cancellation preserves the execution gate" |
| 439 | ); |
| 440 | assert!(controls.lock().unwrap().active.is_none()); |
| 441 | let mut events = handle.rx_event.write().await; |
| 442 | let mut started = 0; |
| 443 | let mut completed = 0; |
| 444 | let mut terminal_was_last = false; |
| 445 | while let Ok(event) = events.try_recv() { |
| 446 | assert!(!terminal_was_last); |
| 447 | match event { |
| 448 | Event::TurnStarted { .. } => started += 1, |
| 449 | Event::TurnComplete { status, usage, .. } => { |
| 450 | completed += 1; |
| 451 | terminal_was_last = true; |
| 452 | assert_eq!(status, TurnOutcomeStatus::Interrupted); |
| 453 | assert_eq!(usage, Usage::default()); |
| 454 | } |
| 455 | _ => {} |
| 456 | } |
| 457 | } |
| 458 | assert_eq!((started, completed), (1, 1)); |
| 459 | assert!(terminal_was_last); |
| 460 | drop(events); |
| 461 | handle.send(Op::Shutdown).await.unwrap(); |
| 462 | tokio::time::timeout(Duration::from_secs(2), task) |
| 463 | .await |
| 464 | .unwrap() |
| 465 | .unwrap(); |
| 466 | } |
| 467 | |
| 468 | #[tokio::test] |
| 469 | async fn event_capacity_cancelled_repl_child_keeps_unknown_cost_and_discards_kernel() { |
| 470 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 471 | struct ReplClient { |
| 472 | inner: MockLlmClient, |
| 473 | stream_requests: std::sync::atomic::AtomicUsize, |
| 474 | child_entered: Arc<tokio::sync::Notify>, |
| 475 | child_dropped: Arc<std::sync::atomic::AtomicBool>, |
| 476 | } |
| 477 | #[async_trait::async_trait] |
| 478 | impl crate::core::model_client::ModelClient for ReplClient { |
| 479 | fn provider_name(&self) -> &str { |
| 480 | "backpressure-repl-fixture" |
| 481 | } |
| 482 | fn model(&self) -> &str { |
| 483 | "mock-model" |
| 484 | } |
| 485 | async fn create_message( |
| 486 | &self, |
| 487 | _: codewhale_models::MessageRequest, |
| 488 | ) -> anyhow::Result<codewhale_models::MessageResponse> { |
| 489 | anyhow::bail!("fixture expects canonical streaming requests") |
| 490 | } |
| 491 | async fn create_message_stream( |
| 492 | &self, |
| 493 | request: codewhale_models::MessageRequest, |
| 494 | ) -> anyhow::Result<crate::llm_client::StreamEventBox> { |
| 495 | if self |
| 496 | .stream_requests |
| 497 | .fetch_add(1, std::sync::atomic::Ordering::SeqCst) |
| 498 | == 0 |
| 499 | { |
| 500 | return crate::core::model_client::ModelClient::create_message_stream( |
| 501 | &self.inner, |
| 502 | request, |
| 503 | ) |
| 504 | .await; |
| 505 | } |
| 506 | let _drop = DropSignal(Arc::clone(&self.child_dropped)); |
| 507 | self.child_entered.notify_one(); |
| 508 | std::future::pending().await |
| 509 | } |
| 510 | async fn health_check(&self) -> anyhow::Result<bool> { |
| 511 | Ok(true) |
| 512 | } |
| 513 | } |
| 514 | let _cost = crate::cost_status::test_scope(); |
| 515 | let workspace = tempdir().unwrap(); |
| 516 | let entered = Arc::new(tokio::sync::Notify::new()); |
| 517 | let dropped = Arc::new(std::sync::atomic::AtomicBool::new(false)); |
| 518 | let mut response = canned::simple_text_turn( |
| 519 | "```repl\nchild = sub_query('hang until cancelled')\nfinalize(child)\n```", |
| 520 | ); |
| 521 | for event in &mut response { |
| 522 | if let codewhale_models::StreamEvent::MessageDelta { usage, .. } = event { |
| 523 | *usage = Some(Usage { |
| 524 | input_tokens: 11, |
| 525 | output_tokens: 5, |
| 526 | ..Usage::default() |
| 527 | }); |
| 528 | } |
| 529 | } |
| 530 | let client = Arc::new(ReplClient { |
| 531 | inner: MockLlmClient::new(vec![response]), |
| 532 | stream_requests: std::sync::atomic::AtomicUsize::new(0), |
| 533 | child_entered: Arc::clone(&entered), |
| 534 | child_dropped: Arc::clone(&dropped), |
| 535 | }); |
| 536 | let api_config = rlm_host::fixture_config("mock-model"); |
| 537 | let (mut engine, handle) = Engine::new_with_model_client( |
| 538 | EngineConfig { |
| 539 | model: "mock-model".into(), |
| 540 | ..deterministic_engine_config(workspace.path()) |
| 541 | }, |
| 542 | &api_config, |
| 543 | client.clone(), |
| 544 | ); |
| 545 | rlm_host::install_fixture_route(&mut engine); |
| 546 | engine.session.auto_approve = true; |
| 547 | engine.session.add_message(Message { |
| 548 | role: Role::User, |
| 549 | content: vec![ContentBlock::Text { |
| 550 | text: "run cancellable REPL".into(), |
| 551 | cache_control: None, |
| 552 | }], |
| 553 | }); |
| 554 | let turn_guard = engine.begin_turn_control(); |
| 555 | let mut turn = TurnContext::new(4); |
| 556 | let registry = crate::tools::ToolRegistry::new(rlm_host::admitted_context(&engine, &turn.id)); |
| 557 | let policy = test_tool_surface( |
| 558 | &engine, |
| 559 | registry, |
| 560 | Some(vec![catalog_tool(CODE_EXECUTION_TOOL_NAME)]), |
| 561 | AppMode::Agent, |
| 562 | ); |
| 563 | let tx = engine.tx_event.clone(); |
| 564 | let mut run = Box::pin(engine.run_turn(&mut turn, policy, None, None)); |
| 565 | tokio::time::timeout(Duration::from_secs(10), async { |
| 566 | tokio::select! { |
| 567 | result = &mut run => panic!("REPL did not enter the child provider request: {result:?}"), |
| 568 | () = entered.notified() => {}, |
| 569 | } |
| 570 | }).await.expect("actual Python kernel dispatches the pending child request"); |
| 571 | while tx |
| 572 | .try_send(Event::status("full during nested REPL")) |
| 573 | .is_ok() |
| 574 | {} |
| 575 | handle.cancel(); |
| 576 | let (status, error) = tokio::time::timeout(Duration::from_secs(2), &mut run) |
| 577 | .await |
| 578 | .expect("cancel drops the pending REPL round even with a full queue"); |
| 579 | drop(run); |
| 580 | assert_eq!(status, TurnOutcomeStatus::Interrupted, "{error:?}"); |
| 581 | assert!(error.is_none()); |
| 582 | assert!( |
| 583 | dropped.load(std::sync::atomic::Ordering::SeqCst), |
| 584 | "nested provider future is released" |
| 585 | ); |
| 586 | assert!( |
| 587 | engine.repl_kernel.is_none(), |
| 588 | "a cancelled round cannot preserve an executing process" |
| 589 | ); |
| 590 | assert_eq!((turn.usage.input_tokens, turn.usage.output_tokens), (11, 5)); |
| 591 | let cost = crate::cost_status::drain(); |
| 592 | // The child provider future never returned a response; its canonical |
| 593 | // dispatch guard records an unknown outcome, not success without usage. |
| 594 | assert!( |
| 595 | cost.unpriced_reasons |
| 596 | .contains(crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown.label()), |
| 597 | "pending child usage stays unknown, never zero" |
| 598 | ); |
| 599 | assert_eq!(cost.priced_turns, 0); |
| 600 | assert_eq!(cost.unpriced_turns, 1); |
| 601 | assert_eq!(cost.missing_usage_sources.len(), 1); |
| 602 | assert!(cost.missing_usage_sources.values().all(|coverage| { |
| 603 | coverage.reason == crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown |
| 604 | && coverage.money_metered |
| 605 | })); |
| 606 | assert!(cost.resolved_missing_usage_sources.is_empty()); |
| 607 | assert_eq!( |
| 608 | client.inner.call_count(), |
| 609 | 1, |
| 610 | "no next root provider request after cancellation" |
| 611 | ); |
| 612 | drop(turn_guard); |
| 613 | } |
| 614 | #[tokio::test] |
| 615 | async fn event_capacity_cancelled_parallel_tool_keeps_completed_span_and_call_when_available() { |
| 616 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 617 | use crate::tools::spec::{ |
| 618 | ApprovalRequirement, PreparedToolCall, ResourceClaim, ToolCapability, ToolSpec, |
| 619 | }; |
| 620 | use codewhale_protocol::engine_owner::OwnerOperationOutcome; |
| 621 | use std::sync::atomic::{AtomicUsize, Ordering}; |
| 622 | |
| 623 | struct CompleteThenCancel { |
| 624 | cancel: tokio_util::sync::CancellationToken, |
| 625 | tx: mpsc::Sender<Event>, |
| 626 | fill_queue: bool, |
| 627 | executed: Arc<AtomicUsize>, |
| 628 | } |
| 629 | #[async_trait::async_trait] |
| 630 | impl ToolSpec for CompleteThenCancel { |
| 631 | // Registered under the canonical read identity so the engine's |
| 632 | // central resource authority (not this fixture) grants the disjoint |
| 633 | // ReadPath claims that form one real parallel chunk. |
| 634 | fn name(&self) -> &str { |
| 635 | "read_file" |
| 636 | } |
| 637 | fn description(&self) -> &str { |
| 638 | "Finish an observed operation before firing its turn cancellation token." |
| 639 | } |
| 640 | fn input_schema(&self) -> Value { |
| 641 | json!({"type": "object"}) |
| 642 | } |
| 643 | fn capabilities(&self) -> Vec<ToolCapability> { |
| 644 | vec![ToolCapability::ReadOnly] |
| 645 | } |
| 646 | fn supports_parallel(&self) -> bool { |
| 647 | true |
| 648 | } |
| 649 | fn prepare( |
| 650 | &self, |
| 651 | input: Value, |
| 652 | context: &ToolContext, |
| 653 | ) -> Result<PreparedToolCall, ToolError> { |
| 654 | let path = input["path"].as_str().expect("fixture path"); |
| 655 | Ok(PreparedToolCall { |
| 656 | name: self.name().to_string(), |
| 657 | description: self.description().to_string(), |
| 658 | read_only: true, |
| 659 | supports_parallel: true, |
| 660 | starts_detached: false, |
| 661 | approval: ApprovalRequirement::Auto, |
| 662 | // Replaced by `registered_resource_claims`; kept honest anyway. |
| 663 | resources: vec![ResourceClaim::ReadPath(context.workspace.join(path))], |
| 664 | input, |
| 665 | }) |
| 666 | } |
| 667 | async fn execute(&self, _: Value, _: &ToolContext) -> Result<ToolResult, ToolError> { |
| 668 | self.executed.fetch_add(1, Ordering::SeqCst); |
| 669 | if self.fill_queue { |
| 670 | while self |
| 671 | .tx |
| 672 | .try_send(Event::status("full at tool completion")) |
| 673 | .is_ok() |
| 674 | {} |
| 675 | } |
| 676 | self.cancel.cancel(); |
| 677 | Ok(ToolResult::success("completed before cancellation")) |
| 678 | } |
| 679 | } |
| 680 | |
| 681 | for fill_queue in [false, true] { |
| 682 | let workspace = tempdir().unwrap(); |
| 683 | for path in ["one", "two"] { |
| 684 | std::fs::write(workspace.path().join(path), path).unwrap(); |
| 685 | } |
| 686 | let mock = Arc::new(MockLlmClient::new(vec![ |
| 687 | tool_batch_turn(&[ |
| 688 | ("one", "read_file", r#"{"path":"one"}"#), |
| 689 | ("two", "read_file", r#"{"path":"two"}"#), |
| 690 | ]), |
| 691 | canned::simple_text_turn("must not run after cancellation"), |
| 692 | ])); |
| 693 | let (mut engine, handle) = Engine::new_with_model_client( |
| 694 | deterministic_engine_config(workspace.path()), |
| 695 | &Config::default(), |
| 696 | mock.clone(), |
| 697 | ); |
| 698 | let turn_guard = engine.begin_turn_control(); |
| 699 | let executed = Arc::new(AtomicUsize::new(0)); |
| 700 | let mut registry = crate::tools::ToolRegistry::new(ToolContext::new(workspace.path())); |
| 701 | registry.register(Arc::new(CompleteThenCancel { |
| 702 | cancel: engine.cancel_token.clone(), |
| 703 | tx: engine.tx_event.clone(), |
| 704 | fill_queue, |
| 705 | executed: executed.clone(), |
| 706 | })); |
| 707 | let tools = Some(registry.to_api_tools_with_cache(true)); |
| 708 | let policy = test_tool_surface(&engine, registry, tools, AppMode::Agent); |
| 709 | let mut turn = TurnContext::new(4); |
| 710 | let (status, error) = tokio::time::timeout( |
| 711 | Duration::from_secs(2), |
| 712 | engine.run_turn(&mut turn, policy, None, None), |
| 713 | ) |
| 714 | .await |
| 715 | .expect("actual parallel tool completion cannot park on a full cancelled queue"); |
| 716 | assert_eq!(status, TurnOutcomeStatus::Interrupted, "{error:?}"); |
| 717 | assert!(error.is_none()); |
| 718 | assert_eq!( |
| 719 | executed.load(Ordering::SeqCst), |
| 720 | 1, |
| 721 | "the peer never executes" |
| 722 | ); |
| 723 | assert_eq!( |
| 724 | mock.call_count(), |
| 725 | 1, |
| 726 | "no provider continuation after cancellation" |
| 727 | ); |
| 728 | let mut rx = handle.rx_event.write().await; |
| 729 | let events: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect(); |
| 730 | assert!(events.iter().any(|event| matches!(event, |
| 731 | Event::Status { message } if message == "Executing 2 read-only tools in 1 parallel chunk(s)" |
| 732 | )), "fixture must exercise one actual two-tool parallel chunk"); |
| 733 | let starts: Vec<_> = events |
| 734 | .iter() |
| 735 | .filter_map(|event| match event { |
| 736 | Event::OperationActivityStarted { |
| 737 | span_id, |
| 738 | activity_kind, |
| 739 | } => Some((span_id, activity_kind)), |
| 740 | _ => None, |
| 741 | }) |
| 742 | .collect(); |
| 743 | assert_eq!(starts.len(), 1, "one actual tool activity was admitted"); |
| 744 | let completed: Vec<_> = events |
| 745 | .iter() |
| 746 | .filter_map(|event| match event { |
| 747 | Event::OperationActivityCompleted { |
| 748 | span_id, |
| 749 | activity_kind, |
| 750 | outcome, |
| 751 | } => Some((span_id, activity_kind, outcome)), |
| 752 | _ => None, |
| 753 | }) |
| 754 | .collect(); |
| 755 | let calls: Vec<_> = events |
| 756 | .iter() |
| 757 | .filter_map(|event| match event { |
| 758 | Event::ToolCallComplete { |
| 759 | id, |
| 760 | model_call, |
| 761 | result, |
| 762 | .. |
| 763 | } => Some((id, model_call, result)), |
| 764 | _ => None, |
| 765 | }) |
| 766 | .collect(); |
| 767 | if fill_queue { |
| 768 | assert!( |
| 769 | completed.is_empty() && calls.is_empty(), |
| 770 | "cancelled full waits release without inventing delivery" |
| 771 | ); |
| 772 | } else { |
| 773 | assert_eq!( |
| 774 | completed.len(), |
| 775 | 1, |
| 776 | "preserve the completed activity after cancellation" |
| 777 | ); |
| 778 | assert_eq!(completed[0].0, starts[0].0, "retain the span relationship"); |
| 779 | assert_eq!(completed[0].1, starts[0].1); |
| 780 | assert_eq!(*completed[0].2, OwnerOperationOutcome::Succeeded); |
| 781 | assert_eq!( |
| 782 | calls.len(), |
| 783 | 2, |
| 784 | "completed call and cancelled peer both settle exactly once" |
| 785 | ); |
| 786 | let succeeded: Vec<_> = calls |
| 787 | .iter() |
| 788 | .filter(|(_, _, result)| result.as_ref().is_ok_and(|result| result.success)) |
| 789 | .collect(); |
| 790 | assert_eq!(succeeded.len(), 1); |
| 791 | let cancelled_peer = calls |
| 792 | .iter() |
| 793 | .find(|(_, _, result)| !result.as_ref().is_ok_and(|result| result.success)) |
| 794 | .expect("cancelled peer has its own completion"); |
| 795 | let peer = cancelled_peer |
| 796 | .2 |
| 797 | .as_ref() |
| 798 | .expect("legacy cancelled peer result"); |
| 799 | assert_eq!(peer.metadata.as_ref().unwrap()["cancelled"], true); |
| 800 | assert_eq!(peer.metadata.as_ref().unwrap()["cleanup_confirmed"], false); |
| 801 | assert_eq!( |
| 802 | succeeded[0].2.as_ref().unwrap().content, |
| 803 | "completed before cancellation" |
| 804 | ); |
| 805 | assert!( |
| 806 | starts[0].0.starts_with(&format!("{}#", succeeded[0].0)), |
| 807 | "span retains its completed execution identity" |
| 808 | ); |
| 809 | let provider_ids: HashSet<_> = calls |
| 810 | .iter() |
| 811 | .map(|(_, model_call, _)| model_call.as_ref().unwrap().provider_id.as_str()) |
| 812 | .collect(); |
| 813 | assert_eq!(provider_ids, HashSet::from(["one", "two"])); |
| 814 | let call_starts: HashSet<_> = events |
| 815 | .iter() |
| 816 | .filter_map(|event| match event { |
| 817 | Event::ToolCallStarted { id, .. } => Some(id), |
| 818 | _ => None, |
| 819 | }) |
| 820 | .collect(); |
| 821 | assert!(calls.iter().all(|(id, _, _)| call_starts.contains(id))); |
| 822 | } |
| 823 | drop(rx); |
| 824 | drop(turn_guard); |
| 825 | } |
| 826 | } |
| 827 | |
| 828 | const REPRESENTATIVE_FIXTURE_ID: &str = "representative-v1"; |
| 829 | const REPRESENTATIVE_PROJECT_AUTHORITY: &str = "REPRESENTATIVE_PROJECT_AUTHORITY"; |
| 830 | const REPRESENTATIVE_PROJECT_AUTHORITY_BODY: &str = concat!( |
| 831 | "# Representative Project Authority\n\n", |
| 832 | "REPRESENTATIVE_PROJECT_AUTHORITY\n\n", |
| 833 | "- Keep all work local to the isolated fixture workspace.\n", |
| 834 | "- Treat the checked-in repository instructions as the authority for edits.\n", |
| 835 | "- Preserve unrelated files and report unsupported checks as unrun.\n", |
| 836 | "- Prefer one owner for each runtime fact and delete duplicated derivations.\n", |
| 837 | "- Use deterministic provider-free tests before claiming a behavior is verified.\n", |
| 838 | "- Keep durable state atomic, recoverable, and explicit about unavailable facts.\n", |
| 839 | "- Do not contact remotes, providers, registries, or production services.\n", |
| 840 | "- Record exact measurements and distinguish source proof from installed proof.\n", |
| 841 | ); |
| 842 | |
| 843 | #[test] |
| 844 | fn snapshot_notice_precedes_first_provider_call_and_is_owned_by_session() { |
| 845 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 846 | let _env = lock_test_env(); |
| 847 | let root = tempdir().unwrap(); |
| 848 | let _home = EnvVarGuard::set("CODEWHALE_HOME", root.path()); |
| 849 | let _user_home = EnvVarGuard::set("HOME", root.path()); |
| 850 | let _user_profile = EnvVarGuard::set("USERPROFILE", root.path()); |
| 851 | let workspace = root.path().join("workspace"); |
| 852 | fs::create_dir(&workspace).unwrap(); |
| 853 | fs::write(workspace.join("large.txt"), vec![b'x'; 4096]).unwrap(); |
| 854 | let runtime = tokio::runtime::Builder::new_current_thread() |
| 855 | .enable_all() |
| 856 | .build() |
| 857 | .unwrap(); |
| 858 | runtime.block_on(async { |
| 859 | // A resumed Engine for session-a must not warn again. Session-b in |
| 860 | // the same process/workspace must receive its own first-turn notice. |
| 861 | for (session_id, expected_notices) in [("session-a", 1), ("session-b", 1), ("session-a", 0)] |
| 862 | { |
| 863 | let config = Config::default(); |
| 864 | let client = std::sync::Arc::new(MockLlmClient::new(Vec::new())); |
| 865 | let (engine, handle) = Engine::new_with_model_client( |
| 866 | EngineConfig { |
| 867 | session_id: Some(session_id.into()), |
| 868 | snapshots_enabled: true, |
| 869 | snapshots_max_workspace_bytes: 1024, |
| 870 | ..deterministic_engine_config(&workspace) |
| 871 | }, |
| 872 | &config, |
| 873 | client.clone(), |
| 874 | ); |
| 875 | let events = std::sync::Arc::clone(&handle.rx_event); |
| 876 | let observations = std::sync::Arc::new(std::sync::Mutex::new(Vec::new())); |
| 877 | let observed = std::sync::Arc::clone(&observations); |
| 878 | client.push_factory(move |_| { |
| 879 | let mut events = events.try_write().expect("fixture owns the event receiver"); |
| 880 | let mut notices = Vec::new(); |
| 881 | while let Ok(event) = events.try_recv() { |
| 882 | if let Event::SnapshotsDisabled { reason, .. } = event { |
| 883 | notices.push(reason); |
| 884 | } |
| 885 | } |
| 886 | observed.lock().unwrap().push(notices); |
| 887 | canned::simple_text_turn("snapshot fixture done") |
| 888 | }); |
| 889 | let run = tokio::spawn(engine.run()); |
| 890 | handle |
| 891 | .send(external_user_message_op( |
| 892 | "check snapshots", |
| 893 | AppMode::Agent, |
| 894 | &config, |
| 895 | )) |
| 896 | .await |
| 897 | .unwrap(); |
| 898 | let snapshot = |
| 899 | tokio::time::timeout(Duration::from_secs(10), handle.get_session_snapshot()) |
| 900 | .await |
| 901 | .unwrap() |
| 902 | .unwrap(); |
| 903 | // The Engine catches provider panics, so assertions inside the |
| 904 | // factory are not a test oracle. Inspect its observations here. |
| 905 | { |
| 906 | let observed = observations.lock().unwrap(); |
| 907 | assert_eq!( |
| 908 | observed.len(), |
| 909 | 1, |
| 910 | "factory must have recorded an observation" |
| 911 | ); |
| 912 | assert_eq!( |
| 913 | observed[0].len(), |
| 914 | expected_notices, |
| 915 | "notice must precede provider dispatch for this session" |
| 916 | ); |
| 917 | assert!(observed[0].iter().all(|reason| { |
| 918 | // One rendered line: consequence, cause, and the remedy |
| 919 | // that lifts this gate, each stated once. |
| 920 | reason.lines().count() == 1 |
| 921 | && reason.contains("Snapshots and /undo are off") |
| 922 | && reason |
| 923 | .matches(crate::core::turn::SNAPSHOTS_CAP_CONFIG_KEY) |
| 924 | .count() |
| 925 | == 1 |
| 926 | })); |
| 927 | } |
| 928 | assert_eq!(client.call_count(), 1); |
| 929 | assert!( |
| 930 | serde_json::to_string(&snapshot.messages) |
| 931 | .unwrap() |
| 932 | .contains("snapshot fixture done") |
| 933 | ); |
| 934 | handle.send(Op::Shutdown).await.unwrap(); |
| 935 | tokio::time::timeout(Duration::from_secs(10), run) |
| 936 | .await |
| 937 | .unwrap() |
| 938 | .unwrap(); |
| 939 | } |
| 940 | }); |
| 941 | // Await the owned blocking post-turn snapshots before restoring test home. |
| 942 | drop(runtime); |
| 943 | } |
| 944 | |
| 945 | /// A recording host (`record_restore_points`) receives every workspace |
| 946 | /// snapshot receipt of a turn before its `TurnComplete`: the pre-turn restore |
| 947 | /// point, a `tool` snapshot naming the file-mutating call and the paths it |
| 948 | /// declared, the `post_tool` snapshot closing it, and the post-turn state. A |
| 949 | /// read-only call takes none. The pre/post-turn trees bracket exactly the |
| 950 | /// turn's write, so a host derives the turn's workspace delta from them. |
| 951 | #[test] |
| 952 | fn recorded_snapshot_receipts_bracket_the_turn_and_its_file_writes() { |
| 953 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 954 | use crate::snapshot::WorkspaceSnapshotKind; |
| 955 | let _env = lock_test_env(); |
| 956 | let root = tempdir().unwrap(); |
| 957 | let _home = EnvVarGuard::set("CODEWHALE_HOME", root.path()); |
| 958 | let _user_home = EnvVarGuard::set("HOME", root.path()); |
| 959 | let _user_profile = EnvVarGuard::set("USERPROFILE", root.path()); |
| 960 | let workspace = root.path().join("workspace"); |
| 961 | fs::create_dir(&workspace).unwrap(); |
| 962 | fs::write(workspace.join("README.md"), "fixture\n").unwrap(); |
| 963 | let runtime = tokio::runtime::Builder::new_current_thread() |
| 964 | .enable_all() |
| 965 | .build() |
| 966 | .unwrap(); |
| 967 | runtime.block_on(async { |
| 968 | let config = Config::default(); |
| 969 | let client = std::sync::Arc::new(MockLlmClient::new(vec![ |
| 970 | canned::tool_call_turn( |
| 971 | "call-write", |
| 972 | "File", |
| 973 | r#"{"action":"write","path":"out.md","content":"out\n"}"#, |
| 974 | ), |
| 975 | canned::tool_call_turn( |
| 976 | "call-read", |
| 977 | "File", |
| 978 | r#"{"action":"read","path":"README.md"}"#, |
| 979 | ), |
| 980 | canned::simple_text_turn("done"), |
| 981 | ])); |
| 982 | let (engine, handle) = Engine::new_with_model_client( |
| 983 | EngineConfig { |
| 984 | session_id: Some("session-restore".into()), |
| 985 | snapshots_enabled: true, |
| 986 | snapshots_max_workspace_bytes: 0, |
| 987 | record_restore_points: true, |
| 988 | ..deterministic_engine_config(&workspace) |
| 989 | }, |
| 990 | &config, |
| 991 | client, |
| 992 | ); |
| 993 | let run = tokio::spawn(engine.run()); |
| 994 | let Op::SendMessage(mut spec) = |
| 995 | external_user_message_op("write out.md", AppMode::Agent, &config) |
| 996 | else { |
| 997 | unreachable!("external_user_message_op builds a SendMessage"); |
| 998 | }; |
| 999 | spec.auto_approve = true; |
| 1000 | spec.trust_mode = true; |
| 1001 | spec.approval_mode = ApprovalMode::Bypass; |
| 1002 | handle.send(Op::SendMessage(spec)).await.unwrap(); |
| 1003 | |
| 1004 | let mut completions = HashMap::new(); |
| 1005 | let mut local_ids = HashMap::new(); |
| 1006 | let mut receipts = Vec::new(); |
| 1007 | let mut rx = handle.rx_event.write().await; |
| 1008 | while let Some(event) = tokio::time::timeout(model_turn_event_timeout(), rx.recv()) |
| 1009 | .await |
| 1010 | .expect("turn events") |
| 1011 | { |
| 1012 | match event { |
| 1013 | Event::ToolCallStarted { |
| 1014 | id, |
| 1015 | model_call: Some(model_call), |
| 1016 | .. |
| 1017 | } => { |
| 1018 | uuid::Uuid::parse_str(&id).expect("host execution id"); |
| 1019 | assert_ne!(id, model_call.provider_id); |
| 1020 | assert!(local_ids.insert(model_call.provider_id, id).is_none()); |
| 1021 | } |
| 1022 | Event::ToolCallComplete { |
| 1023 | id, |
| 1024 | result, |
| 1025 | model_call: Some(model_call), |
| 1026 | .. |
| 1027 | } => { |
| 1028 | assert_eq!(local_ids.get(&model_call.provider_id), Some(&id)); |
| 1029 | completions.insert(model_call.provider_id, result.expect("tool result")); |
| 1030 | } |
| 1031 | Event::WorkspaceSnapshotTaken { snapshot } => receipts.push(snapshot), |
| 1032 | Event::TurnComplete { status, error, .. } => { |
| 1033 | assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}"); |
| 1034 | break; |
| 1035 | } |
| 1036 | _ => {} |
| 1037 | } |
| 1038 | } |
| 1039 | drop(rx); |
| 1040 | |
| 1041 | assert!(completions.get("call-write").expect("write ran").success); |
| 1042 | assert!(completions.get("call-read").expect("read ran").success); |
| 1043 | assert_eq!(local_ids.len(), 2); |
| 1044 | assert_ne!(local_ids["call-write"], local_ids["call-read"]); |
| 1045 | // Every receipt, post-turn included, arrived before TurnComplete. |
| 1046 | assert_eq!( |
| 1047 | receipts |
| 1048 | .iter() |
| 1049 | .map(|receipt| (receipt.kind, receipt.tool_call_id.as_deref())) |
| 1050 | .collect::<Vec<_>>(), |
| 1051 | [ |
| 1052 | (WorkspaceSnapshotKind::PreTurn, None), |
| 1053 | ( |
| 1054 | WorkspaceSnapshotKind::Tool, |
| 1055 | Some(local_ids["call-write"].as_str()) |
| 1056 | ), |
| 1057 | ( |
| 1058 | WorkspaceSnapshotKind::PostTool, |
| 1059 | Some(local_ids["call-write"].as_str()) |
| 1060 | ), |
| 1061 | (WorkspaceSnapshotKind::PostTurn, None), |
| 1062 | ], |
| 1063 | "a read-only call takes no restore point: {receipts:?}" |
| 1064 | ); |
| 1065 | assert!( |
| 1066 | receipts |
| 1067 | .iter() |
| 1068 | .all(|receipt| receipt.session_id == "session-restore") |
| 1069 | ); |
| 1070 | assert_eq!( |
| 1071 | receipts[1].write_paths.as_deref(), |
| 1072 | Some(&["out.md".to_string()][..]) |
| 1073 | ); |
| 1074 | assert_eq!( |
| 1075 | receipts[2].changed_paths.as_deref(), |
| 1076 | Some(&["out.md".to_string()][..]) |
| 1077 | ); |
| 1078 | |
| 1079 | let repo = crate::snapshot::SnapshotRepo::open_existing(&workspace) |
| 1080 | .unwrap() |
| 1081 | .expect("snapshot repo"); |
| 1082 | let listed = repo.list(usize::MAX).unwrap(); |
| 1083 | assert!( |
| 1084 | listed.iter().any(|snapshot| receipts[1].matches(snapshot) |
| 1085 | && snapshot.label == format!("tool:{}", local_ids["call-write"])), |
| 1086 | "the tool receipt names a live snapshot" |
| 1087 | ); |
| 1088 | let delta = repo |
| 1089 | .diff_snapshots( |
| 1090 | &crate::snapshot::SnapshotId::parse(&receipts[0].tree_id).unwrap(), |
| 1091 | &crate::snapshot::SnapshotId::parse(&receipts[3].tree_id).unwrap(), |
| 1092 | 100, |
| 1093 | ) |
| 1094 | .unwrap(); |
| 1095 | assert_eq!( |
| 1096 | delta |
| 1097 | .entries |
| 1098 | .iter() |
| 1099 | .map(|entry| entry.path.as_str()) |
| 1100 | .collect::<Vec<_>>(), |
| 1101 | ["out.md"], |
| 1102 | "the pre/post-turn trees bracket exactly this turn's write" |
| 1103 | ); |
| 1104 | |
| 1105 | handle.send(Op::Shutdown).await.unwrap(); |
| 1106 | tokio::time::timeout(Duration::from_secs(10), run) |
| 1107 | .await |
| 1108 | .unwrap() |
| 1109 | .unwrap(); |
| 1110 | }); |
| 1111 | // Await the owned blocking post-turn snapshots before restoring test home. |
| 1112 | drop(runtime); |
| 1113 | } |
| 1114 | |
| 1115 | #[test] |
| 1116 | fn preview_request_error_preserves_non_semantic_context_chain() { |
| 1117 | let error = anyhow::Error::msg("root cause").context("request preparation failed"); |
| 1118 | assert_eq!( |
| 1119 | preview_request_error_user_message("en", &error), |
| 1120 | "request preparation failed: root cause" |
| 1121 | ); |
| 1122 | assert_eq!( |
| 1123 | initial_stream_error_user_message("en", &error), |
| 1124 | "request preparation failed: root cause" |
| 1125 | ); |
| 1126 | } |
| 1127 | |
| 1128 | #[test] |
| 1129 | fn initial_stream_failure_preserves_sanitized_context_and_typed_category() { |
| 1130 | let error = anyhow::Error::new(crate::llm_client::LlmError::InvalidRequest { |
| 1131 | status: 400, |
| 1132 | message: "image input is unsupported; api_key=fixture-credential-value".to_string(), |
| 1133 | }) |
| 1134 | .context("Responses API request failed"); |
| 1135 | let display = initial_stream_error_user_message("en", &error); |
| 1136 | assert!( |
| 1137 | display.contains("Responses API request failed"), |
| 1138 | "{display}" |
| 1139 | ); |
| 1140 | assert!(display.contains("Invalid request (400)"), "{display}"); |
| 1141 | assert!(display.contains("image input is unsupported"), "{display}"); |
| 1142 | assert!(!display.contains("fixture-credential-value"), "{display}"); |
| 1143 | assert!(display.contains("[redacted]"), "{display}"); |
| 1144 | |
| 1145 | // The real boundary classifies the original error independently of its |
| 1146 | // expanded display text. Preserve the typed terminal invalid-input result. |
| 1147 | let policy_message = error.to_string(); |
| 1148 | let mut envelope = crate::error_taxonomy::envelope_for_llm_error(error, policy_message); |
| 1149 | envelope.message = display; |
| 1150 | assert_eq!( |
| 1151 | envelope.category, |
| 1152 | crate::error_taxonomy::ErrorCategory::InvalidInput |
| 1153 | ); |
| 1154 | assert!(!envelope.recoverable); |
| 1155 | assert_eq!(envelope.code, "llm_invalid_request"); |
| 1156 | } |
| 1157 | const REPRESENTATIVE_INLINE_INSTRUCTIONS: &str = "REPRESENTATIVE_INLINE_INSTRUCTIONS"; |
| 1158 | const REPRESENTATIVE_SKILL_DESCRIPTION: &str = "REPRESENTATIVE_SKILL_DESCRIPTION"; |
| 1159 | const REPRESENTATIVE_MEMORY_CHECKPOINT: &str = "REPRESENTATIVE_MEMORY_CHECKPOINT"; |
| 1160 | const REPRESENTATIVE_GOAL_OBJECTIVE: &str = "REPRESENTATIVE_GOAL_OBJECTIVE"; |
| 1161 | |
| 1162 | #[test] |
| 1163 | fn cancellation_wins_at_the_terminal_child_settlement_seam() { |
| 1164 | assert_eq!( |
| 1165 | terminal_turn_status_at_settlement(TurnOutcomeStatus::Completed, true), |
| 1166 | TurnOutcomeStatus::Interrupted |
| 1167 | ); |
| 1168 | assert_eq!( |
| 1169 | terminal_turn_status_at_settlement(TurnOutcomeStatus::Completed, false), |
| 1170 | TurnOutcomeStatus::Completed |
| 1171 | ); |
| 1172 | assert_eq!( |
| 1173 | terminal_turn_status_at_settlement(TurnOutcomeStatus::Failed, true), |
| 1174 | TurnOutcomeStatus::Failed |
| 1175 | ); |
| 1176 | } |
| 1177 | |
| 1178 | #[tokio::test] |
| 1179 | async fn terminal_barrier_keeps_healthy_child_and_late_completion_alive() { |
| 1180 | use std::sync::atomic::Ordering; |
| 1181 | let turn_token = CancellationToken::new(); |
| 1182 | let (mailbox, _receiver) = Mailbox::new(turn_token.clone()); |
| 1183 | let children = Arc::new(ForegroundChildRegistry::new()); |
| 1184 | let child_token = turn_token.child_token(); |
| 1185 | let registration = children |
| 1186 | .register(child_token.clone(), "agent_healthy") |
| 1187 | .unwrap(); |
| 1188 | let parking = registration.parking_signal(); |
| 1189 | let (complete_tx, mut complete_rx) = tokio::sync::mpsc::unbounded_channel(); |
| 1190 | let (release_tx, release_rx) = tokio::sync::oneshot::channel(); |
| 1191 | let child = tokio::spawn(async move { |
| 1192 | release_rx.await.unwrap(); |
| 1193 | assert!(!child_token.is_cancelled()); |
| 1194 | complete_tx.send("existing completion inbox").unwrap(); |
| 1195 | drop(registration); |
| 1196 | }); |
| 1197 | let (flush_tx, flush_rx) = tokio::sync::oneshot::channel(); |
| 1198 | let drain_handle = tokio::spawn(async { |
| 1199 | let _ = flush_rx.await; |
| 1200 | }); |
| 1201 | let barrier = TurnMailboxBarrier { |
| 1202 | mailbox, |
| 1203 | cancel_token: turn_token.clone(), |
| 1204 | foreground_children: Arc::clone(&children), |
| 1205 | flush_tx, |
| 1206 | drain_handle, |
| 1207 | settle_grace: Duration::from_secs(1), |
| 1208 | }; |
| 1209 | tokio::time::timeout(Duration::from_secs(1), barrier.continue_and_flush()) |
| 1210 | .await |
| 1211 | .unwrap(); |
| 1212 | assert_eq!(children.active_count(), 1); |
| 1213 | assert!(!turn_token.is_cancelled()); |
| 1214 | assert!(!parking.load(Ordering::Acquire)); |
| 1215 | release_tx.send(()).unwrap(); |
| 1216 | assert_eq!(complete_rx.recv().await, Some("existing completion inbox")); |
| 1217 | child.await.unwrap(); |
| 1218 | assert_eq!(children.active_count(), 0); |
| 1219 | } |
| 1220 | |
| 1221 | #[tokio::test] |
| 1222 | async fn terminal_barrier_explicit_cancel_still_joins_owned_child() { |
| 1223 | let turn_token = CancellationToken::new(); |
| 1224 | let (mailbox, _receiver) = Mailbox::new(turn_token.clone()); |
| 1225 | let children = Arc::new(ForegroundChildRegistry::new()); |
| 1226 | let child_token = turn_token.child_token(); |
| 1227 | let registration = children |
| 1228 | .register(child_token.clone(), "agent_owned") |
| 1229 | .unwrap(); |
| 1230 | let child = tokio::spawn(async move { |
| 1231 | child_token.cancelled().await; |
| 1232 | drop(registration); |
| 1233 | }); |
| 1234 | let (flush_tx, flush_rx) = tokio::sync::oneshot::channel(); |
| 1235 | let drain_handle = tokio::spawn(async { |
| 1236 | let _ = flush_rx.await; |
| 1237 | }); |
| 1238 | let barrier = TurnMailboxBarrier { |
| 1239 | mailbox, |
| 1240 | cancel_token: turn_token, |
| 1241 | foreground_children: Arc::clone(&children), |
| 1242 | flush_tx, |
| 1243 | drain_handle, |
| 1244 | settle_grace: Duration::from_secs(1), |
| 1245 | }; |
| 1246 | let unsettled = tokio::time::timeout(Duration::from_secs(1), barrier.cancel_and_flush()) |
| 1247 | .await |
| 1248 | .unwrap(); |
| 1249 | assert!(unsettled.is_empty(), "cooperative children join cleanly"); |
| 1250 | child.await.unwrap(); |
| 1251 | assert_eq!(children.active_count(), 0); |
| 1252 | } |
| 1253 | |
| 1254 | /// Regression for #6184: a foreground child parked on an await that never |
| 1255 | /// observes its cancel token must not withhold the terminal turn event. The |
| 1256 | /// join gives up at `settle_grace` and names the child it left behind. |
| 1257 | #[tokio::test] |
| 1258 | async fn terminal_barrier_cancel_names_child_that_ignores_cancellation() { |
| 1259 | let turn_token = CancellationToken::new(); |
| 1260 | let (mailbox, _receiver) = Mailbox::new(turn_token.clone()); |
| 1261 | let children = Arc::new(ForegroundChildRegistry::new()); |
| 1262 | let child_token = turn_token.child_token(); |
| 1263 | // The fake child keeps its registration for the whole test — it never |
| 1264 | // observes the cancel token, like a task parked on a blocking await. |
| 1265 | let registration = children |
| 1266 | .register(child_token.clone(), "agent_stuck") |
| 1267 | .unwrap(); |
| 1268 | let (flush_tx, flush_rx) = tokio::sync::oneshot::channel(); |
| 1269 | let drain_handle = tokio::spawn(async { |
| 1270 | let _ = flush_rx.await; |
| 1271 | }); |
| 1272 | let barrier = TurnMailboxBarrier { |
| 1273 | mailbox, |
| 1274 | cancel_token: turn_token.clone(), |
| 1275 | foreground_children: Arc::clone(&children), |
| 1276 | flush_tx, |
| 1277 | drain_handle, |
| 1278 | settle_grace: Duration::from_millis(50), |
| 1279 | }; |
| 1280 | // Esc latches the turn token before the barrier runs, so the grace |
| 1281 | // window — not the token — is the bound under test. |
| 1282 | turn_token.cancel(); |
| 1283 | let unsettled = tokio::time::timeout(Duration::from_secs(1), barrier.cancel_and_flush()) |
| 1284 | .await |
| 1285 | .expect("the bounded join must not wait on a parked child"); |
| 1286 | assert_eq!(unsettled, vec!["agent_stuck".to_string()]); |
| 1287 | assert!( |
| 1288 | child_token.is_cancelled(), |
| 1289 | "the bounded join still cancels the child's token" |
| 1290 | ); |
| 1291 | assert_eq!( |
| 1292 | children.active_count(), |
| 1293 | 1, |
| 1294 | "the stuck child is left registered — leaked, not awaited" |
| 1295 | ); |
| 1296 | drop(registration); |
| 1297 | assert_eq!(children.active_count(), 0); |
| 1298 | } |
| 1299 | |
| 1300 | /// Regression for #6184: the mailbox drainer parked in an untimed |
| 1301 | /// `tx_event.send().await` (event channel full, UI not draining) must not |
| 1302 | /// withhold the terminal turn event. `flush` gives up at `settle_grace` and |
| 1303 | /// aborts the drainer, so the turn settles instead of waiting hours. |
| 1304 | #[tokio::test] |
| 1305 | async fn terminal_barrier_flush_bounds_a_drainer_parked_on_a_full_event_channel() { |
| 1306 | let turn_token = CancellationToken::new(); |
| 1307 | let (mailbox, _receiver) = Mailbox::new(turn_token.clone()); |
| 1308 | let children = Arc::new(ForegroundChildRegistry::new()); |
| 1309 | // The wedged shape from the report: a one-slot event channel, already |
| 1310 | // full, with nobody draining — so the drainer's forward parks in `send`. |
| 1311 | let (event_tx, mut event_rx) = tokio::sync::mpsc::channel::<()>(1); |
| 1312 | event_tx.send(()).await.unwrap(); |
| 1313 | let drain_handle = tokio::spawn(async move { |
| 1314 | // Parked forever: the channel is full and the receiver never drains. |
| 1315 | let _ = event_tx.send(()).await; |
| 1316 | }); |
| 1317 | // Give the drainer a moment to park before the flush races it. |
| 1318 | tokio::task::yield_now().await; |
| 1319 | let (flush_tx, _flush_rx) = tokio::sync::oneshot::channel(); |
| 1320 | let barrier = TurnMailboxBarrier { |
| 1321 | mailbox, |
| 1322 | cancel_token: turn_token, |
| 1323 | foreground_children: Arc::clone(&children), |
| 1324 | flush_tx, |
| 1325 | drain_handle, |
| 1326 | settle_grace: Duration::from_millis(50), |
| 1327 | }; |
| 1328 | let started = Instant::now(); |
| 1329 | tokio::time::timeout(Duration::from_secs(2), barrier.continue_and_flush()) |
| 1330 | .await |
| 1331 | .expect("flush must give up at its grace, not park with the drainer"); |
| 1332 | assert!( |
| 1333 | started.elapsed() < Duration::from_secs(2), |
| 1334 | "flush returned in {:?}", |
| 1335 | started.elapsed() |
| 1336 | ); |
| 1337 | // Abort receipt: an aborted drainer drops its sender half, so the |
| 1338 | // buffered item drains and the channel then reads closed. |
| 1339 | assert!(event_rx.recv().await.is_some()); |
| 1340 | assert!( |
| 1341 | event_rx.recv().await.is_none(), |
| 1342 | "aborted drainer must release the event channel" |
| 1343 | ); |
| 1344 | } |
| 1345 | |
| 1346 | /// A turn that fails without user cancellation runs the same barrier with a |
| 1347 | /// live turn token: Esc during the join must still break it rather than sit |
| 1348 | /// out the whole grace period (#6184). |
| 1349 | #[tokio::test] |
| 1350 | async fn terminal_barrier_cancel_join_breaks_on_fresh_esc() { |
| 1351 | let turn_token = CancellationToken::new(); |
| 1352 | let (mailbox, _receiver) = Mailbox::new(turn_token.clone()); |
| 1353 | let children = Arc::new(ForegroundChildRegistry::new()); |
| 1354 | let registration = children |
| 1355 | .register(turn_token.child_token(), "agent_stuck") |
| 1356 | .unwrap(); |
| 1357 | let (flush_tx, flush_rx) = tokio::sync::oneshot::channel(); |
| 1358 | let drain_handle = tokio::spawn(async { |
| 1359 | let _ = flush_rx.await; |
| 1360 | }); |
| 1361 | let barrier = TurnMailboxBarrier { |
| 1362 | mailbox, |
| 1363 | cancel_token: turn_token.clone(), |
| 1364 | foreground_children: Arc::clone(&children), |
| 1365 | flush_tx, |
| 1366 | drain_handle, |
| 1367 | settle_grace: Duration::from_secs(30), |
| 1368 | }; |
| 1369 | let esc = turn_token.clone(); |
| 1370 | let esc_task = tokio::spawn(async move { |
| 1371 | tokio::time::sleep(Duration::from_millis(20)).await; |
| 1372 | esc.cancel(); |
| 1373 | }); |
| 1374 | let unsettled = tokio::time::timeout(Duration::from_secs(1), barrier.cancel_and_flush()) |
| 1375 | .await |
| 1376 | .expect("a fresh Esc must break the join before its grace expires"); |
| 1377 | assert_eq!(unsettled, vec!["agent_stuck".to_string()]); |
| 1378 | esc_task.await.unwrap(); |
| 1379 | drop(registration); |
| 1380 | } |
| 1381 |