| 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®ion=us&refresh-token=abc", |
| 823 | ); |
| 824 | |
| 825 | assert_eq!( |
| 826 | redacted, |
| 827 | "https://***:***@example.com/v1?api_key=***®ion=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 |