| 1 | |
| 2 | |
| 3 | #[tokio::test] |
| 4 | async fn queued_goal_clear_refreshes_prompt_and_cancels_stale_continuation() { |
| 5 | let config = Config::default(); |
| 6 | let engine_config = EngineConfig { |
| 7 | snapshots_enabled: false, |
| 8 | terminal_chrome_enabled: false, |
| 9 | goal_objective: Some("clear this goal".to_string()), |
| 10 | goal_token_budget: Some(42_000), |
| 11 | ..EngineConfig::default() |
| 12 | }; |
| 13 | let (engine, handle) = Engine::new(engine_config, &config); |
| 14 | let goal_state = engine.config.goal_state.clone(); |
| 15 | let run_task = tokio::spawn(engine.run()); |
| 16 | |
| 17 | // Model the mailbox order produced when the user clears a goal while its |
| 18 | // prior turn is still finishing: the control reaches the queue before the |
| 19 | // synthetic continuation that TurnComplete schedules. |
| 20 | handle |
| 21 | .send(Op::SetGoalStatus { |
| 22 | goal_id: None, |
| 23 | status: crate::tools::goal::GoalStatus::Active, |
| 24 | clear: true, |
| 25 | }) |
| 26 | .await |
| 27 | .expect("queue goal clear"); |
| 28 | handle |
| 29 | .send(Op::ContinueGoal { |
| 30 | dynamic_tools: Vec::new(), |
| 31 | engine_schedule_id: None, |
| 32 | }) |
| 33 | .await |
| 34 | .expect("queue stale continuation"); |
| 35 | |
| 36 | // This receipt sits behind both operations. Once it arrives, a stale |
| 37 | // continuation has either incorrectly started a turn or been consumed. |
| 38 | let session = handle |
| 39 | .get_session_snapshot() |
| 40 | .await |
| 41 | .expect("post-clear session snapshot"); |
| 42 | let prompt = match session.system_prompt.expect("post-clear system prompt") { |
| 43 | SystemPrompt::Text(text) => text, |
| 44 | SystemPrompt::Blocks(blocks) => blocks |
| 45 | .into_iter() |
| 46 | .map(|block| block.text) |
| 47 | .collect::<Vec<_>>() |
| 48 | .join("\n"), |
| 49 | }; |
| 50 | assert!( |
| 51 | !prompt.contains("<session_goal>"), |
| 52 | "cleared config fallback must not restore the goal prompt: {prompt}" |
| 53 | ); |
| 54 | let snapshot = goal_state.lock().expect("goal lock").snapshot(); |
| 55 | assert_eq!(snapshot.objective, None); |
| 56 | assert_eq!(snapshot.status, "none"); |
| 57 | assert_eq!(snapshot.token_budget, None); |
| 58 | |
| 59 | let mut saw_clear_session = false; |
| 60 | let mut saw_clear_goal = false; |
| 61 | let mut saw_clear_status = false; |
| 62 | { |
| 63 | let mut events = handle.rx_event.write().await; |
| 64 | while let Ok(event) = events.try_recv() { |
| 65 | match event { |
| 66 | Event::TurnStarted { .. } => { |
| 67 | panic!("queued clear must prevent a stale goal continuation") |
| 68 | } |
| 69 | Event::SessionUpdated { system_prompt, .. } => { |
| 70 | let prompt = match system_prompt.expect("clear SessionUpdated prompt") { |
| 71 | SystemPrompt::Text(text) => text, |
| 72 | SystemPrompt::Blocks(blocks) => blocks |
| 73 | .into_iter() |
| 74 | .map(|block| block.text) |
| 75 | .collect::<Vec<_>>() |
| 76 | .join("\n"), |
| 77 | }; |
| 78 | assert!(!prompt.contains("<session_goal>"), "{prompt}"); |
| 79 | saw_clear_session = true; |
| 80 | } |
| 81 | Event::GoalUpdated { snapshot } => { |
| 82 | assert_eq!(snapshot.objective, None); |
| 83 | assert_eq!(snapshot.status, "none"); |
| 84 | saw_clear_goal = true; |
| 85 | } |
| 86 | Event::Status { message } if message == "Goal cleared." => { |
| 87 | saw_clear_status = true; |
| 88 | } |
| 89 | _ => {} |
| 90 | } |
| 91 | } |
| 92 | } |
| 93 | assert!( |
| 94 | saw_clear_session, |
| 95 | "clear must refresh persisted prompt state" |
| 96 | ); |
| 97 | assert!(saw_clear_goal, "clear must emit a canonical empty snapshot"); |
| 98 | assert!(saw_clear_status, "clear must remain user-visible"); |
| 99 | |
| 100 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 101 | run_task.await.expect("engine task"); |
| 102 | } |
| 103 | |
| 104 | fn goal_custom_route_config() -> Config { |
| 105 | let mut custom = HashMap::new(); |
| 106 | custom.insert( |
| 107 | "custom-a".to_string(), |
| 108 | crate::config::ProviderConfig { |
| 109 | kind: Some("openai-compatible".to_string()), |
| 110 | base_url: Some("http://127.0.0.1:18181/v1".to_string()), |
| 111 | model: Some("local-model".to_string()), |
| 112 | api_key: Some("local-test-key".to_string()), |
| 113 | ..crate::config::ProviderConfig::default() |
| 114 | }, |
| 115 | ); |
| 116 | Config { |
| 117 | provider: Some("custom-a".to_string()), |
| 118 | providers: Some(crate::config::ProvidersConfig { |
| 119 | custom, |
| 120 | ..crate::config::ProvidersConfig::default() |
| 121 | }), |
| 122 | ..Config::default() |
| 123 | } |
| 124 | } |
| 125 | |
| 126 | #[tokio::test] |
| 127 | async fn ordinary_prose_never_activates_a_goal() { |
| 128 | let request_entered = std::sync::Arc::new(tokio::sync::Notify::new()); |
| 129 | let release_request = std::sync::Arc::new(tokio::sync::Notify::new()); |
| 130 | let model = std::sync::Arc::new(FirstRequestGatedGoalModelClient { |
| 131 | calls: std::sync::atomic::AtomicUsize::new(0), |
| 132 | request_entered: std::sync::Arc::clone(&request_entered), |
| 133 | release_request: std::sync::Arc::clone(&release_request), |
| 134 | }); |
| 135 | let config = goal_custom_route_config(); |
| 136 | let client: crate::core::model_client::SharedModelClient = model.clone(); |
| 137 | let (engine, handle) = Engine::new_with_model_client( |
| 138 | EngineConfig { |
| 139 | model: "local-model".to_string(), |
| 140 | max_steps: 1, |
| 141 | snapshots_enabled: false, |
| 142 | terminal_chrome_enabled: false, |
| 143 | ..EngineConfig::default() |
| 144 | }, |
| 145 | &config, |
| 146 | client, |
| 147 | ); |
| 148 | let goal_state = engine.config.goal_state.clone(); |
| 149 | let run_task = tokio::spawn(engine.run()); |
| 150 | |
| 151 | handle |
| 152 | .send(Op::SendMessage(TurnSpec { |
| 153 | max_output_tokens: None, |
| 154 | content: "hello - take over and make it your /goal to solve navier stokes".to_string(), |
| 155 | images: Vec::new(), |
| 156 | mode: AppMode::Agent, |
| 157 | route: resolved_route_for_test(&config, "local-model"), |
| 158 | compaction: Box::new(CompactionConfig::default()), |
| 159 | initial_routed_usage: Box::default(), |
| 160 | goal_objective: None, |
| 161 | goal_token_budget: None, |
| 162 | goal_status: crate::tools::goal::GoalStatus::Active, |
| 163 | reasoning_effort: None, |
| 164 | reasoning_effort_auto: false, |
| 165 | auto_model: false, |
| 166 | allow_shell: false, |
| 167 | trust_mode: false, |
| 168 | auto_approve: false, |
| 169 | approval_mode: ApprovalMode::Suggest, |
| 170 | translation_enabled: false, |
| 171 | allowed_tools: None, |
| 172 | dynamic_tools: Vec::new(), |
| 173 | hook_executor: None, |
| 174 | verbosity: None, |
| 175 | provenance: UserInputProvenance::ExternalUser, |
| 176 | submission_id: None, |
| 177 | })) |
| 178 | .await |
| 179 | .expect("send explicit natural goal turn"); |
| 180 | |
| 181 | // #6290 rework: the natural-language `/goal` prose parser is gone. The |
| 182 | // same wording that used to activate a goal is now an ordinary turn; |
| 183 | // only the model (`create_goal`) or the `/goal` command creates one. |
| 184 | let mut saw_goal = false; |
| 185 | loop { |
| 186 | let event = tokio::time::timeout(model_turn_event_timeout(), async { |
| 187 | handle.rx_event.write().await.recv().await |
| 188 | }) |
| 189 | .await |
| 190 | .expect("prose goal event timeout") |
| 191 | .expect("prose goal event"); |
| 192 | match event { |
| 193 | Event::GoalUpdated { .. } => { |
| 194 | saw_goal = true; |
| 195 | } |
| 196 | Event::TurnStarted { .. } => { |
| 197 | assert!( |
| 198 | !saw_goal, |
| 199 | "ordinary prose must not publish a goal before provider work starts" |
| 200 | ); |
| 201 | break; |
| 202 | } |
| 203 | _ => {} |
| 204 | } |
| 205 | } |
| 206 | |
| 207 | tokio::time::timeout(model_turn_event_timeout(), request_entered.notified()) |
| 208 | .await |
| 209 | .expect("provider request was never entered"); |
| 210 | release_request.notify_one(); |
| 211 | let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot()) |
| 212 | .await |
| 213 | .expect("turn did not settle") |
| 214 | .expect("post-turn session snapshot"); |
| 215 | let snapshot = goal_state.lock().expect("goal lock").snapshot(); |
| 216 | assert_eq!(snapshot.objective.as_deref(), None); |
| 217 | assert!(!snapshot.is_active()); |
| 218 | assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 1); |
| 219 | |
| 220 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 221 | run_task.await.expect("engine task"); |
| 222 | } |
| 223 | |
| 224 | /// Drive one turn in `mode` and report what the goal path did before the |
| 225 | /// provider was called: the objective published by `GoalUpdated` (if any), |
| 226 | /// whether the engine goal state is active, and how many Operate contract |
| 227 | /// messages the session log holds afterwards. |
| 228 | async fn operate_goal_probe(mode: AppMode, prompt: &str) -> (Option<String>, bool, usize) { |
| 229 | let request_entered = std::sync::Arc::new(tokio::sync::Notify::new()); |
| 230 | let release_request = std::sync::Arc::new(tokio::sync::Notify::new()); |
| 231 | let model = std::sync::Arc::new(FirstRequestGatedGoalModelClient { |
| 232 | calls: std::sync::atomic::AtomicUsize::new(0), |
| 233 | request_entered: std::sync::Arc::clone(&request_entered), |
| 234 | release_request: std::sync::Arc::clone(&release_request), |
| 235 | }); |
| 236 | let config = goal_custom_route_config(); |
| 237 | let client: crate::core::model_client::SharedModelClient = model.clone(); |
| 238 | let (engine, handle) = Engine::new_with_model_client( |
| 239 | EngineConfig { |
| 240 | model: "local-model".to_string(), |
| 241 | max_steps: 1, |
| 242 | snapshots_enabled: false, |
| 243 | terminal_chrome_enabled: false, |
| 244 | ..EngineConfig::default() |
| 245 | }, |
| 246 | &config, |
| 247 | client, |
| 248 | ); |
| 249 | let goal_state = engine.config.goal_state.clone(); |
| 250 | let run_task = tokio::spawn(engine.run()); |
| 251 | |
| 252 | handle |
| 253 | .send(Op::SendMessage(TurnSpec { |
| 254 | max_output_tokens: None, |
| 255 | content: prompt.to_string(), |
| 256 | images: Vec::new(), |
| 257 | mode, |
| 258 | route: resolved_route_for_test(&config, "local-model"), |
| 259 | compaction: Box::new(CompactionConfig::default()), |
| 260 | initial_routed_usage: Box::default(), |
| 261 | goal_objective: None, |
| 262 | goal_token_budget: None, |
| 263 | goal_status: crate::tools::goal::GoalStatus::Active, |
| 264 | reasoning_effort: None, |
| 265 | reasoning_effort_auto: false, |
| 266 | auto_model: false, |
| 267 | allow_shell: false, |
| 268 | trust_mode: false, |
| 269 | auto_approve: false, |
| 270 | approval_mode: ApprovalMode::Suggest, |
| 271 | translation_enabled: false, |
| 272 | allowed_tools: None, |
| 273 | dynamic_tools: Vec::new(), |
| 274 | hook_executor: None, |
| 275 | verbosity: None, |
| 276 | provenance: UserInputProvenance::ExternalUser, |
| 277 | submission_id: None, |
| 278 | })) |
| 279 | .await |
| 280 | .expect("send probe turn"); |
| 281 | |
| 282 | let mut published_objective = None; |
| 283 | loop { |
| 284 | let event = tokio::time::timeout(model_turn_event_timeout(), async { |
| 285 | handle.rx_event.write().await.recv().await |
| 286 | }) |
| 287 | .await |
| 288 | .expect("probe event timeout") |
| 289 | .expect("probe event"); |
| 290 | match event { |
| 291 | Event::GoalUpdated { snapshot } if snapshot.status == "active" => { |
| 292 | published_objective = snapshot.objective; |
| 293 | } |
| 294 | Event::TurnStarted { .. } => break, |
| 295 | _ => {} |
| 296 | } |
| 297 | } |
| 298 | tokio::time::timeout(model_turn_event_timeout(), request_entered.notified()) |
| 299 | .await |
| 300 | .expect("provider request was never entered"); |
| 301 | let active = goal_state.lock().expect("goal lock").is_active(); |
| 302 | if active { |
| 303 | // Stop autonomous continuation after the one provider call. |
| 304 | handle |
| 305 | .send(Op::SetGoalStatus { |
| 306 | goal_id: None, |
| 307 | status: crate::tools::goal::GoalStatus::Paused, |
| 308 | clear: false, |
| 309 | }) |
| 310 | .await |
| 311 | .expect("queue goal pause"); |
| 312 | } |
| 313 | release_request.notify_one(); |
| 314 | let snapshot = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot()) |
| 315 | .await |
| 316 | .expect("probe turn did not settle") |
| 317 | .expect("probe session snapshot"); |
| 318 | let contracts = snapshot |
| 319 | .messages |
| 320 | .iter() |
| 321 | .filter(|message| crate::runtime_handoff::is_operate_contract_message(message)) |
| 322 | .count(); |
| 323 | assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 1); |
| 324 | |
| 325 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 326 | run_task.await.expect("engine task"); |
| 327 | (published_objective, active, contracts) |
| 328 | } |
| 329 | |
| 330 | #[tokio::test] |
| 331 | async fn operate_never_promotes_wording_to_a_goal() { |
| 332 | let prompt = |
| 333 | "Migrate the settings loader to the new config crate and keep the old keys readable"; |
| 334 | |
| 335 | // The verb-list promotion is gone: an ordinary work prompt is an ordinary |
| 336 | // turn in every mode, and the model decides goals through `create_goal` |
| 337 | // (docs/design/TUI_DECONSTRUCTION.md, founder clarification 2026-09-09). |
| 338 | let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, prompt).await; |
| 339 | assert_eq!( |
| 340 | objective, None, |
| 341 | "the host must not infer a goal from wording" |
| 342 | ); |
| 343 | assert!(!active); |
| 344 | assert_eq!(contracts, 1, "Operate appends its contract exactly once"); |
| 345 | |
| 346 | let (objective, active, contracts) = operate_goal_probe(AppMode::Agent, prompt).await; |
| 347 | assert_eq!(objective, None, "Work must not promote a prompt to a goal"); |
| 348 | assert!(!active); |
| 349 | assert_eq!(contracts, 0, "Work never sees the Operate contract"); |
| 350 | |
| 351 | // #6290 rework: even an explicit-looking declaration is ordinary |
| 352 | // prose now — the host never parses it, and the model decides goals |
| 353 | // through `create_goal`. |
| 354 | let (objective, active, contracts) = |
| 355 | operate_goal_probe(AppMode::Operate, "Please set /goal to ship the release").await; |
| 356 | assert_eq!( |
| 357 | objective, None, |
| 358 | "prose asking for a goal must not create one host-side" |
| 359 | ); |
| 360 | assert!(!active); |
| 361 | assert_eq!(contracts, 1); |
| 362 | } |
| 363 | |
| 364 | #[tokio::test] |
| 365 | async fn operate_leaves_followup_and_long_questions_as_ordinary_turns() { |
| 366 | let report = "what about like rust or docker builds or something"; |
| 367 | let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, report).await; |
| 368 | assert_eq!( |
| 369 | objective, None, |
| 370 | "conversational followup must remain ordinary turn in Operate" |
| 371 | ); |
| 372 | assert!(!active); |
| 373 | assert_eq!(contracts, 1); |
| 374 | |
| 375 | let long_q = |
| 376 | "why did the build fail on the last step when running under docker on macos with rust 1.80"; |
| 377 | let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, long_q).await; |
| 378 | assert_eq!( |
| 379 | objective, None, |
| 380 | "long question without punctuation must remain ordinary turn" |
| 381 | ); |
| 382 | assert!(!active); |
| 383 | assert_eq!(contracts, 1); |
| 384 | |
| 385 | let zh_followup = "那 rust 或者 docker 构建呢"; |
| 386 | let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, zh_followup).await; |
| 387 | assert_eq!( |
| 388 | objective, None, |
| 389 | "Chinese followup must remain ordinary turn in Operate" |
| 390 | ); |
| 391 | assert!(!active); |
| 392 | assert_eq!(contracts, 1); |
| 393 | } |
| 394 | |
| 395 | #[tokio::test] |
| 396 | async fn operate_does_not_create_a_goal_when_the_work_request_declines_one() { |
| 397 | let prompt = "Run one bounded cancellation check. Do not edit files, inspect other files, create a goal, spawn agents, or start any other tool."; |
| 398 | let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, prompt).await; |
| 399 | assert_eq!(objective, None); |
| 400 | assert!(!active); |
| 401 | assert_eq!( |
| 402 | contracts, 1, |
| 403 | "the ordinary Operate turn still reaches the model" |
| 404 | ); |
| 405 | } |
| 406 | |
| 407 | #[tokio::test] |
| 408 | async fn operate_contract_is_appended_once_and_an_existing_goal_is_never_replaced() { |
| 409 | let first_entered = std::sync::Arc::new(tokio::sync::Notify::new()); |
| 410 | let release_first = std::sync::Arc::new(tokio::sync::Notify::new()); |
| 411 | let model = std::sync::Arc::new(IndexedGatedGoalModelClient { |
| 412 | calls: std::sync::atomic::AtomicUsize::new(0), |
| 413 | gates: HashMap::from([( |
| 414 | 1, |
| 415 | ( |
| 416 | std::sync::Arc::clone(&first_entered), |
| 417 | std::sync::Arc::clone(&release_first), |
| 418 | ), |
| 419 | )]), |
| 420 | max_calls: 2, |
| 421 | }); |
| 422 | let config = goal_custom_route_config(); |
| 423 | let client: crate::core::model_client::SharedModelClient = model.clone(); |
| 424 | let (engine, handle) = Engine::new_with_model_client( |
| 425 | EngineConfig { |
| 426 | model: "local-model".to_string(), |
| 427 | max_steps: 1, |
| 428 | snapshots_enabled: false, |
| 429 | terminal_chrome_enabled: false, |
| 430 | ..EngineConfig::default() |
| 431 | }, |
| 432 | &config, |
| 433 | client, |
| 434 | ); |
| 435 | let goal_state = engine.config.goal_state.clone(); |
| 436 | let run_task = tokio::spawn(engine.run()); |
| 437 | |
| 438 | let send = |content: &str, goal_objective: Option<String>, goal_status| { |
| 439 | Op::SendMessage(TurnSpec { |
| 440 | max_output_tokens: None, |
| 441 | content: content.to_string(), |
| 442 | images: Vec::new(), |
| 443 | mode: AppMode::Operate, |
| 444 | route: resolved_route_for_test(&config, "local-model"), |
| 445 | compaction: Box::new(CompactionConfig::default()), |
| 446 | initial_routed_usage: Box::default(), |
| 447 | goal_objective, |
| 448 | goal_token_budget: None, |
| 449 | goal_status, |
| 450 | reasoning_effort: None, |
| 451 | reasoning_effort_auto: false, |
| 452 | auto_model: false, |
| 453 | allow_shell: false, |
| 454 | trust_mode: false, |
| 455 | auto_approve: false, |
| 456 | approval_mode: ApprovalMode::Suggest, |
| 457 | translation_enabled: false, |
| 458 | allowed_tools: None, |
| 459 | dynamic_tools: Vec::new(), |
| 460 | hook_executor: None, |
| 461 | verbosity: None, |
| 462 | provenance: UserInputProvenance::ExternalUser, |
| 463 | submission_id: None, |
| 464 | }) |
| 465 | }; |
| 466 | |
| 467 | let first_objective = |
| 468 | "Migrate the settings loader to the new config crate and keep the old keys readable"; |
| 469 | // #6290 rework: prose no longer creates goals, so the unfinished goal |
| 470 | // this test needs is seeded directly — the same `GoalState::create` path |
| 471 | // the `/goal` command and the model's `create_goal` tool use. |
| 472 | goal_state |
| 473 | .lock() |
| 474 | .expect("goal lock") |
| 475 | .create(first_objective.to_string(), None) |
| 476 | .expect("seed unfinished goal"); |
| 477 | handle |
| 478 | .send(send( |
| 479 | first_objective, |
| 480 | Some(first_objective.to_string()), |
| 481 | crate::tools::goal::GoalStatus::Active, |
| 482 | )) |
| 483 | .await |
| 484 | .expect("send first Operate turn"); |
| 485 | tokio::time::timeout(model_turn_event_timeout(), first_entered.notified()) |
| 486 | .await |
| 487 | .expect("first provider request was never entered"); |
| 488 | assert_eq!( |
| 489 | goal_state.lock().expect("goal lock").objective(), |
| 490 | Some(first_objective) |
| 491 | ); |
| 492 | handle |
| 493 | .send(Op::SetGoalStatus { |
| 494 | goal_id: None, |
| 495 | status: crate::tools::goal::GoalStatus::Paused, |
| 496 | clear: false, |
| 497 | }) |
| 498 | .await |
| 499 | .expect("queue goal pause"); |
| 500 | release_first.notify_one(); |
| 501 | let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot()) |
| 502 | .await |
| 503 | .expect("first turn did not settle") |
| 504 | .expect("first session snapshot"); |
| 505 | |
| 506 | // Second Operate prompt while the (paused) goal is still unfinished: the |
| 507 | // host reports that goal, so the prompt is ordinary work under it. |
| 508 | handle |
| 509 | .send(send( |
| 510 | "Refactor the provider table so it survives a config reload", |
| 511 | Some(first_objective.to_string()), |
| 512 | crate::tools::goal::GoalStatus::Paused, |
| 513 | )) |
| 514 | .await |
| 515 | .expect("send second Operate turn"); |
| 516 | let snapshot = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot()) |
| 517 | .await |
| 518 | .expect("second turn did not settle") |
| 519 | .expect("second session snapshot"); |
| 520 | assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 2); |
| 521 | let contracts = snapshot |
| 522 | .messages |
| 523 | .iter() |
| 524 | .filter(|message| crate::runtime_handoff::is_operate_contract_message(message)) |
| 525 | .count(); |
| 526 | assert_eq!( |
| 527 | contracts, 1, |
| 528 | "the contract must not repeat on later Operate turns" |
| 529 | ); |
| 530 | let goal = goal_state.lock().expect("goal lock").snapshot(); |
| 531 | assert_eq!(goal.objective.as_deref(), Some(first_objective)); |
| 532 | assert_eq!(goal.status, "paused"); |
| 533 | |
| 534 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 535 | run_task.await.expect("engine task"); |
| 536 | } |
| 537 | |
| 538 | fn without_named_custom_route(mut config: Config) -> Config { |
| 539 | config |
| 540 | .providers |
| 541 | .as_mut() |
| 542 | .expect("custom providers") |
| 543 | .custom |
| 544 | .clear(); |
| 545 | config |
| 546 | } |
| 547 | |
| 548 | #[tokio::test] |
| 549 | async fn exhausted_goal_reaches_route_failure_without_budget_pause() { |
| 550 | let config = goal_custom_route_config(); |
| 551 | let engine_config = EngineConfig { |
| 552 | model: "local-model".to_string(), |
| 553 | snapshots_enabled: false, |
| 554 | terminal_chrome_enabled: false, |
| 555 | goal_objective: Some("stop at the budget".to_string()), |
| 556 | goal_token_budget: Some(10), |
| 557 | ..EngineConfig::default() |
| 558 | }; |
| 559 | let (mut engine, handle) = Engine::new(engine_config, &config); |
| 560 | let goal_state = engine.config.goal_state.clone(); |
| 561 | goal_state.lock().expect("goal lock").record_usage(11, 0); |
| 562 | |
| 563 | let invalid_route_config = without_named_custom_route(config); |
| 564 | engine.authoritative_route_config = |
| 565 | Some(Arc::new(parking_lot::RwLock::new(invalid_route_config))); |
| 566 | assert!( |
| 567 | engine.current_runtime_route().is_err(), |
| 568 | "fixture must prove route resolution cannot succeed" |
| 569 | ); |
| 570 | let run_task = tokio::spawn(engine.run()); |
| 571 | |
| 572 | handle |
| 573 | .send(Op::ContinueGoal { |
| 574 | dynamic_tools: Vec::new(), |
| 575 | engine_schedule_id: None, |
| 576 | }) |
| 577 | .await |
| 578 | .expect("queue exhausted continuation"); |
| 579 | |
| 580 | let mut saw_route_error = false; |
| 581 | let mut saw_blocked_goal = false; |
| 582 | while !(saw_route_error && saw_blocked_goal) { |
| 583 | let event = tokio::time::timeout(model_turn_event_timeout(), async { |
| 584 | handle.rx_event.write().await.recv().await |
| 585 | }) |
| 586 | .await |
| 587 | .expect("budget terminal event timeout") |
| 588 | .expect("budget terminal event"); |
| 589 | match event { |
| 590 | Event::TurnStarted { .. } => { |
| 591 | panic!("an exhausted goal must not start a turn before route resolution") |
| 592 | } |
| 593 | Event::Error { envelope, .. } => { |
| 594 | assert!( |
| 595 | envelope.message.contains("route is no longer valid"), |
| 596 | "the exhausted goal must reach the invalid-route failure, got: {envelope:?}" |
| 597 | ); |
| 598 | saw_route_error = true; |
| 599 | } |
| 600 | Event::GoalUpdated { snapshot } if snapshot.status == "paused" => { |
| 601 | panic!( |
| 602 | "budgets are telemetry-only in unbounded goal mode; the goal must not pause on budget (pause_reason={:?})", |
| 603 | snapshot.pause_reason |
| 604 | ); |
| 605 | } |
| 606 | Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => { |
| 607 | assert_eq!(snapshot.tokens_used, 11); |
| 608 | assert_eq!(snapshot.token_budget, Some(10)); |
| 609 | assert_eq!( |
| 610 | snapshot.pause_reason, None, |
| 611 | "budget must never be the pause reason in unbounded goal mode" |
| 612 | ); |
| 613 | saw_blocked_goal = true; |
| 614 | } |
| 615 | _ => {} |
| 616 | } |
| 617 | } |
| 618 | |
| 619 | let snapshot = goal_state.lock().expect("goal lock").snapshot(); |
| 620 | assert_eq!(snapshot.status, "blocked"); |
| 621 | assert_eq!(snapshot.tokens_used, 11); |
| 622 | assert_eq!(snapshot.token_budget, Some(10)); |
| 623 | assert_eq!(snapshot.pause_reason, None); |
| 624 | |
| 625 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 626 | run_task.await.expect("engine task"); |
| 627 | } |
| 628 | |
| 629 | #[tokio::test] |
| 630 | async fn continuation_circuit_breaker_pauses_with_run_limit_reason() { |
| 631 | // #5052: the backstop is configurable ([goal] max_continuations) and set |
| 632 | // deliberately past the retired hardcoded cap of 10 to prove an operate |
| 633 | // goal is no longer stopped there — only the configured backstop halts a |
| 634 | // pathological loop that never emits a terminal signal. |
| 635 | let backstop = 12u32; |
| 636 | let config = Config::default(); |
| 637 | let (engine, handle) = Engine::new( |
| 638 | EngineConfig { |
| 639 | snapshots_enabled: false, |
| 640 | terminal_chrome_enabled: false, |
| 641 | goal_objective: Some("stop a runaway continuation loop".to_string()), |
| 642 | goal_max_continuations: backstop, |
| 643 | ..EngineConfig::default() |
| 644 | }, |
| 645 | &config, |
| 646 | ); |
| 647 | let goal_state = engine.config.goal_state.clone(); |
| 648 | { |
| 649 | let mut goal = goal_state.lock().expect("goal lock"); |
| 650 | for _ in 0..backstop { |
| 651 | goal.record_continuation(); |
| 652 | } |
| 653 | } |
| 654 | let run_task = tokio::spawn(engine.run()); |
| 655 | |
| 656 | handle |
| 657 | .send(Op::ContinueGoal { |
| 658 | dynamic_tools: Vec::new(), |
| 659 | engine_schedule_id: None, |
| 660 | }) |
| 661 | .await |
| 662 | .expect("queue capped continuation"); |
| 663 | |
| 664 | let mut saw_pause = false; |
| 665 | let mut saw_reason = false; |
| 666 | while !(saw_pause && saw_reason) { |
| 667 | let event = tokio::time::timeout(model_turn_event_timeout(), async { |
| 668 | handle.rx_event.write().await.recv().await |
| 669 | }) |
| 670 | .await |
| 671 | .expect("continuation cap event timeout") |
| 672 | .expect("continuation cap event"); |
| 673 | match event { |
| 674 | Event::TurnStarted { .. } => panic!("capped goal must not start another turn"), |
| 675 | Event::GoalUpdated { snapshot } if snapshot.status == "paused" => { |
| 676 | assert_eq!( |
| 677 | snapshot.pause_reason, |
| 678 | Some(crate::tools::goal::GoalPauseReason::Backoff) |
| 679 | ); |
| 680 | saw_pause = true; |
| 681 | } |
| 682 | Event::Status { message } if message.contains("automatic continuations") => { |
| 683 | assert!(message.contains(&backstop.to_string()), "{message}"); |
| 684 | assert!(message.contains("[goal] max_continuations"), "{message}"); |
| 685 | saw_reason = true; |
| 686 | } |
| 687 | _ => {} |
| 688 | } |
| 689 | } |
| 690 | |
| 691 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 692 | run_task.await.expect("engine task"); |
| 693 | } |
| 694 | |
| 695 | #[tokio::test] |
| 696 | async fn goal_continues_past_legacy_ten_pass_cap_when_budget_remains() { |
| 697 | // #5052 regression: 10 automatic continuations used to be a terminal stop. |
| 698 | // With the default backstop and budget remaining, the loop must keep |
| 699 | // dispatching toward the completion gate. |
| 700 | let config = Config::default(); |
| 701 | let (engine, _handle) = Engine::new( |
| 702 | EngineConfig { |
| 703 | snapshots_enabled: false, |
| 704 | terminal_chrome_enabled: false, |
| 705 | goal_objective: Some("run to the completion gate, not a pass count".to_string()), |
| 706 | goal_token_budget: Some(1_000_000), |
| 707 | ..EngineConfig::default() |
| 708 | }, |
| 709 | &config, |
| 710 | ); |
| 711 | { |
| 712 | let mut goal = engine.config.goal_state.lock().expect("goal lock"); |
| 713 | for _ in 0..10 { |
| 714 | goal.record_continuation(); |
| 715 | } |
| 716 | } |
| 717 | |
| 718 | match engine.goal_continuation_if_active() { |
| 719 | GoalContinuationAction::Dispatch { snapshot, .. } => { |
| 720 | assert_eq!(snapshot.continuation_count, 11); |
| 721 | } |
| 722 | other => panic!("goal must continue past 10 passes, got {other:?}"), |
| 723 | } |
| 724 | |
| 725 | // A backstop of 0 means unlimited-with-budget-stops: even a pathological |
| 726 | // pass count keeps continuing while budget remains. |
| 727 | let (engine, _handle) = Engine::new( |
| 728 | EngineConfig { |
| 729 | snapshots_enabled: false, |
| 730 | terminal_chrome_enabled: false, |
| 731 | goal_objective: Some("unlimited backstop".to_string()), |
| 732 | goal_token_budget: Some(1_000_000), |
| 733 | goal_max_continuations: 0, |
| 734 | ..EngineConfig::default() |
| 735 | }, |
| 736 | &config, |
| 737 | ); |
| 738 | { |
| 739 | let mut goal = engine.config.goal_state.lock().expect("goal lock"); |
| 740 | for _ in 0..500 { |
| 741 | goal.record_continuation(); |
| 742 | } |
| 743 | } |
| 744 | assert!( |
| 745 | matches!( |
| 746 | engine.goal_continuation_if_active(), |
| 747 | GoalContinuationAction::Dispatch { .. } |
| 748 | ), |
| 749 | "backstop 0 must not stop an in-budget goal" |
| 750 | ); |
| 751 | } |
| 752 | |
| 753 | /// T08-03: the cross-turn continuation gate stops at an *enforced* token |
| 754 | /// budget, the same stop the intra-turn gate and the host take, and keeps |
| 755 | /// the default advisory behavior when enforcement is off. |
| 756 | #[tokio::test] |
| 757 | async fn cross_turn_goal_continuation_stops_at_an_enforced_token_budget() { |
| 758 | let config = Config::default(); |
| 759 | for enforce in [true, false] { |
| 760 | let (engine, _handle) = Engine::new( |
| 761 | EngineConfig { |
| 762 | snapshots_enabled: false, |
| 763 | terminal_chrome_enabled: false, |
| 764 | goal_objective: Some("stop at the enforced budget".to_string()), |
| 765 | goal_token_budget: Some(100), |
| 766 | goal_enforce_token_budget: enforce, |
| 767 | ..EngineConfig::default() |
| 768 | }, |
| 769 | &config, |
| 770 | ); |
| 771 | engine |
| 772 | .config |
| 773 | .goal_state |
| 774 | .lock() |
| 775 | .expect("goal lock") |
| 776 | .record_usage(100, 0); |
| 777 | let action = engine.goal_continuation_if_active(); |
| 778 | if enforce { |
| 779 | assert!( |
| 780 | matches!( |
| 781 | action, |
| 782 | GoalContinuationAction::Stopped { |
| 783 | reason: crate::tools::goal::GoalPauseReason::BudgetLimit, |
| 784 | .. |
| 785 | } |
| 786 | ), |
| 787 | "an exhausted enforced budget must not dispatch another turn: {action:?}" |
| 788 | ); |
| 789 | } else { |
| 790 | assert!( |
| 791 | matches!(action, GoalContinuationAction::Dispatch { .. }), |
| 792 | "an advisory budget keeps the goal running: {action:?}" |
| 793 | ); |
| 794 | } |
| 795 | } |
| 796 | } |
| 797 | |
| 798 | #[tokio::test] |
| 799 | async fn invalid_route_blocks_active_goal_and_refreshes_projections() { |
| 800 | let config = goal_custom_route_config(); |
| 801 | let engine_config = EngineConfig { |
| 802 | model: "local-model".to_string(), |
| 803 | snapshots_enabled: false, |
| 804 | terminal_chrome_enabled: false, |
| 805 | goal_objective: Some("keep going across route drift".to_string()), |
| 806 | ..EngineConfig::default() |
| 807 | }; |
| 808 | let (mut engine, handle) = Engine::new(engine_config, &config); |
| 809 | let goal_state = engine.config.goal_state.clone(); |
| 810 | engine.authoritative_route_config = Some(Arc::new(parking_lot::RwLock::new( |
| 811 | without_named_custom_route(config), |
| 812 | ))); |
| 813 | assert!( |
| 814 | engine.current_runtime_route().is_err(), |
| 815 | "fixture must fail before dispatch" |
| 816 | ); |
| 817 | let run_task = tokio::spawn(engine.run()); |
| 818 | |
| 819 | handle |
| 820 | .send(Op::ContinueGoal { |
| 821 | dynamic_tools: Vec::new(), |
| 822 | engine_schedule_id: None, |
| 823 | }) |
| 824 | .await |
| 825 | .expect("queue active continuation"); |
| 826 | let session = handle |
| 827 | .get_session_snapshot() |
| 828 | .await |
| 829 | .expect("post-route-failure session snapshot"); |
| 830 | |
| 831 | let prompt = match session.system_prompt.expect("blocked system prompt") { |
| 832 | SystemPrompt::Text(text) => text, |
| 833 | SystemPrompt::Blocks(blocks) => blocks |
| 834 | .into_iter() |
| 835 | .map(|block| block.text) |
| 836 | .collect::<Vec<_>>() |
| 837 | .join("\n"), |
| 838 | }; |
| 839 | assert!(!prompt.contains("<session_goal>"), "{prompt}"); |
| 840 | let snapshot = goal_state.lock().expect("goal lock").snapshot(); |
| 841 | assert_eq!(snapshot.status, "blocked"); |
| 842 | assert!( |
| 843 | snapshot |
| 844 | .blocker |
| 845 | .as_deref() |
| 846 | .is_some_and(|blocker| blocker.contains("provider route is no longer valid")), |
| 847 | "{snapshot:?}" |
| 848 | ); |
| 849 | |
| 850 | let mut saw_route_error = false; |
| 851 | let mut saw_blocked_session = false; |
| 852 | let mut saw_blocked_goal = false; |
| 853 | let mut saw_blocked_status = false; |
| 854 | { |
| 855 | let mut events = handle.rx_event.write().await; |
| 856 | while let Ok(event) = events.try_recv() { |
| 857 | match event { |
| 858 | Event::TurnStarted { .. } => { |
| 859 | panic!("invalid route must not start a continuation turn") |
| 860 | } |
| 861 | Event::Error { envelope, .. } => { |
| 862 | assert!(format!("{envelope:?}").contains("provider route is no longer valid")); |
| 863 | saw_route_error = true; |
| 864 | } |
| 865 | Event::SessionUpdated { |
| 866 | system_prompt: Some(system_prompt), |
| 867 | .. |
| 868 | } => { |
| 869 | let prompt = match system_prompt { |
| 870 | SystemPrompt::Text(text) => text, |
| 871 | SystemPrompt::Blocks(blocks) => blocks |
| 872 | .into_iter() |
| 873 | .map(|block| block.text) |
| 874 | .collect::<Vec<_>>() |
| 875 | .join("\n"), |
| 876 | }; |
| 877 | assert!(!prompt.contains("<session_goal>"), "{prompt}"); |
| 878 | saw_blocked_session = true; |
| 879 | } |
| 880 | Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => { |
| 881 | saw_blocked_goal = true; |
| 882 | } |
| 883 | Event::Status { message } |
| 884 | if message.contains("provider route is no longer valid") => |
| 885 | { |
| 886 | assert!(message.contains("resume the goal"), "{message}"); |
| 887 | saw_blocked_status = true; |
| 888 | } |
| 889 | _ => {} |
| 890 | } |
| 891 | } |
| 892 | } |
| 893 | assert!(saw_route_error, "route failure must remain visible"); |
| 894 | assert!( |
| 895 | saw_blocked_session, |
| 896 | "session prompt projection must refresh" |
| 897 | ); |
| 898 | assert!(saw_blocked_goal, "sidebar must receive blocked state"); |
| 899 | assert!(saw_blocked_status, "blocked reason must remain visible"); |
| 900 | |
| 901 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 902 | run_task.await.expect("engine task"); |
| 903 | } |
| 904 | |
| 905 | #[tokio::test] |
| 906 | async fn rejected_continuation_dispatch_blocks_goal_after_failed_turn() { |
| 907 | let config = goal_custom_route_config(); |
| 908 | let engine_config = EngineConfig { |
| 909 | model: "local-model".to_string(), |
| 910 | snapshots_enabled: false, |
| 911 | terminal_chrome_enabled: false, |
| 912 | goal_objective: Some("keep going after dispatch".to_string()), |
| 913 | ..EngineConfig::default() |
| 914 | }; |
| 915 | let (mut engine, handle) = Engine::new(engine_config, &config); |
| 916 | let goal_state = engine.config.goal_state.clone(); |
| 917 | assert!(engine.current_runtime_route().is_ok()); |
| 918 | // Exercise the continuation caller's `false` boundary deterministically: |
| 919 | // the route installs, but the injected-model authority has no client with |
| 920 | // which to start the request. |
| 921 | engine.model_client_injected = true; |
| 922 | engine.model_client = None; |
| 923 | let run_task = tokio::spawn(engine.run()); |
| 924 | |
| 925 | handle |
| 926 | .send(Op::ContinueGoal { |
| 927 | dynamic_tools: Vec::new(), |
| 928 | engine_schedule_id: None, |
| 929 | }) |
| 930 | .await |
| 931 | .expect("queue rejected continuation"); |
| 932 | let session = handle |
| 933 | .get_session_snapshot() |
| 934 | .await |
| 935 | .expect("post-rejection session snapshot"); |
| 936 | |
| 937 | let prompt = match session.system_prompt.expect("blocked system prompt") { |
| 938 | SystemPrompt::Text(text) => text, |
| 939 | SystemPrompt::Blocks(blocks) => blocks |
| 940 | .into_iter() |
| 941 | .map(|block| block.text) |
| 942 | .collect::<Vec<_>>() |
| 943 | .join("\n"), |
| 944 | }; |
| 945 | assert!(!prompt.contains("<session_goal>"), "{prompt}"); |
| 946 | let snapshot = goal_state.lock().expect("goal lock").snapshot(); |
| 947 | assert_eq!(snapshot.status, "blocked"); |
| 948 | assert!( |
| 949 | snapshot |
| 950 | .blocker |
| 951 | .as_deref() |
| 952 | .is_some_and(|blocker| blocker.contains("next model turn could not be started")), |
| 953 | "{snapshot:?}" |
| 954 | ); |
| 955 | |
| 956 | let mut starts = 0; |
| 957 | let mut saw_failed_turn = false; |
| 958 | let mut saw_dispatch_error = false; |
| 959 | let mut saw_blocked_session = false; |
| 960 | let mut saw_blocked_goal = false; |
| 961 | let mut saw_blocked_status = false; |
| 962 | { |
| 963 | let mut events = handle.rx_event.write().await; |
| 964 | while let Ok(event) = events.try_recv() { |
| 965 | match event { |
| 966 | Event::TurnStarted { .. } => starts += 1, |
| 967 | Event::TurnComplete { status, .. } => { |
| 968 | assert_eq!(status, TurnOutcomeStatus::Failed); |
| 969 | saw_failed_turn = true; |
| 970 | } |
| 971 | Event::Error { .. } => saw_dispatch_error = true, |
| 972 | Event::SessionUpdated { |
| 973 | system_prompt: Some(system_prompt), |
| 974 | .. |
| 975 | } => { |
| 976 | let prompt = match system_prompt { |
| 977 | SystemPrompt::Text(text) => text, |
| 978 | SystemPrompt::Blocks(blocks) => blocks |
| 979 | .into_iter() |
| 980 | .map(|block| block.text) |
| 981 | .collect::<Vec<_>>() |
| 982 | .join("\n"), |
| 983 | }; |
| 984 | assert!(!prompt.contains("<session_goal>"), "{prompt}"); |
| 985 | saw_blocked_session = true; |
| 986 | } |
| 987 | Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => { |
| 988 | saw_blocked_goal = true; |
| 989 | } |
| 990 | Event::Status { message } |
| 991 | if message.contains("next model turn could not be started") => |
| 992 | { |
| 993 | assert!(message.contains("resume the goal"), "{message}"); |
| 994 | saw_blocked_status = true; |
| 995 | } |
| 996 | _ => {} |
| 997 | } |
| 998 | } |
| 999 | } |
| 1000 | assert_eq!( |
| 1001 | starts, 1, |
| 1002 | "dispatch reached exactly one engine turn boundary" |
| 1003 | ); |
| 1004 | assert!( |
| 1005 | saw_failed_turn, |
| 1006 | "rejected dispatch must surface a failed turn" |
| 1007 | ); |
| 1008 | assert!( |
| 1009 | saw_dispatch_error, |
| 1010 | "rejected dispatch must surface its error" |
| 1011 | ); |
| 1012 | assert!( |
| 1013 | saw_blocked_session, |
| 1014 | "session prompt projection must refresh" |
| 1015 | ); |
| 1016 | assert!(saw_blocked_goal, "sidebar must receive blocked state"); |
| 1017 | assert!(saw_blocked_status, "blocked reason must remain visible"); |
| 1018 | |
| 1019 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 1020 | run_task.await.expect("engine task"); |
| 1021 | } |
| 1022 | |
| 1023 | #[tokio::test] |
| 1024 | async fn started_nonretryable_continuation_failure_blocks_goal_with_bounded_reason() { |
| 1025 | let failure_marker = "HTTP 400 Bad Request: deterministic continuation failure"; |
| 1026 | let leaked_secret = "sk-goal-secret-sentinel-123456"; |
| 1027 | let failure_message = format!( |
| 1028 | "{failure_marker}: {leaked_secret} {}", |
| 1029 | "provider detail ".repeat(80) |
| 1030 | ); |
| 1031 | let model = std::sync::Arc::new(FailingGoalModelClient { |
| 1032 | calls: std::sync::atomic::AtomicUsize::new(0), |
| 1033 | message: failure_message.clone(), |
| 1034 | }); |
| 1035 | let config = goal_custom_route_config(); |
| 1036 | let engine_config = EngineConfig { |
| 1037 | model: "local-model".to_string(), |
| 1038 | snapshots_enabled: false, |
| 1039 | terminal_chrome_enabled: false, |
| 1040 | goal_objective: Some("block a failed continuation truthfully".to_string()), |
| 1041 | ..EngineConfig::default() |
| 1042 | }; |
| 1043 | let client: crate::core::model_client::SharedModelClient = model.clone(); |
| 1044 | let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client); |
| 1045 | let goal_state = engine.config.goal_state.clone(); |
| 1046 | let run_task = tokio::spawn(engine.run()); |
| 1047 | |
| 1048 | handle |
| 1049 | .send(Op::ContinueGoal { |
| 1050 | dynamic_tools: Vec::new(), |
| 1051 | engine_schedule_id: None, |
| 1052 | }) |
| 1053 | .await |
| 1054 | .expect("queue failing continuation"); |
| 1055 | let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot()) |
| 1056 | .await |
| 1057 | .expect("failed continuation did not terminalize") |
| 1058 | .expect("post-failure session snapshot"); |
| 1059 | |
| 1060 | let prompt = match session.system_prompt.expect("blocked system prompt") { |
| 1061 | SystemPrompt::Text(text) => text, |
| 1062 | SystemPrompt::Blocks(blocks) => blocks |
| 1063 | .into_iter() |
| 1064 | .map(|block| block.text) |
| 1065 | .collect::<Vec<_>>() |
| 1066 | .join("\n"), |
| 1067 | }; |
| 1068 | assert!(!prompt.contains("<session_goal>"), "{prompt}"); |
| 1069 | let snapshot = goal_state.lock().expect("goal lock").snapshot(); |
| 1070 | assert_eq!(snapshot.status, "blocked"); |
| 1071 | let blocker = snapshot.blocker.as_deref().expect("failure blocker"); |
| 1072 | assert!(blocker.contains(failure_marker), "{blocker}"); |
| 1073 | assert!(blocker.contains("resume the goal"), "{blocker}"); |
| 1074 | assert!(!blocker.contains(leaked_secret), "{blocker}"); |
| 1075 | assert!( |
| 1076 | blocker.contains(codewhale_config::persistence::REDACTED), |
| 1077 | "{blocker}" |
| 1078 | ); |
| 1079 | assert!( |
| 1080 | blocker.len() <= GOAL_CONTINUATION_FAILURE_DETAIL_MAX_BYTES + 160, |
| 1081 | "failure reason must remain bounded: {} bytes", |
| 1082 | blocker.len() |
| 1083 | ); |
| 1084 | assert!( |
| 1085 | blocker.len() < failure_message.len(), |
| 1086 | "long provider detail must be truncated" |
| 1087 | ); |
| 1088 | assert_eq!( |
| 1089 | model.calls.load(std::sync::atomic::Ordering::SeqCst), |
| 1090 | 1, |
| 1091 | "a nonretryable started failure must not dispatch again" |
| 1092 | ); |
| 1093 | |
| 1094 | let mut starts = 0; |
| 1095 | let mut saw_failed_turn = false; |
| 1096 | let mut saw_provider_error = false; |
| 1097 | let mut saw_blocked_goal = false; |
| 1098 | let mut saw_blocked_status = false; |
| 1099 | { |
| 1100 | let mut events = handle.rx_event.write().await; |
| 1101 | while let Ok(event) = events.try_recv() { |
| 1102 | match event { |
| 1103 | Event::TurnStarted { .. } => starts += 1, |
| 1104 | Event::TurnComplete { status, error, .. } => { |
| 1105 | assert_eq!(status, TurnOutcomeStatus::Failed); |
| 1106 | assert!( |
| 1107 | error |
| 1108 | .as_deref() |
| 1109 | .is_some_and(|message| message.contains(failure_marker)), |
| 1110 | "{error:?}" |
| 1111 | ); |
| 1112 | saw_failed_turn = true; |
| 1113 | } |
| 1114 | Event::Error { envelope, .. } => { |
| 1115 | if envelope.message.contains(failure_marker) { |
| 1116 | saw_provider_error = true; |
| 1117 | } |
| 1118 | } |
| 1119 | Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => { |
| 1120 | saw_blocked_goal = true; |
| 1121 | } |
| 1122 | Event::Status { message } if message.contains(failure_marker) => { |
| 1123 | assert!(message.contains("resume the goal"), "{message}"); |
| 1124 | saw_blocked_status = true; |
| 1125 | } |
| 1126 | _ => {} |
| 1127 | } |
| 1128 | } |
| 1129 | } |
| 1130 | assert_eq!(starts, 1, "exactly one continuation turn must start"); |
| 1131 | assert!(saw_failed_turn, "failed turn receipt must remain visible"); |
| 1132 | assert!( |
| 1133 | saw_provider_error, |
| 1134 | "provider error event must remain visible" |
| 1135 | ); |
| 1136 | assert!(saw_blocked_goal, "goal must publish its blocked snapshot"); |
| 1137 | assert!(saw_blocked_status, "bounded failure must remain visible"); |
| 1138 | |
| 1139 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 1140 | run_task.await.expect("engine task"); |
| 1141 | } |
| 1142 | |
| 1143 | #[tokio::test] |
| 1144 | async fn headless_host_drains_existing_engine_completion_inbox_before_exit() { |
| 1145 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 1146 | |
| 1147 | let workspace = tempdir().unwrap(); |
| 1148 | let config = goal_custom_route_config(); |
| 1149 | let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn( |
| 1150 | "child evidence integrated by the existing Engine", |
| 1151 | )])); |
| 1152 | let (engine, handle) = Engine::new_with_model_client( |
| 1153 | EngineConfig { |
| 1154 | model: "local-model".into(), |
| 1155 | terminal_chrome_enabled: false, |
| 1156 | ..deterministic_engine_config(workspace.path()) |
| 1157 | }, |
| 1158 | &config, |
| 1159 | mock.clone(), |
| 1160 | ); |
| 1161 | assert!(engine.subagent_settlement_snapshot().await.is_settled()); |
| 1162 | // Reproduce the host boundary: the parent already ended, and a terminal |
| 1163 | // child's receipt is waiting for the Engine's normal idle fan-in path. |
| 1164 | engine |
| 1165 | .tx_event |
| 1166 | .send(Event::TurnComplete { |
| 1167 | usage: Usage::default(), |
| 1168 | parent_route_usage: Usage::default(), |
| 1169 | routed_usage_dropped_records: 0, |
| 1170 | status: TurnOutcomeStatus::Completed, |
| 1171 | error: None, |
| 1172 | tool_catalog: None, |
| 1173 | base_url: None, |
| 1174 | }) |
| 1175 | .await |
| 1176 | .unwrap(); |
| 1177 | engine |
| 1178 | .tx_subagent_completion |
| 1179 | .try_send(SubAgentCompletion { |
| 1180 | owner_session_id: engine.session.id.clone(), |
| 1181 | agent_id: "headless-settled-child".into(), |
| 1182 | payload: "bounded local fixture evidence".into(), |
| 1183 | }) |
| 1184 | .unwrap(); |
| 1185 | let pending = engine.subagent_settlement_snapshot().await; |
| 1186 | assert_eq!(pending.running_children, 0); |
| 1187 | assert_eq!(pending.pending_completions, 1); |
| 1188 | assert!( |
| 1189 | !pending.is_settled(), |
| 1190 | "terminal child alone cannot release the host" |
| 1191 | ); |
| 1192 | |
| 1193 | let run = tokio::spawn(engine.run()); |
| 1194 | let mut events = crate::exec_agent::ExecAgentEvents::new( |
| 1195 | handle.clone(), |
| 1196 | Instant::now() + model_turn_event_timeout(), |
| 1197 | ); |
| 1198 | let mut content = String::new(); |
| 1199 | let mut starts = 0; |
| 1200 | tokio::time::timeout(model_turn_event_timeout(), async { |
| 1201 | loop { |
| 1202 | match events.next().await.expect("host event") { |
| 1203 | Event::TurnStarted { .. } => starts += 1, |
| 1204 | Event::MessageDelta { content: delta, .. } => content.push_str(&delta), |
| 1205 | Event::TurnComplete { status, .. } => { |
| 1206 | assert_eq!(status, TurnOutcomeStatus::Completed); |
| 1207 | break; |
| 1208 | } |
| 1209 | _ => {} |
| 1210 | } |
| 1211 | } |
| 1212 | }) |
| 1213 | .await |
| 1214 | .expect("bounded headless settlement"); |
| 1215 | assert_eq!(starts, 1, "fan-in uses exactly one existing Engine turn"); |
| 1216 | assert_eq!(mock.call_count(), 1); |
| 1217 | assert!(content.contains("child evidence integrated")); |
| 1218 | assert!(!handle.is_cancelled()); |
| 1219 | handle.send(Op::Shutdown).await.unwrap(); |
| 1220 | tokio::time::timeout(Duration::from_secs(3), run) |
| 1221 | .await |
| 1222 | .unwrap() |
| 1223 | .unwrap(); |
| 1224 | } |
| 1225 | |
| 1226 | #[tokio::test] |
| 1227 | async fn host_managed_engine_does_not_self_dispatch_goal_continuation() { |
| 1228 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 1229 | |
| 1230 | let mut custom = HashMap::new(); |
| 1231 | custom.insert( |
| 1232 | "custom-a".to_string(), |
| 1233 | crate::config::ProviderConfig { |
| 1234 | kind: Some("openai-compatible".to_string()), |
| 1235 | base_url: Some("http://127.0.0.1:18181/v1".to_string()), |
| 1236 | model: Some("local-model".to_string()), |
| 1237 | api_key: Some("local-test-key".to_string()), |
| 1238 | ..crate::config::ProviderConfig::default() |
| 1239 | }, |
| 1240 | ); |
| 1241 | let config = Config { |
| 1242 | provider: Some("custom-a".to_string()), |
| 1243 | providers: Some(crate::config::ProvidersConfig { |
| 1244 | custom, |
| 1245 | ..crate::config::ProvidersConfig::default() |
| 1246 | }), |
| 1247 | ..Config::default() |
| 1248 | }; |
| 1249 | let runtime_services = crate::tools::spec::RuntimeToolServices { |
| 1250 | active_thread_id: Some("thr_host_managed".to_string()), |
| 1251 | ..crate::tools::spec::RuntimeToolServices::default() |
| 1252 | }; |
| 1253 | let engine_config = EngineConfig { |
| 1254 | max_steps: 0, |
| 1255 | snapshots_enabled: false, |
| 1256 | terminal_chrome_enabled: false, |
| 1257 | goal_objective: Some("keep going".to_string()), |
| 1258 | runtime_services, |
| 1259 | ..EngineConfig::default() |
| 1260 | }; |
| 1261 | let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn( |
| 1262 | "host-managed turn complete", |
| 1263 | )])); |
| 1264 | let client: crate::core::model_client::SharedModelClient = mock.clone(); |
| 1265 | let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client); |
| 1266 | let run_task = tokio::spawn(engine.run()); |
| 1267 | |
| 1268 | handle |
| 1269 | .send(Op::SendMessage(TurnSpec { |
| 1270 | max_output_tokens: None, |
| 1271 | content: "one host-owned turn".to_string(), |
| 1272 | images: Vec::new(), |
| 1273 | mode: AppMode::Agent, |
| 1274 | route: resolved_route_for_test(&config, "local-model"), |
| 1275 | compaction: Box::new(CompactionConfig::default()), |
| 1276 | initial_routed_usage: Box::default(), |
| 1277 | goal_objective: Some("keep going".to_string()), |
| 1278 | goal_token_budget: None, |
| 1279 | goal_status: crate::tools::goal::GoalStatus::Active, |
| 1280 | reasoning_effort: None, |
| 1281 | reasoning_effort_auto: false, |
| 1282 | auto_model: false, |
| 1283 | allow_shell: false, |
| 1284 | trust_mode: false, |
| 1285 | auto_approve: false, |
| 1286 | approval_mode: ApprovalMode::Suggest, |
| 1287 | translation_enabled: false, |
| 1288 | allowed_tools: None, |
| 1289 | dynamic_tools: Vec::new(), |
| 1290 | hook_executor: None, |
| 1291 | verbosity: None, |
| 1292 | provenance: UserInputProvenance::ExternalUser, |
| 1293 | submission_id: None, |
| 1294 | })) |
| 1295 | .await |
| 1296 | .expect("send host-owned goal turn"); |
| 1297 | |
| 1298 | let mut starts = 0; |
| 1299 | loop { |
| 1300 | let event = tokio::time::timeout(Duration::from_secs(3), async { |
| 1301 | handle.rx_event.write().await.recv().await |
| 1302 | }) |
| 1303 | .await |
| 1304 | .expect("host engine event timeout") |
| 1305 | .expect("host engine event"); |
| 1306 | match event { |
| 1307 | Event::TurnStarted { .. } => starts += 1, |
| 1308 | Event::TurnComplete { .. } => break, |
| 1309 | _ => {} |
| 1310 | } |
| 1311 | } |
| 1312 | assert_eq!(starts, 1); |
| 1313 | assert_eq!( |
| 1314 | mock.call_count(), |
| 1315 | 1, |
| 1316 | "the host-owned turn runs exactly once" |
| 1317 | ); |
| 1318 | assert!( |
| 1319 | tokio::time::timeout(Duration::from_millis(200), async { |
| 1320 | handle.rx_event.write().await.recv().await |
| 1321 | }) |
| 1322 | .await |
| 1323 | .is_err(), |
| 1324 | "a hosted engine must wait for an explicit durable turn claim" |
| 1325 | ); |
| 1326 | |
| 1327 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 1328 | run_task.await.expect("engine task"); |
| 1329 | } |
| 1330 | |
| 1331 | #[tokio::test] |
| 1332 | async fn cancellation_during_blocked_idle_handoff_survives_turn_admission() { |
| 1333 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 1334 | |
| 1335 | for child_completion in [true, false] { |
| 1336 | let workspace = tempdir().unwrap(); |
| 1337 | let config = goal_custom_route_config(); |
| 1338 | let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn( |
| 1339 | "must not dispatch after cancellation", |
| 1340 | )])); |
| 1341 | let (mut engine, mut handle) = Engine::new_with_model_client( |
| 1342 | EngineConfig { |
| 1343 | model: "local-model".into(), |
| 1344 | terminal_chrome_enabled: false, |
| 1345 | ..deterministic_engine_config(workspace.path()) |
| 1346 | }, |
| 1347 | &config, |
| 1348 | mock.clone(), |
| 1349 | ); |
| 1350 | // Force the handoff to stop at its status send after its initial |
| 1351 | // cancellation check and before handle_send_message admits a turn. |
| 1352 | let (tx, rx) = tokio::sync::mpsc::channel(1); |
| 1353 | engine.tx_event = tx; |
| 1354 | handle.rx_event = Arc::new(RwLock::new(rx)); |
| 1355 | engine.tx_event.send(Event::status("full")).await.unwrap(); |
| 1356 | let completion = SubAgentCompletion { |
| 1357 | owner_session_id: engine.session.id.clone(), |
| 1358 | agent_id: "cancel-race-child".into(), |
| 1359 | payload: "retained-after-cancel-race".into(), |
| 1360 | }; |
| 1361 | let mut wake: std::pin::Pin<Box<dyn std::future::Future<Output = ()> + '_>> = |
| 1362 | if child_completion { |
| 1363 | Box::pin(engine.handle_idle_subagent_completion(completion)) |
| 1364 | } else { |
| 1365 | Box::pin(engine.handle_idle_shell_completion_wake()) |
| 1366 | }; |
| 1367 | assert!( |
| 1368 | tokio::time::timeout(Duration::from_millis(5), &mut wake) |
| 1369 | .await |
| 1370 | .is_err(), |
| 1371 | "the full event channel must hold the handoff before admission" |
| 1372 | ); |
| 1373 | handle.cancel_with_reason(CancelReason::External); |
| 1374 | handle.rx_event.write().await.try_recv().unwrap(); |
| 1375 | tokio::time::timeout(Duration::from_secs(2), &mut wake) |
| 1376 | .await |
| 1377 | .expect("cancelled handoff returns without a provider call"); |
| 1378 | drop(wake); |
| 1379 | assert_eq!(mock.call_count(), 0); |
| 1380 | assert!(handle.is_cancelled()); |
| 1381 | assert!(engine.delivered_subagent_completion_ids.is_empty()); |
| 1382 | assert_eq!( |
| 1383 | engine.rx_subagent_completion.len(), |
| 1384 | usize::from(child_completion) |
| 1385 | ); |
| 1386 | if child_completion { |
| 1387 | assert!( |
| 1388 | engine |
| 1389 | .rx_subagent_completion |
| 1390 | .try_recv() |
| 1391 | .unwrap() |
| 1392 | .payload |
| 1393 | .contains("retained-after-cancel-race") |
| 1394 | ); |
| 1395 | } |
| 1396 | // A new explicit user action still receives a fresh turn control. |
| 1397 | let _turn = engine.begin_turn_control(); |
| 1398 | assert!(!handle.is_cancelled()); |
| 1399 | let existing_child = engine.cancel_token.child_token(); |
| 1400 | drop(_turn); |
| 1401 | let _automatic = |
| 1402 | engine.begin_turn_control_for_provenance(UserInputProvenance::SubAgentHandoff); |
| 1403 | handle.cancel(); |
| 1404 | assert!( |
| 1405 | existing_child.is_cancelled(), |
| 1406 | "stopping an automatic continuation must also stop earlier request siblings" |
| 1407 | ); |
| 1408 | } |
| 1409 | } |
| 1410 | |
| 1411 | #[tokio::test] |
| 1412 | async fn cancelled_parent_defers_child_receipts_until_an_explicit_turn() { |
| 1413 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 1414 | |
| 1415 | for reason in [CancelReason::User, CancelReason::External] { |
| 1416 | let workspace = tempdir().unwrap(); |
| 1417 | let config = Config::default(); |
| 1418 | let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn( |
| 1419 | "Explicit continuation completed.", |
| 1420 | )])); |
| 1421 | let (mut engine, handle) = Engine::new_with_model_client( |
| 1422 | deterministic_engine_config(workspace.path()), |
| 1423 | &config, |
| 1424 | mock.clone(), |
| 1425 | ); |
| 1426 | handle.cancel_with_reason(reason); |
| 1427 | // Exercise a completion selected just before cancellation arrived. |
| 1428 | engine |
| 1429 | .handle_idle_subagent_completion(SubAgentCompletion { |
| 1430 | owner_session_id: engine.session.id.clone(), |
| 1431 | agent_id: "cancelled-worker".into(), |
| 1432 | payload: "parked-child-evidence".into(), |
| 1433 | }) |
| 1434 | .await; |
| 1435 | assert_eq!( |
| 1436 | mock.call_count(), |
| 1437 | 0, |
| 1438 | "cancellation must forbid a model wake" |
| 1439 | ); |
| 1440 | assert!(engine.delivered_subagent_completion_ids.is_empty()); |
| 1441 | assert_eq!(engine.rx_subagent_completion.len(), 1); |
| 1442 | assert!( |
| 1443 | tokio::time::timeout(Duration::from_millis(50), engine.next_run_input(false)) |
| 1444 | .await |
| 1445 | .is_err(), |
| 1446 | "the idle loop must leave the receipt queued without spinning" |
| 1447 | ); |
| 1448 | |
| 1449 | let run = tokio::spawn(engine.run()); |
| 1450 | handle |
| 1451 | .send(external_user_message_op( |
| 1452 | "Continue explicitly", |
| 1453 | AppMode::Agent, |
| 1454 | &config, |
| 1455 | )) |
| 1456 | .await |
| 1457 | .unwrap(); |
| 1458 | tokio::time::timeout(Duration::from_secs(5), async { |
| 1459 | while let Some(event) = handle.rx_event.write().await.recv().await { |
| 1460 | if matches!(event, Event::TurnComplete { .. }) { |
| 1461 | break; |
| 1462 | } |
| 1463 | } |
| 1464 | }) |
| 1465 | .await |
| 1466 | .expect("explicit turn completes"); |
| 1467 | assert_eq!(mock.call_count(), 1); |
| 1468 | let snapshot = handle.get_session_snapshot().await.unwrap(); |
| 1469 | assert!( |
| 1470 | serde_json::to_string(&snapshot.messages) |
| 1471 | .unwrap() |
| 1472 | .contains("parked-child-evidence"), |
| 1473 | "the next explicit turn must retain the completion receipt" |
| 1474 | ); |
| 1475 | handle.send(Op::Shutdown).await.unwrap(); |
| 1476 | run.await.unwrap(); |
| 1477 | } |
| 1478 | } |