返回 CodeWhale
test_cases_05.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_05.rs
1 #[tokio::test]
2 async fn host_managed_engine_defers_idle_subagent_completion_to_explicit_turn() {
3 use crate::llm_client::mock::{MockLlmClient, canned};
4 use crate::tools::subagent::SubAgentCompletion;
5
6 let mut custom = HashMap::new();
7 custom.insert(
8 "custom-a".to_string(),
9 crate::config::ProviderConfig {
10 kind: Some("openai-compatible".to_string()),
11 base_url: Some("http://127.0.0.1:18181/v1".to_string()),
12 model: Some("local-model".to_string()),
13 api_key: Some("local-test-key".to_string()),
14 ..crate::config::ProviderConfig::default()
15 },
16 );
17 let config = Config {
18 provider: Some("custom-a".to_string()),
19 providers: Some(crate::config::ProvidersConfig {
20 custom,
21 ..crate::config::ProvidersConfig::default()
22 }),
23 ..Config::default()
24 };
25 let runtime_services = crate::tools::spec::RuntimeToolServices {
26 active_thread_id: Some("thr_host_managed".to_string()),
27 ..crate::tools::spec::RuntimeToolServices::default()
28 };
29 let engine_config = EngineConfig {
30 max_steps: 0,
31 snapshots_enabled: false,
32 terminal_chrome_enabled: false,
33 runtime_services,
34 ..EngineConfig::default()
35 };
36 // This branch grants a zero-step host turn exactly one A2 report request
37 // (see host_managed_engine_does_not_self_dispatch_goal_continuation,
38 // which asserts that single dispatch). Script the response the same way
39 // so the turn completes deterministically instead of dialing the loopback
40 // route fixture that nothing serves in this test.
41 let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
42 "drained host turn",
43 )]));
44 let client: crate::core::model_client::SharedModelClient = mock.clone();
45 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
46 let owner_session_id = engine.session.id.clone();
47 let tx_subagent_completion = engine.tx_subagent_completion.clone();
48 let run_task = tokio::spawn(engine.run());
49
50 tx_subagent_completion
51 .try_send(SubAgentCompletion {
52 owner_session_id,
53 agent_id: "agent_deferred".to_string(),
54 payload: "deferred child result".to_string(),
55 })
56 .expect("queue sub-agent completion");
57 assert!(
58 tokio::time::timeout(Duration::from_millis(100), async {
59 handle.rx_event.write().await.recv().await
60 })
61 .await
62 .is_err(),
63 "an idle child completion must not create an unclaimed hosted turn"
64 );
65
66 handle
67 .send(Op::SendMessage(TurnSpec {
68 max_output_tokens: None,
69 content: "claim the next turn".to_string(),
70 images: Vec::new(),
71 mode: AppMode::Agent,
72 route: resolved_route_for_test(&config, "local-model"),
73 compaction: Box::new(CompactionConfig::default()),
74 initial_routed_usage: Box::default(),
75 goal_objective: None,
76 goal_token_budget: None,
77 goal_status: crate::tools::goal::GoalStatus::Active,
78 reasoning_effort: None,
79 reasoning_effort_auto: false,
80 auto_model: false,
81 allow_shell: false,
82 trust_mode: false,
83 auto_approve: false,
84 approval_mode: ApprovalMode::Suggest,
85 translation_enabled: false,
86 allowed_tools: None,
87 dynamic_tools: Vec::new(),
88 hook_executor: None,
89 verbosity: None,
90 provenance: UserInputProvenance::ExternalUser,
91 submission_id: None,
92 }))
93 .await
94 .expect("send explicit host turn");
95
96 let mut starts = 0;
97 let mut drained_completion = false;
98 loop {
99 let event = tokio::time::timeout(Duration::from_secs(3), async {
100 handle.rx_event.write().await.recv().await
101 })
102 .await
103 .expect("host engine event timeout")
104 .expect("host engine event");
105 match event {
106 Event::TurnStarted { .. } => starts += 1,
107 Event::Status { message } => {
108 drained_completion |= message.contains("1 queued sub-agent completion");
109 }
110 Event::TurnComplete { .. } => break,
111 _ => {}
112 }
113 }
114 assert_eq!(starts, 1);
115 assert!(
116 drained_completion,
117 "the next explicit turn must drain the queued child completion"
118 );
119
120 handle.send(Op::Shutdown).await.expect("shutdown engine");
121 run_task.await.expect("engine task");
122 }
123
124 /// A host-submitted turn echoes the correlation token it was submitted with,
125 /// while a runtime self-started turn (idle sub-agent completion resume) never
126 /// carries one. With this contract a host can tell "my pending submission
127 /// started" apart from "an autonomous follow-up overtook it in the event
128 /// stream" and only ever consume a deferred submit-window action on the
129 /// former.
130 #[tokio::test]
131 async fn turn_started_echoes_submission_id_and_self_starts_stay_none() {
132 use crate::llm_client::mock::{MockLlmClient, canned};
133 use crate::tools::subagent::SubAgentCompletion;
134
135 let workspace = tempdir().unwrap();
136 let config = Config::default();
137 let mock = Arc::new(MockLlmClient::new(vec![
138 canned::simple_text_turn("submitted turn finished."),
139 canned::simple_text_turn("self-started continuation finished."),
140 ]));
141 let (engine, handle) = Engine::new_with_model_client(
142 deterministic_engine_config(workspace.path()),
143 &config,
144 mock.clone(),
145 );
146 let owner_session_id = engine.session.id.clone();
147 let completion_tx = engine.tx_subagent_completion.clone();
148 let mut op = external_user_message_op("Watch the child agent.", AppMode::Agent, &config);
149 if let Op::SendMessage(TurnSpec { submission_id, .. }) = &mut op {
150 *submission_id = Some("sub-host-1".to_string());
151 }
152 let run_task = tokio::spawn(engine.run());
153 handle.send(op).await.expect("send correlated turn");
154
155 // The submitted turn's start echoes the token verbatim.
156 let submitted_turn_id = {
157 let mut rx = handle.rx_event.write().await;
158 loop {
159 let event = tokio::time::timeout(model_turn_event_timeout(), rx.recv())
160 .await
161 .expect("timed out waiting for the submitted turn")
162 .expect("engine event");
163 if let Event::TurnStarted {
164 turn_id,
165 submission_id,
166 ..
167 } = event
168 {
169 assert_eq!(
170 submission_id.as_deref(),
171 Some("sub-host-1"),
172 "the submitted turn's TurnStarted must echo its correlation token"
173 );
174 break turn_id;
175 }
176 }
177 };
178 {
179 let mut rx = handle.rx_event.write().await;
180 loop {
181 let event = tokio::time::timeout(model_turn_event_timeout(), rx.recv())
182 .await
183 .expect("timed out waiting for the submitted turn to complete")
184 .expect("engine event");
185 if let Event::TurnComplete { status, error, .. } = event {
186 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
187 break;
188 }
189 }
190 }
191
192 // An idle child completion self-starts the follow-up without any host
193 // submission; its start must not present a token.
194 completion_tx
195 .try_send(SubAgentCompletion {
196 owner_session_id,
197 agent_id: "idle-child".to_string(),
198 payload: "child finished its work".to_string(),
199 })
200 .expect("inject idle sub-agent completion");
201 {
202 let mut rx = handle.rx_event.write().await;
203 loop {
204 let event = tokio::time::timeout(model_turn_event_timeout(), rx.recv())
205 .await
206 .expect("timed out waiting for the self-started turn")
207 .expect("engine event");
208 if let Event::TurnStarted {
209 turn_id,
210 submission_id,
211 ..
212 } = event
213 {
214 assert_ne!(
215 turn_id, submitted_turn_id,
216 "the idle child completion must self-start a new turn"
217 );
218 assert!(
219 submission_id.is_none(),
220 "a runtime self-started turn must not present a submission id"
221 );
222 break;
223 }
224 }
225 }
226 {
227 let mut rx = handle.rx_event.write().await;
228 loop {
229 let event = tokio::time::timeout(model_turn_event_timeout(), rx.recv())
230 .await
231 .expect("timed out waiting for the self-started turn to complete")
232 .expect("engine event");
233 if let Event::TurnComplete { status, error, .. } = event {
234 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
235 break;
236 }
237 }
238 }
239
240 handle.send(Op::Shutdown).await.expect("shutdown engine");
241 run_task.await.expect("engine task");
242 assert_eq!(mock.call_count(), 2);
243 }
244
245 #[test]
246 fn idle_and_in_turn_subagent_delivery_claim_each_completion_once() {
247 use crate::tools::subagent::SubAgentCompletion;
248
249 let mut delivered = HashSet::new();
250 let first = SubAgentCompletion {
251 owner_session_id: "session-a".to_string(),
252 agent_id: "agent_same".to_string(),
253 payload: "first delivery".to_string(),
254 };
255 let duplicate = SubAgentCompletion {
256 owner_session_id: "session-a".to_string(),
257 agent_id: "agent_same".to_string(),
258 payload: "duplicate delivery".to_string(),
259 };
260 let second = SubAgentCompletion {
261 owner_session_id: "session-a".to_string(),
262 agent_id: "agent_other".to_string(),
263 payload: "other delivery".to_string(),
264 };
265
266 assert!(claim_subagent_completion(&mut delivered, first).is_some());
267 assert!(claim_subagent_completion(&mut delivered, duplicate).is_none());
268 assert!(claim_subagent_completion(&mut delivered, second).is_some());
269 assert_eq!(
270 delivered,
271 HashSet::from(["agent_same".to_string(), "agent_other".to_string()])
272 );
273 }
274
275 #[tokio::test]
276 async fn session_switch_drops_old_completion_before_deduplication() {
277 use crate::tools::subagent::SubAgentCompletion;
278
279 let workspace = tempdir().expect("tempdir");
280 let (mut engine, _handle) = Engine::new(
281 deterministic_engine_config(workspace.path()),
282 &Config::default(),
283 );
284 engine.session.id = "session-new".to_string();
285 let messages_before = engine.session.messages.len();
286
287 engine
288 .handle_idle_subagent_completion(SubAgentCompletion {
289 owner_session_id: "session-old".to_string(),
290 agent_id: "agent_same".to_string(),
291 payload: "foreign task state".to_string(),
292 })
293 .await;
294
295 assert_eq!(engine.session.messages.len(), messages_before);
296 assert!(engine.delivered_subagent_completion_ids.is_empty());
297 assert!(
298 claim_subagent_completion_for_session(
299 &mut engine.delivered_subagent_completion_ids,
300 "session-new",
301 SubAgentCompletion {
302 owner_session_id: "session-new".to_string(),
303 agent_id: "agent_same".to_string(),
304 payload: "current task state".to_string(),
305 },
306 )
307 .is_some(),
308 "the rejected foreign completion must not suppress the same id in the active session"
309 );
310 }
311
312 #[tokio::test]
313 async fn idle_subagent_delivery_releases_claim_when_route_fails_before_recording() {
314 use crate::tools::subagent::SubAgentCompletion;
315
316 let workspace = tempdir().expect("tempdir");
317 let api_config = Config {
318 ..Config::default()
319 }
320 .with_legacy_root(
321 Some("test-key".to_string()),
322 Some("http://127.0.0.1:1/v1".to_string()),
323 );
324 let (mut engine, _handle) =
325 Engine::new(deterministic_engine_config(workspace.path()), &api_config);
326 // Make the persisted exact identity structurally unresolvable. The
327 // completion is claimed before route resolution, so this exercises the
328 // early error branch before a transcript record can be written.
329 engine.api_provider = ProviderKind::Custom;
330 engine.api_provider_identity = None;
331
332 engine
333 .handle_idle_subagent_completion(SubAgentCompletion {
334 owner_session_id: engine.session.id.clone(),
335 agent_id: "agent_retryable".to_string(),
336 payload: "completed work".to_string(),
337 })
338 .await;
339
340 assert!(
341 !engine
342 .delivered_subagent_completion_ids
343 .contains("agent_retryable"),
344 "a completion that never reached the transcript must remain retryable"
345 );
346 let owner_session_id = engine.session.id.clone();
347 assert!(
348 claim_subagent_completion(
349 &mut engine.delivered_subagent_completion_ids,
350 SubAgentCompletion {
351 owner_session_id,
352 agent_id: "agent_retryable".to_string(),
353 payload: "retry".to_string(),
354 },
355 )
356 .is_some()
357 );
358 }
359
360 #[test]
361 fn subagent_mailbox_keeps_lifecycle_events_reliable() {
362 use crate::tools::subagent::MailboxMessage;
363 use codewhale_models::Usage;
364
365 assert!(subagent_mailbox_message_is_best_effort(
366 &MailboxMessage::progress("agent_a", "step 1")
367 ));
368 assert!(subagent_mailbox_message_is_best_effort(
369 &MailboxMessage::ToolCallStarted {
370 agent_id: "agent_a".to_string(),
371 tool_name: "read_file".to_string(),
372 step: 1,
373 }
374 ));
375 assert!(subagent_mailbox_message_is_best_effort(
376 &MailboxMessage::ToolCallCompleted {
377 agent_id: "agent_a".to_string(),
378 tool_name: "read_file".to_string(),
379 step: 1,
380 ok: true,
381 }
382 ));
383
384 assert!(!subagent_mailbox_message_is_best_effort(
385 &MailboxMessage::started("agent_a", crate::tools::subagent::FleetRole::Scout)
386 ));
387 assert!(!subagent_mailbox_message_is_best_effort(
388 &MailboxMessage::Completed {
389 agent_id: "agent_a".to_string(),
390 summary: "done".to_string(),
391 }
392 ));
393 assert!(!subagent_mailbox_message_is_best_effort(
394 &MailboxMessage::Failed {
395 agent_id: "agent_a".to_string(),
396 error: "failed".to_string(),
397 }
398 ));
399 assert!(!subagent_mailbox_message_is_best_effort(
400 &MailboxMessage::TokenUsage {
401 agent_id: "agent_a".to_string(),
402 source_id: "response-a".to_string(),
403 route: Box::new(crate::cost_status::EffectiveRouteEnvelope::capture(
404 None,
405 ProviderKind::Deepseek,
406 "deepseek",
407 "model",
408 Some(ProviderKind::Deepseek.provider().default_base_url()),
409 chrono::Utc::now(),
410 )),
411 usage: Usage::default(),
412 }
413 ));
414 }
415
416 #[test]
417 fn subagent_mailbox_samples_best_effort_events_per_agent() {
418 use crate::tools::subagent::MailboxMessage;
419
420 let mut last_sent_at = HashMap::new();
421 let start = Instant::now();
422 let first = MailboxMessage::ToolCallStarted {
423 agent_id: "agent_a".to_string(),
424 tool_name: "exec_shell".to_string(),
425 step: 1,
426 };
427 let second = MailboxMessage::ToolCallCompleted {
428 agent_id: "agent_a".to_string(),
429 tool_name: "exec_shell".to_string(),
430 step: 1,
431 ok: true,
432 };
433 let other_agent = MailboxMessage::ToolCallCompleted {
434 agent_id: "agent_b".to_string(),
435 tool_name: "exec_shell".to_string(),
436 step: 1,
437 ok: true,
438 };
439
440 assert!(subagent_mailbox_best_effort_send_permitted(
441 &mut last_sent_at,
442 &first,
443 start,
444 ));
445 assert!(
446 !subagent_mailbox_best_effort_send_permitted(
447 &mut last_sent_at,
448 &second,
449 start + Duration::from_millis(10),
450 ),
451 "same-agent telemetry inside the sampling window is dropped"
452 );
453 assert!(
454 subagent_mailbox_best_effort_send_permitted(
455 &mut last_sent_at,
456 &other_agent,
457 start + Duration::from_millis(10),
458 ),
459 "sampling is per agent, so one busy child cannot hide another"
460 );
461 assert!(
462 subagent_mailbox_best_effort_send_permitted(
463 &mut last_sent_at,
464 &second,
465 start + SUBAGENT_MAILBOX_BEST_EFFORT_MIN_INTERVAL,
466 ),
467 "the next same-agent update is allowed after the interval"
468 );
469 }
470
471 #[test]
472 fn subagent_mailbox_never_samples_lifecycle_or_usage_events() {
473 use crate::tools::subagent::{FleetRole, MailboxMessage};
474 use codewhale_models::Usage;
475
476 let mut last_sent_at = HashMap::new();
477 let start = Instant::now();
478
479 assert!(subagent_mailbox_best_effort_send_permitted(
480 &mut last_sent_at,
481 &MailboxMessage::started("agent_a", FleetRole::Scout),
482 start,
483 ));
484 assert!(subagent_mailbox_best_effort_send_permitted(
485 &mut last_sent_at,
486 &MailboxMessage::Completed {
487 agent_id: "agent_a".to_string(),
488 summary: "done".to_string(),
489 },
490 start,
491 ));
492 assert!(subagent_mailbox_best_effort_send_permitted(
493 &mut last_sent_at,
494 &MailboxMessage::TokenUsage {
495 agent_id: "agent_a".to_string(),
496 source_id: "response-a".to_string(),
497 route: Box::new(crate::cost_status::EffectiveRouteEnvelope::capture(
498 None,
499 ProviderKind::Deepseek,
500 "deepseek",
501 "model",
502 Some(ProviderKind::Deepseek.provider().default_base_url()),
503 chrono::Utc::now(),
504 )),
505 usage: Usage::default(),
506 },
507 start,
508 ));
509 }
510
511 struct ScopedDeepSeekApiKey {
512 previous: Option<OsString>,
513 }
514
515 impl ScopedDeepSeekApiKey {
516 fn set(value: &str) -> Self {
517 let previous = std::env::var_os("DEEPSEEK_API_KEY");
518 // Safety: tests using this helper serialize with lock_test_env() and
519 // restore the original value in Drop.
520 unsafe {
521 std::env::set_var("DEEPSEEK_API_KEY", value);
522 }
523 Self { previous }
524 }
525 }
526
527 impl Drop for ScopedDeepSeekApiKey {
528 fn drop(&mut self) {
529 // Safety: tests using this helper serialize with lock_test_env().
530 unsafe {
531 if let Some(previous) = self.previous.take() {
532 std::env::set_var("DEEPSEEK_API_KEY", previous);
533 } else {
534 std::env::remove_var("DEEPSEEK_API_KEY");
535 }
536 }
537 }
538 }
539
540 fn catalog_tool(name: &str) -> Tool {
541 Tool {
542 tool_type: None,
543 name: name.to_string(),
544 description: String::new(),
545 input_schema: json!({"type": "object"}),
546 allowed_callers: None,
547 defer_loading: None,
548 input_examples: None,
549 strict: None,
550 cache_control: None,
551 }
552 }
553
554 #[test]
555 fn shell_denial_filters_search_catalog_without_expanding_allow_grants() {
556 let raw_names = [
557 "bash",
558 "Bash",
559 "exec_shell",
560 "task_shell_start",
561 "task_gate_run",
562 "terminal/run",
563 "terminal/send",
564 "terminal/reset",
565 "exec_shell_interact",
566 "exec_interact",
567 "code_execution",
568 "js_execution",
569 "rlm_eval",
570 "start_mcp_server",
571 "start_registry_mcp_server",
572 ];
573 for rule in ["Bash", "eXeC_sHeLl", "baSH*", "exec_shell*"] {
574 let surface = policy_for_catalog(
575 raw_names.into_iter().map(catalog_tool).collect(),
576 None,
577 Some(vec![rule.into()]),
578 );
579 for name in raw_names {
580 assert!(surface.denies_call(name, &json!({})), "{rule}: {name}");
581 assert!(
582 surface.catalog.iter().all(|tool| tool.name != name),
583 "{rule}: {name}"
584 );
585 }
586 }
587 let allowed = policy_for_catalog(
588 raw_names.into_iter().map(catalog_tool).collect(),
589 Some(vec!["Bash".into()]),
590 None,
591 );
592 for name in raw_names.into_iter().skip(3) {
593 assert!(
594 !allowed.passes_allow_list(name),
595 "Bash must not grant {name}"
596 );
597 }
598 }
599
600 #[test]
601 fn shell_denial_preserves_task_reads_and_bounded_verification_actions() {
602 let mut tasks = catalog_tool("tasks");
603 tasks.input_schema =
604 json!({"type":"object", "properties":{"action":{"enum":["list", "read", "gate_run"]}}});
605 let surface = policy_for_catalog(
606 vec![
607 tasks,
608 catalog_tool("Run"),
609 catalog_tool("task_shell_wait"),
610 catalog_tool("terminal/cancel"),
611 ],
612 None,
613 Some(vec!["Bash".into()]),
614 );
615 let tasks = surface
616 .catalog
617 .iter()
618 .find(|tool| tool.name == "tasks")
619 .unwrap();
620 assert_eq!(
621 tasks.input_schema["properties"]["action"]["enum"],
622 json!(["list", "read"])
623 );
624 assert!(!surface.denies_call("tasks", &json!({"action":"list"})));
625 assert!(surface.denies_call("tasks", &json!({"action":"gate_run"})));
626 assert!(!surface.denies_call("task_shell_wait", &json!({})));
627 assert!(!surface.denies_call("terminal/cancel", &json!({})));
628 assert!(!surface.denies_call(
629 "Run",
630 &json!({"action":"tests", "args":"-p fixture selected_test"})
631 ));
632 assert!(!surface.denies_call("Run", &json!({"action":"verifiers", "commands":[]})));
633 assert!(surface.denies_call(
634 "Run",
635 &json!({"action":"tests", "args":"--config build.rustc=malicious"})
636 ));
637 assert!(surface.denies_call(
638 "Run",
639 &json!({"action":"verifiers", "commands":[{"program":"sh"}]})
640 ));
641 }
642
643 #[test]
644 fn shell_denial_applies_to_new_durable_execution_without_hiding_management() {
645 use super::tool_catalog::tool_call_denied;
646 let rules = vec!["Bash".to_string()];
647 for (family, actions) in [
648 ("tasks", vec!["create", "gate_run"]),
649 ("automation", vec!["create", "update", "resume", "run"]),
650 ] {
651 for action in actions {
652 assert!(
653 tool_call_denied(Some(&rules), family, &json!({"action": action})),
654 "{family}/{action}"
655 );
656 }
657 }
658 for (family, actions) in [
659 ("tasks", vec!["list", "read", "cancel"]),
660 ("automation", vec!["list", "read", "pause", "delete"]),
661 ] {
662 for action in actions {
663 assert!(
664 !tool_call_denied(Some(&rules), family, &json!({"action": action})),
665 "{family}/{action}"
666 );
667 }
668 }
669 let fetch_rules = vec!["fetch_url".to_string()];
670 for family in ["rlm", "rlm_open"] {
671 assert!(tool_call_denied(
672 Some(&fetch_rules),
673 family,
674 &json!({"action":"open", "url":"https://example.com/document"})
675 ));
676 assert!(!tool_call_denied(
677 Some(&fetch_rules),
678 family,
679 &json!({"action":"open", "content":"local fixture"})
680 ));
681 }
682 }
683
684 fn policy_for_catalog(
685 catalog: Vec<Tool>,
686 allowed_tools: Option<Vec<String>>,
687 disallowed_tools: Option<Vec<String>>,
688 ) -> ToolSurfacePolicy {
689 ToolSurfacePolicy::new(
690 crate::tools::ToolRegistry::new(crate::tools::ToolContext::new(PathBuf::from("."))),
691 Some(catalog),
692 AppMode::Agent,
693 &HashSet::new(),
694 &[],
695 false,
696 allowed_tools,
697 disallowed_tools,
698 None,
699 crate::core::engine::tool_catalog::ToolMode::Direct,
700 )
701 }
702
703 #[test]
704 fn tool_catalog_scenario() {
705 // Scenario consolidation of: tool_catalog_filter_applies_allow_and_deny_gates, tool_catalog_filter_is_inert_without_gates
706 // from tool_catalog_filter_applies_allow_and_deny_gates
707 {
708 // #3027 AC1: the advertised catalog must not contain tools the execution
709 // gates would deny; deny wins over allow.
710 let catalog = vec![
711 catalog_tool("read_file"),
712 catalog_tool("exec_shell"),
713 catalog_tool("grep_files"),
714 ];
715 let surface = policy_for_catalog(
716 catalog,
717 Some(vec!["read_file".to_string(), "exec_shell".to_string()]),
718 Some(vec!["exec_shell".to_string()]),
719 );
720 let names: Vec<&str> = surface.catalog.iter().map(|t| t.name.as_str()).collect();
721 assert_eq!(names, ["read_file"]);
722 }
723 // from tool_catalog_filter_is_inert_without_gates
724 {
725 let surface = policy_for_catalog(
726 vec![catalog_tool("read_file"), catalog_tool("exec_shell")],
727 None,
728 None,
729 );
730 assert!(surface.catalog.iter().any(|tool| tool.name == "read_file"));
731 assert!(surface.catalog.iter().any(|tool| tool.name == "exec_shell"));
732 }
733 }
734
735 #[test]
736 fn tool_catalog_shell_only_benchmark_surface_hides_native_tools() {
737 let catalog = vec![
738 catalog_tool("exec_shell"),
739 catalog_tool("exec_shell_wait"),
740 catalog_tool("exec_shell_interact"),
741 catalog_tool("read_file"),
742 catalog_tool("write_file"),
743 catalog_tool("list_dir"),
744 catalog_tool("git_status"),
745 catalog_tool("work_update"),
746 ];
747 let shell_only = [
748 "exec_shell".to_string(),
749 "exec_shell_wait".to_string(),
750 "exec_shell_interact".to_string(),
751 ];
752
753 let surface = policy_for_catalog(catalog, Some(shell_only.to_vec()), None);
754
755 let names: Vec<&str> = surface.catalog.iter().map(|t| t.name.as_str()).collect();
756 assert_eq!(
757 names,
758 ["exec_shell", "exec_shell_wait", "exec_shell_interact"]
759 );
760 }
761
762 #[test]
763 fn tool_surface_policy_never_reintroduces_denied_synthetic_tools() {
764 let denied = vec![
765 TOOL_SEARCH_NAME.to_string(),
766 CODE_EXECUTION_TOOL_NAME.to_string(),
767 JS_EXECUTION_TOOL_NAME.to_string(),
768 ];
769 let surface = policy_for_catalog(
770 vec![
771 catalog_tool("read_file"),
772 catalog_tool(CODE_EXECUTION_TOOL_NAME),
773 catalog_tool(JS_EXECUTION_TOOL_NAME),
774 ],
775 Some(vec![
776 TOOL_SEARCH_NAME.to_string(),
777 CODE_EXECUTION_TOOL_NAME.to_string(),
778 JS_EXECUTION_TOOL_NAME.to_string(),
779 ]),
780 Some(denied),
781 );
782
783 for denied_name in [
784 TOOL_SEARCH_NAME,
785 CODE_EXECUTION_TOOL_NAME,
786 JS_EXECUTION_TOOL_NAME,
787 ] {
788 assert!(surface.denies_tool(denied_name));
789 assert!(surface.passes_allow_list(denied_name));
790 assert!(
791 !surface.allows_tool(denied_name),
792 "deny must win over allow for {denied_name}"
793 );
794 assert!(
795 surface.catalog.iter().all(|tool| tool.name != denied_name),
796 "{denied_name} must not reappear after policy narrowing"
797 );
798 assert!(!surface.active_names.contains(denied_name));
799 }
800 }
801
802 #[tokio::test]
803 async fn denied_synthetic_tool_is_blocked_by_the_same_turn_policy_at_execution() {
804 use crate::llm_client::mock::{MockLlmClient, canned};
805
806 let workspace = tempdir().expect("tempdir");
807 let mock = std::sync::Arc::new(MockLlmClient::new(vec![
808 canned::tool_call_turn(
809 "call-denied-search",
810 TOOL_SEARCH_NAME,
811 r#"{"query":"File"}"#,
812 ),
813 canned::simple_text_turn("Denied tool handled."),
814 ]));
815 let client: crate::core::model_client::SharedModelClient = mock;
816 let (mut engine, handle) = Engine::new_with_model_client(
817 deterministic_engine_config(workspace.path()),
818 &Config::default(),
819 client,
820 );
821 let policy = policy_for_catalog(
822 vec![catalog_tool("read_file")],
823 Some(vec![TOOL_SEARCH_NAME.to_string()]),
824 Some(vec![TOOL_SEARCH_NAME.to_string()]),
825 );
826 assert!(!policy.allows_tool(TOOL_SEARCH_NAME));
827 let mut turn = crate::core::turn::TurnContext::new(4);
828
829 let (status, error) = engine.run_turn(&mut turn, policy, None, None).await;
830 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
831
832 let mut events = handle.rx_event.write().await;
833 let denied = std::iter::from_fn(|| events.try_recv().ok()).find_map(|event| match event {
834 Event::ToolCallComplete { name, result, .. } if name == TOOL_SEARCH_NAME => Some(result),
835 _ => None,
836 });
837 let error = denied
838 .expect("denied synthetic tool completion")
839 .expect_err("denied synthetic tool must not execute");
840 assert!(
841 error.to_string().contains("disallowed-tools list"),
842 "{error:?}"
843 );
844 }
845
846 #[tokio::test]
847 async fn healthy_owned_children_do_not_force_another_parent_model_turn() {
848 use crate::llm_client::mock::{MockLlmClient, canned};
849
850 for max_steps in [1, 4] {
851 let workspace = tempdir().expect("tempdir");
852 let mock = Arc::new(MockLlmClient::new(vec![
853 canned::simple_text_turn("The workflow is running; I will report its result."),
854 canned::tool_call_turn("must-not-run", "read_file", r#"{"path":"state.txt"}"#),
855 ]));
856 let client: crate::core::model_client::SharedModelClient = mock.clone();
857 let config = EngineConfig {
858 max_steps,
859 ..deterministic_engine_config(workspace.path())
860 };
861 let (mut engine, _handle) =
862 Engine::new_with_model_client(config, &Config::default(), client);
863 let registry = crate::tools::ToolRegistry::new(crate::tools::ToolContext::new(
864 workspace.path().to_path_buf(),
865 ));
866 let surface = test_tool_surface(&engine, registry, None, AppMode::Agent);
867 let children = Arc::new(ForegroundChildRegistry::new());
868 let child_cancel = CancellationToken::new();
869 let registration = children
870 .register(child_cancel.clone(), "agent_child")
871 .expect("child registered");
872 let mut turn = crate::core::turn::TurnContext::new(max_steps);
873
874 let (status, error) = engine
875 .run_turn(&mut turn, surface, Some(Arc::clone(&children)), None)
876 .await;
877
878 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
879 assert_eq!(mock.call_count(), 1, "no forced coordination request");
880 assert_eq!(mock.remaining_turns(), 1, "extra tool turn remains unused");
881 assert_eq!(children.active_count(), 1);
882 assert!(!child_cancel.is_cancelled());
883 assert!(!engine.session.messages.iter().any(|message| {
884 message.content.iter().any(|block| {
885 matches!(block, ContentBlock::Text { text, .. }
886 if text.contains("turn_owned_children_active"))
887 })
888 }));
889 drop(registration);
890 }
891 }
892
893 #[tokio::test]
894 async fn user_steer_during_parent_answer_still_gets_a_reply_with_healthy_children() {
895 use crate::llm_client::mock::{MockLlmClient, canned};
896
897 let workspace = tempdir().expect("tempdir");
898 fs::write(workspace.path().join("state.txt"), "steer-proof\n").expect("fixture");
899 let mock = Arc::new(MockLlmClient::new(Vec::new()));
900 let client: crate::core::model_client::SharedModelClient = mock.clone();
901 let (mut engine, handle) = Engine::new_with_model_client(
902 deterministic_engine_config(workspace.path()),
903 &Config::default(),
904 client,
905 );
906 mock.push_factory(move |_request| {
907 let turn_id = handle
908 .turn_controls
909 .lock()
910 .unwrap()
911 .active
912 .as_ref()
913 .map(|control| control.id);
914 handle
915 .tx_steer
916 .try_send(handle::SteerInput {
917 replace_pending: false,
918 turn_id,
919 content: "Also read state.txt and include its evidence.".to_string(),
920 outcome: None,
921 })
922 .expect("steer channel open");
923 canned::simple_text_turn("The workflow is still running.")
924 });
925 mock.push_turn(canned::tool_call_turn(
926 "call-steer-read",
927 "read_file",
928 r#"{"path":"state.txt"}"#,
929 ));
930 mock.push_turn(canned::simple_text_turn(
931 "The user-requested evidence is steer-proof.",
932 ));
933 let mut registry = crate::tools::ToolRegistry::new(crate::tools::ToolContext::new(
934 workspace.path().to_path_buf(),
935 ));
936 registry.register(Arc::new(crate::tools::file::ReadFileTool));
937 let tools = Some(registry.to_api_tools_with_cache(true));
938 let surface = test_tool_surface(&engine, registry, tools, AppMode::Agent);
939 let children = Arc::new(ForegroundChildRegistry::new());
940 let child_cancel = CancellationToken::new();
941 let registration = children
942 .register(child_cancel.clone(), "agent_child")
943 .expect("child registered");
944 let mut turn = crate::core::turn::TurnContext::new(8);
945
946 let (status, error) = engine
947 .run_turn(&mut turn, surface, Some(Arc::clone(&children)), None)
948 .await;
949
950 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
951 assert_eq!(mock.call_count(), 3);
952 let requests = mock.captured_requests();
953 assert!(requests[1].messages.iter().any(|message| {
954 message.content.iter().any(|block| {
955 matches!(block, ContentBlock::Text { text, .. }
956 if text.contains("Also read state.txt and include its evidence."))
957 })
958 }));
959 assert!(requests[2].messages.iter().any(|message| {
960 message.content.iter().any(|block| {
961 matches!(block, ContentBlock::ToolResult { tool_use_id, .. }
962 if tool_use_id == "call-steer-read")
963 })
964 }));
965 assert_eq!(children.active_count(), 1);
966 assert!(!child_cancel.is_cancelled());
967 assert_eq!(mock.remaining_turns(), 0);
968 drop(registration);
969 }
970
971 /// Compose one assistant turn that proposes `calls` as a single parallel
972 /// tool-call batch: `(call_id, tool_name, args_json)` per block, in order.
973 fn tool_batch_turn(calls: &[(&str, &str, &str)]) -> Vec<codewhale_models::StreamEvent> {
974 use crate::llm_client::mock::canned;
975
976 let mut events = vec![canned::message_start("mock_tool_batch")];
977 for (index, (call_id, tool_name, args_json)) in calls.iter().enumerate() {
978 let index = u32::try_from(index).expect("test batch index fits u32");
979 events.push(canned::tool_use_block_start(index, call_id, tool_name));
980 events.push(canned::tool_input_delta(index, args_json));
981 events.push(canned::block_stop(index));
982 }
983 events.push(canned::message_delta("tool_use", None));
984 events.push(canned::message_stop());
985 events
986 }
987
988 /// Drive one engine turn against the scripted `mock` turns with a registry
989 /// that only serves `read_file`, collecting every `ToolCallComplete` event as
990 /// `(call_id, result)` in emission order.
991 async fn run_budgeted_read_turn(
992 workspace: &Path,
993 max_tool_calls: Option<u32>,
994 mock: std::sync::Arc<crate::llm_client::mock::MockLlmClient>,
995 ) -> (
996 TurnOutcomeStatus,
997 Option<String>,
998 Vec<(String, Result<ToolResult, ToolError>)>,
999 ) {
1000 let mut engine_config = deterministic_engine_config(workspace);
1001 engine_config.max_tool_calls = max_tool_calls;
1002 let client: crate::core::model_client::SharedModelClient = mock;
1003 let (mut engine, handle) =
1004 Engine::new_with_model_client(engine_config, &Config::default(), client);
1005 let context = crate::tools::ToolContext::new(workspace.to_path_buf());
1006 let mut registry = crate::tools::ToolRegistry::new(context);
1007 registry.register(std::sync::Arc::new(crate::tools::file::ReadFileTool));
1008 let tools = Some(registry.to_api_tools_with_cache(true));
1009 let surface = test_tool_surface(&engine, registry, tools, AppMode::Agent);
1010 let mut turn = crate::core::turn::TurnContext::new(4);
1011
1012 let (status, error) = engine.run_turn(&mut turn, surface, None, None).await;
1013 let mut events = handle.rx_event.write().await;
1014 let completions = std::iter::from_fn(|| events.try_recv().ok())
1015 .filter_map(|event| match event {
1016 Event::ToolCallComplete {
1017 model_call: Some(model_call),
1018 result,
1019 ..
1020 } => Some((model_call.provider_id, result)),
1021 _ => None,
1022 })
1023 .collect::<Vec<_>>();
1024 (status, error, completions)
1025 }
1026
1027 #[tokio::test]
1028 async fn provider_id_reuse_mints_distinct_execution_identity_in_all_dispatch_modes() {
1029 use crate::llm_client::mock::{MockLlmClient, canned};
1030 use crate::tools::spec::{ToolCapability, ToolSpec};
1031
1032 struct IdentityTool {
1033 parallel: bool,
1034 observed: Arc<StdMutex<Vec<String>>>,
1035 }
1036 #[async_trait::async_trait]
1037 impl ToolSpec for IdentityTool {
1038 fn name(&self) -> &str {
1039 "fixture_identity"
1040 }
1041 fn description(&self) -> &str {
1042 "Record the admitted execution identity."
1043 }
1044 fn input_schema(&self) -> Value {
1045 json!({"type": "object"})
1046 }
1047 fn capabilities(&self) -> Vec<ToolCapability> {
1048 vec![ToolCapability::ReadOnly]
1049 }
1050 fn supports_parallel(&self) -> bool {
1051 self.parallel
1052 }
1053 async fn execute(&self, _: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
1054 let id = context
1055 .origin_tool_call_id
1056 .as_ref()
1057 .expect("host origin")
1058 .clone();
1059 self.observed.lock().unwrap().push(id.clone());
1060 Ok(ToolResult::success(id))
1061 }
1062 }
1063 for parallel in [false, true] {
1064 let workspace = tempdir().unwrap();
1065 let calls = [
1066 ("reused", "fixture_identity", "{}"),
1067 ("peer", "fixture_identity", "{}"),
1068 ];
1069 let mock = Arc::new(MockLlmClient::new(vec![
1070 tool_batch_turn(&calls),
1071 tool_batch_turn(&calls),
1072 canned::simple_text_turn("done"),
1073 ]));
1074 let (mut engine, handle) = Engine::new_with_model_client(
1075 deterministic_engine_config(workspace.path()),
1076 &Config::default(),
1077 mock.clone(),
1078 );
1079 let observed = Arc::new(StdMutex::new(Vec::new()));
1080 let mut registry =
1081 crate::tools::ToolRegistry::new(ToolContext::new(workspace.path().to_path_buf()));
1082 registry.register(Arc::new(IdentityTool {
1083 parallel,
1084 observed: observed.clone(),
1085 }));
1086 let tools = Some(registry.to_api_tools_with_cache(true));
1087 let surface = test_tool_surface(&engine, registry, tools, AppMode::Agent);
1088 let mut turn = crate::core::turn::TurnContext::new(4);
1089 let (status, error) = engine.run_turn(&mut turn, surface, None, None).await;
1090 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
1091 let ids: HashSet<_> = observed.lock().unwrap().iter().cloned().collect();
1092 assert_eq!(ids.len(), 4);
1093 assert!(ids.iter().all(|id| uuid::Uuid::parse_str(id).is_ok()));
1094 let mut starts = HashMap::new();
1095 let mut completes = HashMap::new();
1096 let mut events = handle.rx_event.write().await;
1097 while let Ok(event) = events.try_recv() {
1098 match event {
1099 Event::ToolCallStarted {
1100 id,
1101 model_call: Some(correlation),
1102 ..
1103 } => {
1104 assert!(starts.insert(id, correlation).is_none());
1105 }
1106 Event::ToolCallComplete {
1107 id,
1108 model_call: Some(correlation),
1109 result,
1110 ..
1111 } => {
1112 assert_eq!(result.unwrap().content, id);
1113 assert!(completes.insert(id, correlation).is_none());
1114 }
1115 _ => {}
1116 }
1117 }
1118 assert_eq!(starts, completes);
1119 assert_eq!(starts.keys().cloned().collect::<HashSet<_>>(), ids);
1120 let requests = mock.captured_requests();
1121 let history = &requests[2].messages;
1122 let mut pairs = HashMap::<String, (usize, usize)>::new();
1123 for block in history.iter().flat_map(|message| &message.content) {
1124 match block {
1125 ContentBlock::ToolUse {
1126 id,
1127 execution_id: Some(local),
1128 ..
1129 } => {
1130 assert_eq!(starts[local].provider_id, *id);
1131 pairs.entry(local.clone()).or_default().0 += 1;
1132 }
1133 ContentBlock::ToolResult {
1134 tool_use_id,
1135 execution_id: Some(local),
1136 ..
1137 } => {
1138 assert_eq!(starts[local].provider_id, *tool_use_id);
1139 pairs.entry(local.clone()).or_default().1 += 1;
1140 }
1141 _ => {}
1142 }
1143 }
1144 assert_eq!(pairs.len(), 4);
1145 assert!(pairs.values().all(|pair| *pair == (1, 1)));
1146 }
1147 }
1148
1149 #[tokio::test]
1150 async fn invalid_provider_tool_batch_is_rejected_before_observation_or_history() {
1151 use crate::llm_client::mock::MockLlmClient;
1152 for ids in [["duplicate", "duplicate"], ["valid", ""], ["valid", " "]] {
1153 let workspace = tempdir().unwrap();
1154 let mock = Arc::new(MockLlmClient::new(vec![tool_batch_turn(&[
1155 (ids[0], "read_file", r#"{"path":"never-read"}"#),
1156 (ids[1], "read_file", r#"{"path":"never-read"}"#),
1157 ])]));
1158 let (mut engine, handle) = Engine::new_with_model_client(
1159 deterministic_engine_config(workspace.path()),
1160 &Config::default(),
1161 mock.clone(),
1162 );
1163 let mut registry =
1164 crate::tools::ToolRegistry::new(ToolContext::new(workspace.path().to_path_buf()));
1165 registry.register(Arc::new(crate::tools::file::ReadFileTool));
1166 let tools = Some(registry.to_api_tools_with_cache(true));
1167 let surface = test_tool_surface(&engine, registry, tools, AppMode::Agent);
1168 let mut turn = crate::core::turn::TurnContext::new(4);
1169 let (status, error) = engine.run_turn(&mut turn, surface, None, None).await;
1170 assert_eq!(status, TurnOutcomeStatus::Failed);
1171 assert!(error.unwrap().contains("pairing id"));
1172 assert_eq!(mock.call_count(), 1);
1173 let mut events = handle.rx_event.write().await;
1174 while let Ok(event) = events.try_recv() {
1175 assert!(!matches!(
1176 event,
1177 Event::ToolCallStarted { .. }
1178 | Event::ToolCallComplete { .. }
1179 | Event::ApprovalRequired { .. }
1180 | Event::ToolGateDecision { .. }
1181 ));
1182 }
1183 assert!(
1184 !engine
1185 .session
1186 .messages
1187 .iter()
1188 .flat_map(|message| &message.content)
1189 .any(|block| matches!(
1190 block,
1191 ContentBlock::ToolUse { .. } | ContentBlock::ToolResult { .. }
1192 ))
1193 );
1194 }
1195 }
1196
1197 /// B1: tool output is redacted once, as it enters the transcript, so the
1198 /// session messages (and the session JSON built from them) never hold a live
1199 /// credential a tool printed.
1200 #[tokio::test]
1201 async fn tool_output_credentials_are_redacted_when_they_enter_the_transcript() {
1202 use crate::llm_client::mock::{MockLlmClient, canned};
1203
1204 const TOKEN: &str = "sk-ant-oat01-AbCdEfGhIjKlMnOpQrStUvWxYz0123456789abcdefghij";
1205 let workspace = tempdir().expect("tempdir");
1206 fs::write(
1207 workspace.path().join("auth.json"),
1208 format!("{{\n \"access_token\": \"{TOKEN}\",\n \"note\": \"keep me\"\n}}\n"),
1209 )
1210 .expect("write fixture");
1211 let mock = std::sync::Arc::new(MockLlmClient::new(vec![
1212 canned::tool_call_turn("call-read", "read_file", r#"{"path":"auth.json"}"#),
1213 canned::simple_text_turn("done"),
1214 ]));
1215 let client: crate::core::model_client::SharedModelClient = mock;
1216 let (mut engine, _handle) = Engine::new_with_model_client(
1217 deterministic_engine_config(workspace.path()),
1218 &Config::default(),
1219 client,
1220 );
1221 let context = crate::tools::ToolContext::new(workspace.path().to_path_buf());
1222 let mut registry = crate::tools::ToolRegistry::new(context);
1223 registry.register(std::sync::Arc::new(crate::tools::file::ReadFileTool));
1224 let tools = Some(registry.to_api_tools_with_cache(true));
1225 let surface = test_tool_surface(&engine, registry, tools, AppMode::Agent);
1226 let mut turn = crate::core::turn::TurnContext::new(4);
1227 let (status, error) = engine.run_turn(&mut turn, surface, None, None).await;
1228 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
1229
1230 let stored = engine
1231 .session
1232 .messages
1233 .iter()
1234 .flat_map(|message| message.content.iter())
1235 .find_map(|block| match block {
1236 ContentBlock::ToolResult { content, .. } => Some(content.clone()),
1237 _ => None,
1238 })
1239 .expect("the read result is in the transcript");
1240 assert!(!stored.contains(TOKEN), "{stored}");
1241 assert!(stored.contains("keep me"), "ordinary bytes stay: {stored}");
1242 let serialized = serde_json::to_string(&engine.session.messages.iter().collect::<Vec<_>>())
1243 .expect("serialize");
1244 assert!(!serialized.contains(TOKEN));
1245 }
1246
1247 /// #4415 AC(a): an 8-call cap admits exactly 8 calls; the 9th is rejected
1248 /// with the typed reason carrying `remaining=0` and is never executed.
1249 #[tokio::test]
1250 async fn tool_call_budget_admits_exactly_the_cap_and_rejects_the_ninth() {
1251 use crate::llm_client::mock::{MockLlmClient, canned};
1252
1253 let workspace = tempdir().expect("tempdir");
1254 let mut calls = Vec::new();
1255 for index in 1..=9 {
1256 let name = format!("fixture-{index}.txt");
1257 fs::write(workspace.path().join(&name), format!("fixture-{index}\n"))
1258 .expect("write fixture");
1259 calls.push((
1260 format!("call-{index}"),
1261 "read_file".to_string(),
1262 format!(r#"{{"path":"{name}"}}"#),
1263 ));
1264 }
1265 let call_refs = calls
1266 .iter()
1267 .map(|(id, name, args)| (id.as_str(), name.as_str(), args.as_str()))
1268 .collect::<Vec<_>>();
1269 let mock = std::sync::Arc::new(MockLlmClient::new(vec![
1270 tool_batch_turn(&call_refs),
1271 canned::simple_text_turn("done"),
1272 ]));
1273
1274 let (status, error, completions) =
1275 run_budgeted_read_turn(workspace.path(), Some(8), mock.clone()).await;
1276 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
1277 assert_eq!(mock.call_count(), 2, "batch turn then the final text turn");
1278 assert_eq!(
1279 completions.len(),
1280 9,
1281 "every proposed call reports a completion"
1282 );
1283
1284 for (id, result) in &completions {
1285 let index = id.strip_prefix("call-").expect("call id");
1286 if index == "9" {
1287 let rejection = result.as_ref().expect_err("the 9th call must be rejected");
1288 let reason = rejection.to_string();
1289 assert!(reason.contains("budget of 8"), "{reason}");
1290 assert!(reason.contains("remaining=0"), "{reason}");
1291 assert!(reason.contains("not executed"), "{reason}");
1292 } else {
1293 let outcome = result.as_ref().expect("calls within budget execute");
1294 assert!(
1295 outcome.content.contains(&format!("fixture-{index}")),
1296 "call {id} must return its file contents: {outcome:?}"
1297 );
1298 }
1299 }
1300 }
1301
1302 /// #5170: a call stopped by an admission gate never executes, so its
1303 /// debited budget slot is refunded — the cap counts admitted calls only.
1304 /// With a cap of 1, a blocked first proposal must leave room for the
1305 /// second proposal to run.
1306 #[tokio::test]
1307 async fn tool_call_budget_refunds_calls_blocked_by_admission_gates() {
1308 use crate::llm_client::mock::{MockLlmClient, canned};
1309
1310 let workspace = tempdir().expect("tempdir");
1311 fs::write(workspace.path().join("fixture.txt"), "fixture\n").expect("write fixture");
1312 let mock = std::sync::Arc::new(MockLlmClient::new(vec![
1313 tool_batch_turn(&[
1314 (
1315 "call-blocked",
1316 "definitely_not_a_tool",
1317 r#"{"path":"fixture.txt"}"#,
1318 ),
1319 ("call-admitted", "read_file", r#"{"path":"fixture.txt"}"#),
1320 ]),
1321 canned::simple_text_turn("done"),
1322 ]));
1323
1324 let (status, error, completions) =
1325 run_budgeted_read_turn(workspace.path(), Some(1), mock.clone()).await;
1326 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
1327 assert_eq!(
1328 completions.len(),
1329 2,
1330 "every proposed call reports a completion"
1331 );
1332
1333 let blocked = completions[0]
1334 .1
1335 .as_ref()
1336 .expect_err("the unknown tool must be blocked by the missing-tool gate");
1337 assert!(
1338 blocked.to_string().contains("definitely_not_a_tool"),
1339 "{blocked}"
1340 );
1341 let admitted = completions[1]
1342 .1
1343 .as_ref()
1344 .expect("a gate-blocked call refunds its slot, so the second call still fits the cap of 1");
1345 assert!(
1346 admitted.content.contains("fixture"),
1347 "the admitted call must return its file contents: {admitted:?}"
1348 );
1349 }
1350
1351 /// An approval card that expires unanswered is a timeout, not the user's
1352 /// denial: the model is told so, and the call — which never ran — gives its
1353 /// tool-call budget slot back, so a cap of 1 still admits the next call.
1354 #[tokio::test]
1355 async fn approval_timeout_is_reported_as_timeout_and_refunds_the_budget() {
1356 use crate::llm_client::mock::{MockLlmClient, canned};
1357
1358 let workspace = tempdir().expect("tempdir");
1359 fs::write(workspace.path().join("fixture.txt"), "fixture\n").expect("write fixture");
1360 let mock = std::sync::Arc::new(MockLlmClient::new(vec![
1361 canned::tool_call_turn("call-timeout", "bash", r#"{"command":"echo first"}"#),
1362 canned::tool_call_turn("call-after", "read_file", r#"{"path":"fixture.txt"}"#),
1363 canned::simple_text_turn("done"),
1364 ]));
1365 let client: crate::core::model_client::SharedModelClient = mock.clone();
1366 let config = Config::default();
1367 let mut engine_config = deterministic_engine_config(workspace.path());
1368 engine_config.exec_policy_engine = ask_rule_engine("echo first");
1369 engine_config.max_tool_calls = Some(1);
1370 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1371 let task = tokio::spawn(engine.run());
1372 handle
1373 .send(external_user_message_op(
1374 "Run the command, then read the fixture.",
1375 AppMode::Agent,
1376 &config,
1377 ))
1378 .await
1379 .expect("send turn");
1380
1381 let mut timed_out = None;
1382 let mut after = None;
1383 let mut rx = handle.rx_event.write().await;
1384 loop {
1385 let event = tokio::time::timeout(model_turn_event_timeout(), rx.recv())
1386 .await
1387 .expect("timed out waiting for the turn")
1388 .expect("engine event stream closed");
1389 match event {
1390 Event::ApprovalRequired { id, tool_name, .. }
1391 if matches!(tool_name.as_str(), "bash" | "Bash") =>
1392 {
1393 handle
1394 .deny_tool_call_timed_out(&id)
1395 .await
1396 .expect("expire the approval card");
1397 }
1398 Event::ApprovalRequired { id, .. } => {
1399 handle.approve_tool_call(&id).await.expect("approve");
1400 }
1401 Event::ToolCallComplete {
1402 model_call: Some(model_call),
1403 result,
1404 ..
1405 } if model_call.provider_id == "call-timeout" => {
1406 timed_out = Some(result);
1407 }
1408 Event::ToolCallComplete {
1409 model_call: Some(model_call),
1410 result,
1411 ..
1412 } if model_call.provider_id == "call-after" => {
1413 after = Some(result);
1414 }
1415 Event::TurnComplete { status, error, .. } => {
1416 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
1417 break;
1418 }
1419 _ => {}
1420 }
1421 }
1422 drop(rx);
1423
1424 let timeout = timed_out
1425 .expect("the expired call reports a completion")
1426 .expect_err("an expired approval never runs the call")
1427 .to_string();
1428 assert!(timeout.contains("timed out"), "{timeout}");
1429 assert!(timeout.contains("did not deny"), "{timeout}");
1430 assert!(
1431 !timeout.contains("denied by user"),
1432 "a timeout must not read as the user's refusal: {timeout}"
1433 );
1434 let after = after
1435 .expect("the next call reports a completion")
1436 .expect("the expired call refunded its slot, so the cap of 1 admits this call");
1437 assert!(after.content.contains("fixture"), "{after:?}");
1438 handle.send(Op::Shutdown).await.expect("shutdown engine");
1439 task.await.expect("engine task");
1440 }
1441
1441 lines RUST