返回 CodeWhale
test_cases_06.rs
根目录 / crates / tui / src / client / test_cases_06.rs
1
2 #[tokio::test]
3 async fn fetch_catalog_delta_rejects_oversized_bodies_and_rosters() {
4 let server = MockServer::start().await;
5 Mock::given(method("GET"))
6 .and(path("/v1/models"))
7 .respond_with(ResponseTemplate::new(200).set_body_raw(
8 "x".repeat(PROVIDER_CATALOG_MAX_RESPONSE_BYTES + 1),
9 "application/json",
10 ))
11 .mount(&server)
12 .await;
13 assert_eq!(
14 openrouter_client_for(&server)
15 .fetch_catalog_delta()
16 .await
17 .expect_err("oversized body"),
18 CatalogRefreshError::InvalidResponse
19 );
20
21 let server = MockServer::start().await;
22 let rows: Vec<_> = (0..=PROVIDER_CATALOG_MAX_ROWS)
23 .map(|index| json!({"id": format!("synthetic-model-{index}")}))
24 .collect();
25 mount_models_json(&server, 200, json!({"data": rows})).await;
26 assert_eq!(
27 openrouter_client_for(&server)
28 .fetch_catalog_delta()
29 .await
30 .expect_err("oversized roster"),
31 CatalogRefreshError::InvalidResponse
32 );
33 }
34
35 #[tokio::test]
36 async fn refresh_catalog_cache_records_success_then_preserves_rows_on_failure() {
37 // First refresh succeeds and caches live rows.
38 let server = MockServer::start().await;
39 mount_models_json(
40 &server,
41 200,
42 json!({"data": [{"id": "synthetic-model-gamma"}]}),
43 )
44 .await;
45 let client = openrouter_client_for(&server);
46 let mut cache = ProviderCatalogCache::new();
47
48 let status = client.refresh_catalog_cache(&mut cache, 3600).await;
49 assert_eq!(status, CatalogStatus::Fresh);
50 let fp = base_url_fingerprint(&server.uri());
51 let cached = cache.get("openrouter", &fp).expect("cached entry");
52 assert_eq!(cached.offerings.len(), 1);
53 assert_eq!(cached.offerings[0].wire_model_id, "synthetic-model-gamma");
54
55 // A later failing refresh on the same base URL flips status to Failed
56 // but PRESERVES the rows.
57 server.reset().await;
58 mount_models_json(&server, 401, json!({"error": "denied"})).await;
59 let status = client.refresh_catalog_cache(&mut cache, 3600).await;
60 assert!(matches!(
61 status,
62 CatalogStatus::Failed {
63 reason: CatalogRefreshError::Unauthorized,
64 ..
65 }
66 ));
67 let cached = cache.get("openrouter", &fp).expect("entry still present");
68 assert_eq!(
69 cached.offerings.len(),
70 1,
71 "rows from the prior success must survive a failed refresh"
72 );
73 assert!(matches!(cached.status, CatalogStatus::Failed { .. }));
74
75 // #4139: failed/stale rows must still publish into ProviderLake so
76 // pickers keep live coverage instead of dropping back to bundled-only.
77 let visible = cache.all_visible_offerings(now_unix());
78 assert_eq!(visible.len(), 1);
79 assert_eq!(visible[0].wire_model_id, "synthetic-model-gamma");
80 assert!(
81 cache.all_fresh_offerings(now_unix()).is_empty(),
82 "Failed entries are not fresh, but they remain visible"
83 );
84 }
85
86 #[tokio::test]
87 async fn invalid_live_prices_fail_refresh_and_preserve_each_provider_last_known_good() {
88 let openrouter_server = MockServer::start().await;
89 mount_models_json(
90 &openrouter_server,
91 200,
92 json!({"data": [{
93 "id": "synthetic/openrouter-priced",
94 "pricing": {"prompt": "0.000001", "completion": "0.000002"}
95 }]}),
96 )
97 .await;
98 let openrouter = openrouter_client_for(&openrouter_server);
99 let mut openrouter_cache = ProviderCatalogCache::new();
100 assert_eq!(
101 openrouter
102 .refresh_catalog_cache(&mut openrouter_cache, 3_600)
103 .await,
104 CatalogStatus::Fresh
105 );
106 let openrouter_fp = base_url_fingerprint(&openrouter_server.uri());
107 let openrouter_lkg = openrouter_cache
108 .get("openrouter", &openrouter_fp)
109 .expect("OpenRouter LKG")
110 .offerings
111 .clone();
112
113 openrouter_server.reset().await;
114 mount_models_json(
115 &openrouter_server,
116 200,
117 json!({"data": [{
118 "id": "synthetic/openrouter-priced",
119 "pricing": {"prompt": "1e308", "completion": "0.000002"}
120 }]}),
121 )
122 .await;
123 assert!(matches!(
124 openrouter
125 .refresh_catalog_cache(&mut openrouter_cache, 3_600)
126 .await,
127 CatalogStatus::Failed {
128 reason: CatalogRefreshError::InvalidResponse
129 }
130 ));
131 assert_eq!(
132 openrouter_cache
133 .get("openrouter", &openrouter_fp)
134 .expect("preserved OpenRouter LKG")
135 .offerings,
136 openrouter_lkg
137 );
138
139 // Same plumbing for an ordinary custom host: the generic branch
140 // ignores unknown pricing fields, so the failure injector here is a
141 // body that is not a model list at all. Absurd-price rejection stays
142 // covered at the Baseten builder's own fixture tests.
143 let custom_server = MockServer::start().await;
144 mount_models_json(
145 &custom_server,
146 200,
147 json!({"data": [{
148 "id": "synthetic/custom-model"
149 }]}),
150 )
151 .await;
152 let custom = custom_mock_client_for_identity(&custom_server, "custom-lkg");
153 let mut custom_cache = ProviderCatalogCache::new();
154 assert_eq!(
155 custom.refresh_catalog_cache(&mut custom_cache, 3_600).await,
156 CatalogStatus::Fresh
157 );
158 let custom_fp = base_url_fingerprint(&format!("{}/v1", custom_server.uri()));
159 let custom_lkg = custom_cache
160 .get("custom-lkg", &custom_fp)
161 .expect("custom LKG")
162 .offerings
163 .clone();
164
165 custom_server.reset().await;
166 mount_models_json(&custom_server, 200, json!("not a model list")).await;
167 assert!(matches!(
168 custom.refresh_catalog_cache(&mut custom_cache, 3_600).await,
169 CatalogStatus::Failed {
170 reason: CatalogRefreshError::InvalidResponse
171 }
172 ));
173 assert_eq!(
174 custom_cache
175 .get("custom-lkg", &custom_fp)
176 .expect("preserved custom LKG")
177 .offerings,
178 custom_lkg
179 );
180 }
181
182 #[tokio::test]
183 async fn live_catalog_is_scoped_by_base_url_fingerprint() {
184 // Same provider, two different base URLs -> two distinct cache scopes.
185 let server_a = MockServer::start().await;
186 mount_models_json(&server_a, 200, json!({"data": [{"id": "synthetic-a"}]})).await;
187 let server_b = MockServer::start().await;
188 mount_models_json(&server_b, 200, json!({"data": [{"id": "synthetic-b"}]})).await;
189
190 let mut cache = ProviderCatalogCache::new();
191 openrouter_client_for(&server_a)
192 .refresh_catalog_cache(&mut cache, 3600)
193 .await;
194 openrouter_client_for(&server_b)
195 .refresh_catalog_cache(&mut cache, 3600)
196 .await;
197
198 let fp_a = base_url_fingerprint(&server_a.uri());
199 let fp_b = base_url_fingerprint(&server_b.uri());
200 assert_ne!(
201 fp_a, fp_b,
202 "different base URLs must fingerprint differently"
203 );
204 assert_eq!(
205 cache.get("openrouter", &fp_a).expect("a").offerings[0].wire_model_id,
206 "synthetic-a"
207 );
208 assert_eq!(
209 cache.get("openrouter", &fp_b).expect("b").offerings[0].wire_model_id,
210 "synthetic-b"
211 );
212 }
213
214 #[tokio::test]
215 async fn static_rows_survive_a_live_refresh_failure() {
216 // Bundled/static rows compile through even when the live layer is empty
217 // (the state after a failed refresh with no prior success).
218 let server = MockServer::start().await;
219 mount_models_json(&server, 503, json!({"error": "down"})).await;
220 let client = openrouter_client_for(&server);
221 let mut cache = ProviderCatalogCache::new();
222 let status = client.refresh_catalog_cache(&mut cache, 3600).await;
223 assert!(matches!(status, CatalogStatus::Failed { .. }));
224
225 let static_row = CatalogOffering {
226 provider: "openrouter".to_string(),
227 wire_model_id: "synthetic-static".to_string(),
228 endpoint_key: "chat".to_string(),
229 ..CatalogOffering::default()
230 };
231 let fp = base_url_fingerprint(&server.uri());
232 let fresh_live: Vec<CatalogOffering> = cache
233 .get("openrouter", &fp)
234 .filter(|entry| entry.is_fresh(now_unix()))
235 .map(|entry| entry.offerings.clone())
236 .unwrap_or_default();
237 let snapshot = codewhale_config::catalog::CatalogCompiler::new()
238 .with_bundled(vec![static_row])
239 .with_live(fresh_live)
240 .compile();
241 assert!(
242 snapshot
243 .offerings
244 .iter()
245 .any(|offering| offering.wire_model_id == "synthetic-static"),
246 "static fallback row must remain available after a failed refresh"
247 );
248 }
249
250 #[test]
251 fn client_route_envelope_freezes_saved_minimax_billing_mode_and_wire_model() {
252 let _ = rustls::crypto::ring::default_provider().install_default();
253 let config = Config {
254 provider: Some("minimax".to_string()),
255 providers: Some(ProvidersConfig {
256 minimax: ProviderConfig {
257 api_key: Some("test-key".to_string()),
258 mode: Some("pay-as-you-go".to_string()),
259 ..ProviderConfig::default()
260 },
261 ..ProvidersConfig::default()
262 }),
263 ..Config::default()
264 };
265 let client = CodewhaleClient::new(&config).expect("MiniMax client");
266 let dispatched_at =
267 chrono::DateTime::<chrono::Utc>::from_timestamp(1_234, 0).expect("timestamp");
268 let route = client.effective_route_envelope("MiniMax-M3", dispatched_at);
269
270 assert_eq!(route.provider, ProviderKind::Minimax);
271 assert_eq!(route.provider_identity, "minimax");
272 assert_eq!(route.model, "MiniMax-M3");
273 assert_eq!(
274 route.billing_surface.as_deref(),
275 Some(crate::pricing::MINIMAX_PAYG_BILLING_SURFACE)
276 );
277 assert_eq!(
278 route.billing_mode,
279 crate::cost_status::RouteBillingMode::Metered
280 );
281 assert_eq!(route.dispatched_at.timestamp(), 1_234);
282 }
283
284 #[test]
285 fn sanitize_thinking_mode_counts_reasoning_replay_across_assistant_turns() {
286 // Multi-turn body that mimics two prior tool-calling rounds: each
287 // assistant message carries its `reasoning_content`. The sanitizer
288 // should keep all of them and the count helper should tally bytes
289 // across every assistant message.
290 let mut body = json!({
291 "model": "deepseek-v4-pro",
292 "messages": [
293 { "role": "system", "content": "you are helpful" },
294 { "role": "user", "content": "step 1" },
295 {
296 "role": "assistant",
297 "content": "",
298 "reasoning_content": "I need to call tool A first.",
299 "tool_calls": [{ "id": "1", "type": "function" }]
300 },
301 { "role": "tool", "tool_call_id": "1", "content": "ok" },
302 {
303 "role": "assistant",
304 "content": "",
305 "reasoning_content": "Now I call tool B.",
306 "tool_calls": [{ "id": "2", "type": "function" }]
307 },
308 { "role": "tool", "tool_call_id": "2", "content": "ok" },
309 { "role": "user", "content": "step 2" }
310 ]
311 });
312
313 let approx_tokens = sanitize_thinking_mode_messages(
314 &mut body,
315 "deepseek-v4-pro",
316 Some("max"),
317 ProviderKind::Deepseek,
318 )
319 .expect("multi-turn thinking-mode conversation should report replay tokens");
320 // ~4 chars/token; 46 bytes of reasoning -> 11 tokens.
321 assert_eq!(approx_tokens, 11);
322
323 let chars = count_reasoning_replay_chars(&body);
324 // "I need to call tool A first." (28) + "Now I call tool B." (18) = 46
325 assert_eq!(chars, 46);
326
327 // No assistant messages should have lost or had their reasoning_content blanked.
328 let messages = body["messages"].as_array().unwrap();
329 let assistant_with_reasoning: usize = messages
330 .iter()
331 .filter(|m| m["role"] == "assistant")
332 .filter(|m| {
333 m["reasoning_content"]
334 .as_str()
335 .is_some_and(|s| !s.is_empty())
336 })
337 .count();
338 assert_eq!(assistant_with_reasoning, 2);
339 }
340
341 /// Issue #30: when no thinking-mode replay applies (non-thinking model or
342 /// empty conversation), the sanitizer returns `None` so the footer chip
343 /// stays hidden.
344 #[test]
345 fn sanitize_thinking_mode_returns_none_for_non_thinking_model() {
346 let mut body = json!({
347 "model": "deepseek-v4-flash",
348 "messages": [
349 { "role": "user", "content": "hi" }
350 ]
351 });
352 let result = sanitize_thinking_mode_messages(
353 &mut body,
354 "deepseek-v4-flash",
355 None,
356 ProviderKind::Deepseek,
357 );
358 // reasoning_effort is None → no thinking injection, result is None
359 assert!(result.is_none());
360 }
361
362 #[test]
363 fn sanitize_thinking_mode_counts_substituted_placeholder() {
364 // An assistant tool-call message is missing reasoning_content; the
365 // sanitizer must inject the placeholder, and the count helper must
366 // include the placeholder in the total (since it's in the wire
367 // payload that ships to DeepSeek).
368 let mut body = json!({
369 "model": "deepseek-v4-pro",
370 "messages": [
371 { "role": "user", "content": "hi" },
372 {
373 "role": "assistant",
374 "content": "",
375 "tool_calls": [{ "id": "1", "type": "function" }]
376 }
377 ]
378 });
379
380 sanitize_thinking_mode_messages(
381 &mut body,
382 "deepseek-v4-pro",
383 Some("max"),
384 ProviderKind::Deepseek,
385 );
386
387 let chars = count_reasoning_replay_chars(&body);
388 // "(reasoning omitted)" is 19 bytes.
389 assert_eq!(chars, 19);
390 }
391
392 #[test]
393 fn sanitize_thinking_mode_skips_generic_openai_provider() {
394 // #1542 intent (narrowed by #1739/#1694): the sanitizer only skips for
395 // a *genuine non-DeepSeek* model on the generic openai provider. A
396 // DeepSeek reasoning model on the openai provider still gets sanitized
397 // (see chat.rs `deepseek_model_on_openai_provider_still_replays_*`).
398 let mut body = json!({
399 "model": "qwen3-coder",
400 "messages": [
401 { "role": "user", "content": "hi" },
402 {
403 "role": "assistant",
404 "content": "",
405 "tool_calls": [{ "id": "1", "type": "function" }]
406 }
407 ]
408 });
409
410 let result = sanitize_thinking_mode_messages(
411 &mut body,
412 "qwen3-coder",
413 Some("max"),
414 ProviderKind::Openai,
415 );
416
417 assert!(result.is_none());
418 let assistant = body["messages"]
419 .as_array()
420 .and_then(|messages| {
421 messages
422 .iter()
423 .find(|message| message["role"] == "assistant")
424 })
425 .expect("assistant message");
426 assert!(
427 assistant.get("reasoning_content").is_none(),
428 "generic OpenAI-compatible provider payload must not get reasoning_content (#1542)"
429 );
430 }
431
432 #[test]
433 fn sanitize_thinking_mode_keeps_tool_call_placeholder_after_new_user_turn() {
434 let mut body = json!({
435 "model": "deepseek-v4-pro",
436 "messages": [
437 { "role": "user", "content": "step 1" },
438 {
439 "role": "assistant",
440 "content": "",
441 "tool_calls": [{ "id": "1", "type": "function" }]
442 },
443 { "role": "tool", "tool_call_id": "1", "content": "ok" },
444 { "role": "user", "content": "step 2" }
445 ]
446 });
447
448 sanitize_thinking_mode_messages(
449 &mut body,
450 "deepseek-v4-pro",
451 Some("max"),
452 ProviderKind::Deepseek,
453 );
454
455 let messages = body["messages"].as_array().unwrap();
456 let assistant = messages
457 .iter()
458 .find(|m| m["role"] == "assistant")
459 .expect("assistant tool-call message");
460 assert_eq!(
461 assistant.get("reasoning_content").and_then(Value::as_str),
462 Some("(reasoning omitted)")
463 );
464 }
465
466 #[test]
467 fn token_bucket_enforces_delay_when_empty() {
468 let now = Instant::now();
469 let mut bucket = TokenBucket {
470 enabled: true,
471 capacity: 1.0,
472 tokens: 1.0,
473 refill_per_sec: 2.0,
474 last_refill: now,
475 };
476
477 assert!(bucket.delay_until_available(1.0).is_none());
478 let delay = bucket
479 .delay_until_available(1.0)
480 .expect("bucket should require refill delay");
481 assert!(
482 delay >= Duration::from_millis(400) && delay <= Duration::from_millis(600),
483 "unexpected refill delay: {delay:?}"
484 );
485 }
486
487 /// Every queued waiter must be given a *distinct* wake time. `client.rs`
488 /// releases the bucket lock before sleeping (`wait_for_rate_limit`), and a
489 /// clone of the client shares one `Arc<AsyncMutex<TokenBucket>>` across
490 /// sub-agents, so if the bucket hands two waiters the same delay they both
491 /// wake at the same instant and fire together — a burst the configured
492 /// limit was supposed to prevent.
493 #[test]
494 fn token_bucket_queues_concurrent_waiters_instead_of_stacking_them() {
495 let now = Instant::now();
496 let mut bucket = TokenBucket {
497 enabled: true,
498 capacity: 1.0,
499 tokens: 1.0,
500 refill_per_sec: 1.0,
501 last_refill: now,
502 };
503
504 assert!(bucket.delay_until_available(1.0).is_none());
505 let first = bucket
506 .delay_until_available(1.0)
507 .expect("second caller waits for a refill");
508 let second = bucket
509 .delay_until_available(1.0)
510 .expect("third caller waits for a refill");
511
512 assert!(
513 first >= Duration::from_millis(900) && first <= Duration::from_millis(1100),
514 "unexpected first wait: {first:?}"
515 );
516 assert!(
517 second >= Duration::from_millis(1900) && second <= Duration::from_millis(2100),
518 "third caller must queue behind the second, not wake with it: {second:?}"
519 );
520 }
521
522 #[test]
523 fn stream_buffer_pool_reuses_released_buffers() {
524 let mut first = acquire_stream_buffer();
525 first.extend_from_slice(b"hello");
526 let released_capacity = first.capacity();
527 release_stream_buffer(first);
528
529 let second = acquire_stream_buffer();
530 assert!(second.is_empty());
531 assert!(
532 second.capacity() >= released_capacity,
533 "pooled buffer capacity should be reused"
534 );
535 }
536
537 #[test]
538 fn base_url_scenario() {
539 // Scenario consolidation of: base_url_security_rejects_insecure_non_local_http, base_url_security_errors_redact_sensitive_url_parts, base_url_security_allows_localhost_http, base_url_security_allows_non_local_http_with_explicit_opt_in
540 // from base_url_security_rejects_insecure_non_local_http
541 {
542 let _lock = ALLOW_INSECURE_HTTP_ENV_LOCK.lock().unwrap();
543 let _guard = AllowInsecureHttpEnvGuard::capture();
544 unsafe { std::env::remove_var(ALLOW_INSECURE_HTTP_ENV) };
545
546 let err = validate_base_url_security("http://api.deepseek.com", false)
547 .expect_err("non-local insecure HTTP should be rejected");
548 assert!(err.to_string().contains("Refusing insecure base URL"));
549 }
550 // from base_url_security_errors_redact_sensitive_url_parts
551 {
552 let _lock = ALLOW_INSECURE_HTTP_ENV_LOCK.lock().unwrap();
553 let _guard = AllowInsecureHttpEnvGuard::capture();
554 unsafe { std::env::remove_var(ALLOW_INSECURE_HTTP_ENV) };
555
556 let err = validate_base_url_security(
557 "http://user:secret@example.com/v1?api_key=sk-test&ok=1",
558 false,
559 )
560 .expect_err("non-local insecure HTTP should be rejected");
561 let message = err.to_string();
562
563 assert!(message.contains("http://***:***@example.com/v1?api_key=***&ok=1"));
564 assert!(!message.contains("user:secret"));
565 assert!(!message.contains("sk-test"));
566 }
567 // from base_url_security_allows_localhost_http
568 {
569 let _lock = ALLOW_INSECURE_HTTP_ENV_LOCK.lock().unwrap();
570 let _guard = AllowInsecureHttpEnvGuard::capture();
571 unsafe { std::env::remove_var(ALLOW_INSECURE_HTTP_ENV) };
572
573 for url in [
574 "http://localhost:8080",
575 "http://LOCALHOST:8080",
576 "http://127.0.0.1:8080",
577 "http://127.0.0.2:8080",
578 "http://[::1]:8080",
579 "https://provider.example/v1",
580 ] {
581 assert!(validate_base_url_security(url, false).is_ok(), "{url}");
582 }
583 for url in [
584 "http://localhost.attacker.example/v1",
585 "http://127.0.0.1.attacker.example/v1",
586 "http://localhost@attacker.example/v1",
587 "http://127.0.0.1@attacker.example/v1",
588 "HTTP://localhost.attacker.example/v1",
589 ] {
590 assert!(validate_base_url_security(url, false).is_err(), "{url}");
591 assert!(validate_base_url_security(url, true).is_ok(), "{url}");
592 }
593 for url in ["https://", "http://[::1", "file:///tmp/provider"] {
594 assert!(validate_base_url_security(url, true).is_err(), "{url}");
595 }
596 }
597 // from base_url_security_allows_non_local_http_with_explicit_opt_in
598 {
599 let _lock = ALLOW_INSECURE_HTTP_ENV_LOCK.lock().unwrap();
600 let _guard = AllowInsecureHttpEnvGuard::capture();
601 unsafe { std::env::set_var(ALLOW_INSECURE_HTTP_ENV, "1") };
602
603 assert!(validate_base_url_security("http://192.168.0.110:8000/v1", false).is_ok());
604 }
605 // #5991: a provider that opts in via its [providers.<name>] table may
606 // use a plain-HTTP base URL without any env var. This is the
607 // 0.9.11-and-earlier behavior the key silently stopped providing.
608 {
609 let _lock = ALLOW_INSECURE_HTTP_ENV_LOCK.lock().unwrap();
610 let _guard = AllowInsecureHttpEnvGuard::capture();
611 unsafe { std::env::remove_var(ALLOW_INSECURE_HTTP_ENV) };
612
613 assert!(validate_base_url_security("http://192.168.0.110:8000/v1", true).is_ok());
614 // The refusal message now leads with the config key.
615 let err = validate_base_url_security("http://api.deepseek.com", false)
616 .expect_err("still refused without either opt-in");
617 assert!(err.to_string().contains("allow_insecure_http = true"));
618 }
619 }
620
621 /// Serialize tests that mutate `DEEPSEEK_ALLOW_INSECURE_HTTP`; env vars are
622 /// process-global and would otherwise leak across security checks.
623 static ALLOW_INSECURE_HTTP_ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
624
625 struct AllowInsecureHttpEnvGuard {
626 prior: Option<std::ffi::OsString>,
627 prior_legacy: Option<std::ffi::OsString>,
628 }
629 impl AllowInsecureHttpEnvGuard {
630 fn capture() -> Self {
631 let guard = Self {
632 prior: std::env::var_os(ALLOW_INSECURE_HTTP_ENV),
633 prior_legacy: std::env::var_os(LEGACY_ALLOW_INSECURE_HTTP_ENV),
634 };
635 // Clear the legacy alias so ambient shell state cannot satisfy
636 // the CODEWHALE-first fallback chain behind a test's back.
637 unsafe { std::env::remove_var(LEGACY_ALLOW_INSECURE_HTTP_ENV) };
638 guard
639 }
640 }
641 impl Drop for AllowInsecureHttpEnvGuard {
642 fn drop(&mut self) {
643 match &self.prior {
644 Some(v) => unsafe { std::env::set_var(ALLOW_INSECURE_HTTP_ENV, v) },
645 None => unsafe { std::env::remove_var(ALLOW_INSECURE_HTTP_ENV) },
646 }
647 match &self.prior_legacy {
648 Some(v) => unsafe { std::env::set_var(LEGACY_ALLOW_INSECURE_HTTP_ENV, v) },
649 None => unsafe { std::env::remove_var(LEGACY_ALLOW_INSECURE_HTTP_ENV) },
650 }
651 }
652 }
653
654 #[test]
655 fn connection_health_degrades_and_recovers() {
656 let now = Instant::now();
657 let mut health = ConnectionHealth::default();
658 assert_eq!(health.state, ConnectionState::Healthy);
659
660 apply_request_failure(&mut health, now);
661 assert_eq!(health.state, ConnectionState::Healthy);
662
663 apply_request_failure(&mut health, now + Duration::from_millis(1));
664 assert_eq!(health.state, ConnectionState::Degraded);
665 assert_eq!(health.consecutive_failures, 2);
666
667 let recovered = apply_request_success(&mut health, now + Duration::from_secs(1));
668 assert!(recovered);
669 assert_eq!(health.state, ConnectionState::Healthy);
670 assert_eq!(health.consecutive_failures, 0);
671 }
672
673 #[test]
674 fn recovery_probe_respects_cooldown() {
675 let now = Instant::now();
676 let mut health = ConnectionHealth {
677 state: ConnectionState::Degraded,
678 ..ConnectionHealth::default()
679 };
680
681 assert!(mark_recovery_probe_if_due(&mut health, now));
682 assert_eq!(health.state, ConnectionState::Recovering);
683 assert!(!mark_recovery_probe_if_due(
684 &mut health,
685 now + Duration::from_secs(1)
686 ));
687 assert!(mark_recovery_probe_if_due(
688 &mut health,
689 now + RECOVERY_PROBE_COOLDOWN + Duration::from_millis(1)
690 ));
691 }
692
693 // === #103 Phase 2: HTTP/1 escape hatch ===================================
694
695 /// Serialize tests that mutate `DEEPSEEK_FORCE_HTTP1` so they don't race
696 /// against each other — env vars are process-global.
697 pub(super) static FORCE_HTTP1_ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
698
699 struct ForceHttp1EnvGuard {
700 prior: Option<std::ffi::OsString>,
701 }
702 impl ForceHttp1EnvGuard {
703 fn capture() -> Self {
704 Self {
705 prior: std::env::var_os("DEEPSEEK_FORCE_HTTP1"),
706 }
707 }
708 }
709 impl Drop for ForceHttp1EnvGuard {
710 fn drop(&mut self) {
711 // Safety: scoped to test process; reverts to the captured value.
712 match &self.prior {
713 Some(v) => unsafe { std::env::set_var("DEEPSEEK_FORCE_HTTP1", v) },
714 None => unsafe { std::env::remove_var("DEEPSEEK_FORCE_HTTP1") },
715 }
716 }
717 }
718
719 #[tokio::test]
720 async fn configured_http2_keepalive_reaches_real_client_transport() {
721 use tokio::io::{AsyncReadExt, AsyncWriteExt};
722 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
723 let url = format!("http://{}", listener.local_addr().unwrap());
724 let config: Config = toml::from_str(
725 "[stream]\nhttp2_keep_alive_interval_secs=1\nhttp2_keep_alive_timeout_secs=1\n",
726 )
727 .unwrap();
728 let client = CodewhaleClient::http_client_builder_with_auth_mode(
729 "",
730 &HashMap::new(),
731 ProviderKind::Deepseek,
732 &url,
733 WireFormat::ChatCompletions,
734 true,
735 false,
736 &config,
737 )
738 .unwrap()
739 .no_proxy()
740 .http2_prior_knowledge()
741 .build()
742 .unwrap();
743 let request = tokio::spawn(async move { client.get(url).send().await });
744 // Minimal HTTP/2 peer: handshake, leave the request open, and observe
745 // the actual PING. No additional dependency or external provider.
746 let peer = tokio::time::timeout(Duration::from_secs(5), async {
747 let (mut peer, _) = listener.accept().await.unwrap();
748 let mut preface = [0; 24];
749 peer.read_exact(&mut preface).await.unwrap();
750 assert_eq!(&preface, b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n");
751 peer.write_all(&[0, 0, 0, 4, 0, 0, 0, 0, 0]).await.unwrap();
752 loop {
753 let mut header = [0; 9];
754 peer.read_exact(&mut header).await.unwrap();
755 let len = (usize::from(header[0]) << 16)
756 | (usize::from(header[1]) << 8)
757 | usize::from(header[2]);
758 assert!(len <= 65536, "bounded test frame");
759 let mut body = vec![0; len];
760 peer.read_exact(&mut body).await.unwrap();
761 if header[3] == 4 && header[4] & 1 == 0 {
762 peer.write_all(&[0, 0, 0, 4, 1, 0, 0, 0, 0]).await.unwrap();
763 }
764 if header[3] == 6 && header[4] & 1 == 0 {
765 assert_eq!(len, 8);
766 return peer; // Deliberately withhold the PING ACK.
767 }
768 }
769 })
770 .await
771 .expect("configured one-second interval must send a PING before the default 15 seconds");
772 let result = tokio::time::timeout(Duration::from_secs(3), request).await
773 .expect("configured one-second acknowledgement timeout must end the request before default 20 seconds")
774 .unwrap();
775 assert!(
776 result.is_err(),
777 "an unacknowledged PING must fail the request"
778 );
779 drop(peer); // Keep the peer open until the client's own timer fires.
780 }
781
782 #[test]
783 fn force_http1_scenario() {
784 // Scenario consolidation of: force_http1_unset_is_false, force_http1_truthy_values, force_http1_falsy_values
785 // from force_http1_unset_is_false
786 {
787 let _lock = FORCE_HTTP1_ENV_LOCK.lock().unwrap();
788 let _guard = ForceHttp1EnvGuard::capture();
789 unsafe { std::env::remove_var("DEEPSEEK_FORCE_HTTP1") };
790 assert!(!force_http1_from_env());
791 }
792 // from force_http1_truthy_values
793 {
794 let _lock = FORCE_HTTP1_ENV_LOCK.lock().unwrap();
795 let _guard = ForceHttp1EnvGuard::capture();
796 for value in ["1", "true", "True", "YES", "on", " 1 "] {
797 // Safety: serialized by FORCE_HTTP1_ENV_LOCK; reverted by guard.
798 unsafe { std::env::set_var("DEEPSEEK_FORCE_HTTP1", value) };
799 assert!(
800 force_http1_from_env(),
801 "{value:?} should be parsed as truthy",
802 );
803 }
804 }
805 // from force_http1_falsy_values
806 {
807 let _lock = FORCE_HTTP1_ENV_LOCK.lock().unwrap();
808 let _guard = ForceHttp1EnvGuard::capture();
809 for value in ["0", "false", "no", "off", "", "garbage", "2"] {
810 unsafe { std::env::set_var("DEEPSEEK_FORCE_HTTP1", value) };
811 assert!(
812 !force_http1_from_env(),
813 "{value:?} should NOT be parsed as truthy"
814 );
815 }
816 }
817 }
818
819 #[test]
820 fn redact_url_for_display_masks_userinfo_and_sensitive_query_values() {
821 let redacted = redact_url_for_display(
822 "https://user:secret@example.com/v1?api_key=sk-test&region=us&refresh-token=abc",
823 );
824
825 assert_eq!(
826 redacted,
827 "https://***:***@example.com/v1?api_key=***&region=us&refresh-token=***"
828 );
829 }
830
831 /// Build a DeepSeek config with an inline key/base URL plus the resolved
832 /// runtime route for it. `RouteResolver` (reached through
833 /// `resolve_runtime_route`) is the only producer of `ReadyRouteCandidate`,
834 /// so we mint candidates the same way the engine does at switch time.
835 fn deepseek_route_for_test(
836 base_url: &str,
837 model: &str,
838 ) -> (Config, crate::route_runtime::ResolvedRuntimeRoute) {
839 let config = Config {
840 provider: Some("deepseek".to_string()),
841 default_text_model: Some(model.to_string()),
842 ..Config::default()
843 }
844 .with_legacy_root(Some("ds-test".to_string()), Some(base_url.to_string()));
845 let route =
846 crate::route_runtime::resolve_runtime_route(&config, ProviderKind::Deepseek, Some(model))
847 .expect("deepseek route should resolve");
848 (config, route)
849 }
850
851 #[test]
852 fn from_candidate_scenario() {
853 // Scenario consolidation of: from_candidate_uses_candidate_base_url_and_wire_model, from_candidate_matches_new_when_config_agrees
854 // from from_candidate_uses_candidate_base_url_and_wire_model
855 {
856 let (_config, route) =
857 deepseek_route_for_test("https://route.example.com/v1", "deepseek-v4-pro");
858
859 let client = CodewhaleClient::from_candidate(&route.config, &route.candidate)
860 .expect("client should construct from candidate");
861
862 // The transport is bound to the candidate, not re-derived from Config.
863 assert_eq!(client.base_url, route.candidate.endpoint().base_url);
864 assert_eq!(
865 client.default_model,
866 route.candidate.wire_model_id().as_str()
867 );
868 }
869 // from from_candidate_matches_new_when_config_agrees
870 {
871 // For a normal route, the resolver writes the candidate's wire model and
872 // endpoint back into `route.config`, so constructing from the candidate
873 // must be byte-identical to constructing from that config. This pins the
874 // "no behavior change today" guarantee for Slice A.
875 let (_config, route) =
876 deepseek_route_for_test("https://api.deepseek.com/v1", "deepseek-v4-pro");
877
878 let from_new = CodewhaleClient::new(&route.config).expect("new client");
879 let from_candidate = CodewhaleClient::from_candidate(&route.config, &route.candidate)
880 .expect("candidate client");
881
882 assert_eq!(from_candidate.base_url, from_new.base_url);
883 assert_eq!(from_candidate.default_model, from_new.default_model);
884 assert_eq!(from_candidate.api_provider, from_new.api_provider);
885 }
886 }
887
888 fn route_cap_test_client(wire_format: WireFormat, limits: RouteLimits) -> CodewhaleClient {
889 let config = Config {
890 provider: Some("custom".to_string()),
891 default_text_model: Some("DeepSeek-V4-Flash".to_string()),
892 ..Config::default()
893 }
894 .with_legacy_root(
895 Some("route-cap-test".to_string()),
896 Some("https://route-cap.example/v1".to_string()),
897 );
898 CodewhaleClient::from_parts(
899 "https://route-cap.example/v1".to_string(),
900 "DeepSeek-V4-Flash".to_string(),
901 wire_format,
902 Some(limits),
903 &config,
904 )
905 .expect("route cap test client")
906 }
907
908 #[test]
909 fn tool_pairing_admission_uses_the_frozen_wire_protocol() {
910 for wire in [
911 WireFormat::ChatCompletions,
912 WireFormat::AnthropicMessages,
913 WireFormat::Responses,
914 ] {
915 let client = route_cap_test_client(wire, RouteLimits::default());
916 assert!(
917 client
918 .validate_tool_call_ids(["unique-a", "unique-b"])
919 .is_ok()
920 );
921 for ids in [["", "valid"], [" ", "valid"], ["same", "same"]] {
922 assert!(
923 client.validate_tool_call_ids(ids).is_err(),
924 "{wire:?}: {ids:?}"
925 );
926 }
927 let result = client.validate_tool_call_ids(["call|item-a", "call|item-b"]);
928 assert_eq!(result.is_err(), wire == WireFormat::Responses);
929 assert_eq!(
930 client.validate_tool_call_ids(["|item"]).is_err(),
931 wire == WireFormat::Responses
932 );
933 // Each response is independent: provider reuse on another round is valid.
934 assert!(client.validate_tool_call_ids(["reused"]).is_ok());
935 assert!(client.validate_tool_call_ids(["reused"]).is_ok());
936 }
937 }
938
939 #[test]
940 fn outbound_seam_excludes_host_execution_identity_for_every_dialect() {
941 let _lock = crate::test_support::lock_test_env();
942 const IMAGE: &str = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR4nGP4z8DwHwAFAAH/iZk9HQAAAABJRU5ErkJggg==";
943 for (wire, google_route) in [
944 (WireFormat::ChatCompletions, false),
945 (WireFormat::ChatCompletions, true),
946 (WireFormat::Responses, false),
947 (WireFormat::AnthropicMessages, false),
948 ] {
949 let model = if google_route {
950 "gemini-2.5-flash"
951 } else {
952 "DeepSeek-V4-Flash"
953 };
954 let client = if google_route {
955 let base_url = "https://generativelanguage.googleapis.com/v1beta/openai";
956 let config = Config {
957 provider: Some("custom".to_string()),
958 default_text_model: Some(model.to_string()),
959 ..Config::default()
960 }
961 .with_legacy_root(Some("fixture".to_string()), Some(base_url.to_string()));
962 CodewhaleClient::from_parts(
963 base_url.to_string(),
964 model.to_string(),
965 wire,
966 Some(RouteLimits::default()),
967 &config,
968 )
969 .unwrap()
970 } else {
971 route_cap_test_client(wire, RouteLimits::default())
972 };
973 let provider_id = if wire == WireFormat::Responses {
974 "wire-call|provider-item"
975 } else {
976 "wire-call"
977 };
978 let mut request =
979 translation_message_request("inspect", model.to_string(), "English", 1024);
980 request.messages.extend([
981 Message {
982 role: Role::Assistant,
983 content: vec![ContentBlock::ToolUse {
984 id: provider_id.to_string(),
985 execution_id: Some("local-execution-sentinel".to_string()),
986 name: "read".to_string(),
987 input: json!({"path": "shot.png"}),
988 caller: Some(codewhale_models::ToolCaller {
989 caller_type: "code_execution".to_string(),
990 tool_id: Some("parent-wire".to_string()),
991 }),
992 thought_signature: Some("provider-signature".to_string()),
993 }],
994 },
995 Message {
996 role: Role::User,
997 content: vec![ContentBlock::ToolResult {
998 tool_use_id: provider_id.to_string(),
999 execution_id: Some("local-execution-sentinel".to_string()),
1000 content: "captured image".to_string(),
1001 is_error: Some(false),
1002 content_blocks: Some(vec![
1003 json!({"type": "image", "mime_type": "image/png", "data": IMAGE}),
1004 ]),
1005 }],
1006 },
1007 ]);
1008 let mut legacy = request.clone();
1009 for block in legacy
1010 .messages
1011 .iter_mut()
1012 .flat_map(|message| &mut message.content)
1013 {
1014 match block {
1015 ContentBlock::ToolUse { execution_id, .. }
1016 | ContentBlock::ToolResult { execution_id, .. } => *execution_id = None,
1017 _ => {}
1018 }
1019 }
1020 for streaming in [false, true] {
1021 let prepared = client
1022 .prepare_outbound_request(request.clone(), streaming)
1023 .unwrap();
1024 let without_local = client
1025 .prepare_outbound_request(legacy.clone(), streaming)
1026 .unwrap();
1027 assert_eq!(
1028 prepared.body, without_local.body,
1029 "{wire:?}, stream={streaming}"
1030 );
1031 let bytes = prepared.body.to_string();
1032 assert!(!bytes.contains("execution_id") && !bytes.contains("local-execution-sentinel"));
1033 assert!(bytes.contains("wire-call") && bytes.contains(IMAGE));
1034 if wire == WireFormat::ChatCompletions {
1035 assert!(bytes.contains("parent-wire"));
1036 assert_eq!(
1037 bytes.contains("provider-signature"),
1038 google_route,
1039 "Google signatures remain restricted to Google's route"
1040 );
1041 }
1042 if wire == WireFormat::Responses {
1043 let input = prepared.body["input"].as_array().unwrap();
1044 assert!(
1045 input
1046 .iter()
1047 .any(|item| item["type"] == "function_call"
1048 && item["call_id"] == "wire-call")
1049 );
1050 assert!(
1051 input
1052 .iter()
1053 .any(|item| item["type"] == "function_call_output"
1054 && item["call_id"] == "wire-call")
1055 );
1056 }
1057 }
1058 assert_eq!(
1059 request.messages[1].content[0].tool_call_key(),
1060 Some(codewhale_models::ToolCallKey::Execution(
1061 "local-execution-sentinel"
1062 ))
1063 );
1064 }
1065 }
1066
1067 #[test]
1068 fn unresolved_ollama_client_waits_for_catalog_but_probe_can_bootstrap() {
1069 use codewhale_config::catalog::CatalogOffering;
1070 let _env = crate::test_support::lock_test_env();
1071 let _live = crate::provider_lake::lock_live_snapshot();
1072 let home = tempfile::tempdir().unwrap();
1073 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", home.path());
1074 crate::provider_catalog_live::reset_cache_for_test();
1075 crate::provider_lake::clear_live_snapshot();
1076 let endpoint = "http://127.0.0.1:11452/v1";
1077 let mut config = Config {
1078 provider: Some("ollama".into()),
1079 ..Default::default()
1080 };
1081 config
1082 .provider_config_for_mut(&config.test_identity_for_kind(ProviderKind::Ollama))
1083 .unwrap()
1084 .base_url = Some(endpoint.into());
1085 assert!(
1086 CodewhaleClient::new(&config).is_err(),
1087 "unknown must not become a dispatch model"
1088 );
1089 let probe = CodewhaleClient::for_catalog_refresh(&config).expect("catalog bootstrap client");
1090 assert_eq!(probe.base_url, endpoint);
1091 let ticket = crate::provider_catalog_live::begin_refresh_for_identity(
1092 ProviderKind::Ollama,
1093 "ollama",
1094 endpoint,
1095 );
1096 let fingerprint = base_url_fingerprint(endpoint);
1097 let fetched_at = now_unix();
1098 crate::provider_catalog_live::record_success_if_current(
1099 &ticket,
1100 ProviderCatalogDelta {
1101 provider: "ollama".into(),
1102 base_url_fingerprint: fingerprint.clone(),
1103 fetched_at,
1104 offerings: ["zeta:tag", "alpha:tag"]
1105 .into_iter()
1106 .map(|id| CatalogOffering {
1107 provider: "ollama".into(),
1108 wire_model_id: id.into(),
1109 endpoint_key: "chat".into(),
1110 source: CatalogSource::Live {
1111 base_url_fingerprint: fingerprint.clone(),
1112 fetched_at,
1113 },
1114 ..Default::default()
1115 })
1116 .collect(),
1117 },
1118 );
1119 let client = CodewhaleClient::new(&config).expect("fresh local model client");
1120 assert_eq!(client.default_model, "alpha:tag");
1121 let route =
1122 crate::route_runtime::resolve_runtime_route(&config, ProviderKind::Ollama, None).unwrap();
1123 assert_eq!(
1124 client.default_model,
1125 route.candidate.wire_model_id().as_str()
1126 );
1127 config
1128 .set_provider_model_override(
1129 &config.test_identity_for_kind(ProviderKind::Ollama),
1130 Some("saved:tag".into()),
1131 )
1132 .unwrap();
1133 assert_eq!(
1134 CodewhaleClient::new(&config).unwrap().default_model,
1135 "saved:tag"
1136 );
1137 crate::provider_catalog_live::reset_cache_for_test();
1138 crate::provider_lake::clear_live_snapshot();
1139 }
1140
1141 #[test]
1142 fn provider_regression_5820_ollama_config_reaches_the_wire_with_safe_output() {
1143 let _lock = crate::test_support::lock_test_env();
1144 let _canonical = crate::test_support::EnvVarGuard::remove("CODEWHALE_MAX_OUTPUT_TOKENS");
1145 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_MAX_OUTPUT_TOKENS");
1146 let config = Config {
1147 provider: Some("ollama".to_string()),
1148 providers: Some(ProvidersConfig {
1149 ollama: ProviderConfig {
1150 model: Some("qwen2.5:7b".to_string()),
1151 context_window: Some(32_768),
1152 base_url: Some("http://127.0.0.1:11434".to_string()),
1153 ..Default::default()
1154 },
1155 ..Default::default()
1156 }),
1157 ..Default::default()
1158 };
1159 let client = CodewhaleClient::new(&config).unwrap();
1160 assert_eq!(
1161 client
1162 .route_limits()
1163 .and_then(|limits| limits.context_tokens),
1164 Some(32_768)
1165 );
1166 assert_eq!(client.effective_max_output_tokens("qwen2.5:7b"), 8_192);
1167 let prepared = client
1168 .prepare_outbound_request(
1169 MessageRequest {
1170 model: "qwen2.5:7b".to_string(),
1171 messages: vec![Message {
1172 role: Role::User,
1173 content: vec![ContentBlock::Text {
1174 text: "hello".to_string(),
1175 cache_control: None,
1176 }],
1177 }],
1178 max_tokens: 64_000,
1179 system: None,
1180 tools: None,
1181 tool_choice: None,
1182 metadata: None,
1183 thinking: None,
1184 reasoning_effort: None,
1185 stream: Some(false),
1186 temperature: None,
1187 top_p: None,
1188 },
1189 false,
1190 )
1191 .unwrap();
1192 assert_eq!(prepared.body["max_tokens"], 8_192);
1193 }
1194
1195 #[test]
1196 fn outbound_seam_clamps_every_dialect_to_the_exact_route_envelope() {
1197 let _lock = crate::test_support::lock_test_env();
1198 let _canonical = crate::test_support::EnvVarGuard::set("CODEWHALE_MAX_OUTPUT_TOKENS", "384000");
1199 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_MAX_OUTPUT_TOKENS");
1200
1201 for (limits, expected) in [
1202 (
1203 RouteLimits {
1204 context_tokens: Some(327_680),
1205 ..RouteLimits::default()
1206 },
1207 325_632_u64,
1208 ),
1209 (
1210 RouteLimits {
1211 context_tokens: Some(327_680),
1212 output_tokens: Some(100_000),
1213 ..RouteLimits::default()
1214 },
1215 100_000,
1216 ),
1217 (
1218 RouteLimits {
1219 context_tokens: Some(327_680),
1220 output_tokens: Some(128),
1221 ..RouteLimits::default()
1222 },
1223 128,
1224 ),
1225 ] {
1226 for (wire_format, body_field) in [
1227 (WireFormat::ChatCompletions, "max_tokens"),
1228 (WireFormat::Responses, "max_output_tokens"),
1229 (WireFormat::AnthropicMessages, "max_tokens"),
1230 ] {
1231 let client = route_cap_test_client(wire_format, limits);
1232 let prepared = client
1233 .prepare_outbound_request(
1234 MessageRequest {
1235 model: "DeepSeek-V4-Flash".to_string(),
1236 messages: vec![Message {
1237 role: Role::User,
1238 content: vec![ContentBlock::Text {
1239 text: "route cap".to_string(),
1240 cache_control: None,
1241 }],
1242 }],
1243 max_tokens: 384_000,
1244 system: None,
1245 tools: None,
1246 tool_choice: None,
1247 metadata: None,
1248 thinking: None,
1249 reasoning_effort: Some("max".to_string()),
1250 stream: Some(false),
1251 temperature: None,
1252 top_p: None,
1253 },
1254 false,
1255 )
1256 .expect("request prepares through preview/wire seam");
1257 assert_eq!(
1258 prepared.body[body_field].as_u64(),
1259 Some(expected),
1260 "wire={wire_format:?} limits={limits:?}"
1261 );
1262 }
1263 }
1264 }
1265
1266 #[test]
1267 fn same_protocol_model_switch_rebinds_exact_candidate_identity_and_limits() {
1268 let _lock = crate::test_support::lock_test_env();
1269 let _canonical = crate::test_support::EnvVarGuard::set("CODEWHALE_MAX_OUTPUT_TOKENS", "384000");
1270 let config = Config {
1271 provider: Some("openrouter".to_string()),
1272 providers: Some(ProvidersConfig {
1273 openrouter: ProviderConfig {
1274 api_key: Some("openrouter-route-cap-test".to_string()),
1275 base_url: Some("https://openrouter.ai/api/v1".to_string()),
1276 model: Some("deepseek/deepseek-v4-pro".to_string()),
1277 ..ProviderConfig::default()
1278 },
1279 ..ProvidersConfig::default()
1280 }),
1281 ..Config::default()
1282 };
1283 let client = CodewhaleClient::new(&config).expect("OpenRouter client resolves");
1284 assert_eq!(client.wire_format, WireFormat::ChatCompletions);
1285 assert!(client.route_limits.is_some());
1286
1287 let rebound = client
1288 .rebound_for_model_protocol(Some(&config), OPENROUTER_QWEN_3_6_FLASH_MODEL)
1289 .expect("same-protocol alternate route resolves")
1290 .expect("model/limit identity change requires a rebound");
1291 assert_eq!(rebound.wire_format, WireFormat::ChatCompletions);
1292 assert_eq!(rebound.default_model, OPENROUTER_QWEN_3_6_FLASH_MODEL);
1293 assert_ne!(rebound.route_limits, client.route_limits);
1294
1295 let pro_cap = client.effective_max_output_tokens("deepseek/deepseek-v4-pro");
1296 let alternate_cap = client.effective_max_output_tokens(OPENROUTER_QWEN_3_6_FLASH_MODEL);
1297 assert!(
1298 alternate_cap < pro_cap,
1299 "fixture must prove a smaller same-protocol alternate route: pro={pro_cap}, alternate={alternate_cap}"
1300 );
1301 assert_eq!(
1302 alternate_cap,
1303 rebound.effective_max_output_tokens(OPENROUTER_QWEN_3_6_FLASH_MODEL),
1304 "the original bound client and rebound client must resolve the same alternate envelope"
1305 );
1306 let prepared = client
1307 .prepare_outbound_request(
1308 MessageRequest {
1309 model: OPENROUTER_QWEN_3_6_FLASH_MODEL.to_string(),
1310 messages: vec![Message {
1311 role: Role::User,
1312 content: vec![ContentBlock::Text {
1313 text: "alternate route cap".to_string(),
1314 cache_control: None,
1315 }],
1316 }],
1317 max_tokens: 384_000,
1318 system: None,
1319 tools: None,
1320 tool_choice: None,
1321 metadata: None,
1322 thinking: None,
1323 reasoning_effort: Some("max".to_string()),
1324 stream: Some(false),
1325 temperature: None,
1326 top_p: None,
1327 },
1328 false,
1329 )
1330 .expect("same-protocol alternate prepares");
1331 assert_eq!(prepared.body["max_tokens"], json!(alternate_cap));
1332 }
1333
1334 #[allow(clippy::await_holding_lock)]
1335 #[tokio::test(flavor = "current_thread")]
1336 async fn fim_non_message_request_is_clamped_to_bound_route() {
1337 let _lock = crate::test_support::lock_test_env();
1338 let _canonical = crate::test_support::EnvVarGuard::set("CODEWHALE_MAX_OUTPUT_TOKENS", "384000");
1339 let server = MockServer::start().await;
1340 Mock::given(method("POST"))
1341 .and(path("/beta/completions"))
1342 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
1343 "choices": [{"text": "middle"}]
1344 })))
1345 .expect(1)
1346 .mount(&server)
1347 .await;
1348 let base_url = format!("{}/v1", server.uri());
1349 let config = Config {
1350 provider: Some("custom".to_string()),
1351 default_text_model: Some("local-fim".to_string()),
1352 ..Config::default()
1353 }
1354 .with_legacy_root(Some("fim-cap-test".to_string()), Some(base_url.clone()));
1355 let client = CodewhaleClient::from_parts(
1356 base_url,
1357 "local-fim".to_string(),
1358 WireFormat::ChatCompletions,
1359 Some(RouteLimits {
1360 context_tokens: Some(327_680),
1361 output_tokens: Some(128),
1362 ..RouteLimits::default()
1363 }),
1364 &config,
1365 )
1366 .expect("FIM route client");
1367
1368 assert_eq!(
1369 client
1370 .fim_completion("local-fim", "prefix", "suffix", 4_096)
1371 .await
1372 .expect("FIM response"),
1373 "middle"
1374 );
1375 let requests = server
1376 .received_requests()
1377 .await
1378 .expect("recorded FIM request");
1379 let body: Value = serde_json::from_slice(&requests[0].body).expect("FIM JSON");
1380 assert_eq!(body["max_tokens"], json!(128));
1381 }
1382
1383 #[test]
1384 fn official_deepseek_flash_binds_responses_request_and_endpoint() {
1385 let (_config, route) =
1386 deepseek_route_for_test("https://api.deepseek.com/beta", "deepseek-v4-flash");
1387 assert_eq!(route.candidate.protocol(), WireFormat::Responses);
1388
1389 let client = CodewhaleClient::new(&route.config).expect("Flash client resolves");
1390 assert_eq!(client.wire_format, WireFormat::Responses);
1391
1392 let prepared = client
1393 .prepare_outbound_request(
1394 MessageRequest {
1395 model: "deepseek-v4-flash".to_string(),
1396 messages: vec![Message {
1397 role: Role::User,
1398 content: vec![ContentBlock::Text {
1399 text: "hello".to_string(),
1400 cache_control: None,
1401 }],
1402 }],
1403 max_tokens: 64,
1404 system: None,
1405 tools: None,
1406 tool_choice: None,
1407 metadata: None,
1408 thinking: None,
1409 reasoning_effort: Some("max".to_string()),
1410 stream: Some(true),
1411 temperature: None,
1412 top_p: None,
1413 },
1414 true,
1415 )
1416 .expect("Flash Responses request prepares");
1417
1418 assert_eq!(prepared.dialect, WireDialect::OpenAiResponses);
1419 assert_eq!(prepared.endpoint.url, "https://api.deepseek.com/responses");
1420 assert_eq!(prepared.body["model"], "deepseek-v4-flash");
1421 assert_eq!(prepared.body["reasoning"]["effort"], "max");
1422 }
1423
1424 #[test]
1425 fn uncatalogued_deepseek_preview_binds_chat_request_without_model_fallback() {
1426 let model = "deepseek-v4.1-flash-expires-on-0910";
1427 let (_config, route) = deepseek_route_for_test("https://api.deepseek.com", model);
1428 assert_eq!(route.candidate.protocol(), WireFormat::ChatCompletions);
1429 assert_eq!(route.candidate.wire_model_id().as_str(), model);
1430 assert!(route.candidate.canonical_model().is_none());
1431
1432 for client in [
1433 CodewhaleClient::new(&route.config).expect("preview client resolves"),
1434 CodewhaleClient::from_candidate(&route.config, &route.candidate)
1435 .expect("preview client binds the admitted candidate"),
1436 ] {
1437 assert_eq!(client.wire_format, WireFormat::ChatCompletions);
1438 assert_eq!(client.default_model, model);
1439 let prepared = client
1440 .prepare_outbound_request(
1441 MessageRequest {
1442 model: model.to_string(),
1443 messages: vec![Message {
1444 role: Role::User,
1445 content: vec![ContentBlock::Text {
1446 text: "hello".to_string(),
1447 cache_control: None,
1448 }],
1449 }],
1450 max_tokens: 64,
1451 system: None,
1452 tools: None,
1453 tool_choice: None,
1454 metadata: None,
1455 thinking: None,
1456 reasoning_effort: None,
1457 stream: Some(true),
1458 temperature: None,
1459 top_p: None,
1460 },
1461 true,
1462 )
1463 .expect("preview Chat request prepares");
1464 assert_eq!(prepared.dialect, WireDialect::ChatCompletions);
1465 assert_eq!(
1466 prepared.endpoint.url,
1467 "https://api.deepseek.com/v1/chat/completions"
1468 );
1469 assert_eq!(prepared.body["model"], model);
1470 assert_eq!(prepared.body["messages"][0]["content"], "hello");
1471 }
1472 }
1473
1474 #[test]
1475 fn exact_catalog_deepseek_preview_binding_survives_prepare_and_same_model_rebind() {
1476 use codewhale_config::route::{
1477 PricingSku, ProviderId, RouteCapabilities, WireModelId, offering::ProviderModelOffering,
1478 };
1479
1480 let model = "deepseek-v4.1-flash-expires-on-0910";
1481 // Synthetic catalog evidence: the offline catalog has no such row.
1482 // This fixture does not assert that the real preview supports Responses.
1483 let resolver = RouteResolver::from_offerings(vec![ProviderModelOffering {
1484 provider: ProviderId::from("deepseek"),
1485 canonical_model: None,
1486 wire_model_id: WireModelId::from(model),
1487 endpoint_key: "responses".to_string(),
1488 default_for_provider: false,
1489 limits: RouteLimits {
1490 output_tokens: Some(777),
1491 ..Default::default()
1492 },
1493 capabilities: RouteCapabilities::default(),
1494 pricing: PricingSku::UnknownOrStale,
1495 }]);
1496 let candidate = resolver
1497 .resolve(&RouteRequest {
1498 explicit_provider: Some(codewhale_config::ProviderKind::Deepseek),
1499 model_selector: Some(LogicalModelRef::from(model)),
1500 base_url_override: Some("https://api.deepseek.com".to_string()),
1501 ..Default::default()
1502 })
1503 .expect("exact synthetic catalog offering resolves");
1504 assert_eq!(candidate.protocol(), WireFormat::Responses);
1505 assert_eq!(candidate.wire_model_id().as_str(), model);
1506
1507 let config = Config {
1508 provider: Some("deepseek".to_string()),
1509 default_text_model: Some(model.to_string()),
1510 ..Default::default()
1511 }
1512 .with_legacy_root(
1513 Some("ds-test".to_string()),
1514 Some("https://api.deepseek.com".to_string()),
1515 );
1516 let client = CodewhaleClient::from_candidate(&config, &candidate)
1517 .expect("client binds exact synthetic catalog offering");
1518 assert!(
1519 client
1520 .rebound_for_model_protocol(None, model)
1521 .expect("the same admitted model needs no offline lookup")
1522 .is_none()
1523 );
1524 let prepared = client
1525 .prepare_outbound_request(
1526 translation_message_request("hello", model.to_string(), "English", 4_096),
1527 true,
1528 )
1529 .expect("same-model request preserves the exact catalog binding");
1530 assert_eq!(prepared.dialect, WireDialect::OpenAiResponses);
1531 assert_eq!(prepared.endpoint.url, "https://api.deepseek.com/responses");
1532 assert_eq!(prepared.body["model"], model);
1533 assert_eq!(prepared.body["max_output_tokens"], 777);
1534
1535 let other_model = "deepseek-v4-pro";
1536 assert!(
1537 client
1538 .prepare_outbound_request(
1539 translation_message_request("hello", other_model.to_string(), "English", 4_096,),
1540 true,
1541 )
1542 .is_err(),
1543 "a different model must still obey the protocol switch guard"
1544 );
1545 let rebound = client
1546 .rebound_for_model_protocol(Some(&config), other_model)
1547 .expect("different model resolves independently")
1548 .expect("Pro requires a Chat client");
1549 let prepared = rebound
1550 .prepare_outbound_request(
1551 translation_message_request("hello", other_model.to_string(), "English", 4_096),
1552 true,
1553 )
1554 .expect("rebound Pro prepares Chat");
1555 assert_eq!(prepared.dialect, WireDialect::ChatCompletions);
1556 assert_eq!(
1557 prepared.endpoint.url,
1558 "https://api.deepseek.com/v1/chat/completions"
1559 );
1560 assert_eq!(prepared.body["model"], other_model);
1561 }
1562
1563 #[test]
1564 fn rebinding_a_chat_bound_client_for_flash_switches_to_responses() {
1565 // #5042: fleet dispatch binds the child client before the profile
1566 // model is resolved; a chat-bound DeepSeek client asked to run flash
1567 // must be rebuilt on the Responses protocol by the central resolver
1568 // instead of failing deterministically at first send.
1569 let (_config, route) =
1570 deepseek_route_for_test("https://api.deepseek.com/beta", "deepseek-v4-pro");
1571 let client = CodewhaleClient::new(&route.config).expect("pro client resolves");
1572 assert_eq!(client.wire_format, WireFormat::ChatCompletions);
1573
1574 let rebound = client
1575 .rebound_for_model_protocol(Some(&route.config), "deepseek-v4-flash")
1576 .expect("flash rebind resolves")
1577 .expect("flash requires a different protocol");
1578 assert_eq!(rebound.wire_format, WireFormat::Responses);
1579 assert_eq!(rebound.default_model, "deepseek-v4-flash");
1580
1581 assert!(
1582 client
1583 .rebound_for_model_protocol(Some(&route.config), "deepseek-v4-pro")
1584 .expect("pro rebind resolves")
1585 .is_none(),
1586 "a matching protocol must not rebuild the client"
1587 );
1588 }
1589
1589 lines RUST