返回 CodeWhale
test_cases_02.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_02.rs
1 const REPRESENTATIVE_HANDOFF_RELAY: &str = "REPRESENTATIVE_HANDOFF_RELAY";
2
3 #[test]
4 fn ordinary_engine_default_does_not_install_a_step_budget() {
5 assert_eq!(
6 DEFAULT_MODEL_STEPS,
7 crate::core::engine::turn_budget::DEFAULT_MAX_MODEL_STEPS
8 );
9 assert_eq!(EngineConfig::default().max_steps, DEFAULT_MODEL_STEPS);
10 let mut turn = TurnContext::new(EngineConfig::default().max_steps);
11 assert_eq!(turn.step_limit(), None);
12 assert_eq!(turn.stop_diagnostics.effective_max_steps, None);
13 // No substitute ceiling or overflow may stop an uncapped turn.
14 turn.step = u32::MAX - 1;
15 assert!(turn.next_step());
16 assert!(turn.next_step());
17 assert!(!turn.at_max_steps());
18 assert_eq!(turn.steps_used(), u32::MAX);
19 }
20
21 #[test]
22 fn registry_instruction_is_in_the_initial_prompt_only_when_mcp_is_enabled() {
23 let enabled = EngineConfig::default();
24 let (engine, _handle) = Engine::new(enabled, &Config::default());
25 let prompt = crate::prompts::system_prompt_flat_text(
26 engine
27 .session
28 .system_prompt
29 .as_ref()
30 .expect("system prompt"),
31 );
32 assert!(prompt.contains(MCP_REGISTRY_FIRST_INSTRUCTION_SOURCE));
33 assert!(prompt.contains("registry_sync"));
34 assert!(prompt.contains("start_registry_mcp_server"));
35
36 let mut disabled = EngineConfig::default();
37 disabled.features.disable(Feature::Mcp);
38 let (engine, _handle) = Engine::new(disabled, &Config::default());
39 let prompt = crate::prompts::system_prompt_flat_text(
40 engine
41 .session
42 .system_prompt
43 .as_ref()
44 .expect("system prompt"),
45 );
46 assert!(!prompt.contains(MCP_REGISTRY_FIRST_INSTRUCTION_SOURCE));
47 }
48
49 /// The engine hands every agent the live posture cell, not a copy: a
50 /// posture the person publishes after the agent's context was built is what
51 /// the agent's next call runs under: approval, shell, sandbox.
52 #[test]
53 fn agent_tool_contexts_follow_the_published_posture() {
54 let (engine, handle) = Engine::new(EngineConfig::default(), &Config::default());
55 let context = engine.build_tool_context(AppMode::Agent, false);
56 let mut before = context.clone();
57 before
58 .refresh_live_posture()
59 .expect("engine contexts carry the live cell");
60 assert!(!before.auto_approve);
61 assert_ne!(
62 before.elevated_sandbox_policy,
63 Some(crate::sandbox::SandboxPolicy::DangerFullAccess)
64 );
65 handle.publish_turn_authority(
66 AppMode::Agent,
67 true,
68 false,
69 true,
70 ApprovalMode::Bypass,
71 None,
72 );
73 let mut after = context;
74 after.refresh_live_posture();
75 assert!(after.auto_approve);
76 assert_eq!(after.approval_mode, ApprovalMode::Bypass);
77 assert_eq!(
78 after.elevated_sandbox_policy,
79 Some(crate::sandbox::SandboxPolicy::DangerFullAccess)
80 );
81 // Narrowing reaches it too: Plan drops shell and makes the sandbox read-only.
82 handle.publish_turn_authority(
83 AppMode::Plan,
84 true,
85 false,
86 false,
87 ApprovalMode::Suggest,
88 None,
89 );
90 after.refresh_live_posture();
91 assert!(!after.auto_approve);
92 assert_eq!(after.shell_policy, crate::worker_profile::ShellPolicy::None);
93 assert_eq!(
94 after.elevated_sandbox_policy,
95 Some(crate::sandbox::SandboxPolicy::ReadOnly)
96 );
97 }
98
99 /// "Gets stuck": an agent's approval answer only reached the agent when the
100 /// engine was idle or itself awaiting approval. While the parent turn
101 /// streamed or ran tools nobody read it, and the agent waited on. The handle
102 /// now hands it to the agent directly — here with no engine running at all.
103 #[tokio::test]
104 async fn an_agent_approval_answer_reaches_the_agent_while_the_engine_is_busy() {
105 let mock = mock_engine_handle();
106 let id = format!("agent:agent_busy:approval:{}", uuid::Uuid::new_v4());
107 let (_, waiting) = mock
108 .handle
109 .subagent_manager
110 .write()
111 .await
112 .register_child_approval("agent_busy", &id, "bash", "held")
113 .expect("register");
114 mock.handle.approve_tool_call(id).await.expect("send");
115 let outcome = tokio::time::timeout(std::time::Duration::from_secs(2), waiting)
116 .await
117 .expect("the waiting agent is answered without the engine")
118 .expect("answer");
119 assert_eq!(
120 outcome,
121 crate::tools::subagent::ChildApprovalOutcome::Approved
122 );
123 }
124
125 /// The regression this test exists for. A real DeepSeek turn — "build a
126 /// self-contained HTML focus timer, read a local fixture, verify it" — spent
127 /// its steps on five `tool_search` calls for `registry_sync`, a deferred-schema
128 /// retry, and reasoning about starting a browser MCP server, because the
129 /// always-visible instruction ordered Registry discovery *before* code
130 /// execution or a manual implementation and named two tools that are not in the
131 /// catalog head.
132 ///
133 /// Two properties keep that from coming back, and both are about this prompt,
134 /// not about a second policy surface: discovery is never ordered ahead of
135 /// ordinary local work, and the deferred-tool cost of reaching the Registry
136 /// tools is stated where the model reads about them.
137 #[test]
138 fn registry_instruction_does_not_gate_ordinary_local_work() {
139 let (engine, _handle) = Engine::new(EngineConfig::default(), &Config::default());
140 let prompt = crate::prompts::system_prompt_flat_text(
141 engine
142 .session
143 .system_prompt
144 .as_ref()
145 .expect("system prompt"),
146 );
147 assert!(prompt.contains("## MCP Registry"));
148
149 for ordered in [
150 "must call `registry_sync`",
151 "before `exec_shell`",
152 "you must call `start_registry_mcp_server`",
153 "not a reason to skip Registry discovery",
154 ] {
155 assert!(
156 !prompt.contains(ordered),
157 "Registry guidance must not order discovery ahead of ordinary work: {ordered:?}"
158 );
159 }
160
161 // Available capability comes first, and the deferred cost is disclosed
162 // where the model reads about the two tools it would have to hunt for.
163 assert!(prompt.contains("Prefer what is already available"));
164 assert!(prompt.contains("not a step before ordinary work"));
165 assert!(prompt.contains("Both Registry tools are deferred"));
166 assert!(prompt.contains("load one with `tool_search`"));
167
168 // The trust boundary the instruction is actually load-bearing for: a
169 // Registry package is started through the approved host path, never
170 // installed or run through the shell.
171 assert!(prompt.contains("rather than installing or running its package command"));
172 }
173
174 /// A provider's input bill describes one route's tokenization of one prompt.
175 /// The compaction gate and the preflight guard lift the honest estimate to
176 /// it, so a bill carried across a route switch would measure the next
177 /// request with the previous route's tokenizer and prefix. Re-installing the
178 /// same route keeps the carry-over #5577 relies on; a different route drops
179 /// it.
180 /// A named custom provider keeps its name, model string, and (absent) limits
181 /// across a config reload that points it at a different server. A different
182 /// server is a different tokenizer, so the bill from the old one must not
183 /// measure the first request to the new one (post-merge finding on #6380).
184 #[test]
185 fn custom_route_endpoint_change_forgets_the_previous_bill() {
186 let mut custom = HashMap::new();
187 custom.insert(
188 "lm-studio".to_string(),
189 crate::config::ProviderConfig {
190 kind: Some("openai-compatible".to_string()),
191 base_url: Some("http://127.0.0.1:18181/v1".to_string()),
192 model: Some("local-model".to_string()),
193 api_key: Some("local-test-key".to_string()),
194 ..crate::config::ProviderConfig::default()
195 },
196 );
197 let config = Config {
198 provider: Some("lm-studio".to_string()),
199 providers: Some(crate::config::ProvidersConfig {
200 custom,
201 ..crate::config::ProvidersConfig::default()
202 }),
203 ..Config::default()
204 };
205 let (mut engine, _handle) = Engine::new(EngineConfig::default(), &config);
206 let install = |engine: &mut Engine, config: &Config| {
207 let route = resolve_runtime_route(config, ProviderKind::Custom, Some("local-model"))
208 .expect("resolve lm-studio")
209 .validate()
210 .expect("preflight lm-studio");
211 engine.install_validated_runtime_route(route);
212 };
213 install(&mut engine, &config);
214 engine.session.latest_parent_input_tokens = Some(150_000);
215
216 install(&mut engine, &config);
217 assert_eq!(
218 engine.session.latest_parent_input_tokens,
219 Some(150_000),
220 "the same endpoint keeps the last bill"
221 );
222
223 let mut reloaded = config;
224 reloaded
225 .providers
226 .as_mut()
227 .and_then(|providers| providers.custom.get_mut("lm-studio"))
228 .expect("named custom provider")
229 .base_url = Some("http://127.0.0.1:18182/v1".to_string());
230 install(&mut engine, &reloaded);
231 assert_eq!(
232 engine.api_provider_identity.as_ref().unwrap().key.as_str(),
233 "lm-studio"
234 );
235 assert_eq!(
236 engine.session.latest_parent_input_tokens, None,
237 "a new endpoint under the same name drops the old server's bill"
238 );
239 }
240
241 /// A catalog refresh can keep a route's name, base URL, model, and limits and
242 /// still move it to another endpoint key or wire protocol. Chat Completions
243 /// and Responses serialize a prompt differently, so the bill from one must
244 /// not measure the first request on the other (post-merge finding on #6381).
245 #[test]
246 fn route_protocol_change_forgets_the_previous_bill() {
247 use codewhale_config::route::{RequestProtocol, ResolvedEndpoint};
248 let (mut engine, _handle) = Engine::new(EngineConfig::default(), &Config::default());
249 let chat = ResolvedEndpoint {
250 base_url: "https://gateway.example/v1".to_string(),
251 endpoint_key: "chat".to_string(),
252 protocol: RequestProtocol::ChatCompletions,
253 };
254 engine.active_route_endpoint = Some(chat.clone());
255 let identity = engine.api_provider_identity.clone().unwrap();
256 let model = engine.session.model.clone();
257 let limits = engine.active_route_limits;
258
259 engine.session.latest_parent_input_tokens = Some(150_000);
260 engine.forget_input_bill_if_route_changes(
261 identity.key.as_str(),
262 identity.persisted_id(),
263 Some(&chat),
264 &model,
265 limits,
266 );
267 assert_eq!(
268 engine.session.latest_parent_input_tokens,
269 Some(150_000),
270 "the same endpoint keeps the last bill"
271 );
272
273 let responses = ResolvedEndpoint {
274 endpoint_key: "responses".to_string(),
275 protocol: RequestProtocol::Responses,
276 ..chat
277 };
278 engine.forget_input_bill_if_route_changes(
279 identity.key.as_str(),
280 identity.persisted_id(),
281 Some(&responses),
282 &model,
283 limits,
284 );
285 assert_eq!(
286 engine.session.latest_parent_input_tokens, None,
287 "a new endpoint key or protocol at the same URL drops the bill"
288 );
289 }
290
291 #[test]
292 fn route_switch_forgets_the_previous_routes_input_bill() {
293 let mut custom = HashMap::new();
294 for (name, base_url, model) in [
295 ("custom-a", "http://127.0.0.1:18181/v1", "model-a"),
296 ("custom-b", "http://127.0.0.1:18182/v1", "model-b"),
297 ] {
298 custom.insert(
299 name.to_string(),
300 crate::config::ProviderConfig {
301 kind: Some("openai-compatible".to_string()),
302 base_url: Some(base_url.to_string()),
303 model: Some(model.to_string()),
304 api_key: Some("local-test-key".to_string()),
305 ..crate::config::ProviderConfig::default()
306 },
307 );
308 }
309 let config = Config {
310 provider: Some("custom-a".to_string()),
311 providers: Some(crate::config::ProvidersConfig {
312 custom,
313 ..crate::config::ProvidersConfig::default()
314 }),
315 ..Config::default()
316 };
317 let (mut engine, _handle) = Engine::new(EngineConfig::default(), &config);
318 let route_a = || {
319 resolve_runtime_route(&config, ProviderKind::Custom, Some("model-a"))
320 .expect("resolve custom A")
321 .validate()
322 .expect("preflight custom A")
323 };
324 engine.install_validated_runtime_route(route_a());
325 engine.session.latest_parent_input_tokens = Some(150_000);
326
327 engine.install_validated_runtime_route(route_a());
328 assert_eq!(
329 engine.session.latest_parent_input_tokens,
330 Some(150_000),
331 "re-installing the same route keeps the last bill"
332 );
333
334 let mut target = config.clone();
335 target.provider = Some("custom-b".to_string());
336 let route_b = resolve_runtime_route(&target, ProviderKind::Custom, Some("model-b"))
337 .expect("resolve custom B")
338 .validate()
339 .expect("preflight custom B");
340 engine.install_validated_runtime_route(route_b);
341 assert_eq!(
342 engine.session.latest_parent_input_tokens, None,
343 "a different route drops the previous route's bill"
344 );
345 }
346
347 #[test]
348 fn custom_route_identity_change_rebuilds_client_for_new_named_endpoint() {
349 let mut custom = HashMap::new();
350 for (name, base_url, model) in [
351 ("custom-a", "http://127.0.0.1:18181/v1", "model-a"),
352 ("custom-b", "http://127.0.0.1:18182/v1", "model-b"),
353 ] {
354 custom.insert(
355 name.to_string(),
356 crate::config::ProviderConfig {
357 kind: Some("openai-compatible".to_string()),
358 base_url: Some(base_url.to_string()),
359 model: Some(model.to_string()),
360 api_key: Some("local-test-key".to_string()),
361 ..crate::config::ProviderConfig::default()
362 },
363 );
364 }
365 let config = Config {
366 provider: Some("custom-a".to_string()),
367 providers: Some(crate::config::ProvidersConfig {
368 custom,
369 ..crate::config::ProvidersConfig::default()
370 }),
371 ..Config::default()
372 };
373 let (mut engine, _handle) = Engine::new(EngineConfig::default(), &config);
374 assert_eq!(
375 engine.api_provider_identity.as_ref().unwrap().key.as_str(),
376 "custom-a"
377 );
378 assert_eq!(
379 engine
380 .codewhale_client
381 .as_ref()
382 .expect("custom A client")
383 .base_url(),
384 "http://127.0.0.1:18181/v1"
385 );
386
387 let mut target = config.clone();
388 target.provider = Some("custom-b".to_string());
389 let route = resolve_runtime_route(&target, ProviderKind::Custom, Some("model-b"))
390 .expect("resolve custom B")
391 .validate()
392 .expect("preflight custom B");
393 engine.install_validated_runtime_route(route);
394
395 assert_eq!(
396 engine.api_provider_identity.as_ref().unwrap().key.as_str(),
397 "custom-b"
398 );
399 assert_eq!(
400 engine
401 .codewhale_client
402 .as_ref()
403 .expect("custom B client")
404 .base_url(),
405 "http://127.0.0.1:18182/v1"
406 );
407 }
408
409 #[test]
410 fn custom_route_config_reload_rebuilds_client_when_identity_is_unchanged() {
411 let mut custom = HashMap::new();
412 custom.insert(
413 "lm-studio".to_string(),
414 crate::config::ProviderConfig {
415 kind: Some("openai-compatible".to_string()),
416 base_url: Some("http://127.0.0.1:18181/v1".to_string()),
417 model: Some("local-model".to_string()),
418 api_key: Some("old-local-test-key".to_string()),
419 ..crate::config::ProviderConfig::default()
420 },
421 );
422 let config = Config {
423 provider: Some("lm-studio".to_string()),
424 providers: Some(crate::config::ProvidersConfig {
425 custom,
426 ..crate::config::ProvidersConfig::default()
427 }),
428 ..Config::default()
429 };
430 let (mut engine, _handle) = Engine::new(EngineConfig::default(), &config);
431
432 let mut reloaded = config;
433 let provider = reloaded
434 .providers
435 .as_mut()
436 .and_then(|providers| providers.custom.get_mut("lm-studio"))
437 .expect("named custom provider");
438 provider.base_url = Some("http://127.0.0.1:18182/v1".to_string());
439 provider.api_key = Some("new-local-test-key".to_string());
440
441 let route = resolve_runtime_route(&reloaded, ProviderKind::Custom, Some("local-model"))
442 .expect("resolve reloaded route")
443 .validate()
444 .expect("preflight reloaded route");
445 engine.install_validated_runtime_route(route);
446
447 assert_eq!(
448 engine.api_provider_identity.as_ref().unwrap().key.as_str(),
449 "lm-studio"
450 );
451 assert_eq!(
452 engine
453 .codewhale_client
454 .as_ref()
455 .expect("reloaded custom client")
456 .base_url(),
457 "http://127.0.0.1:18182/v1"
458 );
459 assert_eq!(
460 engine.api_config.active_route_base_url(),
461 "http://127.0.0.1:18182/v1"
462 );
463 }
464
465 #[test]
466 fn failed_same_identity_route_preflight_leaves_old_client_untouched() {
467 let mut custom = HashMap::new();
468 custom.insert(
469 "lm-studio".to_string(),
470 crate::config::ProviderConfig {
471 kind: Some("openai-compatible".to_string()),
472 base_url: Some("http://127.0.0.1:18181/v1".to_string()),
473 model: Some("local-model".to_string()),
474 api_key: Some("old-local-test-key".to_string()),
475 ..crate::config::ProviderConfig::default()
476 },
477 );
478 let config = Config {
479 provider: Some("lm-studio".to_string()),
480 providers: Some(crate::config::ProvidersConfig {
481 custom,
482 ..crate::config::ProvidersConfig::default()
483 }),
484 ..Config::default()
485 };
486 let (engine, _handle) = Engine::new(EngineConfig::default(), &config);
487 assert!(engine.codewhale_client.is_some());
488
489 let mut invalid = config;
490 invalid
491 .providers
492 .as_mut()
493 .and_then(|providers| providers.custom.get_mut("lm-studio"))
494 .expect("named custom provider")
495 .base_url = Some("ftp://invalid.example/v1".to_string());
496 let err = crate::route_runtime::resolve_runtime_route_for_identity(
497 &invalid,
498 engine
499 .api_provider_identity
500 .as_ref()
501 .expect("captured named identity"),
502 Some("local-model"),
503 )
504 .expect_err("invalid route must fail before installation");
505
506 assert!(err.contains("must be an http(s) URL with a host"), "{err}");
507 assert_eq!(
508 engine.api_provider_identity.as_ref().unwrap().key.as_str(),
509 "lm-studio"
510 );
511 assert!(engine.codewhale_client.is_some());
512 assert!(engine.model_client.is_some());
513 assert!(engine.codewhale_client_error.is_none());
514 }
515
516 #[tokio::test]
517 async fn exact_turn_snapshot_restores_custom_endpoint_and_turn_receipt_after_builtin_route() {
518 use wiremock::matchers::{method, path};
519 use wiremock::{Mock, MockServer, ResponseTemplate};
520
521 let custom_server = MockServer::start().await;
522 let custom_base_url = format!("{}/v1", custom_server.uri());
523 let done_sse = concat!(
524 "data: {\"id\":\"chatcmpl-exact-route\",\"choices\":[{\"index\":0,",
525 "\"delta\":{\"content\":\"exact route\"},\"finish_reason\":null}]}\n\n",
526 "data: {\"id\":\"chatcmpl-exact-route\",\"choices\":[{\"index\":0,",
527 "\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n",
528 "data: [DONE]\n\n",
529 );
530 Mock::given(method("POST"))
531 .and(path("/v1/chat/completions"))
532 .respond_with(
533 ResponseTemplate::new(200)
534 .insert_header("content-type", "text/event-stream")
535 .set_body_string(done_sse),
536 )
537 .expect(1)
538 .mount(&custom_server)
539 .await;
540
541 let mut custom = HashMap::new();
542 custom.insert(
543 "custom-a".to_string(),
544 crate::config::ProviderConfig {
545 kind: Some("openai-compatible".to_string()),
546 base_url: Some(custom_base_url.clone()),
547 model: Some("local-model".to_string()),
548 api_key: Some("local-test-key".to_string()),
549 ..crate::config::ProviderConfig::default()
550 },
551 );
552 let config = Config {
553 provider: Some("custom-a".to_string()),
554 providers: Some(crate::config::ProvidersConfig {
555 openai: crate::config::ProviderConfig {
556 base_url: Some("http://127.0.0.1:18182/v1".to_string()),
557 model: Some("gpt-5.5".to_string()),
558 api_key: Some("builtin-test-key".to_string()),
559 ..crate::config::ProviderConfig::default()
560 },
561 custom,
562 ..crate::config::ProvidersConfig::default()
563 }),
564 ..Config::default()
565 };
566 let engine_config = EngineConfig {
567 max_steps: 1,
568 snapshots_enabled: false,
569 ..EngineConfig::default()
570 };
571 let (mut engine, handle) = Engine::new(engine_config, &config);
572
573 let mut builtin_config = config.clone();
574 builtin_config.provider = Some("openai".to_string());
575 let builtin_route =
576 resolve_runtime_route(&builtin_config, ProviderKind::Openai, Some("gpt-5.5"))
577 .expect("resolve intervening builtin route")
578 .validate()
579 .expect("preflight intervening builtin route");
580 engine.install_validated_runtime_route(builtin_route);
581 assert_eq!(engine.api_provider, ProviderKind::Openai);
582 assert_eq!(
583 engine
584 .codewhale_client
585 .as_ref()
586 .expect("builtin client")
587 .base_url(),
588 "http://127.0.0.1:18182/v1"
589 );
590
591 let run_task = tokio::spawn(engine.run());
592 handle
593 .send(Op::SendMessage(TurnSpec {
594 max_output_tokens: None,
595 content: "verify exact route".to_string(),
596 images: Vec::new(),
597 mode: AppMode::Agent,
598 route: Box::new(
599 resolve_runtime_route(&config, ProviderKind::Custom, Some("local-model"))
600 .expect("resolve exact custom route"),
601 ),
602 compaction: Box::new(CompactionConfig::default()),
603 initial_routed_usage: Box::default(),
604 goal_objective: None,
605 goal_token_budget: None,
606 goal_status: crate::tools::goal::GoalStatus::Active,
607 reasoning_effort: None,
608 reasoning_effort_auto: false,
609 auto_model: true,
610 allow_shell: false,
611 trust_mode: false,
612 auto_approve: false,
613 approval_mode: ApprovalMode::Suggest,
614 translation_enabled: false,
615 allowed_tools: None,
616 dynamic_tools: Vec::new(),
617 hook_executor: None,
618 verbosity: None,
619 provenance: UserInputProvenance::ExternalUser,
620 submission_id: None,
621 }))
622 .await
623 .expect("send exact custom turn");
624
625 // This test runs alongside more than ten thousand TUI tests in the release
626 // parity job. Keep the assertion bounded, but leave enough headroom for a
627 // saturated shared runner to schedule the loopback SSE response.
628 let deadline = tokio::time::Instant::now() + Duration::from_secs(15);
629 let mut lifecycle_stage = 0u8;
630 let mut diagnostics = Vec::new();
631 let mut rx = handle.rx_event.write().await;
632 loop {
633 let event = tokio::time::timeout_at(deadline, rx.recv())
634 .await
635 .unwrap_or_else(|_| {
636 panic!("timed out waiting for semantic route sequence: {diagnostics:?}")
637 })
638 .expect("engine event channel closed before terminal route receipt");
639 diagnostics.push(match &event {
640 Event::TurnStarted { .. } => "turn_started",
641 Event::RouteDispatched { .. } => "route_dispatched",
642 Event::TurnComplete { .. } => "turn_complete",
643 Event::SessionUpdated { .. } => "session_updated",
644 Event::PrefixCacheChange { .. } => "prefix_cache",
645 Event::Status { .. } => "status",
646 _ => "other",
647 });
648 match event {
649 Event::TurnStarted { route, .. } => {
650 assert_eq!(
651 lifecycle_stage, 0,
652 "duplicate/reordered start: {diagnostics:?}"
653 );
654 // Lifecycle start still carries the installed-route receipt
655 // hosts authorize follow-up work against, but it must carry no
656 // billing envelope: nothing has been dispatched yet, and an
657 // undispatched route has no metering surface or billing time.
658 assert!(
659 route.as_ref().is_none_or(|route| route.billing.is_none()),
660 "billing route must not be stamped at lifecycle start"
661 );
662 lifecycle_stage = 1;
663 }
664 Event::RouteDispatched { route, .. } => {
665 assert_eq!(
666 lifecycle_stage, 1,
667 "dispatch missing, duplicated, or reordered: {diagnostics:?}"
668 );
669 assert_eq!(route.provider, ProviderKind::Custom);
670 assert_eq!(route.provider_identity, "custom-a");
671 assert_eq!(route.model, "local-model");
672 assert_eq!(
673 route
674 .billing
675 .as_ref()
676 .and_then(|billing| billing.endpoint_fingerprint.clone()),
677 crate::cost_status::endpoint_fingerprint(&custom_base_url),
678 "dispatch receipt borrowed the later ambient route"
679 );
680 lifecycle_stage = 2;
681 }
682 Event::TurnComplete { base_url, .. } => {
683 assert_eq!(
684 lifecycle_stage, 2,
685 "terminal arrived without an ordered dispatch receipt: {diagnostics:?}"
686 );
687 assert_eq!(base_url.as_deref(), Some(custom_base_url.as_str()));
688 lifecycle_stage = 3;
689 break;
690 }
691 _ => {}
692 }
693 }
694 drop(rx);
695 assert_eq!(lifecycle_stage, 3);
696 assert_eq!(
697 custom_server
698 .received_requests()
699 .await
700 .expect("recorded custom-route request")
701 .len(),
702 1,
703 "semantic dispatch sequence must bracket one real provider request"
704 );
705 handle.send(Op::Shutdown).await.expect("shutdown engine");
706 run_task.await.expect("engine task");
707 }
708
709 /// #6690: the main interactive turn froze its dispatch quote only from the
710 /// provider lake, so an operator's `[[custom_models]]` rate never priced it
711 /// even though background/review envelopes honored the same row.
712 #[tokio::test]
713 async fn main_turn_dispatch_freezes_declared_custom_model_rate() {
714 use wiremock::matchers::{method, path};
715 use wiremock::{Mock, MockServer, ResponseTemplate};
716
717 let custom_server = MockServer::start().await;
718 let custom_base_url = format!("{}/v1", custom_server.uri());
719 let done_sse = concat!(
720 "data: {\"id\":\"chatcmpl-declared-rate\",\"choices\":[{\"index\":0,",
721 "\"delta\":{\"content\":\"priced\"},\"finish_reason\":null}]}\n\n",
722 "data: {\"id\":\"chatcmpl-declared-rate\",\"choices\":[{\"index\":0,",
723 "\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n",
724 "data: [DONE]\n\n",
725 );
726 Mock::given(method("POST"))
727 .and(path("/v1/chat/completions"))
728 .respond_with(
729 ResponseTemplate::new(200)
730 .insert_header("content-type", "text/event-stream")
731 .set_body_string(done_sse),
732 )
733 .mount(&custom_server)
734 .await;
735
736 let mut custom = HashMap::new();
737 custom.insert(
738 "custom-a".to_string(),
739 crate::config::ProviderConfig {
740 kind: Some("openai-compatible".to_string()),
741 base_url: Some(custom_base_url.clone()),
742 model: Some("local-model".to_string()),
743 api_key: Some("local-test-key".to_string()),
744 ..crate::config::ProviderConfig::default()
745 },
746 );
747 let declared: codewhale_config::catalog::configured::ConfiguredModel =
748 toml::from_str(&format!(
749 "provider = \"custom-a\"\nbase_url = \"{custom_base_url}\"\n\
750 id = \"local-model\"\ncost = {{ input = 1.0, output = 2.0 }}\n"
751 ))
752 .expect("declared model row");
753 let config = Config {
754 provider: Some("custom-a".to_string()),
755 providers: Some(crate::config::ProvidersConfig {
756 custom,
757 ..crate::config::ProvidersConfig::default()
758 }),
759 custom_models: Some(vec![declared]),
760 ..Config::default()
761 };
762 let engine_config = EngineConfig {
763 max_steps: 1,
764 snapshots_enabled: false,
765 ..EngineConfig::default()
766 };
767 let (engine, handle) = Engine::new(engine_config, &config);
768 let run_task = tokio::spawn(engine.run());
769 handle
770 .send(Op::SendMessage(TurnSpec {
771 max_output_tokens: None,
772 content: "price this turn".to_string(),
773 images: Vec::new(),
774 mode: AppMode::Agent,
775 route: Box::new(
776 resolve_runtime_route(&config, ProviderKind::Custom, Some("local-model"))
777 .expect("resolve declared custom route"),
778 ),
779 compaction: Box::new(CompactionConfig::default()),
780 initial_routed_usage: Box::default(),
781 goal_objective: None,
782 goal_token_budget: None,
783 goal_status: crate::tools::goal::GoalStatus::Active,
784 reasoning_effort: None,
785 reasoning_effort_auto: false,
786 auto_model: false,
787 allow_shell: false,
788 trust_mode: false,
789 auto_approve: false,
790 approval_mode: ApprovalMode::Suggest,
791 translation_enabled: false,
792 allowed_tools: None,
793 dynamic_tools: Vec::new(),
794 hook_executor: None,
795 verbosity: None,
796 provenance: UserInputProvenance::ExternalUser,
797 submission_id: None,
798 }))
799 .await
800 .expect("send declared custom turn");
801
802 let deadline = tokio::time::Instant::now() + Duration::from_secs(15);
803 let mut rx = handle.rx_event.write().await;
804 let quote = loop {
805 let event = tokio::time::timeout_at(deadline, rx.recv())
806 .await
807 .expect("timed out waiting for route dispatch")
808 .expect("engine event channel closed before dispatch");
809 if let Event::RouteDispatched { route, .. } = event {
810 break route
811 .billing
812 .and_then(|billing| billing.provider_live_pricing)
813 .expect("declared custom_models rate must be frozen at main-turn dispatch");
814 }
815 };
816 drop(rx);
817 assert_eq!(
818 quote.provenance,
819 codewhale_config::pricing::PricingProvenance::UserOverride
820 );
821 assert_eq!(quote.input_per_million.as_deref(), Some("1"));
822 assert_eq!(quote.output_per_million.as_deref(), Some("2"));
823 assert_eq!(quote.cache_read_per_million, None);
824 handle.send(Op::Shutdown).await.expect("shutdown engine");
825 run_task.await.expect("engine task");
826 }
827
828 struct GatedGoalModelClient {
829 calls: std::sync::atomic::AtomicUsize,
830 requests: std::sync::Mutex<Vec<codewhale_models::MessageRequest>>,
831 second_request_entered: std::sync::Arc<tokio::sync::Notify>,
832 release_second_request: std::sync::Arc<tokio::sync::Notify>,
833 first_usage: Option<Usage>,
834 }
835
836 struct FirstRequestGatedGoalModelClient {
837 calls: std::sync::atomic::AtomicUsize,
838 request_entered: std::sync::Arc<tokio::sync::Notify>,
839 release_request: std::sync::Arc<tokio::sync::Notify>,
840 }
841
842 struct IndexedGatedGoalModelClient {
843 calls: std::sync::atomic::AtomicUsize,
844 gates: HashMap<
845 usize,
846 (
847 std::sync::Arc<tokio::sync::Notify>,
848 std::sync::Arc<tokio::sync::Notify>,
849 ),
850 >,
851 max_calls: usize,
852 }
853
854 #[async_trait::async_trait]
855 impl crate::core::model_client::ModelClient for IndexedGatedGoalModelClient {
856 fn provider_name(&self) -> &str {
857 "deterministic-goal"
858 }
859
860 fn model(&self) -> &str {
861 "local-model"
862 }
863
864 async fn create_message(
865 &self,
866 _request: codewhale_models::MessageRequest,
867 ) -> anyhow::Result<codewhale_models::MessageResponse> {
868 anyhow::bail!("indexed gate regression uses the streaming model boundary")
869 }
870
871 async fn create_message_stream(
872 &self,
873 _request: codewhale_models::MessageRequest,
874 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
875 let call = self
876 .calls
877 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
878 .saturating_add(1);
879 if call > self.max_calls {
880 anyhow::bail!("unexpected indexed goal model request #{call}");
881 }
882 if let Some((entered, release)) = self.gates.get(&call).cloned() {
883 entered.notify_one();
884 release.notified().await;
885 }
886
887 let events = crate::llm_client::mock::canned::simple_text_turn("still working")
888 .into_iter()
889 .map(Ok);
890 Ok(Box::pin(futures_util::stream::iter(events)))
891 }
892
893 async fn health_check(&self) -> anyhow::Result<bool> {
894 Ok(true)
895 }
896 }
897
898 #[async_trait::async_trait]
899 impl crate::core::model_client::ModelClient for FirstRequestGatedGoalModelClient {
900 fn provider_name(&self) -> &str {
901 "deterministic-goal"
902 }
903
904 fn model(&self) -> &str {
905 "local-model"
906 }
907
908 async fn create_message(
909 &self,
910 _request: codewhale_models::MessageRequest,
911 ) -> anyhow::Result<codewhale_models::MessageResponse> {
912 anyhow::bail!("mailbox regression uses the streaming model boundary")
913 }
914
915 async fn create_message_stream(
916 &self,
917 _request: codewhale_models::MessageRequest,
918 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
919 let call = self
920 .calls
921 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
922 .saturating_add(1);
923 if call > 1 {
924 anyhow::bail!("unexpected mailbox regression model request #{call}");
925 }
926 self.request_entered.notify_one();
927 self.release_request.notified().await;
928
929 let events = crate::llm_client::mock::canned::simple_text_turn("still working")
930 .into_iter()
931 .map(Ok);
932 Ok(Box::pin(futures_util::stream::iter(events)))
933 }
934
935 async fn health_check(&self) -> anyhow::Result<bool> {
936 Ok(true)
937 }
938 }
939
940 struct FailingGoalModelClient {
941 calls: std::sync::atomic::AtomicUsize,
942 message: String,
943 }
944
945 #[async_trait::async_trait]
946 impl crate::core::model_client::ModelClient for FailingGoalModelClient {
947 fn provider_name(&self) -> &str {
948 "deterministic-goal"
949 }
950
951 fn model(&self) -> &str {
952 "local-model"
953 }
954
955 async fn create_message(
956 &self,
957 _request: codewhale_models::MessageRequest,
958 ) -> anyhow::Result<codewhale_models::MessageResponse> {
959 anyhow::bail!("failure regression uses the streaming model boundary")
960 }
961
962 async fn create_message_stream(
963 &self,
964 _request: codewhale_models::MessageRequest,
965 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
966 self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
967 anyhow::bail!(self.message.clone())
968 }
969
970 async fn health_check(&self) -> anyhow::Result<bool> {
971 Ok(true)
972 }
973 }
974
975 impl GatedGoalModelClient {
976 fn captured_requests(&self) -> Vec<codewhale_models::MessageRequest> {
977 self.requests
978 .lock()
979 .expect("goal model request lock")
980 .clone()
981 }
982 }
983
984 #[async_trait::async_trait]
985 impl crate::core::model_client::ModelClient for GatedGoalModelClient {
986 fn provider_name(&self) -> &str {
987 "deterministic-goal"
988 }
989
990 fn model(&self) -> &str {
991 "local-model"
992 }
993
994 async fn create_message(
995 &self,
996 _request: codewhale_models::MessageRequest,
997 ) -> anyhow::Result<codewhale_models::MessageResponse> {
998 anyhow::bail!("goal regression uses the streaming model boundary")
999 }
1000
1001 async fn create_message_stream(
1002 &self,
1003 request: codewhale_models::MessageRequest,
1004 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
1005 let call = self
1006 .calls
1007 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
1008 .saturating_add(1);
1009 self.requests
1010 .lock()
1011 .expect("goal model request lock")
1012 .push(request);
1013 if call == 2 {
1014 self.second_request_entered.notify_one();
1015 self.release_second_request.notified().await;
1016 } else if call > 2 {
1017 anyhow::bail!("unexpected goal model request #{call}");
1018 }
1019
1020 let mut events = crate::llm_client::mock::canned::simple_text_turn("still working");
1021 if call == 1
1022 && let Some(usage) = self.first_usage.clone()
1023 && let Some(codewhale_models::StreamEvent::MessageDelta { usage: slot, .. }) = events
1024 .iter_mut()
1025 .find(|event| matches!(event, codewhale_models::StreamEvent::MessageDelta { .. }))
1026 {
1027 *slot = Some(usage);
1028 }
1029 let events = events.into_iter().map(Ok);
1030 Ok(Box::pin(futures_util::stream::iter(events)))
1031 }
1032
1033 async fn health_check(&self) -> anyhow::Result<bool> {
1034 Ok(true)
1035 }
1036 }
1037
1038 #[tokio::test]
1039 async fn goal_continuation_preserves_goal_and_resolves_updated_authoritative_route() {
1040 let first_base_url = "http://127.0.0.1:18181/v1".to_string();
1041 let second_base_url = "http://127.0.0.1:18182/v1".to_string();
1042 let second_request_entered = std::sync::Arc::new(tokio::sync::Notify::new());
1043 let release_second_request = std::sync::Arc::new(tokio::sync::Notify::new());
1044 let model = std::sync::Arc::new(GatedGoalModelClient {
1045 calls: std::sync::atomic::AtomicUsize::new(0),
1046 requests: std::sync::Mutex::new(Vec::new()),
1047 second_request_entered: std::sync::Arc::clone(&second_request_entered),
1048 release_second_request: std::sync::Arc::clone(&release_second_request),
1049 first_usage: Some(Usage {
1050 input_tokens: 3,
1051 output_tokens: 2,
1052 ..Usage::default()
1053 }),
1054 });
1055 let mut custom = HashMap::new();
1056 custom.insert(
1057 "custom-a".to_string(),
1058 crate::config::ProviderConfig {
1059 kind: Some("openai-compatible".to_string()),
1060 base_url: Some(first_base_url.clone()),
1061 model: Some("local-model".to_string()),
1062 api_key: Some("local-test-key".to_string()),
1063 ..crate::config::ProviderConfig::default()
1064 },
1065 );
1066 let config = Config {
1067 provider: Some("custom-a".to_string()),
1068 providers: Some(crate::config::ProvidersConfig {
1069 custom,
1070 ..crate::config::ProvidersConfig::default()
1071 }),
1072 ..Config::default()
1073 };
1074 let engine_config = EngineConfig {
1075 max_steps: 1,
1076 snapshots_enabled: false,
1077 terminal_chrome_enabled: false,
1078 goal_objective: Some("keep going".to_string()),
1079 goal_token_budget: Some(50_000),
1080 ..EngineConfig::default()
1081 };
1082 let authoritative = Arc::new(parking_lot::RwLock::new(config.clone()));
1083 let client: crate::core::model_client::SharedModelClient = model.clone();
1084 let (mut engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1085 engine.authoritative_route_config = Some(Arc::clone(&authoritative));
1086 let goal_state = engine.config.goal_state.clone();
1087
1088 handle
1089 .send(Op::SendMessage(TurnSpec {
1090 max_output_tokens: None,
1091 content: "first turn".to_string(),
1092 images: Vec::new(),
1093 mode: AppMode::Agent,
1094 route: resolved_route_for_test(&config, "local-model"),
1095 compaction: Box::new(CompactionConfig::default()),
1096 initial_routed_usage: Box::default(),
1097 goal_objective: Some("keep going".to_string()),
1098 goal_token_budget: Some(50_000),
1099 goal_status: crate::tools::goal::GoalStatus::Active,
1100 reasoning_effort: None,
1101 reasoning_effort_auto: false,
1102 auto_model: false,
1103 allow_shell: false,
1104 trust_mode: false,
1105 auto_approve: false,
1106 approval_mode: ApprovalMode::Suggest,
1107 translation_enabled: false,
1108 allowed_tools: None,
1109 dynamic_tools: Vec::new(),
1110 hook_executor: None,
1111 verbosity: None,
1112 provenance: UserInputProvenance::ExternalUser,
1113 submission_id: None,
1114 }))
1115 .await
1116 .expect("send first goal turn");
1117
1118 let mut reloaded = config;
1119 reloaded
1120 .providers
1121 .as_mut()
1122 .and_then(|providers| providers.custom.get_mut("custom-a"))
1123 .expect("custom route")
1124 .base_url = Some(second_base_url.clone());
1125 *authoritative.write() = reloaded;
1126 let refreshed_route = engine
1127 .current_runtime_route()
1128 .expect("resolve the updated authoritative route");
1129 assert_eq!(
1130 refreshed_route.candidate.endpoint().base_url,
1131 second_base_url,
1132 "the synthetic continuation must resolve the latest authoritative endpoint"
1133 );
1134 let run_task = tokio::spawn(engine.run());
1135
1136 let mut lifecycle_starts = 0;
1137 let mut dispatches = 0;
1138 let mut completes = 0;
1139 let mut awaiting_second_sync = false;
1140 let mut verified_synthetic_goal = false;
1141 while completes < 2 {
1142 let event = tokio::time::timeout(Duration::from_secs(3), async {
1143 handle.rx_event.write().await.recv().await
1144 })
1145 .await
1146 .expect("goal engine event timeout")
1147 .expect("goal engine event");
1148 match event {
1149 Event::TurnStarted { route, .. } => {
1150 assert!(
1151 route.as_ref().is_none_or(|route| route.billing.is_none()),
1152 "lifecycle start must not carry billing time"
1153 );
1154 lifecycle_starts += 1;
1155 if lifecycle_starts == 2 {
1156 awaiting_second_sync = true;
1157 }
1158 }
1159 Event::RouteDispatched { route, .. } => {
1160 dispatches += 1;
1161 assert_eq!(route.provider_identity, "custom-a");
1162 let expected_base_url = if dispatches == 1 {
1163 first_base_url.as_str()
1164 } else {
1165 second_base_url.as_str()
1166 };
1167 assert_eq!(
1168 route
1169 .billing
1170 .as_ref()
1171 .and_then(|billing| billing.endpoint_fingerprint.clone()),
1172 crate::cost_status::endpoint_fingerprint(expected_base_url),
1173 "goal continuation dispatch borrowed the wrong authoritative route"
1174 );
1175 }
1176 Event::SessionUpdated {
1177 messages,
1178 system_prompt,
1179 ..
1180 } if awaiting_second_sync => {
1181 awaiting_second_sync = false;
1182 let snapshot = goal_state.lock().expect("goal lock").snapshot();
1183 assert_eq!(snapshot.objective.as_deref(), Some("keep going"));
1184 assert_eq!(snapshot.token_budget, Some(50_000));
1185 // The first turn records one bounded intra-turn pass, then the
1186 // synthetic boundary records the second pass before dispatch.
1187 assert_eq!(snapshot.continuation_count, 2);
1188 assert!(snapshot.is_active(), "synthetic turn must retain the goal");
1189
1190 let continuation = messages
1191 .last()
1192 .expect("synthetic continuation message")
1193 .content
1194 .iter()
1195 .find_map(|block| match block {
1196 ContentBlock::Text { text, .. }
1197 if text.contains("## Active Goal State") =>
1198 {
1199 Some(text.as_str())
1200 }
1201 _ => None,
1202 })
1203 .expect("durable goal state in synthetic message");
1204 assert!(continuation.contains("\"objective\": \"keep going\""));
1205 assert!(continuation.contains("\"token_budget\": 50000"));
1206 assert!(continuation.contains("\"continuation_count\": 2"));
1207 assert!(continuation.contains("Continuation pass #2."));
1208
1209 let system_prompt = match system_prompt.expect("synthetic system prompt") {
1210 SystemPrompt::Text(text) => text,
1211 SystemPrompt::Blocks(blocks) => blocks
1212 .into_iter()
1213 .map(|block| block.text)
1214 .collect::<Vec<_>>()
1215 .join("\n"),
1216 };
1217 assert!(system_prompt.contains("<session_goal>"));
1218 assert!(system_prompt.contains("keep going"));
1219 verified_synthetic_goal = true;
1220
1221 tokio::time::timeout(
1222 model_turn_event_timeout(),
1223 second_request_entered.notified(),
1224 )
1225 .await
1226 .expect("second goal model request was never entered");
1227 handle
1228 .send(Op::SetGoalStatus {
1229 goal_id: None,
1230 status: crate::tools::goal::GoalStatus::Paused,
1231 clear: false,
1232 })
1233 .await
1234 .expect("queue goal pause");
1235 // The model future cannot finish until the pause operation is
1236 // already in the engine mailbox, making the queue-order proof
1237 // deterministic under arbitrarily loaded CI runners.
1238 release_second_request.notify_one();
1239 }
1240 Event::TurnComplete { base_url, .. } => {
1241 completes += 1;
1242 assert!(
1243 base_url.is_none(),
1244 "an injected provider-neutral transport must not claim the auxiliary route's endpoint"
1245 );
1246 }
1247 _ => {}
1248 }
1249 }
1250 assert_eq!(lifecycle_starts, 2);
1251 assert_eq!(dispatches, 2);
1252 assert!(verified_synthetic_goal);
1253 let requests = model.captured_requests();
1254 assert_eq!(requests.len(), 2);
1255 let first_intra_turn_prompt = requests[1]
1256 .messages
1257 .iter()
1258 .flat_map(|message| &message.content)
1259 .find_map(|block| match block {
1260 ContentBlock::Text { text, .. } if text.contains("Continuation pass #1.") => {
1261 Some(text.as_str())
1262 }
1263 _ => None,
1264 })
1265 .expect("first intra-turn goal snapshot must survive into the next request");
1266 assert!(
1267 first_intra_turn_prompt.contains("\"tokens_used\": 5"),
1268 "current-turn usage must be rendered without waiting for durable recording: {first_intra_turn_prompt}"
1269 );
1270
1271 // The pause was queued while the second turn was still running. Wait for
1272 // that control operation, then put a snapshot receipt behind the already
1273 // queued continuation. Receiving the receipt proves the continuation was
1274 // consumed; no third TurnStarted may have been emitted.
1275 let mut saw_paused_prompt = false;
1276 let mut saw_paused_goal = false;
1277 loop {
1278 let event = tokio::time::timeout(Duration::from_secs(3), async {
1279 handle.rx_event.write().await.recv().await
1280 })
1281 .await
1282 .expect("goal pause event timeout")
1283 .expect("goal pause event");
1284 match event {
1285 Event::SessionUpdated {
1286 system_prompt: Some(system_prompt),
1287 ..
1288 } => {
1289 let prompt = match system_prompt {
1290 SystemPrompt::Text(text) => text,
1291 SystemPrompt::Blocks(blocks) => blocks
1292 .into_iter()
1293 .map(|block| block.text)
1294 .collect::<Vec<_>>()
1295 .join("\n"),
1296 };
1297 if !prompt.contains("<session_goal>") {
1298 saw_paused_prompt = true;
1299 }
1300 }
1301 Event::GoalUpdated { snapshot } if snapshot.status == "paused" => {
1302 assert_eq!(snapshot.objective.as_deref(), Some("keep going"));
1303 saw_paused_goal = true;
1304 }
1305 Event::Status { ref message } if message == "Goal paused." => {
1306 assert!(
1307 saw_paused_prompt,
1308 "pause status must follow the persisted prompt refresh"
1309 );
1310 assert!(
1311 saw_paused_goal,
1312 "pause status must follow the visible goal snapshot"
1313 );
1314 break;
1315 }
1316 Event::TurnStarted { .. } => {
1317 panic!("queued pause must prevent an additional goal turn")
1318 }
1319 _ => {}
1320 }
1321 }
1322
1323 let (snapshot_tx, snapshot_rx) = tokio::sync::oneshot::channel();
1324 handle
1325 .send(Op::GetSessionSnapshot {
1326 tx: std::sync::Arc::new(std::sync::Mutex::new(Some(snapshot_tx))),
1327 })
1328 .await
1329 .expect("queue post-continuation receipt");
1330 tokio::time::timeout(Duration::from_secs(3), snapshot_rx)
1331 .await
1332 .expect("post-continuation receipt timeout")
1333 .expect("post-continuation receipt");
1334 {
1335 let mut events = handle.rx_event.write().await;
1336 while let Ok(event) = events.try_recv() {
1337 assert!(
1338 !matches!(event, Event::TurnStarted { .. }),
1339 "paused goal continuation started a stale turn"
1340 );
1341 }
1342 }
1343
1344 handle.send(Op::Shutdown).await.expect("queue shutdown");
1345 run_task.await.expect("engine task");
1346 }
1347
1348 #[tokio::test]
1349 async fn saturated_mailbox_does_not_deadlock_goal_continuation_self_dispatch() {
1350 let request_entered = std::sync::Arc::new(tokio::sync::Notify::new());
1351 let release_request = std::sync::Arc::new(tokio::sync::Notify::new());
1352 let model = std::sync::Arc::new(FirstRequestGatedGoalModelClient {
1353 calls: std::sync::atomic::AtomicUsize::new(0),
1354 request_entered: std::sync::Arc::clone(&request_entered),
1355 release_request: std::sync::Arc::clone(&release_request),
1356 });
1357 let config = goal_custom_route_config();
1358 let engine_config = EngineConfig {
1359 model: "local-model".to_string(),
1360 max_steps: 1,
1361 snapshots_enabled: false,
1362 terminal_chrome_enabled: false,
1363 goal_objective: Some("survive a saturated mailbox".to_string()),
1364 ..EngineConfig::default()
1365 };
1366 let client: crate::core::model_client::SharedModelClient = model.clone();
1367 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1368 let goal_state = engine.config.goal_state.clone();
1369 let run_task = tokio::spawn(engine.run());
1370
1371 handle
1372 .send(Op::SendMessage(TurnSpec {
1373 max_output_tokens: None,
1374 content: "start the saturated goal turn".to_string(),
1375 images: Vec::new(),
1376 mode: AppMode::Agent,
1377 route: resolved_route_for_test(&config, "local-model"),
1378 compaction: Box::new(CompactionConfig::default()),
1379 initial_routed_usage: Box::default(),
1380 goal_objective: Some("survive a saturated mailbox".to_string()),
1381 goal_token_budget: None,
1382 goal_status: crate::tools::goal::GoalStatus::Active,
1383 reasoning_effort: None,
1384 reasoning_effort_auto: false,
1385 auto_model: false,
1386 allow_shell: false,
1387 trust_mode: false,
1388 auto_approve: false,
1389 approval_mode: ApprovalMode::Suggest,
1390 translation_enabled: false,
1391 allowed_tools: None,
1392 dynamic_tools: Vec::new(),
1393 hook_executor: None,
1394 verbosity: None,
1395 provenance: UserInputProvenance::ExternalUser,
1396 submission_id: None,
1397 }))
1398 .await
1399 .expect("send saturated goal turn");
1400 tokio::time::timeout(model_turn_event_timeout(), request_entered.notified())
1401 .await
1402 .expect("first goal request was never entered");
1403
1404 // The engine has consumed the SendMessage and is gated inside the model
1405 // request, so every slot below belongs to a queued control operation. The
1406 // final pause must remain ahead of the synthetic continuation.
1407 for index in 0..ENGINE_OP_CHANNEL_CAPACITY {
1408 let status = if index + 1 == ENGINE_OP_CHANNEL_CAPACITY {
1409 crate::tools::goal::GoalStatus::Paused
1410 } else {
1411 crate::tools::goal::GoalStatus::Active
1412 };
1413 handle
1414 .tx_op
1415 .try_send(Op::SetGoalStatus {
1416 goal_id: None,
1417 status,
1418 clear: false,
1419 })
1420 .unwrap_or_else(|error| panic!("fill op mailbox slot {index}: {error}"));
1421 }
1422 assert_eq!(handle.tx_op.capacity(), 0, "fixture must saturate mailbox");
1423
1424 release_request.notify_one();
1425 let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
1426 .await
1427 .expect("saturated mailbox deadlocked the engine")
1428 .expect("post-saturation session snapshot");
1429
1430 let prompt = match session.system_prompt.expect("paused system prompt") {
1431 SystemPrompt::Text(text) => text,
1432 SystemPrompt::Blocks(blocks) => blocks
1433 .into_iter()
1434 .map(|block| block.text)
1435 .collect::<Vec<_>>()
1436 .join("\n"),
1437 };
1438 assert!(!prompt.contains("<session_goal>"), "{prompt}");
1439 let goal = goal_state.lock().expect("goal lock").snapshot();
1440 assert_eq!(goal.status, "paused");
1441 assert_eq!(
1442 model.calls.load(std::sync::atomic::Ordering::SeqCst),
1443 1,
1444 "the queued pause must suppress the stale continuation"
1445 );
1446
1447 let mut starts = 0;
1448 {
1449 let mut events = handle.rx_event.write().await;
1450 while let Ok(event) = events.try_recv() {
1451 if matches!(event, Event::TurnStarted { .. }) {
1452 starts += 1;
1453 }
1454 }
1455 }
1456 assert_eq!(starts, 1, "only the original goal turn may start");
1457
1458 handle.send(Op::Shutdown).await.expect("shutdown engine");
1459 tokio::time::timeout(model_turn_event_timeout(), run_task)
1460 .await
1461 .expect("engine did not shut down after mailbox saturation")
1462 .expect("engine task");
1463 }
1464
1465 #[tokio::test]
1466 async fn queued_ordinary_turn_does_not_multiply_engine_goal_continuations() {
1467 let first_entered = std::sync::Arc::new(tokio::sync::Notify::new());
1468 let release_first = std::sync::Arc::new(tokio::sync::Notify::new());
1469 let third_entered = std::sync::Arc::new(tokio::sync::Notify::new());
1470 let release_third = std::sync::Arc::new(tokio::sync::Notify::new());
1471 let model = std::sync::Arc::new(IndexedGatedGoalModelClient {
1472 calls: std::sync::atomic::AtomicUsize::new(0),
1473 gates: HashMap::from([
1474 (
1475 1,
1476 (
1477 std::sync::Arc::clone(&first_entered),
1478 std::sync::Arc::clone(&release_first),
1479 ),
1480 ),
1481 (
1482 3,
1483 (
1484 std::sync::Arc::clone(&third_entered),
1485 std::sync::Arc::clone(&release_third),
1486 ),
1487 ),
1488 ]),
1489 max_calls: 3,
1490 });
1491 let config = goal_custom_route_config();
1492 let engine_config = EngineConfig {
1493 model: "local-model".to_string(),
1494 max_steps: 1,
1495 snapshots_enabled: false,
1496 terminal_chrome_enabled: false,
1497 goal_objective: Some("coalesce queued goal turns".to_string()),
1498 ..EngineConfig::default()
1499 };
1500 let client: crate::core::model_client::SharedModelClient = model.clone();
1501 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1502 let goal_state = engine.config.goal_state.clone();
1503 let run_task = tokio::spawn(engine.run());
1504 let send_message = |content: &str| {
1505 Op::SendMessage(TurnSpec {
1506 max_output_tokens: None,
1507 content: content.to_string(),
1508 images: Vec::new(),
1509 mode: AppMode::Agent,
1510 route: resolved_route_for_test(&config, "local-model"),
1511 compaction: Box::new(CompactionConfig::default()),
1512 initial_routed_usage: Box::default(),
1513 goal_objective: Some("coalesce queued goal turns".to_string()),
1514 goal_token_budget: None,
1515 goal_status: crate::tools::goal::GoalStatus::Active,
1516 reasoning_effort: None,
1517 reasoning_effort_auto: false,
1518 auto_model: false,
1519 allow_shell: false,
1520 trust_mode: false,
1521 auto_approve: false,
1522 approval_mode: ApprovalMode::Suggest,
1523 translation_enabled: false,
1524 allowed_tools: None,
1525 dynamic_tools: Vec::new(),
1526 hook_executor: None,
1527 verbosity: None,
1528 provenance: UserInputProvenance::ExternalUser,
1529 submission_id: None,
1530 })
1531 };
1532
1533 handle
1534 .send(send_message("start the goal turn"))
1535 .await
1536 .expect("send first goal turn");
1537 tokio::time::timeout(model_turn_event_timeout(), first_entered.notified())
1538 .await
1539 .expect("first goal request was never entered");
1540
1541 // This ordinary user turn is already ahead of the first synthetic token
1542 // when the gated turn completes. It may refresh that token's tools, but it
1543 // must not create a second autonomous continuation.
1544 handle
1545 .send(send_message("queued ordinary follow-up"))
1546 .await
1547 .expect("queue ordinary follow-up");
1548 release_first.notify_one();
1549
1550 tokio::time::timeout(model_turn_event_timeout(), third_entered.notified())
1551 .await
1552 .expect("coalesced synthetic continuation was never entered");
1553 handle
1554 .send(Op::SetGoalStatus {
1555 goal_id: None,
1556 status: crate::tools::goal::GoalStatus::Paused,
1557 clear: false,
1558 })
1559 .await
1560 .expect("queue goal pause behind synthetic turn");
1561 release_third.notify_one();
1562
1563 let _session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
1564 .await
1565 .expect("queued-turn coalescing did not settle")
1566 .expect("post-coalescing session snapshot");
1567 assert_eq!(
1568 model.calls.load(std::sync::atomic::Ordering::SeqCst),
1569 3,
1570 "one initial turn, one queued user turn, and one synthetic continuation are expected"
1571 );
1572 assert_eq!(
1573 goal_state.lock().expect("goal lock").snapshot().status,
1574 "paused"
1575 );
1576
1577 handle.send(Op::Shutdown).await.expect("shutdown engine");
1578 tokio::time::timeout(model_turn_event_timeout(), run_task)
1579 .await
1580 .expect("engine did not shut down after queued-turn coalescing")
1581 .expect("engine task");
1582 }
1583
1583 lines RUST