返回 CodeWhale
test_cases_01.rs
根目录 / crates / tui / src / client / test_cases_01.rs
1
2 use super::*;
3 use crate::client::chat::{
4 build_chat_messages, build_chat_messages_for_request,
5 build_chat_messages_for_request_and_provider, count_reasoning_replay_chars,
6 parse_chat_message, parse_sse_chunk, sanitize_thinking_mode_messages, tool_to_chat,
7 tool_to_chat_for_base_url,
8 };
9 use crate::client::responses::build_responses_body;
10 use crate::config::{
11 DEFAULT_CONCENTRATE_BASE_URL, DEFAULT_CONCENTRATE_MODEL, DEFAULT_EDENAI_MODEL,
12 DEFAULT_TELECOMJS_MODEL, OPENROUTER_QWEN_3_6_FLASH_MODEL, ProviderConfig, ProvidersConfig,
13 };
14 use crate::tools::apply_patch::ApplyPatchTool;
15 use crate::tools::spec::ToolSpec;
16 use crate::tools::{ToolContext, ToolRegistryBuilder};
17 use codewhale_models::{
18 ContentBlock, ContentBlockStart, Delta, Message, MessageRequest, MessageResponse,
19 StreamEvent, Tool,
20 };
21 use codewhale_protocol::runtime::DynamicToolSpec;
22 use serde_json::json;
23 use wiremock::matchers::{header, method, path};
24 use wiremock::{Mock, MockServer, ResponseTemplate};
25
26 /// OpenRouter app attribution: the headers that put an app on
27 /// openrouter.ai's rankings are sent for OpenRouter routes only, and a
28 /// user-configured header of the same name wins.
29 #[test]
30 fn openrouter_routes_carry_app_attribution_headers() {
31 let headers = build_default_headers(
32 "sk-or-key",
33 &HashMap::new(),
34 ProviderKind::Openrouter,
35 "https://openrouter.ai/api/v1",
36 WireFormat::ChatCompletions,
37 false,
38 )
39 .expect("headers");
40 assert_eq!(
41 headers.get("http-referer").and_then(|v| v.to_str().ok()),
42 Some("https://codewhale.net")
43 );
44 assert_eq!(
45 headers.get("x-title").and_then(|v| v.to_str().ok()),
46 Some("Codewhale")
47 );
48 // The current display-name header, alongside the legacy one.
49 assert_eq!(
50 headers
51 .get("x-openrouter-title")
52 .and_then(|v| v.to_str().ok()),
53 Some("Codewhale")
54 );
55 // Marketplace categories must match OpenRouter's spelling exactly:
56 // unrecognised values are dropped silently, so a typo here is
57 // invisible in production and only shows up as a missing listing.
58 assert_eq!(
59 headers
60 .get("x-openrouter-categories")
61 .and_then(|v| v.to_str().ok()),
62 Some("cli-agent,personal-agent")
63 );
64
65 // Attribution identifies this app to OpenRouter's rankings and has
66 // no business on any other route, whichever provider that is.
67 for provider in [
68 ProviderKind::Deepseek,
69 ProviderKind::Moonshot,
70 ProviderKind::Zai,
71 ProviderKind::Openai,
72 ProviderKind::Custom,
73 ] {
74 let other = build_default_headers(
75 "sk-key",
76 &HashMap::new(),
77 provider,
78 "https://example.invalid/v1",
79 WireFormat::ChatCompletions,
80 false,
81 )
82 .expect("headers");
83 assert!(
84 other.get("http-referer").is_none()
85 && other.get("x-title").is_none()
86 && other.get("x-openrouter-title").is_none(),
87 "{provider:?} must not receive OpenRouter attribution headers"
88 );
89 }
90
91 let overridden = build_default_headers(
92 "sk-or-key",
93 &HashMap::from([("X-Title".to_string(), "My Fork".to_string())]),
94 ProviderKind::Openrouter,
95 "https://openrouter.ai/api/v1",
96 WireFormat::ChatCompletions,
97 false,
98 )
99 .expect("headers");
100 // A user title on the legacy header alone must follow onto the
101 // current one, or OpenRouter would still show "Codewhale".
102 for name in ["x-title", "x-openrouter-title"] {
103 assert_eq!(
104 overridden.get(name).and_then(|v| v.to_str().ok()),
105 Some("My Fork"),
106 "{name} must carry the user's title"
107 );
108 }
109
110 // And the other way round: setting only the current header also
111 // renames the legacy one, so the two never disagree.
112 let current_only = build_default_headers(
113 "sk-or-key",
114 &HashMap::from([("X-OpenRouter-Title".to_string(), "My Fork".to_string())]),
115 ProviderKind::Openrouter,
116 "https://openrouter.ai/api/v1",
117 WireFormat::ChatCompletions,
118 false,
119 )
120 .expect("headers");
121 for name in ["x-title", "x-openrouter-title"] {
122 assert_eq!(
123 current_only.get(name).and_then(|v| v.to_str().ok()),
124 Some("My Fork"),
125 "{name} must carry the user's title"
126 );
127 }
128 }
129
130 #[test]
131 fn openrouter_pricing_maps_cache_write_per_token_to_per_million() {
132 let payload = r#"{"data":[{
133 "id":"anthropic/claude-sonnet-4-6",
134 "pricing":{
135 "prompt":"0.000003",
136 "completion":"0.000015",
137 "input_cache_read":"0.0000003",
138 "input_cache_write":"0.00000375"
139 }
140 },{
141 "id":"some/no-write-row",
142 "pricing":{"prompt":"0.000001","completion":"0.000002"}
143 }]}"#;
144
145 let items = parse_openrouter_models_response(payload).expect("parses");
146 let priced = openrouter_to_catalog_offering(&items[0], "openrouter", "fp", 42)
147 .expect("valid priced row");
148 let cost = priced.cost.as_ref().expect("pricing row");
149 assert_eq!(cost.input, Some(3.0));
150 assert_eq!(cost.output, Some(15.0));
151 assert_eq!(cost.cache_read, Some(0.3));
152 assert_eq!(cost.cache_write, Some(3.75));
153
154 // A cache-write premium must actually reach the estimator: the same
155 // tokens cost more when they are cache-creation rather than cache-read.
156 let pricing = codewhale_config::pricing::OfferingPricing::from_catalog_offering(&priced)
157 .expect("priced offering");
158 let write = codewhale_config::pricing::TokenUsage {
159 cache_write: 1_000_000,
160 ..Default::default()
161 };
162 assert_eq!(pricing.estimate_cost(&write), Some(3.75));
163 assert!(pricing.unpriced_used_classes(&write).is_empty());
164
165 // A row without a published write rate stays unknown, not zero, and
166 // fails closed for cache-creation turns.
167 let unwritten = openrouter_to_catalog_offering(&items[1], "openrouter", "fp", 42)
168 .expect("valid row without cache-write rate");
169 assert_eq!(
170 unwritten.cost.as_ref().and_then(|cost| cost.cache_write),
171 None
172 );
173 let unwritten =
174 codewhale_config::pricing::OfferingPricing::from_catalog_offering(&unwritten)
175 .expect("priced offering");
176 assert_eq!(unwritten.estimate_cost(&write), None);
177 assert_eq!(
178 unwritten.unpriced_used_classes(&write),
179 vec![codewhale_config::pricing::TokenClass::CacheWrite]
180 );
181 }
182
183 #[test]
184 fn baseten_catalog_maps_provider_stated_prices_limits_and_features() {
185 // Exact current Baseten shape: pricing is captured from the official
186 // baseten-switch repository; Model APIs publishes context_length,
187 // max_completion_tokens, and supported_features including `vision`.
188 let payload = r#"{"data":[{
189 "id":"deepseek-ai/DeepSeek-V4-Pro",
190 "context_length":"1048576",
191 "max_completion_tokens":262144,
192 "pricing":{
193 "prompt":0.0000014,
194 "completion":"0.0000044",
195 "input_cache_read":0.00000014
196 },
197 "supported_features":["reasoning","vision"],
198 "reasoning_options":[{"type":"toggle"}]
199 }]}"#;
200
201 let items = parse_baseten_models_response(payload).expect("Baseten catalog");
202 let offering =
203 baseten_to_catalog_offering(&items[0], "baseten", "baseten-fp", 42).expect("row");
204 assert_eq!(offering.provider, "baseten");
205 assert_eq!(
206 offering.wire_model_id,
207 codewhale_config::catalog::BASETEN_DEFAULT_MODEL
208 );
209 assert!(offering.default_for_provider);
210 let limit = offering.limit.expect("published limits");
211 assert_eq!(limit.context, Some(1_048_576));
212 assert_eq!(limit.input, Some(1_048_576));
213 assert_eq!(limit.output, Some(262_144));
214 let cost = offering.cost.expect("published pricing");
215 assert_eq!(cost.input, Some(1.4));
216 assert_eq!(cost.output, Some(4.4));
217 assert_eq!(cost.cache_read, Some(0.14));
218 assert_eq!(cost.cache_write, None);
219 assert_eq!(offering.reasoning, Some(true));
220 assert_eq!(offering.tool_call, Some(true));
221 assert_eq!(offering.structured_output, Some(true));
222 assert_eq!(offering.attachment, Some(true));
223 let modalities = offering.modalities.expect("vision feature modalities");
224 assert_eq!(modalities.input, vec!["text", "image"]);
225 assert_eq!(modalities.output, vec!["text"]);
226 assert_eq!(offering.reasoning_options, vec![json!({"type":"toggle"})]);
227 assert!(matches!(offering.source, CatalogSource::Live { .. }));
228 }
229
230 #[test]
231 fn baseten_catalog_rejects_duplicate_ids_and_invalid_known_numbers() {
232 let duplicates = r#"{"data":[{"id":"same/model"},{"id":"same/model"}]}"#;
233 assert_eq!(
234 parse_baseten_models_response(duplicates).unwrap_err(),
235 CatalogRefreshError::InvalidResponse
236 );
237
238 let negative = r#"{"data":[{
239 "id":"synthetic/model",
240 "pricing":{"prompt":-0.000001,"completion":0.000002}
241 }]}"#;
242 let items = parse_baseten_models_response(negative).expect("shape parses");
243 assert_eq!(
244 baseten_to_catalog_offering(&items[0], "baseten", "fp", 1).unwrap_err(),
245 CatalogRefreshError::InvalidResponse
246 );
247 }
248
249 #[test]
250 fn provider_live_price_parsers_reject_present_bad_rates_but_keep_zero_and_omission() {
251 for invalid in ["not-a-number", "-0.1", "NaN", "inf", "1e308", "0.100001"] {
252 let openrouter = json!({
253 "data": [{
254 "id": "synthetic/openrouter-invalid-price",
255 "pricing": { "prompt": invalid, "completion": "0.000001" }
256 }]
257 })
258 .to_string();
259 let items = parse_openrouter_models_response(&openrouter).expect("OpenRouter shape");
260 let offering = openrouter_to_catalog_offering(&items[0], "openrouter", "fp", 1);
261 if invalid.starts_with('-') {
262 // OpenRouter's negative sentinel means "no fixed rate" (#6690).
263 assert_eq!(offering.expect("variable-price row").cost, None);
264 } else {
265 assert_eq!(
266 offering.unwrap_err(),
267 CatalogRefreshError::InvalidResponse,
268 "OpenRouter must reject {invalid:?}"
269 );
270 }
271
272 let baseten = json!({
273 "data": [{
274 "id": "synthetic/baseten-invalid-price",
275 "pricing": { "prompt": invalid, "completion": "0.000001" }
276 }]
277 })
278 .to_string();
279 let items = parse_baseten_models_response(&baseten).expect("Baseten shape");
280 assert_eq!(
281 baseten_to_catalog_offering(&items[0], "baseten", "fp", 1).unwrap_err(),
282 CatalogRefreshError::InvalidResponse,
283 "Baseten must reject {invalid:?}"
284 );
285 }
286
287 let openrouter = parse_openrouter_models_response(
288 r#"{"data":[{"id":"synthetic/openrouter-free","pricing":{"prompt":"0","completion":"0"}}]}"#,
289 )
290 .expect("OpenRouter zero row");
291 let openrouter = openrouter_to_catalog_offering(&openrouter[0], "openrouter", "fp", 1)
292 .expect("explicit zero is a valid published price");
293 let cost = openrouter.cost.expect("published zero cost");
294 assert_eq!(cost.input, Some(0.0));
295 assert_eq!(cost.output, Some(0.0));
296 assert_eq!(cost.cache_read, None);
297 assert_eq!(cost.cache_write, None);
298
299 let baseten = parse_baseten_models_response(
300 r#"{"data":[{"id":"synthetic/baseten-free","pricing":{"prompt":"0","completion":0}}]}"#,
301 )
302 .expect("Baseten zero row");
303 let baseten = baseten_to_catalog_offering(&baseten[0], "baseten", "fp", 1)
304 .expect("explicit zero is a valid published price");
305 let cost = baseten.cost.expect("published zero cost");
306 assert_eq!(cost.input, Some(0.0));
307 assert_eq!(cost.output, Some(0.0));
308 assert_eq!(cost.cache_read, None);
309 assert_eq!(cost.cache_write, None);
310 }
311
312 #[test]
313 fn baseten_feature_names_require_exact_normalized_aliases() {
314 let payload = r#"{"data":[{
315 "id":"synthetic/text-only",
316 "supported_features":["revision","pre_reasoning_filter"]
317 }]}"#;
318 let items = parse_baseten_models_response(payload).expect("Baseten catalog");
319 let offering =
320 baseten_to_catalog_offering(&items[0], "baseten", "fp", 1).expect("valid row");
321 assert_eq!(offering.reasoning, Some(false));
322 assert_eq!(offering.modalities, None);
323 assert_eq!(offering.attachment, None);
324 assert_eq!(offering.tool_call, Some(true));
325 assert_eq!(offering.structured_output, Some(true));
326 }
327
328 fn test_tool(name: &str) -> Tool {
329 Tool {
330 tool_type: None,
331 name: name.to_string(),
332 description: format!("{name} test tool"),
333 input_schema: json!({
334 "type": "object",
335 "properties": {},
336 }),
337 allowed_callers: None,
338 defer_loading: Some(false),
339 input_examples: None,
340 strict: Some(true),
341 cache_control: None,
342 }
343 }
344
345 fn apply_patch_request_tool() -> Tool {
346 let spec = ApplyPatchTool;
347 Tool {
348 tool_type: None,
349 name: spec.name().to_string(),
350 description: spec.description().to_string(),
351 input_schema: spec.input_schema(),
352 allowed_callers: None,
353 defer_loading: Some(false),
354 input_examples: None,
355 strict: None,
356 cache_control: None,
357 }
358 }
359
360 fn deferred_dynamic_request_tool() -> Tool {
361 let registry = ToolRegistryBuilder::new()
362 .with_dynamic_tools(&[DynamicToolSpec {
363 namespace: Some("capture".to_string()),
364 name: "deferred_lookup".to_string(),
365 description: "Look up a record after deferred loading".to_string(),
366 input_schema: json!({
367 "type": "object",
368 "properties": {
369 "mode": {"type": "string", "const": "fast"},
370 "query": {
371 "anyOf": [
372 {"type": "string"},
373 {"type": "null"}
374 ]
375 }
376 },
377 "required": ["mode"]
378 }),
379 defer_loading: true,
380 }])
381 .build(ToolContext::new(
382 std::env::temp_dir().join("codewhale-k3-deferred-capture"),
383 ));
384 registry
385 .to_api_tools()
386 .into_iter()
387 .find(|tool| tool.name == "deferred_lookup")
388 .expect("dynamic tool remains model-visible")
389 }
390
391 fn value_contains_key(value: &Value, needle: &str) -> bool {
392 match value {
393 Value::Object(object) => {
394 object.contains_key(needle)
395 || object
396 .values()
397 .any(|child| value_contains_key(child, needle))
398 }
399 Value::Array(values) => values.iter().any(|child| value_contains_key(child, needle)),
400 _ => false,
401 }
402 }
403
404 fn captured_function<'a>(body: &'a Value, name: &str) -> &'a Value {
405 body["tools"]
406 .as_array()
407 .and_then(|tools| tools.iter().find(|tool| tool["function"]["name"] == name))
408 .map(|tool| &tool["function"])
409 .unwrap_or_else(|| panic!("captured tool catalog is missing {name}: {body}"))
410 }
411
412 fn moonshot_request_boundary_client(
413 route_base_url: &str,
414 model: &str,
415 transport_base_url: String,
416 ) -> CodewhaleClient {
417 let mut client = CodewhaleClient::new(&Config {
418 provider: Some("moonshot".to_string()),
419 providers: Some(ProvidersConfig {
420 moonshot: ProviderConfig {
421 api_key: Some("moonshot-request-boundary-key".to_string()),
422 base_url: Some(route_base_url.to_string()),
423 model: Some(model.to_string()),
424 ..ProviderConfig::default()
425 },
426 ..ProvidersConfig::default()
427 }),
428 ..Config::default()
429 })
430 .expect("Moonshot request-boundary client");
431 assert_eq!(client.base_url, route_base_url);
432 client.test_chat_transport_base_url = Some(transport_base_url);
433 client
434 }
435
436 fn zai_request_boundary_client(
437 route_base_url: &str,
438 model: &str,
439 transport_base_url: String,
440 ) -> CodewhaleClient {
441 let _ = rustls::crypto::ring::default_provider().install_default();
442 let mut client = CodewhaleClient::new(&Config {
443 provider: Some("zai".to_string()),
444 providers: Some(ProvidersConfig {
445 zai: ProviderConfig {
446 api_key: Some("zai-request-boundary-key".to_string()),
447 base_url: Some(route_base_url.to_string()),
448 model: Some(model.to_string()),
449 ..ProviderConfig::default()
450 },
451 ..ProvidersConfig::default()
452 }),
453 ..Config::default()
454 })
455 .expect("Z.ai request-boundary client");
456 assert_eq!(client.base_url, route_base_url);
457 client.test_chat_transport_base_url = Some(transport_base_url);
458 client
459 }
460
461 fn minimax_request_boundary_client(
462 route_base_url: &str,
463 model: &str,
464 transport_base_url: String,
465 ) -> CodewhaleClient {
466 let _ = rustls::crypto::ring::default_provider().install_default();
467 let mut client = CodewhaleClient::new(&Config {
468 provider: Some("minimax".to_string()),
469 providers: Some(ProvidersConfig {
470 minimax: ProviderConfig {
471 api_key: Some("minimax-request-boundary-key".to_string()),
472 base_url: Some(route_base_url.to_string()),
473 model: Some(model.to_string()),
474 ..ProviderConfig::default()
475 },
476 ..ProvidersConfig::default()
477 }),
478 ..Config::default()
479 })
480 .expect("MiniMax request-boundary client");
481 assert_eq!(client.base_url, route_base_url);
482 client.test_chat_transport_base_url = Some(transport_base_url);
483 client
484 }
485
486 fn deepseek_request_boundary_client(
487 route_base_url: &str,
488 transport_base_url: String,
489 ) -> CodewhaleClient {
490 let mut client = CodewhaleClient::new(
491 &Config {
492 provider: Some("deepseek".to_string()),
493 default_text_model: Some("deepseek-v4-pro".to_string()),
494 ..Config::default()
495 }
496 .with_legacy_root(
497 Some("deepseek-request-boundary-key".to_string()),
498 Some(route_base_url.to_string()),
499 ),
500 )
501 .expect("DeepSeek request-boundary client");
502 client.test_chat_transport_base_url = Some(transport_base_url);
503 client
504 }
505
506 fn ollama_cloud_request_boundary_client(transport_base_url: String) -> CodewhaleClient {
507 let mut client = CodewhaleClient::new(&Config {
508 provider: Some("ollama-cloud".to_string()),
509 providers: Some(ProvidersConfig {
510 ollama_cloud: ProviderConfig {
511 api_key: Some("ollama-cloud-request-boundary-key".to_string()),
512 base_url: Some(crate::config::DEFAULT_OLLAMA_CLOUD_BASE_URL.to_string()),
513 model: Some("gpt-oss:120b".to_string()),
514 ..ProviderConfig::default()
515 },
516 ..ProvidersConfig::default()
517 }),
518 ..Config::default()
519 })
520 .expect("Ollama Cloud request-boundary client");
521 assert_eq!(
522 client.base_url,
523 crate::config::DEFAULT_OLLAMA_CLOUD_BASE_URL
524 );
525 client.test_chat_transport_base_url = Some(transport_base_url);
526 client
527 }
528
529 /// Swap in a short non-streaming envelope for the test's duration,
530 /// restoring the previous value on drop.
531 struct NonStreamingEnvelopeGuard(u64);
532
533 impl NonStreamingEnvelopeGuard {
534 fn millis(ms: u64) -> Self {
535 Self(TEST_NON_STREAMING_ENVELOPE_MS.swap(ms, std::sync::atomic::Ordering::SeqCst))
536 }
537 }
538
539 impl Drop for NonStreamingEnvelopeGuard {
540 fn drop(&mut self) {
541 TEST_NON_STREAMING_ENVELOPE_MS.store(self.0, std::sync::atomic::Ordering::SeqCst);
542 }
543 }
544
545 /// The provider accepts the connection but stalls far past the budgeted
546 /// envelope before answering: the non-streaming request must be cut off
547 /// with a timeout instead of wedging the caller (mid-turn compaction,
548 /// translate, provider-native search) indefinitely.
549 #[tokio::test]
550 async fn non_streaming_envelope_bounds_a_stalled_provider() {
551 // The injected budget is process-global; serialize against other tests
552 // (which may issue non-streaming requests with their own timing
553 // assumptions) through the shared test-env lock.
554 let _env_lock = crate::test_support::lock_test_env();
555 let _envelope = NonStreamingEnvelopeGuard::millis(2000);
556 let server = MockServer::start().await;
557 Mock::given(method("POST"))
558 .respond_with(
559 ResponseTemplate::new(200)
560 .set_body_json(json!({
561 "id": "chatcmpl_envelope",
562 "object": "chat.completion",
563 "model": "deepseek-v4-pro",
564 "choices": [{
565 "index": 0,
566 "message": {"role": "assistant", "content": "late"},
567 "finish_reason": "stop"
568 }],
569 "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}
570 }))
571 .set_delay(Duration::from_secs(5)),
572 )
573 .mount(&server)
574 .await;
575
576 let client = deepseek_request_boundary_client(&server.uri(), server.uri());
577 let err = client
578 .create_message(k3_request_fixture("deepseek-v4-pro", Some("off"), false))
579 .await
580 .expect_err("a provider that never answers must hit the envelope");
581 assert!(
582 err.to_string().to_lowercase().contains("timed out"),
583 "envelope timeout must be reported as such; got {err:#}"
584 );
585 }
586
587 /// Every attempt answers 429 with an hour-long Retry-After. Honoring the
588 /// header must not hand the total budget to the server: the retry loop is
589 /// capped by the envelope, so a gateway answering 429 + Retry-After: 3600
590 /// forever fails the turn in bounded time instead of wedging it for hours.
591 #[tokio::test]
592 async fn retry_after_honoring_cannot_extend_the_non_streaming_envelope() {
593 let _env_lock = crate::test_support::lock_test_env();
594 let _envelope = NonStreamingEnvelopeGuard::millis(2000);
595 let server = MockServer::start().await;
596 Mock::given(method("POST"))
597 .respond_with(
598 ResponseTemplate::new(429)
599 .insert_header("retry-after", "3600")
600 .set_body_string("rate limited"),
601 )
602 .mount(&server)
603 .await;
604
605 let client = deepseek_request_boundary_client(&server.uri(), server.uri());
606 let started = std::time::Instant::now();
607 let err = client
608 .create_message(k3_request_fixture("deepseek-v4-pro", Some("off"), false))
609 .await
610 .expect_err("unbounded Retry-After honoring must still hit the envelope");
611 assert!(
612 err.to_string().to_lowercase().contains("timed out"),
613 "the envelope must cut off the Retry-After wait; got {err:#}"
614 );
615 assert!(
616 started.elapsed() < Duration::from_secs(30),
617 "a 3600s Retry-After must not run past the envelope; took {:?}",
618 started.elapsed()
619 );
620 crate::retry_status::clear();
621 crate::retry_status::clear_rate_limit();
622 }
623
624 /// The open answers past the injected non-streaming envelope. The
625 /// streaming-open path must not inherit any total: reqwest's per-request
626 /// timeout wraps the response body, so a total set on the open would ride
627 /// on the returned body and hard-cut a live stream mid-generation.
628 #[tokio::test]
629 async fn stream_open_retry_path_sets_no_total_deadline() {
630 let _env_lock = crate::test_support::lock_test_env();
631 let _envelope = NonStreamingEnvelopeGuard::millis(2000);
632 let server = MockServer::start().await;
633 Mock::given(method("POST"))
634 .respond_with(
635 ResponseTemplate::new(200)
636 .set_body_string("data: [DONE]\n\n")
637 .set_delay(Duration::from_millis(2500)),
638 )
639 .mount(&server)
640 .await;
641
642 let client = deepseek_request_boundary_client(&server.uri(), server.uri());
643 let response = client
644 .send_stream_open_with_retry(|| {
645 client
646 .http_client
647 .post(format!("{}/chat/completions", server.uri()))
648 .header(reqwest::header::CONTENT_TYPE, "application/json")
649 .body("{}".to_string())
650 })
651 .await
652 .expect("stream open must not carry the non-streaming envelope");
653 assert!(response.status().is_success());
654 let text = response.text().await.expect("read stream-open body");
655 assert_eq!(text, "data: [DONE]\n\n");
656 }
657
658 /// A caller that pins its own, larger per-attempt total (`list_models`
659 /// pins 30s) must keep it: `.timeout()` on the builder is a pure
660 /// overwrite, so an unconditional envelope would silently replace the
661 /// pinned budget with the injected 2s and fail this 2.5s-late response.
662 #[tokio::test]
663 async fn pinned_request_total_survives_the_shared_retry_envelope() {
664 let _env_lock = crate::test_support::lock_test_env();
665 let _envelope = NonStreamingEnvelopeGuard::millis(2000);
666 let server = MockServer::start().await;
667 Mock::given(method("GET"))
668 .respond_with(
669 ResponseTemplate::new(200)
670 .set_body_string("{}")
671 .set_delay(Duration::from_millis(2500)),
672 )
673 .mount(&server)
674 .await;
675
676 let client = deepseek_request_boundary_client(&server.uri(), server.uri());
677 let response = client
678 .send_with_retry_total_error_body(
679 Duration::from_secs(5),
680 || client.http_client.get(format!("{}/models", server.uri())),
681 &ErrorBodyDisclosure::Full,
682 )
683 .await
684 .expect("caller-pinned total must not be overwritten by the envelope");
685 assert!(response.status().is_success());
686 }
687
688 /// The per-chunk line cap is backpressure relief, not a data budget. When
689 /// one transport chunk carries more SSE lines than the cap, the drain loop
690 /// stops mid-buffer and the outer loop waits for the *next* chunk before
691 /// draining any more — so whatever is still buffered when the stream ends
692 /// never reaches the decoder. `flush_sse_line` cannot rescue it: it treats
693 /// the whole remainder as one unterminated line. A long stream of small
694 /// deltas (or provider heartbeats) therefore loses its tail — the last
695 /// tokens, `finish_reason`, and usage — silently.
696 #[tokio::test]
697 async fn chat_stream_drains_chunks_carrying_more_lines_than_the_per_chunk_cap() {
698 // Comment lines are counted by the drain loop and are cheap enough
699 // that one transport read holds far more than SSE_MAX_LINES_PER_CHUNK.
700 let heartbeats = ": ping\n".repeat(20 * SSE_MAX_LINES_PER_CHUNK * 4);
701 let body = format!(
702 "{heartbeats}data: {}\n\ndata: [DONE]\n\n",
703 json!({
704 "choices": [{
705 "index": 0,
706 "delta": {"content": "pong"},
707 "finish_reason": "stop"
708 }]
709 })
710 );
711
712 let server = MockServer::start().await;
713 Mock::given(method("POST"))
714 .respond_with(
715 ResponseTemplate::new(200)
716 .insert_header("content-type", "text/event-stream")
717 .set_body_string(body),
718 )
719 .expect(1)
720 .mount(&server)
721 .await;
722
723 let client = deepseek_request_boundary_client("https://api.deepseek.com/v1", server.uri());
724 let request = MessageRequest {
725 model: "deepseek-v4-pro".to_string(),
726 messages: vec![Message {
727 role: Role::User,
728 content: vec![ContentBlock::Text {
729 text: "sse drain regression".to_string(),
730 cache_control: None,
731 }],
732 }],
733 max_tokens: 64,
734 system: None,
735 tools: None,
736 tool_choice: None,
737 metadata: None,
738 thinking: None,
739 reasoning_effort: Some("off".to_string()),
740 stream: Some(true),
741 temperature: None,
742 top_p: None,
743 };
744
745 let mut stream = client
746 .create_message_stream(request)
747 .await
748 .expect("streaming request succeeds");
749 let mut text = String::new();
750 while let Some(event) = stream.next().await {
751 if let StreamEvent::ContentBlockDelta {
752 delta: Delta::TextDelta { text: chunk },
753 ..
754 } = event.expect("heartbeat-padded SSE stays valid")
755 {
756 text.push_str(&chunk);
757 }
758 }
759
760 assert_eq!(
761 text, "pong",
762 "the data frame after the heartbeat flood must still be decoded"
763 );
764 }
765
766 #[tokio::test]
767 async fn chat_stream_eof_without_done_or_finish_reason_is_error() {
768 let server = MockServer::start().await;
769 Mock::given(method("POST"))
770 .respond_with(
771 ResponseTemplate::new(200)
772 .insert_header("content-type", "text/event-stream")
773 .set_body_string(": provider heartbeat\n\n"),
774 )
775 .expect(1)
776 .mount(&server)
777 .await;
778
779 let client = deepseek_request_boundary_client("https://api.deepseek.com/v1", server.uri());
780 let request = MessageRequest {
781 model: "deepseek-v4-pro".to_string(),
782 messages: vec![Message {
783 role: Role::User,
784 content: vec![ContentBlock::Text {
785 text: "premature EOF regression".to_string(),
786 cache_control: None,
787 }],
788 }],
789 max_tokens: 64,
790 system: None,
791 tools: None,
792 tool_choice: None,
793 metadata: None,
794 thinking: None,
795 reasoning_effort: Some("off".to_string()),
796 stream: Some(true),
797 temperature: None,
798 top_p: None,
799 };
800
801 let mut stream = client
802 .create_message_stream(request)
803 .await
804 .expect("HTTP request succeeds before the stream closes");
805 let mut saw_stop = false;
806 let mut failure = None;
807 while let Some(event) = stream.next().await {
808 match event {
809 Ok(StreamEvent::MessageStop) => saw_stop = true,
810 Ok(_) => {}
811 Err(error) => failure = Some(error.to_string()),
812 }
813 }
814
815 assert!(
816 !saw_stop,
817 "premature EOF must not be reported as MessageStop"
818 );
819 assert!(
820 failure
821 .as_deref()
822 .is_some_and(|message| message.contains("before [DONE] or finish_reason")),
823 "premature EOF must remain a typed stream failure: {failure:?}"
824 );
825 }
826
827 #[tokio::test]
828 async fn chat_stream_finish_reason_without_done_is_terminal() {
829 let server = MockServer::start().await;
830 Mock::given(method("POST"))
831 .respond_with(
832 ResponseTemplate::new(200)
833 .insert_header("content-type", "text/event-stream")
834 .set_body_string(concat!(
835 "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"ok\"},\"finish_reason\":null}]}\n\n",
836 "data: {\"choices\":[{\"index\":0,\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n",
837 )),
838 )
839 .expect(1)
840 .mount(&server)
841 .await;
842
843 let client = deepseek_request_boundary_client("https://api.deepseek.com/v1", server.uri());
844 let request = MessageRequest {
845 model: "deepseek-v4-pro".to_string(),
846 messages: vec![Message {
847 role: Role::User,
848 content: vec![ContentBlock::Text {
849 text: "finish reason regression".to_string(),
850 cache_control: None,
851 }],
852 }],
853 max_tokens: 64,
854 system: None,
855 tools: None,
856 tool_choice: None,
857 metadata: None,
858 thinking: None,
859 reasoning_effort: Some("off".to_string()),
860 stream: Some(true),
861 temperature: None,
862 top_p: None,
863 };
864
865 let mut stream = client
866 .create_message_stream(request)
867 .await
868 .expect("streaming request succeeds");
869 let mut saw_stop = false;
870 while let Some(event) = stream.next().await {
871 if matches!(
872 event.expect("terminal stream stays valid"),
873 StreamEvent::MessageStop
874 ) {
875 saw_stop = true;
876 }
877 }
878 assert!(
879 saw_stop,
880 "finish_reason is valid terminal proof without [DONE]"
881 );
882 }
883
884 async fn capture_deepseek_chat_request(
885 route_base_url: &str,
886 strict: bool,
887 streaming: bool,
888 ) -> (String, Value) {
889 let server = MockServer::start().await;
890 let response = if streaming {
891 ResponseTemplate::new(200)
892 .insert_header("content-type", "text/event-stream")
893 .set_body_string("data: [DONE]\n\n")
894 } else {
895 ResponseTemplate::new(200).set_body_json(json!({
896 "id": "chatcmpl-deepseek-request-boundary",
897 "object": "chat.completion",
898 "model": "deepseek-v4-pro",
899 "choices": [{
900 "index": 0,
901 "message": {"role": "assistant", "content": "ok"},
902 "finish_reason": "stop"
903 }],
904 "usage": {
905 "prompt_tokens": 1,
906 "completion_tokens": 1,
907 "total_tokens": 2
908 }
909 }))
910 };
911 Mock::given(method("POST"))
912 .respond_with(response)
913 .expect(1)
914 .mount(&server)
915 .await;
916
917 let mut tool = test_tool("lookup");
918 if !strict {
919 tool.strict = None;
920 }
921 let request = MessageRequest {
922 model: "deepseek-v4-pro".to_string(),
923 messages: vec![Message {
924 role: Role::User,
925 content: vec![ContentBlock::Text {
926 text: "provider-free DeepSeek route fixture".to_string(),
927 cache_control: None,
928 }],
929 }],
930 max_tokens: 64,
931 system: None,
932 tools: Some(vec![tool]),
933 tool_choice: Some(json!(if strict { "required" } else { "auto" })),
934 metadata: None,
935 thinking: None,
936 reasoning_effort: Some("off".to_string()),
937 stream: Some(streaming),
938 temperature: None,
939 top_p: None,
940 };
941 let client = deepseek_request_boundary_client(route_base_url, server.uri());
942
943 if streaming {
944 let mut stream = client
945 .create_message_stream(request)
946 .await
947 .expect("streaming request succeeds");
948 while let Some(event) = stream.next().await {
949 event.expect("captured SSE response remains valid");
950 }
951 } else {
952 client
953 .create_message(request)
954 .await
955 .expect("non-streaming request succeeds");
956 }
957
958 let requests = server.received_requests().await.expect("recorded request");
959 assert_eq!(requests.len(), 1);
960 let path = requests[0].url.path().to_string();
961 let body = serde_json::from_slice(&requests[0].body).expect("captured request JSON");
962 (path, body)
963 }
964
965 /// #6540: the compaction summary is the parent turn plus one trailing
966 /// user instruction. Rendered through the production wire builders, its
967 /// prompt must be byte-for-byte the parent's prompt followed by that one
968 /// item, with the same model, system/instructions, tools and reasoning
969 /// controls — otherwise the provider cache misses the whole history.
970 #[test]
971 fn compaction_summary_request_extends_the_parent_turn_prompt_bytes() {
972 fn text(role: Role, text: &str) -> Message {
973 Message {
974 role,
975 content: vec![ContentBlock::Text {
976 text: text.to_string(),
977 cache_control: None,
978 }],
979 }
980 }
981 fn assert_extends(label: &str, parent: &Value, summary: &Value) {
982 let parent_items = parent.as_array().expect("parent prompt items");
983 let summary_items = summary.as_array().expect("summary prompt items");
984 assert_eq!(summary_items.len(), parent_items.len() + 1, "{label}");
985 let parent_bytes = serde_json::to_string(parent).expect("serialize parent");
986 let summary_bytes = serde_json::to_string(summary).expect("serialize summary");
987 let open_prefix = parent_bytes.strip_suffix(']').expect("JSON array");
988 assert!(
989 summary_bytes.starts_with(open_prefix),
990 "{label}: summary prompt diverges from the parent prefix"
991 );
992 }
993
994 let model = "deepseek-v4-pro";
995 let history = vec![
996 text(Role::User, "fix the failing session_store test"),
997 text(Role::Assistant, "Reading the test first."),
998 text(Role::User, "keep the branch name"),
999 ];
1000 let system = SystemPrompt::Text("pinned system prompt".to_string());
1001 let tools = vec![test_tool("read"), test_tool("bash")];
1002 let effort = "high";
1003 let parent = codewhale_core::request::prepare_primary_turn_request(
1004 codewhale_core::request::PrimaryTurnRequest {
1005 model: model.to_string(),
1006 messages: history.clone(),
1007 max_tokens: 64_000,
1008 system: Some(system.clone()),
1009 tools: Some(tools.clone()),
1010 tool_choice: Some(json!({"type": "auto"})),
1011 reasoning_effort: Some(effort.to_string()),
1012 },
1013 );
1014 let config = crate::compaction::CompactionConfig {
1015 model: model.to_string(),
1016 ..Default::default()
1017 };
1018 let mut summary_history = history;
1019 summary_history.push(text(Role::User, "write the handoff summary"));
1020 let summary = crate::compaction::compaction_summary_request(
1021 summary_history,
1022 &config,
1023 Some(&system),
1024 Some(&tools),
1025 Some(effort),
1026 8_192,
1027 );
1028
1029 // Chat Completions: system and history share one `messages` array.
1030 let client = deepseek_request_boundary_client(
1031 "https://api.deepseek.com/v1",
1032 "http://127.0.0.1:9".into(),
1033 );
1034 let parent_body = client
1035 .prepare_outbound_request(parent.clone(), true)
1036 .expect("parent prepares")
1037 .body;
1038 let summary_body = client
1039 .prepare_outbound_request(summary.clone(), false)
1040 .expect("summary prepares")
1041 .body;
1042 assert_extends("chat", &parent_body["messages"], &summary_body["messages"]);
1043 for key in ["model", "tools", "reasoning_effort", "thinking"] {
1044 assert_eq!(parent_body.get(key), summary_body.get(key), "chat {key}");
1045 }
1046
1047 // Responses (the Codex route where the 0% hit was recorded).
1048 let parent_body =
1049 responses::build_responses_body_for_provider(&parent, ProviderKind::OpenaiCodex, None);
1050 let summary_body =
1051 responses::build_responses_body_for_provider(&summary, ProviderKind::OpenaiCodex, None);
1052 assert_extends("responses", &parent_body["input"], &summary_body["input"]);
1053 for key in ["model", "instructions", "tools", "reasoning", "include"] {
1054 assert_eq!(
1055 parent_body.get(key),
1056 summary_body.get(key),
1057 "responses {key}"
1058 );
1059 }
1060 assert!(parent_body.get("reasoning").is_some());
1061 }
1062
1063 #[tokio::test]
1064 async fn core_primary_request_preparation_matches_captured_transport_bytes() {
1065 let server = MockServer::start().await;
1066 Mock::given(method("POST"))
1067 .respond_with(
1068 ResponseTemplate::new(200)
1069 .insert_header("content-type", "text/event-stream")
1070 .set_body_string("data: [DONE]\n\n"),
1071 )
1072 .expect(1)
1073 .mount(&server)
1074 .await;
1075
1076 let request = codewhale_core::request::prepare_primary_turn_request(
1077 codewhale_core::request::PrimaryTurnRequest {
1078 model: "deepseek-v4-pro".to_string(),
1079 messages: vec![Message {
1080 role: Role::User,
1081 content: vec![ContentBlock::Text {
1082 text: "core request boundary".to_string(),
1083 cache_control: None,
1084 }],
1085 }],
1086 max_tokens: 64,
1087 system: None,
1088 tools: Some(vec![Tool {
1089 input_schema: json!({
1090 "zeta": {"type": "string"},
1091 "alpha": {"type": "number"},
1092 "type": "object",
1093 }),
1094 ..test_tool("lookup")
1095 }]),
1096 tool_choice: Some(json!({"type": "auto"})),
1097 reasoning_effort: Some("off".to_string()),
1098 },
1099 );
1100 let client = deepseek_request_boundary_client("https://api.deepseek.com/v1", server.uri());
1101 let prepared = client
1102 .prepare_outbound_request(request.clone(), true)
1103 .expect("core request prepares through the production seam");
1104 let prepared_bytes = serde_json::to_vec(&prepared.body).expect("prepared body serializes");
1105
1106 let mut stream = client
1107 .create_message_stream(request)
1108 .await
1109 .expect("production transport accepts the core request");
1110 while let Some(event) = stream.next().await {
1111 event.expect("captured SSE response remains valid");
1112 }
1113
1114 let requests = server.received_requests().await.expect("recorded request");
1115 assert_eq!(requests.len(), 1);
1116 assert_eq!(requests[0].body, prepared_bytes);
1117 let captured = std::str::from_utf8(&requests[0].body).expect("request body is UTF-8 JSON");
1118 assert!(
1119 captured.contains(
1120 r#""parameters":{"zeta":{"type":"string"},"alpha":{"type":"number"},"type":"object"}"#
1121 ),
1122 "nested core-owned JSON order drifted: {captured}"
1123 );
1124 }
1125
1126 #[tokio::test]
1127 async fn ollama_cloud_uses_authenticated_openai_compatible_v1_wire() {
1128 let server = MockServer::start().await;
1129 Mock::given(method("POST"))
1130 .and(path("/v1/chat/completions"))
1131 .and(header(
1132 "authorization",
1133 "Bearer ollama-cloud-request-boundary-key",
1134 ))
1135 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
1136 "id": "chatcmpl-ollama-cloud-request-boundary",
1137 "object": "chat.completion",
1138 "model": "gpt-oss:120b",
1139 "choices": [{
1140 "index": 0,
1141 "message": {"role": "assistant", "content": "ok"},
1142 "finish_reason": "stop"
1143 }],
1144 "usage": {
1145 "prompt_tokens": 1,
1146 "completion_tokens": 1,
1147 "total_tokens": 2
1148 }
1149 })))
1150 .expect(5)
1151 .mount(&server)
1152 .await;
1153
1154 let client = ollama_cloud_request_boundary_client(server.uri());
1155 for requested in ["off", "low", "medium", "high", "max"] {
1156 client
1157 .create_message(MessageRequest {
1158 model: "gpt-oss:120b".to_string(),
1159 messages: vec![Message {
1160 role: Role::User,
1161 content: vec![ContentBlock::Text {
1162 text: "Ollama Cloud request boundary".to_string(),
1163 cache_control: None,
1164 }],
1165 }],
1166 max_tokens: 64,
1167 system: None,
1168 tools: None,
1169 tool_choice: None,
1170 metadata: None,
1171 thinking: None,
1172 reasoning_effort: Some(requested.to_string()),
1173 stream: Some(false),
1174 temperature: None,
1175 top_p: None,
1176 })
1177 .await
1178 .expect("Ollama Cloud request succeeds");
1179 }
1180
1181 let requests = server.received_requests().await.expect("recorded request");
1182 assert_eq!(requests.len(), 5);
1183 for (request, expected) in requests
1184 .iter()
1185 .zip(["none", "low", "medium", "high", "max"])
1186 {
1187 let body: Value = serde_json::from_slice(&request.body).expect("captured request JSON");
1188 assert_eq!(body["model"], "gpt-oss:120b");
1189 assert_eq!(body["reasoning_effort"], expected);
1190 assert!(
1191 body.get("think").is_none(),
1192 "native Ollama field leaked: {body}"
1193 );
1194 assert!(
1195 body.get("thinking").is_none(),
1196 "foreign field leaked: {body}"
1197 );
1198 }
1199 }
1200
1201 // This synchronous guard deliberately spans every await: the assertions
1202 // require exclusive access to process-global retry state for the full call.
1203 #[allow(clippy::await_holding_lock)]
1204 #[tokio::test(flavor = "current_thread")]
1205 async fn cache_free_message_call_neither_reads_nor_writes_global_cache() {
1206 let _retry_guard = crate::retry_status::test_guard();
1207 crate::retry_status::clear();
1208 crate::retry_status::clear_rate_limit();
1209 crate::retry_status::start(7, Duration::from_secs(60), "foreground sentinel");
1210 let server = MockServer::start().await;
1211 Mock::given(method("POST"))
1212 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
1213 "id": "chatcmpl-cache-free-provider",
1214 "object": "chat.completion",
1215 "model": "deepseek-v4-pro",
1216 "choices": [{
1217 "index": 0,
1218 "message": {"role": "assistant", "content": "provider result"},
1219 "finish_reason": "stop"
1220 }],
1221 "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}
1222 })))
1223 .expect(1)
1224 .mount(&server)
1225 .await;
1226
1227 let client = deepseek_request_boundary_client("https://api.deepseek.com/v1", server.uri());
1228 crate::retry_status::note_rate_limit(&client.rate_limit_scope(), Duration::from_secs(60));
1229 let request = MessageRequest {
1230 model: "deepseek-v4-pro".to_string(),
1231 messages: vec![Message {
1232 role: Role::User,
1233 content: vec![ContentBlock::Text {
1234 text: "preview-router-cache-isolation-regression".to_string(),
1235 cache_control: None,
1236 }],
1237 }],
1238 max_tokens: 128,
1239 system: None,
1240 tools: None,
1241 tool_choice: None,
1242 metadata: None,
1243 thinking: None,
1244 reasoning_effort: Some("off".to_string()),
1245 stream: Some(false),
1246 temperature: Some(0.0),
1247 top_p: None,
1248 };
1249 let prepared = client
1250 .prepare_outbound_request(request.clone(), false)
1251 .expect("request prepares");
1252 let wire_body = serde_json::to_vec(&prepared.body).expect("wire body serializes");
1253 let cache_key = crate::llm_response_cache::ResponseCache::make_key(
1254 client.api_provider.as_str(),
1255 &client.base_url,
1256 client.path_suffix.as_deref(),
1257 &client.api_key,
1258 &wire_body,
1259 );
1260 crate::llm_response_cache::response_cache().put(
1261 cache_key,
1262 MessageResponse {
1263 id: "cached-sentinel-must-survive".to_string(),
1264 r#type: "message".to_string(),
1265 role: "assistant".to_string(),
1266 content: Vec::new(),
1267 model: "deepseek-v4-pro".to_string(),
1268 stop_reason: Some("end_turn".to_string()),
1269 stop_sequence: None,
1270 container: None,
1271 usage: Default::default(),
1272 },
1273 );
1274
1275 let response = client
1276 .create_message_without_response_cache(request)
1277 .await
1278 .expect("cache-free provider call succeeds");
1279 assert_eq!(response.id, "chatcmpl-cache-free-provider");
1280 assert_eq!(
1281 crate::llm_response_cache::response_cache()
1282 .get(&cache_key)
1283 .expect("sentinel remains")
1284 .id,
1285 "cached-sentinel-must-survive"
1286 );
1287 match crate::retry_status::snapshot() {
1288 crate::retry_status::RetryState::Active(banner) => {
1289 assert_eq!(banner.attempt, 7);
1290 assert_eq!(banner.reason, "foreground sentinel");
1291 }
1292 state => panic!("isolated success mutated retry state: {state:?}"),
1293 }
1294 assert!(
1295 crate::retry_status::rate_limit_remaining(&client.rate_limit_scope()).is_some(),
1296 "isolated success must not clear the foreground provider pause"
1297 );
1298 crate::retry_status::clear();
1299 crate::retry_status::clear_rate_limit();
1300 }
1301
1302 // This synchronous guard deliberately spans every await: the assertions
1303 // require exclusive access to process-global retry state for the full call.
1304 #[allow(clippy::await_holding_lock)]
1305 #[tokio::test(flavor = "current_thread")]
1306 async fn cache_free_classifier_429_does_not_publish_global_retry_or_rate_limit_state() {
1307 let _retry_guard = crate::retry_status::test_guard();
1308 crate::retry_status::clear();
1309 crate::retry_status::clear_rate_limit();
1310 crate::retry_status::start(9, Duration::from_secs(60), "foreground sentinel 429");
1311
1312 let server = MockServer::start().await;
1313 Mock::given(method("POST"))
1314 .respond_with(
1315 ResponseTemplate::new(429)
1316 .insert_header("retry-after", "120")
1317 .set_body_string("rate limited"),
1318 )
1319 .expect(1)
1320 .mount(&server)
1321 .await;
1322 let mut client =
1323 deepseek_request_boundary_client("https://api.deepseek.com/v1", server.uri());
1324 client.retry.enabled = false;
1325 client.retry.max_retries = 0;
1326 crate::retry_status::note_rate_limit(&client.rate_limit_scope(), Duration::from_secs(60));
1327 let request = MessageRequest {
1328 model: "deepseek-v4-pro".to_string(),
1329 messages: vec![Message {
1330 role: Role::User,
1331 content: vec![ContentBlock::Text {
1332 text: "preview-router-429-isolation-regression".to_string(),
1333 cache_control: None,
1334 }],
1335 }],
1336 max_tokens: 64,
1337 system: None,
1338 tools: None,
1339 tool_choice: None,
1340 metadata: None,
1341 thinking: None,
1342 reasoning_effort: Some("off".to_string()),
1343 stream: Some(false),
1344 temperature: Some(0.0),
1345 top_p: None,
1346 };
1347 let error = client
1348 .create_message_without_response_cache(request)
1349 .await
1350 .expect_err("429 must fail when isolated retries are disabled");
1351 assert!(
1352 matches!(
1353 error.downcast_ref::<LlmError>(),
1354 Some(LlmError::RateLimited { .. })
1355 ),
1356 "{error:#}"
1357 );
1358 match crate::retry_status::snapshot() {
1359 crate::retry_status::RetryState::Active(banner) => {
1360 assert_eq!(banner.attempt, 9);
1361 assert_eq!(banner.reason, "foreground sentinel 429");
1362 }
1363 state => panic!("isolated 429 mutated retry state: {state:?}"),
1364 }
1365 let remaining = crate::retry_status::rate_limit_remaining(&client.rate_limit_scope())
1366 .expect("foreground provider pause remains");
1367 assert!(
1368 remaining < Duration::from_secs(70),
1369 "classifier Retry-After must not extend the global pause: {remaining:?}"
1370 );
1371 crate::retry_status::clear();
1372 crate::retry_status::clear_rate_limit();
1373 }
1374
1375 async fn assert_deepseek_strict_request_route_boundary(streaming: bool) {
1376 for (route_base_url, strict, expected_path, expected_wire_strict) in [
1377 (
1378 "https://api.deepseek.com/beta",
1379 false,
1380 "/v1/chat/completions",
1381 None,
1382 ),
1383 (
1384 "https://api.deepseek.com/beta",
1385 true,
1386 "/beta/chat/completions",
1387 Some(true),
1388 ),
1389 (
1390 "https://api.deepseek.com/v1",
1391 true,
1392 "/v1/chat/completions",
1393 None,
1394 ),
1395 ] {
1396 let (captured_path, body) =
1397 capture_deepseek_chat_request(route_base_url, strict, streaming).await;
1398 assert_eq!(captured_path, expected_path, "{route_base_url} {body}");
1399 assert_eq!(
1400 body.pointer("/tools/0/function/strict")
1401 .and_then(Value::as_bool),
1402 expected_wire_strict,
1403 "{route_base_url} {body}"
1404 );
1405 }
1406 }
1407
1408 fn k3_request_fixture(model: &str, effort: Option<&str>, stream: bool) -> MessageRequest {
1409 MessageRequest {
1410 model: model.to_string(),
1411 messages: vec![Message {
1412 role: Role::User,
1413 content: vec![ContentBlock::Text {
1414 text: "request-boundary fixture".to_string(),
1415 cache_control: None,
1416 }],
1417 }],
1418 max_tokens: 64,
1419 system: None,
1420 tools: None,
1421 tool_choice: None,
1422 metadata: None,
1423 thinking: None,
1424 reasoning_effort: effort.map(str::to_string),
1425 stream: Some(stream),
1426 temperature: Some(0.25),
1427 top_p: Some(0.75),
1428 }
1429 }
1430
1431 async fn capture_moonshot_chat_request(
1432 route_base_url: &str,
1433 model: &str,
1434 effort: Option<&str>,
1435 streaming: bool,
1436 ) -> Value {
1437 let request = k3_request_fixture(model, effort, streaming);
1438 capture_moonshot_chat_request_body(route_base_url, model, request).await
1439 }
1440
1441 async fn capture_moonshot_chat_request_body(
1442 route_base_url: &str,
1443 model: &str,
1444 request: MessageRequest,
1445 ) -> Value {
1446 let streaming = request.stream == Some(true);
1447 let server = MockServer::start().await;
1448 let response = if streaming {
1449 ResponseTemplate::new(200)
1450 .insert_header("content-type", "text/event-stream")
1451 .set_body_string("data: [DONE]\n\n")
1452 } else {
1453 ResponseTemplate::new(200).set_body_json(json!({
1454 "id": "chatcmpl-k3-request-boundary",
1455 "object": "chat.completion",
1456 "model": model,
1457 "choices": [{
1458 "index": 0,
1459 "message": {"role": "assistant", "content": "ok"},
1460 "finish_reason": "stop"
1461 }],
1462 "usage": {
1463 "prompt_tokens": 1,
1464 "completion_tokens": 1,
1465 "total_tokens": 2
1466 }
1467 }))
1468 };
1469 Mock::given(method("POST"))
1470 .and(path("/v1/chat/completions"))
1471 .respond_with(response)
1472 .expect(1)
1473 .mount(&server)
1474 .await;
1475
1476 let client = moonshot_request_boundary_client(route_base_url, model, server.uri());
1477
1478 if streaming {
1479 let mut stream = client
1480 .create_message_stream(request)
1481 .await
1482 .expect("streaming request succeeds");
1483 while let Some(event) = stream.next().await {
1484 event.expect("captured SSE response remains valid");
1485 }
1486 } else {
1487 client
1488 .create_message(request)
1489 .await
1490 .expect("non-streaming request succeeds");
1491 }
1492
1493 let requests = server.received_requests().await.expect("recorded request");
1494 assert_eq!(requests.len(), 1);
1495 serde_json::from_slice(&requests[0].body).expect("captured request JSON")
1496 }
1497
1498 async fn capture_route_chat_request_body(
1499 model: &str,
1500 request: MessageRequest,
1501 client_for_transport: impl FnOnce(String) -> CodewhaleClient,
1502 ) -> (String, Value) {
1503 let streaming = request.stream == Some(true);
1504 let server = MockServer::start().await;
1505 let response = if streaming {
1506 ResponseTemplate::new(200)
1507 .insert_header("content-type", "text/event-stream")
1508 .set_body_string("data: [DONE]\n\n")
1509 } else {
1510 ResponseTemplate::new(200).set_body_json(json!({
1511 "id": "chatcmpl-provider-request-boundary",
1512 "object": "chat.completion",
1513 "model": model,
1514 "choices": [{
1515 "index": 0,
1516 "message": {"role": "assistant", "content": "ok"},
1517 "finish_reason": "stop"
1518 }],
1519 "usage": {
1520 "prompt_tokens": 1,
1521 "completion_tokens": 1,
1522 "total_tokens": 2
1523 }
1524 }))
1525 };
1526 Mock::given(method("POST"))
1527 .and(path("/v1/chat/completions"))
1528 .respond_with(response)
1529 .expect(1)
1530 .mount(&server)
1531 .await;
1532
1533 let client = client_for_transport(server.uri());
1534 if streaming {
1535 let mut stream = client
1536 .create_message_stream(request)
1537 .await
1538 .expect("streaming request succeeds");
1539 while let Some(event) = stream.next().await {
1540 event.expect("captured SSE response remains valid");
1541 }
1542 } else {
1543 client
1544 .create_message(request)
1545 .await
1546 .expect("non-streaming request succeeds");
1547 }
1548
1549 let requests = server.received_requests().await.expect("recorded request");
1550 assert_eq!(requests.len(), 1);
1551 (
1552 requests[0].url.path().to_string(),
1553 serde_json::from_slice(&requests[0].body).expect("captured request JSON"),
1554 )
1555 }
1556
1557 async fn capture_zai_chat_request(
1558 route_base_url: &str,
1559 model: &str,
1560 effort: Option<&str>,
1561 streaming: bool,
1562 ) -> (String, Value) {
1563 capture_route_chat_request_body(
1564 model,
1565 k3_request_fixture(model, effort, streaming),
1566 |uri| zai_request_boundary_client(route_base_url, model, uri),
1567 )
1568 .await
1569 }
1570
1571 async fn capture_minimax_chat_request(
1572 route_base_url: &str,
1573 model: &str,
1574 effort: Option<&str>,
1575 streaming: bool,
1576 ) -> (String, Value) {
1577 capture_route_chat_request_body(
1578 model,
1579 k3_request_fixture(model, effort, streaming),
1580 |uri| minimax_request_boundary_client(route_base_url, model, uri),
1581 )
1582 .await
1583 }
1584
1585 fn modelstudio_request_boundary_client(
1586 route_base_url: &str,
1587 model: &str,
1588 transport_base_url: String,
1589 ) -> CodewhaleClient {
1590 let _ = rustls::crypto::ring::default_provider().install_default();
1591 let mut client = CodewhaleClient::new(&Config {
1592 provider: Some("modelstudio-token-plan".to_string()),
1593 providers: Some(ProvidersConfig {
1594 modelstudio_token_plan: ProviderConfig {
1595 api_key: Some("modelstudio-request-boundary-key".to_string()),
1596 base_url: Some(route_base_url.to_string()),
1597 model: Some(model.to_string()),
1598 ..ProviderConfig::default()
1599 },
1600 ..ProvidersConfig::default()
1601 }),
1602 ..Config::default()
1603 })
1604 .expect("Model Studio request-boundary client");
1605 assert_eq!(client.base_url, route_base_url);
1606 client.test_chat_transport_base_url = Some(transport_base_url);
1607 client
1608 }
1609
1610 async fn capture_modelstudio_chat_request(
1611 route_base_url: &str,
1612 model: &str,
1613 effort: Option<&str>,
1614 streaming: bool,
1615 ) -> (String, Value) {
1616 capture_route_chat_request_body(
1617 model,
1618 k3_request_fixture(model, effort, streaming),
1619 |uri| modelstudio_request_boundary_client(route_base_url, model, uri),
1620 )
1621 .await
1622 }
1623
1624 async fn assert_modelstudio_request_truth(streaming: bool) {
1625 // Token Plan and Coding Plan share DashScope's reasoning controls on
1626 // their OpenAI-compatible Chat Completions endpoints — but the fields
1627 // are model-specific, not provider-wide.
1628 for base_url in [
1629 crate::config::DEFAULT_MODELSTUDIO_TOKEN_PLAN_BASE_URL,
1630 crate::config::DEFAULT_MODELSTUDIO_CODING_PLAN_BASE_URL,
1631 ] {
1632 // The default model, qwen3.8-max, is thinking-only: the bundled
1633 // catalog records it as `thinking: always_on`, and
1634 // qwen3.8-max-preview has effort/budget options with no toggle.
1635 // Neither accepts an enable/disable switch, so CodeWhale must not
1636 // send one — not even `false` for an explicit `off`. This assertion
1637 // used to pin the opposite; PR #5233 caught it.
1638 for effort in [None, Some("off"), Some("high"), Some("max")] {
1639 let (path, body) = capture_modelstudio_chat_request(
1640 base_url,
1641 crate::config::DEFAULT_MODELSTUDIO_TOKEN_PLAN_MODEL,
1642 effort,
1643 streaming,
1644 )
1645 .await;
1646 assert_eq!(path, "/v1/chat/completions");
1647 assert!(
1648 body.get("enable_thinking").is_none(),
1649 "{base_url} {effort:?}: {body}"
1650 );
1651 assert!(
1652 body.get("thinking").is_none(),
1653 "{base_url} {effort:?}: {body}"
1654 );
1655 assert!(
1656 body.get("reasoning_effort").is_none(),
1657 "{base_url} {effort:?}: {body}"
1658 );
1659 }
1660
1661 // A hybrid model does get the documented switch, plus
1662 // `preserve_thinking` so the next turn keeps its trace.
1663 for (effort, enabled) in [(None, true), (Some("high"), true), (Some("off"), false)] {
1664 let (_, body) =
1665 capture_modelstudio_chat_request(base_url, "qwen3.7-plus", effort, streaming)
1666 .await;
1667 assert_eq!(
1668 body["enable_thinking"],
1669 json!(enabled),
1670 "{base_url} {effort:?}: {body}"
1671 );
1672 assert_eq!(
1673 body["preserve_thinking"],
1674 json!(enabled),
1675 "{base_url} {effort:?}: {body}"
1676 );
1677 // The hybrid Qwen families have no effort ladder on the wire.
1678 assert!(
1679 body.get("reasoning_effort").is_none(),
1680 "{base_url} {effort:?}: {body}"
1681 );
1682 }
1683
1684 // DeepSeek-V4 is one of the two families with a documented effort
1685 // ladder (`high` / `max`).
1686 let (_, deepseek) = capture_modelstudio_chat_request(
1687 base_url,
1688 "deepseek-v4-pro",
1689 Some("xhigh"),
1690 streaming,
1691 )
1692 .await;
1693 assert_eq!(
1694 deepseek["enable_thinking"],
1695 json!(true),
1696 "{base_url}: {deepseek}"
1697 );
1698 assert_eq!(
1699 deepseek["reasoning_effort"],
1700 json!("max"),
1701 "{base_url}: {deepseek}"
1702 );
1703 }
1704
1705 // Fail closed: the same provider identity pointed at a custom gateway
1706 // must not be handed Alibaba's dialect.
1707 let (_, proxied) = capture_modelstudio_chat_request(
1708 "https://proxy.example/v1",
1709 "qwen3.7-plus",
1710 Some("high"),
1711 streaming,
1712 )
1713 .await;
1714 assert!(proxied.get("enable_thinking").is_none(), "{proxied}");
1715 assert!(proxied.get("preserve_thinking").is_none(), "{proxied}");
1716 assert!(proxied.get("reasoning_effort").is_none(), "{proxied}");
1717 }
1717 lines RUST