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