返回 CodeWhale
test_cases_18.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_18.rs
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
1852 lines RUST