| 1 | /// A thinking-only mid-stream drop must recover without ever claiming a |
| 2 | /// visible partial reply was preserved, and must leave exactly one |
| 3 | /// authoritative assistant answer in the persisted conversation. |
| 4 | #[tokio::test] |
| 5 | async fn interactive_thinking_only_drop_preserves_nothing_and_never_claims_it_did() { |
| 6 | let model = std::sync::Arc::new(ThinkingOnlyDropModelClient { |
| 7 | calls: std::sync::atomic::AtomicUsize::new(0), |
| 8 | failures: 1, |
| 9 | }); |
| 10 | let client: crate::core::model_client::SharedModelClient = model.clone(); |
| 11 | let config = Config::default(); |
| 12 | let engine_config = EngineConfig { |
| 13 | max_steps: 1, |
| 14 | snapshots_enabled: false, |
| 15 | subagents_enabled: false, |
| 16 | terminal_chrome_enabled: true, |
| 17 | ..EngineConfig::default() |
| 18 | }; |
| 19 | let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client); |
| 20 | let run_task = tokio::spawn(engine.run()); |
| 21 | |
| 22 | handle |
| 23 | .send(Op::SendMessage(TurnSpec { |
| 24 | max_output_tokens: None, |
| 25 | content: "solve the task".to_string(), |
| 26 | images: Vec::new(), |
| 27 | mode: AppMode::Agent, |
| 28 | route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL), |
| 29 | compaction: Box::new(CompactionConfig::default()), |
| 30 | initial_routed_usage: Box::default(), |
| 31 | goal_objective: None, |
| 32 | goal_token_budget: None, |
| 33 | goal_status: crate::tools::goal::GoalStatus::Active, |
| 34 | reasoning_effort: None, |
| 35 | reasoning_effort_auto: false, |
| 36 | auto_model: false, |
| 37 | allow_shell: false, |
| 38 | trust_mode: false, |
| 39 | auto_approve: false, |
| 40 | approval_mode: ApprovalMode::Suggest, |
| 41 | translation_enabled: false, |
| 42 | allowed_tools: None, |
| 43 | dynamic_tools: Vec::new(), |
| 44 | hook_executor: None, |
| 45 | verbosity: None, |
| 46 | provenance: UserInputProvenance::ExternalUser, |
| 47 | submission_id: None, |
| 48 | })) |
| 49 | .await |
| 50 | .expect("send thinking-only drop turn"); |
| 51 | |
| 52 | let mut events = Vec::new(); |
| 53 | loop { |
| 54 | let event = tokio::time::timeout(model_turn_event_timeout(), async { |
| 55 | handle.rx_event.write().await.recv().await |
| 56 | }) |
| 57 | .await |
| 58 | .expect("thinking-only drop event timeout") |
| 59 | .expect("thinking-only drop event"); |
| 60 | let terminal = matches!(event, Event::TurnComplete { .. }); |
| 61 | events.push(event); |
| 62 | if terminal { |
| 63 | break; |
| 64 | } |
| 65 | } |
| 66 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 67 | run_task.await.expect("engine task"); |
| 68 | |
| 69 | assert_eq!( |
| 70 | model.calls.load(std::sync::atomic::Ordering::SeqCst), |
| 71 | 2, |
| 72 | "the thinking-only drop must be re-issued exactly once" |
| 73 | ); |
| 74 | let status = events |
| 75 | .iter() |
| 76 | .find_map(|event| match event { |
| 77 | Event::TurnComplete { status, .. } => Some(*status), |
| 78 | _ => None, |
| 79 | }) |
| 80 | .expect("terminal TurnComplete"); |
| 81 | assert_eq!(status, TurnOutcomeStatus::Completed); |
| 82 | |
| 83 | assert!( |
| 84 | events.iter().any(|event| matches!( |
| 85 | event, |
| 86 | Event::Status { message } if message.starts_with("Retry attempt: stream-resume 1/") |
| 87 | )), |
| 88 | "the first retry must be visible without claiming hidden reasoning was preserved" |
| 89 | ); |
| 90 | |
| 91 | // The persisted conversation keeps the operator's turn and exactly one |
| 92 | // authoritative assistant answer — no synthetic `[runtime]` user message, |
| 93 | // no duplicated answer, no orphaned thinking-only assistant cell. |
| 94 | let transcript = events |
| 95 | .iter() |
| 96 | .rev() |
| 97 | .find_map(|event| match event { |
| 98 | Event::SessionUpdated { messages, .. } => Some(messages.clone()), |
| 99 | _ => None, |
| 100 | }) |
| 101 | .expect("final SessionUpdated"); |
| 102 | assert_eq!( |
| 103 | transcript |
| 104 | .iter() |
| 105 | .filter(|message| message.role == "user") |
| 106 | .count(), |
| 107 | 1, |
| 108 | "the operator's own turn must be the only user message: {transcript:?}" |
| 109 | ); |
| 110 | assert_eq!( |
| 111 | transcript |
| 112 | .iter() |
| 113 | .filter(|message| message.role == "assistant") |
| 114 | .count(), |
| 115 | 1, |
| 116 | "exactly one authoritative assistant answer after recovery: {transcript:?}" |
| 117 | ); |
| 118 | let transcript_text = transcript |
| 119 | .iter() |
| 120 | .flat_map(|message| message.content.iter()) |
| 121 | .filter_map(|block| match block { |
| 122 | ContentBlock::Text { text, .. } => Some(text.as_str()), |
| 123 | _ => None, |
| 124 | }) |
| 125 | .collect::<Vec<_>>() |
| 126 | .join("\n"); |
| 127 | assert!( |
| 128 | !transcript_text.contains("[runtime]"), |
| 129 | "a retried turn must not insert a synthetic user message: {transcript_text}" |
| 130 | ); |
| 131 | assert!( |
| 132 | !transcript_text.contains("hidden reasoning that no operator ever saw"), |
| 133 | "an invisible thinking-only fragment must not be persisted as reply text: {transcript_text}" |
| 134 | ); |
| 135 | assert_eq!( |
| 136 | transcript_text |
| 137 | .matches("the one authoritative answer") |
| 138 | .count(), |
| 139 | 1, |
| 140 | "the recovered answer must be persisted exactly once: {transcript_text}" |
| 141 | ); |
| 142 | } |
| 143 | |
| 144 | /// A model client that answers with ONLY hidden reasoning and a clean stop for |
| 145 | /// its first `reasoning_only` calls, then a real text answer. No transport |
| 146 | /// error: the stream completes normally but carries no sendable content — the |
| 147 | /// reasoning-model failure mode #5546-adjacent that used to dead-end the turn. |
| 148 | struct ReasoningOnlyCleanFinishModelClient { |
| 149 | calls: std::sync::atomic::AtomicUsize, |
| 150 | reasoning_only: usize, |
| 151 | stop_reason: &'static str, |
| 152 | /// Every outbound request's messages, in order, so a test can tell a |
| 153 | /// request-scoped nudge from one written into the session. |
| 154 | requests: std::sync::Mutex<Vec<Vec<codewhale_models::Message>>>, |
| 155 | } |
| 156 | |
| 157 | #[async_trait::async_trait] |
| 158 | impl crate::core::model_client::ModelClient for ReasoningOnlyCleanFinishModelClient { |
| 159 | fn provider_name(&self) -> &str { |
| 160 | "reasoning-only" |
| 161 | } |
| 162 | |
| 163 | fn model(&self) -> &str { |
| 164 | "local-model" |
| 165 | } |
| 166 | |
| 167 | async fn create_message( |
| 168 | &self, |
| 169 | _request: codewhale_models::MessageRequest, |
| 170 | ) -> anyhow::Result<codewhale_models::MessageResponse> { |
| 171 | anyhow::bail!("reasoning-only recovery uses the streaming model boundary") |
| 172 | } |
| 173 | |
| 174 | async fn create_message_stream( |
| 175 | &self, |
| 176 | _request: codewhale_models::MessageRequest, |
| 177 | ) -> anyhow::Result<crate::llm_client::StreamEventBox> { |
| 178 | use crate::llm_client::mock::canned; |
| 179 | if let Ok(mut requests) = self.requests.lock() { |
| 180 | requests.push(_request.messages.clone()); |
| 181 | } |
| 182 | let call = self |
| 183 | .calls |
| 184 | .fetch_add(1, std::sync::atomic::Ordering::SeqCst) |
| 185 | .saturating_add(1); |
| 186 | if call <= self.reasoning_only { |
| 187 | // A protocol-complete response that opened and closed only a |
| 188 | // thinking block: no text, no tool call, and a clean stop reason. |
| 189 | let events: Vec<anyhow::Result<codewhale_models::StreamEvent>> = vec![ |
| 190 | Ok(canned::message_start("reasoning_only_msg")), |
| 191 | Ok(StreamEvent::ContentBlockStart { |
| 192 | index: 0, |
| 193 | content_block: codewhale_models::ContentBlockStart::Thinking { |
| 194 | thinking: String::new(), |
| 195 | }, |
| 196 | }), |
| 197 | Ok(canned::thinking_delta(0, "reasoning with no final channel")), |
| 198 | Ok(canned::block_stop(0)), |
| 199 | Ok(canned::message_delta(self.stop_reason, None)), |
| 200 | Ok(canned::message_stop()), |
| 201 | ]; |
| 202 | return Ok(Box::pin(futures_util::stream::iter(events))); |
| 203 | } |
| 204 | let events = canned::simple_text_turn("the recovered answer") |
| 205 | .into_iter() |
| 206 | .map(Ok); |
| 207 | Ok(Box::pin(futures_util::stream::iter(events))) |
| 208 | } |
| 209 | |
| 210 | async fn health_check(&self) -> anyhow::Result<bool> { |
| 211 | Ok(true) |
| 212 | } |
| 213 | } |
| 214 | |
| 215 | async fn run_reasoning_only_turn( |
| 216 | reasoning_only: usize, |
| 217 | stop_reason: &'static str, |
| 218 | ) -> ( |
| 219 | std::sync::Arc<ReasoningOnlyCleanFinishModelClient>, |
| 220 | Vec<Event>, |
| 221 | ) { |
| 222 | run_reasoning_only_turn_with_reprompts( |
| 223 | reasoning_only, |
| 224 | stop_reason, |
| 225 | crate::config::DEFAULT_REASONING_ONLY_REPROMPTS, |
| 226 | ) |
| 227 | .await |
| 228 | } |
| 229 | |
| 230 | async fn run_reasoning_only_turn_with_reprompts( |
| 231 | reasoning_only: usize, |
| 232 | stop_reason: &'static str, |
| 233 | max_reprompts: u32, |
| 234 | ) -> ( |
| 235 | std::sync::Arc<ReasoningOnlyCleanFinishModelClient>, |
| 236 | Vec<Event>, |
| 237 | ) { |
| 238 | let model = std::sync::Arc::new(ReasoningOnlyCleanFinishModelClient { |
| 239 | calls: std::sync::atomic::AtomicUsize::new(0), |
| 240 | reasoning_only, |
| 241 | stop_reason, |
| 242 | requests: std::sync::Mutex::new(Vec::new()), |
| 243 | }); |
| 244 | let client: crate::core::model_client::SharedModelClient = model.clone(); |
| 245 | let config = Config::default(); |
| 246 | let engine_config = EngineConfig { |
| 247 | max_steps: 1, |
| 248 | snapshots_enabled: false, |
| 249 | subagents_enabled: false, |
| 250 | terminal_chrome_enabled: true, |
| 251 | reasoning_only_max_reprompts: max_reprompts, |
| 252 | ..EngineConfig::default() |
| 253 | }; |
| 254 | let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client); |
| 255 | let run_task = tokio::spawn(engine.run()); |
| 256 | handle |
| 257 | .send(Op::SendMessage(TurnSpec { |
| 258 | max_output_tokens: None, |
| 259 | content: "solve the task".to_string(), |
| 260 | images: Vec::new(), |
| 261 | mode: AppMode::Agent, |
| 262 | route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL), |
| 263 | compaction: Box::new(CompactionConfig::default()), |
| 264 | initial_routed_usage: Box::default(), |
| 265 | goal_objective: None, |
| 266 | goal_token_budget: None, |
| 267 | goal_status: crate::tools::goal::GoalStatus::Active, |
| 268 | reasoning_effort: None, |
| 269 | reasoning_effort_auto: false, |
| 270 | auto_model: false, |
| 271 | allow_shell: false, |
| 272 | trust_mode: false, |
| 273 | auto_approve: false, |
| 274 | approval_mode: ApprovalMode::Suggest, |
| 275 | translation_enabled: false, |
| 276 | allowed_tools: None, |
| 277 | dynamic_tools: Vec::new(), |
| 278 | hook_executor: None, |
| 279 | verbosity: None, |
| 280 | provenance: UserInputProvenance::ExternalUser, |
| 281 | submission_id: None, |
| 282 | })) |
| 283 | .await |
| 284 | .expect("send reasoning-only turn"); |
| 285 | let mut events = Vec::new(); |
| 286 | loop { |
| 287 | let event = tokio::time::timeout(model_turn_event_timeout(), async { |
| 288 | handle.rx_event.write().await.recv().await |
| 289 | }) |
| 290 | .await |
| 291 | .expect("reasoning-only event timeout") |
| 292 | .expect("reasoning-only event"); |
| 293 | let terminal = matches!(event, Event::TurnComplete { .. }); |
| 294 | events.push(event); |
| 295 | if terminal { |
| 296 | break; |
| 297 | } |
| 298 | } |
| 299 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 300 | run_task.await.expect("engine task"); |
| 301 | (model, events) |
| 302 | } |
| 303 | |
| 304 | /// A reasoning-only clean-stop response is re-requested and the turn recovers |
| 305 | /// with the real answer. Local fixtures make no cache-hit or billing claim. |
| 306 | #[tokio::test] |
| 307 | async fn reasoning_only_clean_stop_is_retried_and_recovers() { |
| 308 | let (model, events) = run_reasoning_only_turn(1, "stop").await; |
| 309 | |
| 310 | let terminal = events |
| 311 | .iter() |
| 312 | .find_map(|event| match event { |
| 313 | Event::ToolRequestSnapshot { snapshot } => snapshot.terminal.as_ref(), |
| 314 | _ => None, |
| 315 | }) |
| 316 | .expect("terminal diagnostics through existing event authority"); |
| 317 | assert_eq!(terminal.model_requests_started, 2); |
| 318 | assert_eq!(terminal.reasoning_only_reprompts, 1); |
| 319 | assert_eq!(terminal.transparent_stream_retries, 0); |
| 320 | assert_eq!(terminal.status, Some(TurnOutcomeStatus::Completed)); |
| 321 | |
| 322 | assert_eq!( |
| 323 | model.calls.load(std::sync::atomic::Ordering::SeqCst), |
| 324 | 2, |
| 325 | "the reasoning-only response must be re-requested exactly once" |
| 326 | ); |
| 327 | let status = events |
| 328 | .iter() |
| 329 | .find_map(|event| match event { |
| 330 | Event::TurnComplete { status, .. } => Some(*status), |
| 331 | _ => None, |
| 332 | }) |
| 333 | .expect("terminal TurnComplete"); |
| 334 | assert_eq!(status, TurnOutcomeStatus::Completed); |
| 335 | |
| 336 | let recovery = events |
| 337 | .iter() |
| 338 | .filter_map(|event| match event { |
| 339 | Event::Status { message } if message.contains("re-requesting the answer") => { |
| 340 | Some(message.clone()) |
| 341 | } |
| 342 | _ => None, |
| 343 | }) |
| 344 | .collect::<Vec<_>>(); |
| 345 | assert_eq!(recovery.len(), 1, "exactly one recovery notice: {events:?}"); |
| 346 | assert!( |
| 347 | recovery[0].starts_with("Retry attempt: reasoning-only 1/"), |
| 348 | "attempt announced: {recovery:?}" |
| 349 | ); |
| 350 | |
| 351 | // No hard failure surfaced. |
| 352 | assert!( |
| 353 | !events.iter().any(|event| matches!( |
| 354 | event, |
| 355 | Event::Error { envelope, .. } if envelope.message.contains("no answer or tool call") |
| 356 | )), |
| 357 | "a recovered turn must not surface the incomplete-response error: {events:?}" |
| 358 | ); |
| 359 | assert_eq!(events.iter().filter(|event| matches!(event, Event::Status { message } if message.starts_with("Retry recovery: reasoning-only used 1/") && message.ends_with("turn completed"))).count(), 1); |
| 360 | } |
| 361 | |
| 362 | /// An output-length stop is a real budget hit, not a transient — it must NOT |
| 363 | /// be retried and must fail honestly. |
| 364 | #[tokio::test] |
| 365 | async fn reasoning_only_length_stop_fails_without_retry() { |
| 366 | let (model, events) = run_reasoning_only_turn(1, "length").await; |
| 367 | |
| 368 | let diagnostic = events |
| 369 | .iter() |
| 370 | .find_map(|event| match event { |
| 371 | Event::ToolRequestSnapshot { snapshot } => snapshot.terminal.as_ref(), |
| 372 | _ => None, |
| 373 | }) |
| 374 | .expect("terminal request diagnostics"); |
| 375 | assert!( |
| 376 | diagnostic |
| 377 | .last_prepared_output_limit_tokens |
| 378 | .is_some_and(|tokens| tokens > 0) |
| 379 | ); |
| 380 | assert!( |
| 381 | events.iter().any(|event| matches!( |
| 382 | event, |
| 383 | Event::Error { envelope, .. } |
| 384 | if envelope.message.contains("response output limit") |
| 385 | && envelope.message.contains("including reasoning") |
| 386 | )), |
| 387 | "a length stop must explain the actual output constraint" |
| 388 | ); |
| 389 | |
| 390 | assert_eq!( |
| 391 | model.calls.load(std::sync::atomic::Ordering::SeqCst), |
| 392 | 1, |
| 393 | "a length stop must not be retried" |
| 394 | ); |
| 395 | assert!( |
| 396 | !events.iter().any(|event| matches!( |
| 397 | event, |
| 398 | Event::Status { message } if message.contains("re-requesting the answer") |
| 399 | )), |
| 400 | "a length stop must not announce a retry: {events:?}" |
| 401 | ); |
| 402 | let status = events |
| 403 | .iter() |
| 404 | .find_map(|event| match event { |
| 405 | Event::TurnComplete { status, .. } => Some(*status), |
| 406 | _ => None, |
| 407 | }) |
| 408 | .expect("terminal TurnComplete"); |
| 409 | assert_eq!(status, TurnOutcomeStatus::Failed); |
| 410 | } |
| 411 | |
| 412 | /// The reasoning-only nudge rides one request and is never written to the |
| 413 | /// session. |
| 414 | /// |
| 415 | /// This is the distinction that matters: a nudge added with |
| 416 | /// `add_session_message` would persist into the transcript, the exports, and |
| 417 | /// every later turn's context — a message the user never sent. A |
| 418 | /// request-scoped nudge appears in exactly one outbound request and leaves the |
| 419 | /// conversation as it found it. |
| 420 | /// |
| 421 | /// The two are told apart by message counts across successive requests. With |
| 422 | /// a ceiling of 3 the model is asked four times. Persisted, the counts would |
| 423 | /// grow cumulatively (n, n, n+1, n+2); request-scoped, the nudged requests |
| 424 | /// each carry exactly one extra message over the same baseline. |
| 425 | #[tokio::test] |
| 426 | async fn the_reasoning_only_nudge_rides_one_request_and_never_joins_the_session() { |
| 427 | let (model, events) = run_reasoning_only_turn_with_reprompts(usize::MAX, "stop", 3).await; |
| 428 | |
| 429 | let requests = model.requests.lock().expect("captured requests").clone(); |
| 430 | assert_eq!(requests.len(), 4, "one initial request plus three retries"); |
| 431 | |
| 432 | let baseline = requests[0].len(); |
| 433 | assert_eq!( |
| 434 | requests[1].len(), |
| 435 | baseline, |
| 436 | "the first retry is a bare cached-prefix re-request, with no nudge" |
| 437 | ); |
| 438 | assert_eq!( |
| 439 | requests[2].len(), |
| 440 | baseline + 1, |
| 441 | "the second retry carries the nudge" |
| 442 | ); |
| 443 | assert_eq!( |
| 444 | requests[3].len(), |
| 445 | baseline + 1, |
| 446 | "the nudge did not accumulate: it was spent on the previous request, \ |
| 447 | not added to the session" |
| 448 | ); |
| 449 | |
| 450 | let nudge = crate::config::DEFAULT_REASONING_ONLY_REPROMPT_MESSAGE; |
| 451 | let carries_nudge = |messages: &Vec<codewhale_models::Message>| { |
| 452 | serde_json::to_string(messages) |
| 453 | .expect("messages serialize") |
| 454 | .contains(nudge) |
| 455 | }; |
| 456 | assert!(!carries_nudge(&requests[0]), "no nudge before any failure"); |
| 457 | assert!(!carries_nudge(&requests[1]), "no nudge on the first retry"); |
| 458 | assert!( |
| 459 | carries_nudge(&requests[2]), |
| 460 | "nudge present once retrying again" |
| 461 | ); |
| 462 | |
| 463 | // C02-04: model-visible means logged. Every request the nudge rides |
| 464 | // leaves a durable (internal) receipt carrying its exact text. |
| 465 | let receipts: Vec<&str> = events |
| 466 | .iter() |
| 467 | .filter_map(|event| match event { |
| 468 | Event::Status { message } |
| 469 | if message.starts_with(super::turn_loop::REQUEST_NUDGE_RECEIPT_PREFIX) => |
| 470 | { |
| 471 | Some(message.as_str()) |
| 472 | } |
| 473 | _ => None, |
| 474 | }) |
| 475 | .collect(); |
| 476 | assert_eq!( |
| 477 | receipts.len(), |
| 478 | requests |
| 479 | .iter() |
| 480 | .filter(|request| carries_nudge(request)) |
| 481 | .count(), |
| 482 | "one receipt per nudged request" |
| 483 | ); |
| 484 | assert!(receipts.iter().all(|receipt| receipt.ends_with(nudge))); |
| 485 | assert_eq!( |
| 486 | crate::core::events::status_visibility(receipts[0]), |
| 487 | crate::core::events::StatusVisibility::Internal, |
| 488 | "durable clients keep the receipt, collapsed" |
| 489 | ); |
| 490 | } |
| 491 | |
| 492 | /// A model that only ever returns reasoning is bounded: it retries up to the |
| 493 | /// ceiling and then fails honestly rather than looping forever. |
| 494 | #[tokio::test] |
| 495 | async fn reasoning_only_forever_is_bounded_then_fails() { |
| 496 | let (model, events) = run_reasoning_only_turn(usize::MAX, "stop").await; |
| 497 | |
| 498 | assert_eq!( |
| 499 | model.calls.load(std::sync::atomic::Ordering::SeqCst), |
| 500 | 1 + crate::config::DEFAULT_REASONING_ONLY_REPROMPTS as usize, |
| 501 | "reasoning-only retries are bounded by [reasoning_only] max_reprompts" |
| 502 | ); |
| 503 | let status = events |
| 504 | .iter() |
| 505 | .find_map(|event| match event { |
| 506 | Event::TurnComplete { status, .. } => Some(*status), |
| 507 | _ => None, |
| 508 | }) |
| 509 | .expect("terminal TurnComplete"); |
| 510 | assert_eq!(status, TurnOutcomeStatus::Failed); |
| 511 | let attempts = events |
| 512 | .iter() |
| 513 | .filter_map(|event| match event { |
| 514 | Event::Status { message } if message.starts_with("Retry attempt: reasoning-only ") => { |
| 515 | Some(message) |
| 516 | } |
| 517 | _ => None, |
| 518 | }) |
| 519 | .collect::<Vec<_>>(); |
| 520 | let max = crate::config::DEFAULT_REASONING_ONLY_REPROMPTS; |
| 521 | assert_eq!(attempts.len(), max as usize); |
| 522 | for (index, message) in attempts.iter().enumerate() { |
| 523 | assert!(message.starts_with(&format!( |
| 524 | "Retry attempt: reasoning-only {}/{};", |
| 525 | index + 1, |
| 526 | max |
| 527 | ))); |
| 528 | } |
| 529 | assert_eq!(events.iter().filter(|event| matches!(event, Event::Status { message } if message == &format!("Retry stopped: reasoning-only used {max}/{max} retries; turn failed"))).count(), 1); |
| 530 | } |
| 531 | |
| 532 | #[tokio::test] |
| 533 | async fn headless_turn_fails_with_real_error_after_network_drop_budget_exhausted() { |
| 534 | let (model, events) = |
| 535 | run_headless_turn_with_flaky_network(1 + super::MAX_STREAM_RETRIES as usize).await; |
| 536 | |
| 537 | assert_eq!( |
| 538 | model.calls.load(std::sync::atomic::Ordering::SeqCst), |
| 539 | 1 + super::MAX_STREAM_RETRIES as usize, |
| 540 | "initial attempt plus the bounded resume budget, then the turn fails" |
| 541 | ); |
| 542 | let (status, error) = events |
| 543 | .iter() |
| 544 | .find_map(|event| match event { |
| 545 | Event::TurnComplete { status, error, .. } => Some((status, error)), |
| 546 | _ => None, |
| 547 | }) |
| 548 | .expect("terminal TurnComplete"); |
| 549 | assert_eq!(*status, TurnOutcomeStatus::Failed); |
| 550 | let error = error |
| 551 | .as_deref() |
| 552 | .expect("exhausted network-drop retries must report the real error"); |
| 553 | assert!( |
| 554 | error.contains("Provider stream connection dropped"), |
| 555 | "the surfaced error must name the network drop: {error}" |
| 556 | ); |
| 557 | assert!( |
| 558 | error.contains("error decoding response body"), |
| 559 | "the underlying provider error must stay attached: {error}" |
| 560 | ); |
| 561 | assert_eq!( |
| 562 | crate::error_taxonomy::classify_error_message(error), |
| 563 | crate::error_taxonomy::ErrorCategory::Network, |
| 564 | "the terminal failure must classify as retryable infra (network)" |
| 565 | ); |
| 566 | let error_events = events |
| 567 | .iter() |
| 568 | .filter(|event| matches!(event, Event::Error { .. })) |
| 569 | .count(); |
| 570 | assert_eq!( |
| 571 | error_events, 1, |
| 572 | "only the final, budget-exhausted attempt may emit an error event: {events:?}" |
| 573 | ); |
| 574 | assert_eq!( |
| 575 | events |
| 576 | .iter() |
| 577 | .filter(|event| matches!(event, |
| 578 | Event::Status { message } if message.starts_with("Retry attempt: stream-resume ") |
| 579 | )) |
| 580 | .count(), |
| 581 | super::MAX_STREAM_RETRIES as usize, |
| 582 | "every admitted resume must have its own numbered progress receipt" |
| 583 | ); |
| 584 | for attempt in 1..=super::MAX_STREAM_RETRIES { |
| 585 | assert!(events.iter().any(|event| matches!(event, |
| 586 | Event::Status { message } if message.starts_with(&format!("Retry attempt: stream-resume {attempt}/")) |
| 587 | ))); |
| 588 | } |
| 589 | } |
| 590 | |
| 591 | // === Issue #66: error taxonomy wired through engine + audit + capacity === |
| 592 | |
| 593 | /// A failed-tool audit entry must carry the typed `category` and `severity` |
| 594 | /// fields derived from the underlying `ToolError`. This is what makes |
| 595 | /// downstream tooling able to bucket failures without scraping the message |
| 596 | /// string. |
| 597 | #[test] |
| 598 | fn tool_failure_audit_payload_carries_category_and_severity() { |
| 599 | use crate::error_taxonomy::ErrorEnvelope; |
| 600 | use crate::tools::spec::ToolError; |
| 601 | |
| 602 | let error = ToolError::Timeout { seconds: 30 }; |
| 603 | let envelope: ErrorEnvelope = error.clone().into(); |
| 604 | let payload = json!({ |
| 605 | "event": "tool.result", |
| 606 | "tool_id": "tool-1", |
| 607 | "tool_name": "exec_shell", |
| 608 | "status": ToolExecutionOutcome::from_legacy(Err(error.clone())).status.as_str(), |
| 609 | "success": false, |
| 610 | "error": error.to_string(), |
| 611 | "category": envelope.category.to_string(), |
| 612 | "severity": envelope.severity.to_string(), |
| 613 | }); |
| 614 | |
| 615 | assert_eq!(payload["category"], "timeout"); |
| 616 | assert_eq!(payload["severity"], "warning"); |
| 617 | assert_eq!(payload["status"], "timed_out"); |
| 618 | assert_eq!(payload["success"], false); |
| 619 | } |
| 620 | |
| 621 | // ── #136: post-edit LSP diagnostics hook ───────────────────────────────── |
| 622 | |
| 623 | #[test] |
| 624 | fn edited_paths_scenario() { |
| 625 | // Scenario consolidation of: edited_paths_for_edit_file_returns_path, edited_paths_for_write_file_returns_path, edited_paths_for_apply_patch_with_replace_returns_each_path, edited_paths_for_apply_patch_with_legacy_changes_returns_each_path, edited_paths_for_apply_patch_with_diff_text_extracts_paths, edited_paths_for_apply_patch_with_invalid_diff_returns_empty, edited_paths_for_unknown_tool_returns_empty |
| 626 | // from edited_paths_for_edit_file_returns_path |
| 627 | { |
| 628 | let input = json!({ "path": "src/foo.rs", "search": "x", "replace": "y" }); |
| 629 | let paths = edited_paths_for_tool("edit_file", &input); |
| 630 | assert_eq!(paths, vec![PathBuf::from("src/foo.rs")]); |
| 631 | } |
| 632 | // from edited_paths_for_write_file_returns_path |
| 633 | { |
| 634 | let input = json!({ "path": "src/bar.rs", "content": "fn main() {}" }); |
| 635 | let paths = edited_paths_for_tool("write_file", &input); |
| 636 | assert_eq!(paths, vec![PathBuf::from("src/bar.rs")]); |
| 637 | } |
| 638 | // from edited_paths_for_apply_patch_with_replace_returns_each_path |
| 639 | { |
| 640 | let input = json!({ |
| 641 | "replace": [ |
| 642 | { "path": "a.rs", "content": "" }, |
| 643 | { "path": "b.rs", "content": "" } |
| 644 | ] |
| 645 | }); |
| 646 | let paths = edited_paths_for_tool("apply_patch", &input); |
| 647 | assert_eq!(paths, vec![PathBuf::from("a.rs"), PathBuf::from("b.rs")]); |
| 648 | } |
| 649 | // from edited_paths_for_apply_patch_with_legacy_changes_returns_each_path |
| 650 | { |
| 651 | let input = json!({ |
| 652 | "changes": [ |
| 653 | { "path": "a.rs", "content": "" }, |
| 654 | { "path": "b.rs", "content": "" } |
| 655 | ] |
| 656 | }); |
| 657 | let paths = edited_paths_for_tool("apply_patch", &input); |
| 658 | assert_eq!(paths, vec![PathBuf::from("a.rs"), PathBuf::from("b.rs")]); |
| 659 | } |
| 660 | // from edited_paths_for_apply_patch_with_diff_text_extracts_paths |
| 661 | { |
| 662 | let input = json!({ |
| 663 | "patch": "--- a/foo.rs\n+++ b/foo.rs\n@@ -1 +1 @@\n-let x: i32 = 0;\n+let x: i32 = \"oops\";\n" |
| 664 | }); |
| 665 | let paths = edited_paths_for_tool("apply_patch", &input); |
| 666 | assert_eq!(paths, vec![PathBuf::from("foo.rs")]); |
| 667 | } |
| 668 | // from edited_paths_for_apply_patch_with_invalid_diff_returns_empty |
| 669 | { |
| 670 | let input = json!({ |
| 671 | "patch": "@@ -1 +1 @@\n-old\n+new\n" |
| 672 | }); |
| 673 | let paths = edited_paths_for_tool("apply_patch", &input); |
| 674 | assert!(paths.is_empty()); |
| 675 | } |
| 676 | // from edited_paths_for_unknown_tool_returns_empty |
| 677 | { |
| 678 | let input = json!({ "path": "irrelevant.rs" }); |
| 679 | let paths = edited_paths_for_tool("read_file", &input); |
| 680 | assert!(paths.is_empty()); |
| 681 | let paths = edited_paths_for_tool("grep_files", &input); |
| 682 | assert!(paths.is_empty()); |
| 683 | } |
| 684 | } |
| 685 | |
| 686 | #[test] |
| 687 | fn parse_patch_paths_skips_dev_null() { |
| 688 | let patch = "--- a/keep.rs\n+++ b/keep.rs\n@@ -1 +1 @@\n-old\n+new\n--- a/deleted.rs\n+++ /dev/null\n@@ -1 +0,0 @@\n-delete me\n"; |
| 689 | let paths = edited_paths_for_tool("apply_patch", &json!({ "patch": patch })); |
| 690 | assert_eq!(paths, vec![PathBuf::from("keep.rs")]); |
| 691 | } |
| 692 | |
| 693 | #[tokio::test] |
| 694 | async fn post_edit_hook_injects_diagnostics_message_before_next_request() { |
| 695 | use crate::lsp::{Diagnostic, Language, Severity}; |
| 696 | use std::sync::Arc; |
| 697 | |
| 698 | let tmp = tempdir().expect("tempdir"); |
| 699 | let workspace = tmp.path().to_path_buf(); |
| 700 | let target = workspace.join("src").join("main.rs"); |
| 701 | fs::create_dir_all(workspace.join("src")).unwrap(); |
| 702 | fs::write(&target, "let x: i32 = \"not a number\";").unwrap(); |
| 703 | |
| 704 | let lsp_config = crate::lsp::LspConfig::default(); |
| 705 | let engine_config = EngineConfig { |
| 706 | workspace: workspace.clone(), |
| 707 | lsp_config: Some(lsp_config), |
| 708 | ..Default::default() |
| 709 | }; |
| 710 | let (mut engine, _handle) = Engine::new(engine_config, &Config::default()); |
| 711 | |
| 712 | // Install a fake transport that always reports a type error. |
| 713 | let fake = Arc::new(crate::lsp::tests::FakeTransport::new(vec![Diagnostic { |
| 714 | line: 1, |
| 715 | column: 14, |
| 716 | severity: Severity::Error, |
| 717 | message: "expected i32, found &str".to_string(), |
| 718 | }])); |
| 719 | engine |
| 720 | .lsp_manager |
| 721 | .install_test_transport(Language::Rust, fake) |
| 722 | .await; |
| 723 | |
| 724 | // Simulate the success path of an edit_file tool call. |
| 725 | let input = json!({ "path": "src/main.rs", "search": "0", "replace": "\"not a number\"" }); |
| 726 | engine.run_post_edit_lsp_hook("edit_file", &input).await; |
| 727 | assert_eq!(engine.pending_lsp_blocks.len(), 1); |
| 728 | |
| 729 | // Flush prepares the synthetic message. |
| 730 | let messages_before = engine.session.messages.len(); |
| 731 | engine.flush_pending_lsp_diagnostics().await; |
| 732 | assert_eq!(engine.session.messages.len(), messages_before + 1); |
| 733 | |
| 734 | let last = engine.session.messages.last().expect("message appended"); |
| 735 | assert_eq!(last.role, "user"); |
| 736 | // turn_meta is now at the tail of the content array (PR #2517). |
| 737 | let meta = match last.content.last() { |
| 738 | Some(codewhale_models::ContentBlock::Text { text, .. }) => text.clone(), |
| 739 | other => panic!("expected text block at tail, got {other:?}"), |
| 740 | }; |
| 741 | assert!(meta.starts_with("<turn_meta>\n")); |
| 742 | let diagnostic_text = last |
| 743 | .content |
| 744 | .iter() |
| 745 | .find_map(|block| match block { |
| 746 | codewhale_models::ContentBlock::Text { text, .. } |
| 747 | if text.contains("<diagnostics file=\"") => |
| 748 | { |
| 749 | Some(text) |
| 750 | } |
| 751 | _ => None, |
| 752 | }) |
| 753 | .expect("diagnostics text block"); |
| 754 | assert!(diagnostic_text.contains("ERROR [1:14] expected i32, found &str")); |
| 755 | } |
| 756 | |
| 757 | #[tokio::test] |
| 758 | async fn post_edit_hook_is_silent_when_lsp_disabled() { |
| 759 | let tmp = tempdir().expect("tempdir"); |
| 760 | let workspace = tmp.path().to_path_buf(); |
| 761 | let target = workspace.join("src").join("main.rs"); |
| 762 | fs::create_dir_all(workspace.join("src")).unwrap(); |
| 763 | fs::write(&target, "fn main() {}").unwrap(); |
| 764 | |
| 765 | let lsp_config = crate::lsp::LspConfig { |
| 766 | enabled: false, |
| 767 | ..Default::default() |
| 768 | }; |
| 769 | let engine_config = EngineConfig { |
| 770 | workspace: workspace.clone(), |
| 771 | lsp_config: Some(lsp_config), |
| 772 | ..Default::default() |
| 773 | }; |
| 774 | let (mut engine, _handle) = Engine::new(engine_config, &Config::default()); |
| 775 | |
| 776 | let input = json!({ "path": "src/main.rs", "search": "x", "replace": "y" }); |
| 777 | engine.run_post_edit_lsp_hook("edit_file", &input).await; |
| 778 | assert!(engine.pending_lsp_blocks.is_empty()); |
| 779 | |
| 780 | let messages_before = engine.session.messages.len(); |
| 781 | engine.flush_pending_lsp_diagnostics().await; |
| 782 | assert_eq!(engine.session.messages.len(), messages_before); |
| 783 | } |
| 784 | |
| 785 | #[tokio::test] |
| 786 | async fn post_edit_hook_skips_unknown_tool_names() { |
| 787 | use crate::lsp::{Diagnostic, Language, Severity}; |
| 788 | use std::sync::Arc; |
| 789 | |
| 790 | let tmp = tempdir().expect("tempdir"); |
| 791 | let engine_config = EngineConfig { |
| 792 | workspace: tmp.path().to_path_buf(), |
| 793 | lsp_config: Some(crate::lsp::LspConfig::default()), |
| 794 | ..Default::default() |
| 795 | }; |
| 796 | let (mut engine, _handle) = Engine::new(engine_config, &Config::default()); |
| 797 | let fake = Arc::new(crate::lsp::tests::FakeTransport::new(vec![Diagnostic { |
| 798 | line: 1, |
| 799 | column: 1, |
| 800 | severity: Severity::Error, |
| 801 | message: "should not be reported".to_string(), |
| 802 | }])); |
| 803 | engine |
| 804 | .lsp_manager |
| 805 | .install_test_transport(Language::Rust, fake.clone()) |
| 806 | .await; |
| 807 | |
| 808 | let input = json!({ "path": "src/main.rs" }); |
| 809 | engine.run_post_edit_lsp_hook("read_file", &input).await; |
| 810 | assert!(engine.pending_lsp_blocks.is_empty()); |
| 811 | assert_eq!(fake.call_count(), 0); |
| 812 | } |
| 813 | |
| 814 | // ── #3802: non-blocking send for ListSubAgents refresh events ───────────── |
| 815 | |
| 816 | #[test] |
| 817 | fn agent_list_event_carries_the_typed_coordination_projection() { |
| 818 | use crate::tools::subagent::coord::{DecisionRecord, DecisionStatus}; |
| 819 | |
| 820 | let mut manager = SubAgentManager::new(PathBuf::from("."), 1); |
| 821 | let recorded = manager |
| 822 | .record_coordination_decision(DecisionRecord { |
| 823 | decision_id: "decision-event".to_string(), |
| 824 | subject: "typed event".to_string(), |
| 825 | status: DecisionStatus::Accepted, |
| 826 | owner: "root".to_string(), |
| 827 | scope: Vec::new(), |
| 828 | constraints: Vec::new(), |
| 829 | evidence_handles: Vec::new(), |
| 830 | version: 1, |
| 831 | sequence: 0, |
| 832 | }) |
| 833 | .expect("record decision"); |
| 834 | manager |
| 835 | .stamp_coordination_sequence_for_session(recorded.sequence, "session-a") |
| 836 | .expect("stamp decision owner"); |
| 837 | |
| 838 | let Event::AgentList { |
| 839 | owner_session_id, |
| 840 | agents, |
| 841 | coordination, |
| 842 | .. |
| 843 | } = agent_list_event(&manager, "session-a") |
| 844 | else { |
| 845 | panic!("expected AgentList event"); |
| 846 | }; |
| 847 | assert!(agents.is_empty()); |
| 848 | assert_eq!(owner_session_id, "session-a"); |
| 849 | assert_eq!(coordination.decisions.len(), 1); |
| 850 | assert_eq!(coordination.decisions[0].decision_id, "decision-event"); |
| 851 | assert_eq!(coordination.decisions[0].status, DecisionStatus::Accepted); |
| 852 | assert!(coordination.bounded); |
| 853 | assert_eq!(coordination.limit, 24); |
| 854 | } |
| 855 | |
| 856 | #[test] |
| 857 | fn engine_handle_try_send_does_not_block_when_op_channel_is_full() { |
| 858 | use tokio::sync::mpsc; |
| 859 | |
| 860 | // Create a channel with the smallest possible capacity. |
| 861 | let (tx_op, rx_op) = mpsc::channel::<Op>(1); |
| 862 | |
| 863 | // Construct a minimal EngineHandle with the tiny channel. |
| 864 | let cancel_token = CancellationToken::new(); |
| 865 | let handle = EngineHandle { |
| 866 | goal_state: new_shared_goal_state(), |
| 867 | tx_op, |
| 868 | rx_event: Arc::new(RwLock::new(mpsc::channel::<Event>(1).1)), |
| 869 | cancel_token: Arc::new(StdMutex::new(cancel_token)), |
| 870 | cancel_reason: Arc::new(StdMutex::new(None)), |
| 871 | tx_approval: mpsc::channel(1).0, |
| 872 | tx_user_input: mpsc::channel(1).0, |
| 873 | tx_steer: mpsc::channel(1).0, |
| 874 | turn_controls: Arc::new(StdMutex::new(handle::TurnControls::default())), |
| 875 | shared_paused: Arc::new(StdMutex::new(false)), |
| 876 | client_preflight_required: true, |
| 877 | live_runtime_authority: Arc::new(StdMutex::new(LiveRuntimeAuthorityState::new( |
| 878 | LiveRuntimeAuthority::from_fields( |
| 879 | AppMode::Agent, |
| 880 | false, |
| 881 | false, |
| 882 | false, |
| 883 | ApprovalMode::Suggest, |
| 884 | None, |
| 885 | ), |
| 886 | ))), |
| 887 | compaction_cancellation: Arc::new(StdMutex::new(CompactionCancellationState::default())), |
| 888 | turn_heartbeat: turn_heartbeat::TurnHeartbeat::new(), |
| 889 | subagent_manager: crate::tools::subagent::new_shared_subagent_manager( |
| 890 | std::env::temp_dir(), |
| 891 | 1, |
| 892 | ), |
| 893 | }; |
| 894 | |
| 895 | // Fill the op channel with one message (capacity = 1). |
| 896 | handle |
| 897 | .tx_op |
| 898 | .try_send(Op::ListSubAgents) |
| 899 | .expect("first send should succeed"); |
| 900 | |
| 901 | // A live posture update must publish immediately even though its wake-up |
| 902 | // cannot fit. The already-queued operation will wake the engine, which |
| 903 | // applies this pending authority before handling it. |
| 904 | let result = handle.try_send(Op::ChangeMode { |
| 905 | mode: AppMode::Operate, |
| 906 | allow_shell: true, |
| 907 | trust_mode: false, |
| 908 | auto_approve: false, |
| 909 | approval_mode: ApprovalMode::Auto, |
| 910 | configured_sandbox_mode: None, |
| 911 | }); |
| 912 | let error = result.expect_err("try_send should fail when channel is full"); |
| 913 | assert!(matches!( |
| 914 | error.downcast_ref::<mpsc::error::TrySendError<Op>>(), |
| 915 | Some(mpsc::error::TrySendError::Full(Op::ChangeMode { .. })) |
| 916 | )); |
| 917 | let authority = handle.runtime_permission_authority(); |
| 918 | assert_eq!(authority.approval_mode, ApprovalMode::Auto); |
| 919 | assert!(!authority.auto_approve); |
| 920 | |
| 921 | handle |
| 922 | .cancel_compaction("compact-full-mailbox") |
| 923 | .expect("full mailbox must not block compaction cancellation"); |
| 924 | assert!( |
| 925 | handle |
| 926 | .compaction_cancellation |
| 927 | .lock() |
| 928 | .expect("cancellation state") |
| 929 | .claim("compact-full-mailbox") |
| 930 | .is_none(), |
| 931 | "cancellation authority remains visible even when its wake-up op cannot fit" |
| 932 | ); |
| 933 | drop(rx_op); |
| 934 | let error = handle.try_send(Op::ListSubAgents).unwrap_err(); |
| 935 | assert!(matches!( |
| 936 | error.downcast_ref::<mpsc::error::TrySendError<Op>>(), |
| 937 | Some(mpsc::error::TrySendError::Closed(Op::ListSubAgents)) |
| 938 | )); |
| 939 | } |
| 940 | |
| 941 | #[tokio::test] |
| 942 | async fn full_mailbox_posture_update_supersedes_queued_change_mode() { |
| 943 | use ApprovalMode; |
| 944 | |
| 945 | let tmp = tempdir().expect("tempdir"); |
| 946 | let config = EngineConfig { |
| 947 | workspace: tmp.path().to_path_buf(), |
| 948 | ..Default::default() |
| 949 | }; |
| 950 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 951 | |
| 952 | handle |
| 953 | .try_send(Op::ChangeMode { |
| 954 | mode: AppMode::Plan, |
| 955 | allow_shell: false, |
| 956 | trust_mode: false, |
| 957 | auto_approve: false, |
| 958 | approval_mode: ApprovalMode::Suggest, |
| 959 | configured_sandbox_mode: None, |
| 960 | }) |
| 961 | .expect("queue older posture"); |
| 962 | for _ in 1..ENGINE_OP_CHANNEL_CAPACITY { |
| 963 | handle |
| 964 | .try_send(Op::ListSubAgents) |
| 965 | .expect("fill operation mailbox"); |
| 966 | } |
| 967 | |
| 968 | let result = handle.try_send(Op::ChangeMode { |
| 969 | mode: AppMode::Operate, |
| 970 | allow_shell: true, |
| 971 | trust_mode: false, |
| 972 | auto_approve: false, |
| 973 | approval_mode: ApprovalMode::Auto, |
| 974 | configured_sandbox_mode: Some("read-only".to_string()), |
| 975 | }); |
| 976 | assert!( |
| 977 | result.is_err(), |
| 978 | "latest posture wake-up must see a full mailbox" |
| 979 | ); |
| 980 | |
| 981 | let run = tokio::spawn(engine.run()); |
| 982 | let snapshot = tokio::time::timeout( |
| 983 | std::time::Duration::from_secs(2), |
| 984 | handle.get_session_snapshot(), |
| 985 | ) |
| 986 | .await |
| 987 | .expect("snapshot after mailbox drain") |
| 988 | .expect("session snapshot"); |
| 989 | |
| 990 | assert_eq!(snapshot.mode, "operate"); |
| 991 | let authority = handle.runtime_permission_authority(); |
| 992 | assert_eq!(authority.approval_mode, ApprovalMode::Auto); |
| 993 | assert!(!authority.auto_approve); |
| 994 | |
| 995 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 996 | run.await.expect("engine task"); |
| 997 | } |
| 998 | |
| 999 | #[tokio::test] |
| 1000 | async fn reload_mcp_op_recovers_from_invalid_initial_config_in_process() { |
| 1001 | let tmp = tempdir().expect("tempdir"); |
| 1002 | let workspace = tmp.path().join("workspace"); |
| 1003 | std::fs::create_dir_all(&workspace).expect("workspace"); |
| 1004 | let config_path = tmp.path().join("mcp.json"); |
| 1005 | let secret = "mcp-op-secret-must-not-escape"; |
| 1006 | std::fs::write( |
| 1007 | &config_path, |
| 1008 | format!(r#"{{"servers":{{"bad":{{"token":"{secret}"}} trailing}}}}"#), |
| 1009 | ) |
| 1010 | .expect("invalid config"); |
| 1011 | let engine_config = EngineConfig { |
| 1012 | workspace, |
| 1013 | mcp_config_path: config_path.clone(), |
| 1014 | ..Default::default() |
| 1015 | }; |
| 1016 | let (engine, handle) = Engine::new(engine_config, &Config::default()); |
| 1017 | let task = tokio::spawn(async move { engine.run().await }); |
| 1018 | |
| 1019 | let error = handle |
| 1020 | .reload_mcp(config_path.clone()) |
| 1021 | .await |
| 1022 | .expect_err("invalid config must fail closed"); |
| 1023 | assert!(!error.to_string().contains(secret)); |
| 1024 | std::fs::write( |
| 1025 | &config_path, |
| 1026 | r#"{"servers":{"ready":{"command":"node","disabled":true}}}"#, |
| 1027 | ) |
| 1028 | .expect("fixed config"); |
| 1029 | |
| 1030 | let snapshot = handle |
| 1031 | .reload_mcp(config_path.clone()) |
| 1032 | .await |
| 1033 | .expect("fixed config reloads without restarting the engine") |
| 1034 | .snapshot; |
| 1035 | assert!(!snapshot.reload_required); |
| 1036 | assert_eq!(snapshot.servers.len(), 1); |
| 1037 | assert_eq!(snapshot.servers[0].name, "ready"); |
| 1038 | assert!(!snapshot.servers[0].enabled); |
| 1039 | |
| 1040 | let alternate_path = tmp.path().join("alternate-mcp.json"); |
| 1041 | std::fs::write( |
| 1042 | &alternate_path, |
| 1043 | r#"{"servers":{"alternate":{"command":"node","disabled":true}}}"#, |
| 1044 | ) |
| 1045 | .expect("alternate config"); |
| 1046 | let alternate = handle |
| 1047 | .reload_mcp(alternate_path.clone()) |
| 1048 | .await |
| 1049 | .expect("a changed config path replaces the engine pool in process") |
| 1050 | .snapshot; |
| 1051 | assert_eq!(alternate.config_path, alternate_path); |
| 1052 | assert_eq!(alternate.servers.len(), 1); |
| 1053 | assert_eq!(alternate.servers[0].name, "alternate"); |
| 1054 | |
| 1055 | handle.send(Op::Shutdown).await.expect("shutdown"); |
| 1056 | task.await.expect("engine task"); |
| 1057 | } |
| 1058 | |
| 1059 | #[tokio::test] |
| 1060 | async fn mcp_boot_reports_ready_server_before_stalled_server_finishes() { |
| 1061 | assert_incremental_mcp_boot(false).await; |
| 1062 | } |
| 1063 | |
| 1064 | #[tokio::test] |
| 1065 | async fn first_turn_waits_for_explicit_mcp_schema_without_waiting_for_unrelated_server() { |
| 1066 | let Some(node) = crate::dependencies::resolve_node() else { |
| 1067 | return; |
| 1068 | }; |
| 1069 | let tmp = tempdir().expect("tempdir"); |
| 1070 | let server = tmp.path().join("server.mjs"); |
| 1071 | let release = tmp.path().join("release-slow"); |
| 1072 | let release_fast = tmp.path().join("release-fast"); |
| 1073 | fs::write(&server, r#"import fs from 'node:fs'; |
| 1074 | import path from 'node:path'; |
| 1075 | import readline from 'node:readline'; |
| 1076 | readline.createInterface({ input: process.stdin }).on('line', async line => { |
| 1077 | const request = JSON.parse(line); |
| 1078 | if (request.id === undefined) return; |
| 1079 | if (request.method === 'initialize') { |
| 1080 | fs.writeFileSync(path.join(process.argv[3], 'started-' + process.argv[2]), 'ready'); |
| 1081 | while (!fs.existsSync(path.join(process.argv[3], 'release-' + process.argv[2]))) await new Promise(r => setTimeout(r, 10)); |
| 1082 | } |
| 1083 | const result = request.method === 'initialize' |
| 1084 | ? { protocolVersion: '2024-11-05', capabilities: { tools: {} }, serverInfo: { name: process.argv[2], version: '1' } } |
| 1085 | : { tools: ['ready', 'denied', 'hidden'].map(name => ({ name, inputSchema: { type: 'object' } })) }; |
| 1086 | process.stdout.write(JSON.stringify({ jsonrpc: '2.0', id: request.id, result }) + '\n'); |
| 1087 | });"#).expect("fixture"); |
| 1088 | let config_path = tmp.path().join("mcp.json"); |
| 1089 | fs::write( |
| 1090 | &config_path, |
| 1091 | serde_json::to_vec(&json!({ |
| 1092 | // `slow` must still be connecting when the fast build completes, |
| 1093 | // and `fast` must not be declared dead while a cold Windows runner |
| 1094 | // spawns Node. Ordering here is proven by the release files below, |
| 1095 | // never by a timeout, so this bound only has to outlast the test. |
| 1096 | "timeouts": { "connect_timeout": 120 }, |
| 1097 | "servers": { |
| 1098 | "fast": { "command": node, "args": [server, "fast", tmp.path()] }, |
| 1099 | // `slow` is deliberately unselected: `required` keeps it in |
| 1100 | // the eager boot set under lazy boot (#6033) so it can stand |
| 1101 | // in for "an unrelated server still connecting". |
| 1102 | "slow": { "command": node, "args": [server, "slow", tmp.path()], "required": true }, |
| 1103 | "failed": { "command": "codewhale-missing-mcp-fixture-38911" } |
| 1104 | } |
| 1105 | })) |
| 1106 | .unwrap(), |
| 1107 | ) |
| 1108 | .unwrap(); |
| 1109 | let api_config = Config::default(); |
| 1110 | let (mut engine, _handle) = Engine::new( |
| 1111 | EngineConfig { |
| 1112 | workspace: tmp.path().to_path_buf(), |
| 1113 | mcp_config_path: config_path, |
| 1114 | tools_always_load: HashSet::from(["mcp_fast_ready".to_string()]), |
| 1115 | ..Default::default() |
| 1116 | }, |
| 1117 | &api_config, |
| 1118 | ); |
| 1119 | engine |
| 1120 | .start_mcp_session_boot(McpConnectRefresh::IfChanged) |
| 1121 | .await |
| 1122 | .expect("session boot starts"); |
| 1123 | assert!( |
| 1124 | engine.mcp_tools().await.is_empty(), |
| 1125 | "ordinary startup remains nonblocking" |
| 1126 | ); |
| 1127 | // Separate Windows/CI process startup from the schema-wait assertion. |
| 1128 | // Both children have received initialize, but neither can answer until |
| 1129 | // this test releases its own gate. No fixed delay stands in for readiness. |
| 1130 | // The budget is generous because it covers two cold Node spawns on a |
| 1131 | // windows-latest runner that has just finished a ~15 min compile; a tight |
| 1132 | // bound here fails the setup, not the behavior under test. |
| 1133 | tokio::time::timeout(Duration::from_secs(60), async { |
| 1134 | while !tmp.path().join("started-fast").exists() || !tmp.path().join("started-slow").exists() |
| 1135 | { |
| 1136 | engine.drain_mcp_boot_updates().await; |
| 1137 | for name in ["fast", "slow"] { |
| 1138 | assert!( |
| 1139 | !engine.mcp_connection_errors.contains_key(name), |
| 1140 | "{name} fixture failed before initialize: {:?}", |
| 1141 | engine.mcp_connection_errors.get(name) |
| 1142 | ); |
| 1143 | } |
| 1144 | tokio::time::sleep(Duration::from_millis(10)).await; |
| 1145 | } |
| 1146 | }) |
| 1147 | .await |
| 1148 | .expect("both MCP fixtures must reach initialize before checking first-turn ordering"); |
| 1149 | let route = TurnRouteContext { |
| 1150 | provider: ProviderKind::Deepseek, |
| 1151 | model: DEFAULT_TEXT_MODEL.to_string(), |
| 1152 | capabilities: codewhale_config::route::RouteCapabilities::default(), |
| 1153 | limits: None, |
| 1154 | client: engine.codewhale_client.clone(), |
| 1155 | api_config: Box::new(api_config), |
| 1156 | locale_tag: engine.config.locale_tag.clone(), |
| 1157 | role_models: engine.subagent_role_models(), |
| 1158 | auto_model: false, |
| 1159 | reasoning_effort: None, |
| 1160 | reasoning_effort_auto: false, |
| 1161 | }; |
| 1162 | let policy = crate::core::authority::TurnAuthority::from_effective_fields( |
| 1163 | AppMode::Agent, |
| 1164 | false, |
| 1165 | false, |
| 1166 | false, |
| 1167 | ApprovalMode::Suggest, |
| 1168 | ); |
| 1169 | let build = { |
| 1170 | let build = engine.build_turn_tool_registry_and_catalog( |
| 1171 | &policy, |
| 1172 | &[], |
| 1173 | Some(vec![ |
| 1174 | "mcp_fast_ready".to_string(), |
| 1175 | "mcp_failed_ready".to_string(), |
| 1176 | ]), |
| 1177 | SubAgentWiring::Inert, |
| 1178 | McpAccess::Connect, |
| 1179 | route, |
| 1180 | "", |
| 1181 | ); |
| 1182 | tokio::pin!(build); |
| 1183 | std::future::poll_fn(|cx| { |
| 1184 | assert!( |
| 1185 | std::future::Future::poll(build.as_mut(), cx).is_pending(), |
| 1186 | "the first turn must wait for the explicitly selected fast schema" |
| 1187 | ); |
| 1188 | std::task::Poll::Ready(()) |
| 1189 | }) |
| 1190 | .await; |
| 1191 | fs::write(&release_fast, "release").unwrap(); |
| 1192 | tokio::time::timeout(Duration::from_secs(5), build).await |
| 1193 | }; |
| 1194 | let unrelated_pending = !release.exists() |
| 1195 | && engine.mcp_boot_in_flight |
| 1196 | && !engine.mcp_connection_errors.contains_key("slow"); |
| 1197 | let connected = engine |
| 1198 | .mcp_pool |
| 1199 | .as_ref() |
| 1200 | .unwrap() |
| 1201 | .lock() |
| 1202 | .await |
| 1203 | .connected_servers() |
| 1204 | .into_iter() |
| 1205 | .map(str::to_owned) |
| 1206 | .collect::<Vec<_>>(); |
| 1207 | engine.cancel_token.cancel(); |
| 1208 | tokio::time::timeout( |
| 1209 | Duration::from_millis(100), |
| 1210 | engine.wait_for_explicit_mcp_boot(Some(&["mcp_slow_ready".to_string()])), |
| 1211 | ) |
| 1212 | .await |
| 1213 | .expect("stop interrupts explicit schema wait"); |
| 1214 | let _turn_control = engine.begin_turn_control(); |
| 1215 | fs::write(&release, "release").unwrap(); |
| 1216 | let build = build.unwrap_or_else(|error| { |
| 1217 | panic!( |
| 1218 | "explicit fast/failed selections must not wait for slow: {error:?}; connected={connected:?}; errors={:?}", |
| 1219 | engine.mcp_connection_errors |
| 1220 | ) |
| 1221 | }); |
| 1222 | assert!( |
| 1223 | unrelated_pending, |
| 1224 | "success must precede the unrelated server's release, completion, or timeout" |
| 1225 | ); |
| 1226 | let active = build.surface.active.unwrap_or_default(); |
| 1227 | assert_eq!( |
| 1228 | active |
| 1229 | .iter() |
| 1230 | .map(|tool| tool.name.as_str()) |
| 1231 | .collect::<Vec<_>>(), |
| 1232 | ["mcp_fast_ready"] |
| 1233 | ); |
| 1234 | assert!(engine.mcp_connection_errors.contains_key("failed")); |
| 1235 | |
| 1236 | // The unrelated connection finishes during the same turn. Refresh into a |
| 1237 | // narrowed policy, then execute the actual tool-search activation path. |
| 1238 | let policy = ToolSurfacePolicy::new( |
| 1239 | ToolRegistryBuilder::new().build(ToolContext::for_empty_registry()), |
| 1240 | Some(vec![api_tool("read")]), |
| 1241 | AppMode::Agent, |
| 1242 | &HashSet::new(), |
| 1243 | &[], |
| 1244 | false, |
| 1245 | Some(vec![ |
| 1246 | "tool_search".into(), |
| 1247 | "mcp_slow_ready".into(), |
| 1248 | "mcp_slow_denied".into(), |
| 1249 | ]), |
| 1250 | Some(vec!["mcp_slow_denied".into()]), |
| 1251 | None, |
| 1252 | crate::core::engine::tool_catalog::ToolMode::Direct, |
| 1253 | ); |
| 1254 | let mut catalog = policy.catalog.clone(); |
| 1255 | let mut active = policy.active_names.clone(); |
| 1256 | catalog.push(api_tool("mcp_removed_ready")); |
| 1257 | active.insert("mcp_removed_ready".to_string()); |
| 1258 | tokio::time::timeout(Duration::from_secs(5), async { |
| 1259 | while !catalog.iter().any(|tool| tool.name == "mcp_slow_ready") { |
| 1260 | engine |
| 1261 | .refresh_boot_mcp_catalog(&policy, &mut catalog, &mut active) |
| 1262 | .await; |
| 1263 | tokio::task::yield_now().await; |
| 1264 | } |
| 1265 | }) |
| 1266 | .await |
| 1267 | .expect("completed tools join this turn"); |
| 1268 | assert!( |
| 1269 | !active.contains("mcp_slow_ready"), |
| 1270 | "fresh MCP tools stay deferred" |
| 1271 | ); |
| 1272 | assert!( |
| 1273 | !active.contains("mcp_removed_ready"), |
| 1274 | "removed authority leaves active tools" |
| 1275 | ); |
| 1276 | assert!( |
| 1277 | catalog |
| 1278 | .iter() |
| 1279 | .all(|tool| !tool.name.starts_with("mcp_") || tool.name == "mcp_slow_ready") |
| 1280 | ); |
| 1281 | let result = tool_catalog::execute_tool_search_with_cache( |
| 1282 | "tool_search", |
| 1283 | &json!({"query":"mcp_slow_ready", "match":"regex"}), |
| 1284 | &catalog, |
| 1285 | &mut active, |
| 1286 | &mut engine.session.tool_activation_cache, |
| 1287 | ) |
| 1288 | .expect("real search"); |
| 1289 | assert!(result.success); |
| 1290 | assert!(active.contains("mcp_slow_ready")); |
| 1291 | engine.wait_for_mcp_boot().await; |
| 1292 | } |
| 1293 | |
| 1294 | #[tokio::test] |
| 1295 | async fn mcp_boot_does_not_restore_servers_removed_during_handshake() { |
| 1296 | assert_incremental_mcp_boot(true).await; |
| 1297 | } |
| 1298 | |
| 1299 | async fn assert_incremental_mcp_boot(invalidate_config: bool) { |
| 1300 | if std::process::Command::new("node") |
| 1301 | .arg("--version") |
| 1302 | .output() |
| 1303 | .is_err() |
| 1304 | { |
| 1305 | tracing::warn!("skipping MCP stdio fixture because node is unavailable"); |
| 1306 | return; |
| 1307 | } |
| 1308 | let tmp = tempdir().expect("tempdir"); |
| 1309 | let server = tmp.path().join("server.mjs"); |
| 1310 | let release = tmp.path().join("release-slow"); |
| 1311 | std::fs::write( |
| 1312 | &server, |
| 1313 | r#"import fs from 'node:fs'; |
| 1314 | import readline from 'node:readline'; |
| 1315 | const lines = readline.createInterface({ input: process.stdin }); |
| 1316 | lines.on('line', async (line) => { |
| 1317 | const request = JSON.parse(line); |
| 1318 | if (request.id === undefined) return; |
| 1319 | if (process.argv[2] === 'slow' && request.method === 'initialize') { |
| 1320 | while (!fs.existsSync(process.argv[3])) { |
| 1321 | await new Promise(resolve => setTimeout(resolve, 10)); |
| 1322 | } |
| 1323 | } |
| 1324 | const result = request.method === 'initialize' |
| 1325 | ? { protocolVersion: '2024-11-05', capabilities: { tools: {} }, |
| 1326 | serverInfo: { name: process.argv[2], version: '1' } } |
| 1327 | : { tools: [{ name: 'ready', inputSchema: { type: 'object' } }] }; |
| 1328 | process.stdout.write(JSON.stringify({ jsonrpc: '2.0', id: request.id, result }) + '\n'); |
| 1329 | }); |
| 1330 | "#, |
| 1331 | ) |
| 1332 | .expect("server fixture"); |
| 1333 | let config_path = tmp.path().join("mcp.json"); |
| 1334 | std::fs::write( |
| 1335 | &config_path, |
| 1336 | serde_json::to_vec(&serde_json::json!({ |
| 1337 | "timeouts": { "connect_timeout": 30 }, |
| 1338 | "servers": { |
| 1339 | // Both marked `required` so lazy boot (#6033) still starts |
| 1340 | // them eagerly — this test proves progress ordering, not the |
| 1341 | // lazy/eligible split. |
| 1342 | "fast": { "command": "node", "args": [server, "fast", release], "required": true }, |
| 1343 | "slow": { "command": "node", "args": [server, "slow", release], "required": true } |
| 1344 | } |
| 1345 | })) |
| 1346 | .expect("config JSON"), |
| 1347 | ) |
| 1348 | .expect("MCP config"); |
| 1349 | let (mut engine, handle) = Engine::new( |
| 1350 | EngineConfig { |
| 1351 | workspace: tmp.path().to_path_buf(), |
| 1352 | mcp_config_path: config_path.clone(), |
| 1353 | ..Default::default() |
| 1354 | }, |
| 1355 | &Config::default(), |
| 1356 | ); |
| 1357 | let pool = engine.ensure_mcp_pool().await.expect("engine pool"); |
| 1358 | let task = tokio::spawn(async move { engine.run().await }); |
| 1359 | let mut events = handle.rx_event.write().await; |
| 1360 | // This proves ordering, not Node cold-start speed on a loaded runner. |
| 1361 | let progress = tokio::time::timeout(Duration::from_secs(30), async { |
| 1362 | while let Some(event) = events.recv().await { |
| 1363 | if let Event::McpSessionBoot { |
| 1364 | snapshot, |
| 1365 | connecting, |
| 1366 | finished: false, |
| 1367 | .. |
| 1368 | } = event |
| 1369 | && connecting == ["slow"] |
| 1370 | { |
| 1371 | return snapshot; |
| 1372 | } |
| 1373 | } |
| 1374 | panic!("engine event channel closed"); |
| 1375 | }) |
| 1376 | .await; |
| 1377 | let ready_tools = pool.lock().await.to_api_tools(); |
| 1378 | if invalidate_config { |
| 1379 | std::fs::write( |
| 1380 | &config_path, |
| 1381 | r#"{"servers":{"slow":{"command":"node","disabled":true}}}"#, |
| 1382 | ) |
| 1383 | .expect("remove servers"); |
| 1384 | pool.lock() |
| 1385 | .await |
| 1386 | .reload_if_config_changed() |
| 1387 | .await |
| 1388 | .expect("reload config"); |
| 1389 | } |
| 1390 | // Release and shut down even when testing the old batch-buffered behavior. |
| 1391 | std::fs::write(&release, "continue").expect("release stalled fixture"); |
| 1392 | let finished = tokio::time::timeout(Duration::from_secs(30), async { |
| 1393 | while let Some(event) = events.recv().await { |
| 1394 | if let Event::McpSessionBoot { |
| 1395 | snapshot, |
| 1396 | finished: true, |
| 1397 | .. |
| 1398 | } = event |
| 1399 | { |
| 1400 | return snapshot; |
| 1401 | } |
| 1402 | } |
| 1403 | panic!("engine event channel closed"); |
| 1404 | }) |
| 1405 | .await; |
| 1406 | drop(events); |
| 1407 | handle.send(Op::Shutdown).await.expect("shutdown"); |
| 1408 | task.await.expect("engine task"); |
| 1409 | let progress = progress.expect("fast server must be visible before slow server is released"); |
| 1410 | assert!( |
| 1411 | progress |
| 1412 | .servers |
| 1413 | .iter() |
| 1414 | .any(|row| row.name == "fast" && row.connected) |
| 1415 | ); |
| 1416 | assert!( |
| 1417 | progress |
| 1418 | .servers |
| 1419 | .iter() |
| 1420 | .any(|row| row.name == "slow" && !row.connected) |
| 1421 | ); |
| 1422 | assert!(ready_tools.iter().any(|tool| tool.name == "mcp_fast_ready")); |
| 1423 | assert!(!ready_tools.iter().any(|tool| tool.name == "mcp_slow_ready")); |
| 1424 | let finished = finished.expect("finished boot"); |
| 1425 | if invalidate_config { |
| 1426 | assert_eq!(finished.servers.len(), 1); |
| 1427 | assert!(!finished.servers[0].enabled); |
| 1428 | assert!(!finished.servers[0].connected); |
| 1429 | assert!(pool.lock().await.to_api_tools().is_empty()); |
| 1430 | } else { |
| 1431 | assert!(finished.servers.iter().all(|row| row.connected)); |
| 1432 | } |
| 1433 | } |
| 1434 | |
| 1435 | /// Lazy boot (#6033): a configured server nobody selected and nobody marked |
| 1436 | /// `required` must not be spawned at session start. The fixture writes a |
| 1437 | /// `started-<name>` marker when it receives `initialize`, so the lazy |
| 1438 | /// server's absence is proven by the file that never appears — not by a |
| 1439 | /// timeout on "it would have started by now". |
| 1440 | #[tokio::test] |
| 1441 | async fn lazy_boot_leaves_unselected_servers_unspawned() { |
| 1442 | let Some(node) = crate::dependencies::resolve_node() else { |
| 1443 | return; |
| 1444 | }; |
| 1445 | let tmp = tempdir().expect("tempdir"); |
| 1446 | let server = tmp.path().join("server.mjs"); |
| 1447 | fs::write( |
| 1448 | &server, |
| 1449 | r#"import fs from 'node:fs'; |
| 1450 | import path from 'node:path'; |
| 1451 | import readline from 'node:readline'; |
| 1452 | readline.createInterface({ input: process.stdin }).on('line', line => { |
| 1453 | const request = JSON.parse(line); |
| 1454 | if (request.id === undefined) return; |
| 1455 | if (request.method === 'initialize') { |
| 1456 | fs.writeFileSync(path.join(process.argv[3], 'started-' + process.argv[2]), 'ready'); |
| 1457 | } |
| 1458 | const result = request.method === 'initialize' |
| 1459 | ? { protocolVersion: '2024-11-05', capabilities: { tools: {} }, serverInfo: { name: process.argv[2], version: '1' } } |
| 1460 | : { tools: [{ name: 'ready', inputSchema: { type: 'object' } }] }; |
| 1461 | process.stdout.write(JSON.stringify({ jsonrpc: '2.0', id: request.id, result }) + '\n'); |
| 1462 | });"#, |
| 1463 | ) |
| 1464 | .expect("fixture"); |
| 1465 | let config_path = tmp.path().join("mcp.json"); |
| 1466 | fs::write( |
| 1467 | &config_path, |
| 1468 | serde_json::to_vec(&serde_json::json!({ |
| 1469 | "servers": { |
| 1470 | "eager": { "command": node, "args": [server, "eager", tmp.path()], "required": true }, |
| 1471 | "lazy": { "command": node, "args": [server, "lazy", tmp.path()] } |
| 1472 | } |
| 1473 | })) |
| 1474 | .unwrap(), |
| 1475 | ) |
| 1476 | .expect("MCP config"); |
| 1477 | let (engine, handle) = Engine::new( |
| 1478 | EngineConfig { |
| 1479 | workspace: tmp.path().to_path_buf(), |
| 1480 | mcp_config_path: config_path, |
| 1481 | ..Default::default() |
| 1482 | }, |
| 1483 | &Config::default(), |
| 1484 | ); |
| 1485 | let task = tokio::spawn(async move { engine.run().await }); |
| 1486 | let mut events = handle.rx_event.write().await; |
| 1487 | let finished = tokio::time::timeout(Duration::from_secs(30), async { |
| 1488 | while let Some(event) = events.recv().await { |
| 1489 | if let Event::McpSessionBoot { |
| 1490 | snapshot, |
| 1491 | connecting, |
| 1492 | finished: true, |
| 1493 | .. |
| 1494 | } = event |
| 1495 | { |
| 1496 | return (snapshot, connecting); |
| 1497 | } |
| 1498 | } |
| 1499 | panic!("engine event channel closed"); |
| 1500 | }) |
| 1501 | .await |
| 1502 | .expect("boot must finish"); |
| 1503 | drop(events); |
| 1504 | handle.send(Op::Shutdown).await.expect("shutdown"); |
| 1505 | task.await.expect("engine task"); |
| 1506 | |
| 1507 | let (snapshot, connecting) = finished; |
| 1508 | assert!( |
| 1509 | tmp.path().join("started-eager").exists(), |
| 1510 | "a required server is still eager" |
| 1511 | ); |
| 1512 | assert!( |
| 1513 | !tmp.path().join("started-lazy").exists(), |
| 1514 | "an unselected, unrequired server must not be spawned at boot" |
| 1515 | ); |
| 1516 | assert!(connecting.is_empty()); |
| 1517 | let lazy = snapshot |
| 1518 | .servers |
| 1519 | .iter() |
| 1520 | .find(|server| server.name == "lazy") |
| 1521 | .expect("configured lazy server still appears in the snapshot"); |
| 1522 | assert!(!lazy.connected); |
| 1523 | assert!(lazy.error.is_none(), "lazy is a state, not a failure"); |
| 1524 | let eager = snapshot |
| 1525 | .servers |
| 1526 | .iter() |
| 1527 | .find(|server| server.name == "eager") |
| 1528 | .expect("required server in snapshot"); |
| 1529 | assert!(eager.connected); |
| 1530 | } |
| 1531 | |
| 1532 | #[tokio::test] |
| 1533 | async fn mcp_boot_updates_preserve_authority_errors_and_replace_ordinary_errors() { |
| 1534 | let tmp = tempdir().expect("tempdir"); |
| 1535 | let engine_config = EngineConfig { |
| 1536 | workspace: tmp.path().to_path_buf(), |
| 1537 | ..Default::default() |
| 1538 | }; |
| 1539 | let (mut engine, _handle) = Engine::new(engine_config, &Config::default()); |
| 1540 | engine.mcp_event_generation = 3; |
| 1541 | engine.mcp_boot_generation = Some(3); |
| 1542 | engine.mcp_boot_in_flight = true; |
| 1543 | engine.mcp_connection_errors = HashMap::from([( |
| 1544 | "stale-transport".to_string(), |
| 1545 | "obsolete connection failure".to_string(), |
| 1546 | )]); |
| 1547 | let authority_errors = Arc::new(HashMap::from([( |
| 1548 | "revoked-plugin".to_string(), |
| 1549 | "plugin authority revoked or changed".to_string(), |
| 1550 | )])); |
| 1551 | |
| 1552 | engine |
| 1553 | .apply_mcp_boot_update(McpBootUpdate::Progress { |
| 1554 | generation: 3, |
| 1555 | authority_errors: Arc::clone(&authority_errors), |
| 1556 | connection_errors: HashMap::from([( |
| 1557 | "current-transport".to_string(), |
| 1558 | "current connection failure".to_string(), |
| 1559 | )]), |
| 1560 | connecting: Vec::new(), |
| 1561 | }) |
| 1562 | .await; |
| 1563 | assert_eq!( |
| 1564 | engine.mcp_connection_errors, |
| 1565 | HashMap::from([ |
| 1566 | ( |
| 1567 | "revoked-plugin".to_string(), |
| 1568 | "plugin authority revoked or changed".to_string(), |
| 1569 | ), |
| 1570 | ( |
| 1571 | "current-transport".to_string(), |
| 1572 | "current connection failure".to_string(), |
| 1573 | ), |
| 1574 | ]) |
| 1575 | ); |
| 1576 | assert!(!engine.mcp_connection_errors.contains_key("stale-transport")); |
| 1577 | assert_eq!( |
| 1578 | engine.session.pending_prefix_change_reason.as_deref(), |
| 1579 | Some("mcp-session-boot") |
| 1580 | ); |
| 1581 | |
| 1582 | engine.mcp_connection_errors.insert( |
| 1583 | "stale-between-updates".to_string(), |
| 1584 | "must not survive finished".to_string(), |
| 1585 | ); |
| 1586 | engine |
| 1587 | .apply_mcp_boot_update(McpBootUpdate::Finished { |
| 1588 | generation: 3, |
| 1589 | authority_errors, |
| 1590 | connection_errors: HashMap::from([( |
| 1591 | "final-transport".to_string(), |
| 1592 | "final connection failure".to_string(), |
| 1593 | )]), |
| 1594 | }) |
| 1595 | .await; |
| 1596 | assert_eq!( |
| 1597 | engine.mcp_connection_errors, |
| 1598 | HashMap::from([ |
| 1599 | ( |
| 1600 | "revoked-plugin".to_string(), |
| 1601 | "plugin authority revoked or changed".to_string(), |
| 1602 | ), |
| 1603 | ( |
| 1604 | "final-transport".to_string(), |
| 1605 | "final connection failure".to_string(), |
| 1606 | ), |
| 1607 | ]) |
| 1608 | ); |
| 1609 | } |
| 1610 | // Actual optional stdio servers; no provider request or invented tool catalogue. |
| 1611 | async fn mcp_search_fixture( |
| 1612 | allowed: Option<Vec<String>>, |
| 1613 | ) -> ( |
| 1614 | Engine, |
| 1615 | tool_catalog::ToolSurfacePolicy, |
| 1616 | tempfile::TempDir, |
| 1617 | Arc<crate::llm_client::mock::MockLlmClient>, |
| 1618 | ) { |
| 1619 | let node = |
| 1620 | crate::dependencies::resolve_node().expect("MCP discovery qualification requires Node"); |
| 1621 | let tmp = tempdir().expect("tempdir"); |
| 1622 | let server = tmp.path().join("discovery.mjs"); |
| 1623 | fs::write(&server, r#"import fs from 'node:fs'; |
| 1624 | import path from 'node:path'; |
| 1625 | import readline from 'node:readline'; |
| 1626 | readline.createInterface({input:process.stdin}).on('line', line => { |
| 1627 | const r = JSON.parse(line); if (r.id === undefined) return; |
| 1628 | const name = process.argv[2], root = process.argv[3]; |
| 1629 | if (r.method === 'initialize') { |
| 1630 | fs.writeFileSync(path.join(root, 'started-' + name), 'started'); |
| 1631 | if (name === 'stalled') return; |
| 1632 | } |
| 1633 | if (r.method === 'tools/list') fs.writeFileSync(path.join(root, 'listed-' + name), 'listed'); |
| 1634 | const result = r.method === 'initialize' |
| 1635 | ? {protocolVersion:'2024-11-05',capabilities:{tools:{}},serverInfo:{name,version:'1'}} |
| 1636 | : {tools:[{name:'actual',description:'Actual MCP fixture',inputSchema:{type:'object',properties:{verified:{type:'boolean'}},required:['verified']}}]}; |
| 1637 | process.stdout.write(JSON.stringify({jsonrpc:'2.0',id:r.id,result}) + '\n'); |
| 1638 | });"#).unwrap(); |
| 1639 | let config_path = tmp.path().join("mcp.json"); |
| 1640 | fs::write( |
| 1641 | &config_path, |
| 1642 | serde_json::to_vec(&json!({"servers": { |
| 1643 | "engram": {"command":node,"args":[server,"engram",tmp.path()],"required":false}, |
| 1644 | "unrelated": {"command":node,"args":[server,"unrelated",tmp.path()],"required":false}, |
| 1645 | "disabled": {"command":node,"args":[server,"disabled",tmp.path()],"enabled":false}, |
| 1646 | "stalled": {"command":node,"args":[server,"stalled",tmp.path()],"required":false}, |
| 1647 | "failed": {"command":"codewhale-missing-mcp-discovery-6828"} |
| 1648 | }})) |
| 1649 | .unwrap(), |
| 1650 | ) |
| 1651 | .unwrap(); |
| 1652 | let api = Config::default(); |
| 1653 | let model = Arc::new(crate::llm_client::mock::MockLlmClient::new(Vec::new())); |
| 1654 | let client: crate::core::model_client::SharedModelClient = model.clone(); |
| 1655 | let (mut engine, _) = Engine::new_with_model_client( |
| 1656 | EngineConfig { |
| 1657 | workspace: tmp.path().to_path_buf(), |
| 1658 | mcp_config_path: config_path, |
| 1659 | ..Default::default() |
| 1660 | }, |
| 1661 | &api, |
| 1662 | client, |
| 1663 | ); |
| 1664 | engine |
| 1665 | .start_mcp_session_boot(McpConnectRefresh::IfChanged) |
| 1666 | .await |
| 1667 | .unwrap(); |
| 1668 | assert!( |
| 1669 | !tmp.path().join("started-engram").exists(), |
| 1670 | "optional boot must remain lazy" |
| 1671 | ); |
| 1672 | let authority = crate::core::authority::TurnAuthority::from_effective_fields( |
| 1673 | AppMode::Agent, |
| 1674 | true, |
| 1675 | false, |
| 1676 | false, |
| 1677 | ApprovalMode::Suggest, |
| 1678 | ); |
| 1679 | let route = TurnRouteContext { |
| 1680 | provider: ProviderKind::Deepseek, |
| 1681 | model: DEFAULT_TEXT_MODEL.to_string(), |
| 1682 | capabilities: codewhale_config::route::RouteCapabilities::default(), |
| 1683 | limits: None, |
| 1684 | client: engine.codewhale_client.clone(), |
| 1685 | api_config: Box::new(api), |
| 1686 | locale_tag: engine.config.locale_tag.clone(), |
| 1687 | role_models: engine.subagent_role_models(), |
| 1688 | auto_model: false, |
| 1689 | reasoning_effort: None, |
| 1690 | reasoning_effort_auto: false, |
| 1691 | }; |
| 1692 | let build = engine |
| 1693 | .build_turn_tool_registry_and_catalog( |
| 1694 | &authority, |
| 1695 | &[], |
| 1696 | allowed, |
| 1697 | SubAgentWiring::Inert, |
| 1698 | McpAccess::Connect, |
| 1699 | route, |
| 1700 | "discovery", |
| 1701 | ) |
| 1702 | .await; |
| 1703 | (engine, build.surface, tmp, model) |
| 1704 | } |
| 1705 | |
| 1706 | #[tokio::test] |
| 1707 | async fn mcp_tool_search_discovers_optional_server_real_schema_without_eager_siblings() { |
| 1708 | let (mut engine, policy, tmp, model) = mcp_search_fixture(None).await; |
| 1709 | let mut surface = crate::core::engine::ChildSurfaceProbe { |
| 1710 | policy, |
| 1711 | cache: crate::core::session::ToolActivationCache::default(), |
| 1712 | }; |
| 1713 | let result = engine |
| 1714 | .probe_child_tool_batch( |
| 1715 | &mut surface, |
| 1716 | crate::core::engine::ChildProbeCall { |
| 1717 | id: "mcp-discovery".into(), |
| 1718 | execution_id: "mcp-discovery".into(), |
| 1719 | name: "tool_search".into(), |
| 1720 | input: json!({"query":"engram","match":"bm25"}), |
| 1721 | }, |
| 1722 | ) |
| 1723 | .await |
| 1724 | .expect("actual Core planner and executor search"); |
| 1725 | assert!(result.result.success); |
| 1726 | let catalog = &surface.policy.catalog; |
| 1727 | assert!( |
| 1728 | tmp.path().join("listed-engram").exists(), |
| 1729 | "schema must come from tools/list" |
| 1730 | ); |
| 1731 | assert!(!tmp.path().join("started-unrelated").exists()); |
| 1732 | let tool = catalog |
| 1733 | .iter() |
| 1734 | .find(|tool| tool.name == "mcp_engram_actual") |
| 1735 | .expect("real advertised tool"); |
| 1736 | assert_eq!( |
| 1737 | tool.input_schema["properties"]["verified"]["type"], |
| 1738 | "boolean" |
| 1739 | ); |
| 1740 | assert_eq!(tool.input_schema["required"], json!(["verified"])); |
| 1741 | assert!(surface.policy.active_names.contains("mcp_engram_actual")); |
| 1742 | assert!( |
| 1743 | result.result.metadata.as_ref().unwrap()["tool_references"] |
| 1744 | .as_array() |
| 1745 | .unwrap() |
| 1746 | .iter() |
| 1747 | .any(|name| name == "mcp_engram_actual") |
| 1748 | ); |
| 1749 | assert!(!catalog.iter().any(|tool| tool.name == "mcp_engram_guessed")); |
| 1750 | assert_eq!( |
| 1751 | model.call_count(), |
| 1752 | 0, |
| 1753 | "MCP discovery must not call a provider" |
| 1754 | ); |
| 1755 | } |
| 1756 | |
| 1757 | #[tokio::test] |
| 1758 | async fn mcp_tool_search_respects_captured_turn_ceiling_and_disabled_servers() { |
| 1759 | let (mut engine, policy, tmp, model) = |
| 1760 | mcp_search_fixture(Some(vec!["tool_search".into()])).await; |
| 1761 | let mut catalog = policy.catalog.clone(); |
| 1762 | let mut active = policy.active_names.clone(); |
| 1763 | engine |
| 1764 | .discover_mcp_for_tool_search( |
| 1765 | ("tool_search", &json!({"query":"mcp_.*","match":"regex"})), |
| 1766 | &policy, |
| 1767 | &mut catalog, |
| 1768 | &mut active, |
| 1769 | None, |
| 1770 | ) |
| 1771 | .await |
| 1772 | .unwrap(); |
| 1773 | for name in ["engram", "unrelated", "disabled", "stalled"] { |
| 1774 | assert!(!tmp.path().join(format!("started-{name}")).exists()); |
| 1775 | } |
| 1776 | assert!(!catalog.iter().any(|tool| tool.name.starts_with("mcp_"))); |
| 1777 | assert_eq!( |
| 1778 | model.call_count(), |
| 1779 | 0, |
| 1780 | "MCP discovery must not call a provider" |
| 1781 | ); |
| 1782 | } |
| 1783 | |
| 1784 | #[tokio::test] |
| 1785 | async fn mcp_tool_search_failed_connection_is_bounded_without_guessed_tools() { |
| 1786 | let (mut engine, policy, _tmp, model) = mcp_search_fixture(None).await; |
| 1787 | let mut catalog = policy.catalog.clone(); |
| 1788 | let mut active = policy.active_names.clone(); |
| 1789 | tokio::time::timeout( |
| 1790 | Duration::from_secs(6), |
| 1791 | engine.discover_mcp_for_tool_search( |
| 1792 | ("tool_search", &json!({"query":"failed"})), |
| 1793 | &policy, |
| 1794 | &mut catalog, |
| 1795 | &mut active, |
| 1796 | None, |
| 1797 | ), |
| 1798 | ) |
| 1799 | .await |
| 1800 | .expect("existing five-second bound") |
| 1801 | .unwrap(); |
| 1802 | assert!(engine.mcp_connection_errors.contains_key("failed")); |
| 1803 | assert!( |
| 1804 | !catalog |
| 1805 | .iter() |
| 1806 | .any(|tool| tool.name.starts_with("mcp_failed_")) |
| 1807 | ); |
| 1808 | assert_eq!( |
| 1809 | model.call_count(), |
| 1810 | 0, |
| 1811 | "MCP discovery must not call a provider" |
| 1812 | ); |
| 1813 | } |
| 1814 | |
| 1815 | #[tokio::test] |
| 1816 | async fn mcp_tool_search_cancel_aborts_handshake_and_clears_connecting() { |
| 1817 | let (mut engine, policy, tmp, model) = mcp_search_fixture(None).await; |
| 1818 | let pool = engine.mcp_pool.clone().unwrap(); |
| 1819 | let cancel = engine.cancel_token.clone(); |
| 1820 | let mut catalog = policy.catalog.clone(); |
| 1821 | let mut active = policy.active_names.clone(); |
| 1822 | let input = json!({"query":"stalled"}); |
| 1823 | let discovery = engine.discover_mcp_for_tool_search( |
| 1824 | ("tool_search", &input), |
| 1825 | &policy, |
| 1826 | &mut catalog, |
| 1827 | &mut active, |
| 1828 | None, |
| 1829 | ); |
| 1830 | tokio::pin!(discovery); |
| 1831 | let result = tokio::time::timeout(Duration::from_secs(6), async { |
| 1832 | loop { |
| 1833 | tokio::select! { |
| 1834 | result = &mut discovery => break result, |
| 1835 | () = tokio::time::sleep(Duration::from_millis(10)) => { |
| 1836 | if tmp.path().join("started-stalled").exists() { cancel.cancel(); } |
| 1837 | } |
| 1838 | } |
| 1839 | } |
| 1840 | }) |
| 1841 | .await |
| 1842 | .expect("cancel remains bounded"); |
| 1843 | assert!(matches!(result, Err(ToolError::PermissionDenied { .. }))); |
| 1844 | assert!(pool.lock().await.connecting_servers().is_empty()); |
| 1845 | assert!(pool.lock().await.to_api_tools().is_empty()); |
| 1846 | assert_eq!( |
| 1847 | model.call_count(), |
| 1848 | 0, |
| 1849 | "MCP discovery must not call a provider" |
| 1850 | ); |
| 1851 | } |
| 1852 |