返回 CodeWhale
client.rs
根目录 / crates / tui / src / client.rs
1 //! HTTP client for the resolved provider route.
2 //!
3 //! Routes reach the provider through its OpenAI-compatible or native wire
4 //! surface; `/chat/completions` is the common primary endpoint.
5
6 use std::collections::HashMap;
7 use std::sync::atomic::{AtomicUsize, Ordering};
8 use std::sync::{Arc, Mutex as StdMutex, OnceLock};
9 use std::time::{Duration, Instant};
10
11 use anyhow::{Context, Result, bail};
12 use base64::{Engine as _, engine::general_purpose};
13 use futures_util::StreamExt;
14 use reqwest::header::{AUTHORIZATION, CONTENT_TYPE, HeaderMap, HeaderName, HeaderValue};
15 use serde::{Deserialize, Serialize};
16 use serde_json::{Value, json};
17 use tokio::sync::{
18 Mutex as AsyncMutex, OwnedRwLockReadGuard, OwnedRwLockWriteGuard, OwnedSemaphorePermit, RwLock,
19 Semaphore,
20 };
21
22 use codewhale_config::catalog::{
23 CatalogOffering, CatalogRefreshError, CatalogSource, ProviderCatalogDelta,
24 base_url_fingerprint, now_unix,
25 };
26 #[cfg(test)]
27 use codewhale_config::catalog::{CatalogSnapshot, CatalogStatus, ProviderCatalogCache};
28 use codewhale_config::provider::WireFormat;
29 use codewhale_config::route::{
30 LogicalModelRef, ReadyRouteCandidate, RouteLimits, RouteRequest, RouteResolver,
31 };
32 use codewhale_config::{auth_mode_disables_api_key, is_upstream_auth_header};
33
34 use crate::config::{
35 Config, ProviderIdentity, ProviderKind, RetryPolicy, validate_route,
36 wire_model_for_provider_route,
37 };
38 use crate::llm_client::{
39 LlmClient, LlmError, RetryConfig as LlmRetryConfig, extract_retry_after,
40 sanitize_http_error_body, with_retry,
41 };
42 #[cfg(test)]
43 #[path = "client/catalog_tests.rs"]
44 mod catalog_tests;
45
46 use crate::logging;
47 use codewhale_models::Role;
48 use codewhale_models::{
49 ContentBlock, Message, MessageRequest, MessageResponse, SystemPrompt, Usage,
50 };
51
52 /// Every provider request that can feed the interactive TUI's attached CWC run
53 /// takes a shared permit at this lowest common dispatch seam. Runtime Chat holds
54 /// the exclusive permit from native admission through durable terminal
55 /// acknowledgement. This covers ordinary turns, auto-route classification,
56 /// advisor calls, detached subagents, compaction, purge, and streaming without
57 /// serializing independent RuntimeThreadManager stores.
58 #[cfg(not(test))]
59 static RUNTIME_CHAT_INFERENCE_GATE: OnceLock<Arc<RwLock<()>>> = OnceLock::new();
60
61 /// Unit tests run many independent Tokio runtimes in one process. Keying the
62 /// otherwise-identical gate by runtime keeps unrelated libtest cases from
63 /// manufacturing contention while preserving exact read/write behavior among
64 /// tasks on the same current-thread or multi-thread runtime.
65 #[cfg(test)]
66 static RUNTIME_CHAT_INFERENCE_TEST_GATES: OnceLock<
67 StdMutex<HashMap<RuntimeChatInferenceTestScope, std::sync::Weak<RwLock<()>>>>,
68 > = OnceLock::new();
69
70 #[cfg(test)]
71 #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
72 enum RuntimeChatInferenceTestScope {
73 Runtime(tokio::runtime::Id),
74 Thread(std::thread::ThreadId),
75 }
76
77 #[cfg(not(test))]
78 fn runtime_chat_inference_gate() -> Arc<RwLock<()>> {
79 Arc::clone(RUNTIME_CHAT_INFERENCE_GATE.get_or_init(|| Arc::new(RwLock::new(()))))
80 }
81
82 #[cfg(test)]
83 fn runtime_chat_inference_gate() -> Arc<RwLock<()>> {
84 let scope = tokio::runtime::Handle::try_current()
85 .map(|handle| RuntimeChatInferenceTestScope::Runtime(handle.id()))
86 .unwrap_or_else(|_| RuntimeChatInferenceTestScope::Thread(std::thread::current().id()));
87 let mut gates = RUNTIME_CHAT_INFERENCE_TEST_GATES
88 .get_or_init(|| StdMutex::new(HashMap::new()))
89 .lock()
90 .unwrap_or_else(std::sync::PoisonError::into_inner);
91 if let Some(gate) = gates.get(&scope).and_then(std::sync::Weak::upgrade) {
92 return gate;
93 }
94 let gate = Arc::new(RwLock::new(()));
95 gates.insert(scope, Arc::downgrade(&gate));
96 gate
97 }
98
99 /// Exclusive ownership retained by isolated Runtime Chat.
100 pub(crate) struct RuntimeChatInferenceOwnership {
101 _gate: OwnedRwLockWriteGuard<()>,
102 }
103
104 /// Shared participation retained for the full provider-output lifecycle.
105 pub(crate) struct RemoteControlInferencePermit {
106 _gate: OwnedRwLockReadGuard<()>,
107 }
108
109 #[cfg(test)]
110 pub(crate) async fn acquire_runtime_chat_inference_ownership() -> RuntimeChatInferenceOwnership {
111 let gate = runtime_chat_inference_gate().write_owned().await;
112 RuntimeChatInferenceOwnership { _gate: gate }
113 }
114
115 pub(crate) fn try_acquire_runtime_chat_inference_ownership() -> Option<RuntimeChatInferenceOwnership>
116 {
117 let gate = runtime_chat_inference_gate().try_write_owned().ok()?;
118 Some(RuntimeChatInferenceOwnership { _gate: gate })
119 }
120
121 /// Join the attached interactive CWC run as a provider-output participant.
122 /// Standalone background adapters that bypass `CodewhaleClient::create_message`
123 /// use this guard and retain it through response decoding.
124 pub(crate) async fn acquire_remote_control_inference_participant() -> RemoteControlInferencePermit {
125 let gate = runtime_chat_inference_gate().read_owned().await;
126 RemoteControlInferencePermit { _gate: gate }
127 }
128
129 pub(super) fn to_api_tool_name(name: &str) -> String {
130 let mut out = String::new();
131 for ch in name.chars() {
132 if ch.is_ascii_alphanumeric() || ch == '_' {
133 out.push(ch);
134 } else if ch == '-' {
135 out.push_str("--");
136 } else {
137 out.push_str("-x");
138 out.push_str(&format!("{:06X}", ch as u32));
139 out.push('-');
140 }
141 }
142 out
143 }
144
145 pub(super) fn from_api_tool_name(name: &str) -> String {
146 let mut out = String::new();
147 let mut iter = name.chars().peekable();
148 while let Some(ch) = iter.next() {
149 if ch != '-' {
150 out.push(ch);
151 continue;
152 }
153 if let Some('-') = iter.peek().copied() {
154 iter.next();
155 out.push('-');
156 continue;
157 }
158 if iter.peek().copied() == Some('x') {
159 iter.next();
160 let mut hex = String::new();
161 for _ in 0..6 {
162 if let Some(h) = iter.next() {
163 hex.push(h);
164 } else {
165 break;
166 }
167 }
168 // Only decode if we got exactly 6 hex digits (matching encoder output).
169 // Fewer digits means a truncated/malformed sequence — pass through as-is.
170 if hex.len() == 6
171 && let Ok(code) = u32::from_str_radix(&hex, 16)
172 && let Some(decoded) = std::char::from_u32(code)
173 {
174 if let Some('-') = iter.peek().copied() {
175 iter.next();
176 }
177 out.push(decoded);
178 continue;
179 }
180 out.push('-');
181 out.push('x');
182 out.push_str(&hex);
183 continue;
184 }
185 out.push('-');
186 }
187
188 // Second pass: decode bare hex escapes (e.g. `x00002E`) that the model
189 // may produce when it mangles the `-x00002E-` delimiter form. Only
190 // decode when the resulting character is one that `to_api_tool_name`
191 // would have encoded (not alphanumeric, not `_`, not `-`).
192 decode_bare_hex_escapes(&out)
193 }
194
195 /// Decode bare `x[0-9A-Fa-f]{6}` sequences (optionally followed by `-`)
196 /// that survive the standard delimiter-based pass. This handles cases
197 /// where the model strips or replaces the leading `-` of `-x00002E-`.
198 pub(super) fn decode_bare_hex_escapes(input: &str) -> String {
199 use regex::Regex;
200 use std::sync::OnceLock;
201
202 static RE: OnceLock<Regex> = OnceLock::new();
203 let re = RE.get_or_init(|| Regex::new(r"x([0-9A-Fa-f]{6})-?").unwrap());
204
205 let result = re.replace_all(input, |caps: &regex::Captures| {
206 let hex = &caps[1];
207 if let Ok(code) = u32::from_str_radix(hex, 16)
208 && let Some(decoded) = std::char::from_u32(code)
209 {
210 // Only decode characters that to_api_tool_name would have encoded
211 if !decoded.is_ascii_alphanumeric() && decoded != '_' && decoded != '-' {
212 return decoded.to_string();
213 }
214 }
215 // Not a character we'd encode — leave as-is
216 caps[0].to_string()
217 });
218 result.into_owned()
219 }
220
221 // === Types ===
222
223 /// Model descriptor returned by the provider's `/v1/models` endpoint.
224 #[derive(Debug, Clone, Serialize, PartialEq, Eq)]
225 pub struct AvailableModel {
226 pub id: String,
227 pub owned_by: Option<String>,
228 pub created: Option<u64>,
229 #[serde(skip_serializing_if = "Option::is_none")]
230 pub display_name: Option<String>,
231 }
232
233 /// Request payload for Xiaomi MiMo speech synthesis models.
234 ///
235 /// MiMo-V2.5-TTS / MiMo-V2-TTS use the OpenAI-compatible
236 /// `/v1/chat/completions` endpoint: the optional style/voice instruction is
237 /// sent as a `user` message, while the text to synthesize is sent as an
238 /// `assistant` message.
239 #[derive(Debug, Clone)]
240 pub struct SpeechSynthesisRequest {
241 pub model: String,
242 pub text: String,
243 pub instruction: Option<String>,
244 pub audio_format: String,
245 pub voice: Option<String>,
246 }
247
248 /// Decoded speech synthesis result.
249 #[derive(Debug, Clone)]
250 pub struct SpeechSynthesisResponse {
251 pub model: String,
252 pub audio_format: String,
253 pub audio_bytes: Vec<u8>,
254 pub transcript: Option<String>,
255 pub voice: Option<String>,
256 }
257
258 /// One decoded provider response from the auxiliary translation path.
259 ///
260 /// The immutable route and provider-reported usage travel beside the semantic
261 /// translation result so callers can account for a successful provider call
262 /// before rejecting an incomplete, empty, or otherwise unusable translation.
263 /// `usage == None` is distinct from a transport failure: the provider returned
264 /// a response, but omitted the receipt needed to price it exactly.
265 pub(crate) struct TranslationProviderResponse {
266 pub(crate) translated: Result<String>,
267 pub(crate) route: crate::cost_status::EffectiveRouteEnvelope,
268 pub(crate) usage: Option<Usage>,
269 }
270
271 /// Universal client for the resolved provider route.
272 #[must_use]
273 pub struct CodewhaleClient {
274 pub(super) http_client: reqwest::Client,
275 // Catalogs and probes must never forward frozen custom auth headers to a
276 // provider-supplied redirect destination. Inference keeps its own policy.
277 models_http_client: reqwest::Client,
278 /// HTTP/1.1-only twin of [`Self::http_client`], used for automatic
279 /// stream-header fallback when H2 stalls. Same auth and headers.
280 pub(super) http1_client: reqwest::Client,
281 api_key: String,
282 /// Where `api_key` came from (secret store slot, config file, env var
283 /// name, CLI, OAuth, …), named in authentication errors (#6528).
284 api_key_source: String,
285 /// For a subscription sign-in route (ChatGPT, xAI OAuth): which account
286 /// is signed in and how to switch, appended to plan-quota errors. Holds
287 /// an account label only, never token material.
288 subscription_limit_guidance: Option<String>,
289 /// Exact configured credential values removed from model-bound tool
290 /// results. Structural redaction handles config/JSON assignments, while
291 /// this list closes the gap for bare provider tokens with no recognizable
292 /// prefix (for example token-plan and provider-specific keys).
293 model_bound_secret_values: Arc<Vec<String>>,
294 /// Exact values a catalog endpoint could echo back at us: the active API
295 /// key and every user-configured HTTP header value, plus everything in
296 /// `model_bound_secret_values`.
297 ///
298 /// Deliberately a second, wider list rather than a widening of that one:
299 /// they answer different questions at different trust boundaries. That
300 /// list is "what must never reach a *model*", which is why it covers only
301 /// auth-shaped headers and values of at least
302 /// `MIN_EXACT_SECRET_CHARS`. This one is "what must never come back out
303 /// of an endpoint the user typed during setup" (#6173), where a
304 /// three-character key is still a key and a custom header the user
305 /// configured is still theirs.
306 catalog_error_secret_values: Arc<Vec<String>>,
307 /// Whether credential-shaped tool output is masked before it is sent to an
308 /// upstream model. The safe default is `true`; it is `false` only after the
309 /// user disabled `[redaction] model_bound` and confirmed the opt-out on the
310 /// startup gate (see [`codewhale_config::redaction`]). Routing/classification
311 /// summaries and durable goal-state text keep their own always-on redaction
312 /// regardless of this flag.
313 model_bound_masking: bool,
314 pub(super) base_url: String,
315 pub(super) api_provider: ProviderKind,
316 /// Exact configured provider identity and billing mode frozen when this
317 /// client is built. Child/tool calls only carry the client at dispatch, so
318 /// these route facts must travel with it instead of being reconstructed
319 /// from the mutable parent session at completion time.
320 admitted_identity: ProviderIdentity,
321 openrouter_vendor: Option<String>,
322 billing_surface: Option<String>,
323 billing_mode: crate::cost_status::RouteBillingMode,
324 configured_models: Arc<Vec<codewhale_config::catalog::configured::ConfiguredModel>>,
325 /// Non-secret limits frozen from the same resolved candidate as the
326 /// endpoint and wire model. Auxiliary calls carry only this client, so
327 /// they must not reconstruct output caps with `None` and discard a custom
328 /// route's context window.
329 route_limits: Option<RouteLimits>,
330 /// Verified ChatGPT subject captured with Codewhale's own selected grant.
331 pub(super) codex_account_id: Option<String>,
332 /// Opaque reasoning is bound to this verified grant identity, stable across
333 /// token refresh and absent on custom API-key routes.
334 pub(super) chatgpt_reasoning_api: Option<String>,
335 wire_format: WireFormat,
336 retry: RetryPolicy,
337 /// Auxiliary inspection calls use the normal bounded retry schedule but
338 /// never publish retry/rate-limit state into process-global UI cells.
339 isolated_request_state: bool,
340 /// Whether this concrete client can contribute provider output to the
341 /// interactive TUI's attached CWC run. Isolated Runtime Chat executes under
342 /// the host's exclusive permit; independent RuntimeThreadManager stores do
343 /// not share that run and remain concurrent.
344 remote_control_inference_participant: bool,
345 default_model: String,
346 connection_health: Arc<AsyncMutex<ConnectionHealth>>,
347 rate_limiter: Arc<AsyncMutex<TokenBucket>>,
348 request_concurrency: Option<ProviderConcurrencyLimiter>,
349 path_suffix: Option<String>,
350 /// Unit tests keep the semantic route exact while sending the actual
351 /// production request through a local capture server. This field is
352 /// compiled out of release builds.
353 #[cfg(test)]
354 test_chat_transport_base_url: Option<String>,
355 /// Messages equivalent of `test_chat_transport_base_url`; keeps exact
356 /// route shaping bound to the semantic endpoint while tests capture on a
357 /// local server.
358 #[cfg(test)]
359 test_messages_transport_base_url: Option<String>,
360 pub(super) reasoning_stream_style: Option<String>,
361 pub(super) stream_idle_timeout: Duration,
362 /// Bounded wait for SSE response headers, resolved once from
363 /// `[stream].open_timeout_secs`, legacy `[tui]`, or the environment fallback.
364 pub(super) stream_open_timeout: Duration,
365 /// HTTP/1.1 pin resolved once from `Config::force_http1` (#6700); the
366 /// single source every client builder and stream open reads.
367 pub(super) force_http1: bool,
368 }
369
370 const CONNECTION_FAILURE_THRESHOLD: u32 = 2;
371 const RECOVERY_PROBE_COOLDOWN: Duration = Duration::from_secs(15);
372
373 const DEFAULT_CLIENT_RATE_LIMIT_RPS: f64 = 8.0;
374 const DEFAULT_CLIENT_RATE_LIMIT_BURST: f64 = 16.0;
375 const ALLOW_INSECURE_HTTP_ENV: &str = "CODEWHALE_ALLOW_INSECURE_HTTP";
376 /// Legacy alias for [`ALLOW_INSECURE_HTTP_ENV`].
377 const LEGACY_ALLOW_INSECURE_HTTP_ENV: &str = "DEEPSEEK_ALLOW_INSECURE_HTTP";
378
379 fn client_user_agent(_api_provider: ProviderKind) -> &'static str {
380 concat!(
381 "Mozilla/5.0 (compatible; codewhale/",
382 env!("CARGO_PKG_VERSION"),
383 "; +https://github.com/codewhale-hq/CodeWhale)"
384 )
385 }
386
387 /// Upper bound on a single sleep inside the provider-wide rate-limit pause
388 /// loop in `send_with_retry`. The pause window lives in process-global state
389 /// (`retry_status`), so waiting requests re-poll it on this cadence instead
390 /// of committing to the full remaining window up front.
391 const RATE_LIMIT_PAUSE_RECHECK_INTERVAL: Duration = Duration::from_millis(250);
392
393 /// Total budget for one non-streaming request. Two layers use it: each
394 /// attempt carries it as a reqwest per-request total (connect through body
395 /// end, so a trickling body cannot extend forever), and the retry loop
396 /// through `send_with_retry` is wrapped in one outer envelope of the same
397 /// length (all attempts, backoff, and honored Retry-After included). The
398 /// shared client intentionally has no client-level total timeout, so without
399 /// these nothing bounds a non-streaming completion: a provider that accepts
400 /// the connection and then stalls — or a gateway answering 429 +
401 /// `Retry-After: 3600` forever — wedged the caller indefinitely.
402 ///
403 /// Streaming paths never carry it: their opens go through
404 /// `send_stream_open_with_retry`, which sets no per-request total (a total
405 /// would ride on the returned body and hard-cut a live stream), so a stream
406 /// stays bounded by its open cap and per-chunk idle checks only.
407 pub(super) const NON_STREAMING_REQUEST_ENVELOPE: Duration = Duration::from_secs(1800);
408
409 #[cfg(test)]
410 static TEST_NON_STREAMING_ENVELOPE_MS: std::sync::atomic::AtomicU64 =
411 std::sync::atomic::AtomicU64::new(0);
412
413 fn non_streaming_request_envelope() -> Duration {
414 #[cfg(test)]
415 {
416 let ms = TEST_NON_STREAMING_ENVELOPE_MS.load(std::sync::atomic::Ordering::SeqCst);
417 if ms > 0 {
418 return Duration::from_millis(ms);
419 }
420 }
421 NON_STREAMING_REQUEST_ENVELOPE
422 }
423
424 pub(super) const SSE_BACKPRESSURE_HIGH_WATERMARK: usize = 1024 * 1024; // 1 MB
425 pub(super) const SSE_BACKPRESSURE_SLEEP_MS: u64 = 10;
426 pub(super) const SSE_MAX_LINES_PER_CHUNK: usize = 256;
427 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
428 enum ConnectionState {
429 Healthy,
430 Degraded,
431 Recovering,
432 }
433
434 #[derive(Debug)]
435 struct ConnectionHealth {
436 state: ConnectionState,
437 consecutive_failures: u32,
438 last_failure: Option<Instant>,
439 last_success: Option<Instant>,
440 last_probe: Option<Instant>,
441 }
442
443 impl Default for ConnectionHealth {
444 fn default() -> Self {
445 Self {
446 state: ConnectionState::Healthy,
447 consecutive_failures: 0,
448 last_failure: None,
449 last_success: None,
450 last_probe: None,
451 }
452 }
453 }
454
455 #[derive(Debug)]
456 struct TokenBucket {
457 enabled: bool,
458 capacity: f64,
459 tokens: f64,
460 refill_per_sec: f64,
461 last_refill: Instant,
462 }
463
464 #[derive(Debug, Clone)]
465 struct ProviderConcurrencyLimiter {
466 semaphore: Arc<Semaphore>,
467 active: Arc<AtomicUsize>,
468 limit: usize,
469 }
470
471 struct ProviderRequestPermit {
472 _permit: OwnedSemaphorePermit,
473 active: Arc<AtomicUsize>,
474 }
475
476 impl ProviderConcurrencyLimiter {
477 fn new(limit: usize) -> Self {
478 let limit = limit.max(1);
479 Self {
480 semaphore: Arc::new(Semaphore::new(limit)),
481 active: Arc::new(AtomicUsize::new(0)),
482 limit,
483 }
484 }
485
486 async fn acquire(&self) -> Option<ProviderRequestPermit> {
487 let permit = Arc::clone(&self.semaphore).acquire_owned().await.ok()?;
488 self.active.fetch_add(1, Ordering::AcqRel);
489 Some(ProviderRequestPermit {
490 _permit: permit,
491 active: Arc::clone(&self.active),
492 })
493 }
494
495 fn active(&self) -> usize {
496 self.active.load(Ordering::Acquire)
497 }
498
499 fn limit(&self) -> usize {
500 self.limit
501 }
502 }
503
504 impl Drop for ProviderRequestPermit {
505 fn drop(&mut self) {
506 self.active.fetch_sub(1, Ordering::AcqRel);
507 }
508 }
509
510 impl TokenBucket {
511 fn from_env() -> Self {
512 let rps = std::env::var("CODEWHALE_RATE_LIMIT_RPS")
513 .or_else(|_| std::env::var("DEEPSEEK_RATE_LIMIT_RPS"))
514 .ok()
515 .and_then(|v| v.parse::<f64>().ok())
516 .unwrap_or(DEFAULT_CLIENT_RATE_LIMIT_RPS)
517 .max(0.0);
518 let burst = std::env::var("CODEWHALE_RATE_LIMIT_BURST")
519 .or_else(|_| std::env::var("DEEPSEEK_RATE_LIMIT_BURST"))
520 .ok()
521 .and_then(|v| v.parse::<f64>().ok())
522 .unwrap_or(DEFAULT_CLIENT_RATE_LIMIT_BURST)
523 .max(1.0);
524 let enabled = rps > 0.0;
525 Self {
526 enabled,
527 capacity: burst,
528 tokens: burst,
529 refill_per_sec: rps,
530 last_refill: Instant::now(),
531 }
532 }
533
534 fn refill(&mut self, now: Instant) {
535 if !self.enabled {
536 return;
537 }
538 let elapsed = now.duration_since(self.last_refill).as_secs_f64();
539 self.last_refill = now;
540 self.tokens = (self.tokens + elapsed * self.refill_per_sec).min(self.capacity);
541 }
542
543 /// Reserve `tokens` and report how long the caller must sleep first.
544 ///
545 /// The debt is *kept* (the balance is allowed to go negative) rather than
546 /// floored at zero. Callers release the bucket lock before sleeping, so a
547 /// floored balance would hand every queued waiter the same short delay and
548 /// they would all wake — and fire — at the same instant, which is the
549 /// burst the configured limit exists to prevent. Carrying the deficit
550 /// spaces successive waiters one refill interval apart, and `refill`'s
551 /// clamp to `capacity` still caps how much credit an idle bucket banks.
552 fn delay_until_available(&mut self, tokens: f64) -> Option<Duration> {
553 if !self.enabled {
554 return None;
555 }
556 let now = Instant::now();
557 self.refill(now);
558 self.tokens -= tokens;
559 if self.tokens >= 0.0 {
560 return None;
561 }
562 if self.refill_per_sec <= 0.0 {
563 return Some(Duration::from_secs(1));
564 }
565 Some(Duration::from_secs_f64(-self.tokens / self.refill_per_sec))
566 }
567 }
568
569 fn apply_request_success(health: &mut ConnectionHealth, now: Instant) -> bool {
570 let recovered = health.state != ConnectionState::Healthy;
571 health.state = ConnectionState::Healthy;
572 health.consecutive_failures = 0;
573 health.last_success = Some(now);
574 recovered
575 }
576
577 fn apply_request_failure(health: &mut ConnectionHealth, now: Instant) {
578 health.consecutive_failures = health.consecutive_failures.saturating_add(1);
579 health.last_failure = Some(now);
580 if health.consecutive_failures >= CONNECTION_FAILURE_THRESHOLD {
581 health.state = ConnectionState::Degraded;
582 }
583 }
584
585 fn mark_recovery_probe_if_due(health: &mut ConnectionHealth, now: Instant) -> bool {
586 if health.state == ConnectionState::Healthy {
587 return false;
588 }
589 if health
590 .last_probe
591 .is_some_and(|last| now.duration_since(last) < RECOVERY_PROBE_COOLDOWN)
592 {
593 return false;
594 }
595 health.last_probe = Some(now);
596 health.state = ConnectionState::Recovering;
597 true
598 }
599
600 /// The one command that replaces a rejected key, by where it came from.
601 fn auth_fix_hint(route: &str, key_source: &str) -> String {
602 if key_source.starts_with("--api-key") {
603 "pass a valid --api-key".to_string()
604 } else if let Some(rest) = key_source.strip_prefix("env var ")
605 && rest.contains("api_key_env")
606 {
607 let name = rest.split_whitespace().next().unwrap_or(rest);
608 format!("set {name} to a valid key")
609 } else if key_source.contains("OAuth") || key_source.contains("login") {
610 format!("sign in again for {route}")
611 } else {
612 format!("codewhale auth set --provider {route}")
613 }
614 }
615
616 fn buffer_pool() -> &'static StdMutex<Vec<Vec<u8>>> {
617 static POOL: OnceLock<StdMutex<Vec<Vec<u8>>>> = OnceLock::new();
618 POOL.get_or_init(|| StdMutex::new(Vec::new()))
619 }
620
621 fn acquire_stream_buffer() -> Vec<u8> {
622 if let Ok(mut pool) = buffer_pool().lock() {
623 pool.pop().unwrap_or_else(|| Vec::with_capacity(8192))
624 } else {
625 Vec::with_capacity(8192)
626 }
627 }
628
629 fn release_stream_buffer(mut buf: Vec<u8>) {
630 buf.clear();
631 if buf.capacity() > 256 * 1024 {
632 buf.shrink_to(256 * 1024);
633 }
634 if let Ok(mut pool) = buffer_pool().lock()
635 && pool.len() < 8
636 {
637 pool.push(buf);
638 }
639 }
640
641 impl Clone for CodewhaleClient {
642 fn clone(&self) -> Self {
643 Self {
644 http_client: self.http_client.clone(),
645 models_http_client: self.models_http_client.clone(),
646 http1_client: self.http1_client.clone(),
647 api_key: self.api_key.clone(),
648 api_key_source: self.api_key_source.clone(),
649 subscription_limit_guidance: self.subscription_limit_guidance.clone(),
650 model_bound_secret_values: Arc::clone(&self.model_bound_secret_values),
651 catalog_error_secret_values: Arc::clone(&self.catalog_error_secret_values),
652 model_bound_masking: self.model_bound_masking,
653 base_url: self.base_url.clone(),
654 api_provider: self.api_provider,
655 admitted_identity: self.admitted_identity.clone(),
656 openrouter_vendor: self.openrouter_vendor.clone(),
657 billing_surface: self.billing_surface.clone(),
658 billing_mode: self.billing_mode,
659 configured_models: Arc::clone(&self.configured_models),
660 route_limits: self.route_limits,
661 codex_account_id: self.codex_account_id.clone(),
662 chatgpt_reasoning_api: self.chatgpt_reasoning_api.clone(),
663 wire_format: self.wire_format,
664 retry: self.retry.clone(),
665 isolated_request_state: self.isolated_request_state,
666 remote_control_inference_participant: self.remote_control_inference_participant,
667 default_model: self.default_model.clone(),
668 connection_health: self.connection_health.clone(),
669 rate_limiter: self.rate_limiter.clone(),
670 request_concurrency: self.request_concurrency.clone(),
671 path_suffix: self.path_suffix.clone(),
672 #[cfg(test)]
673 test_chat_transport_base_url: self.test_chat_transport_base_url.clone(),
674 #[cfg(test)]
675 test_messages_transport_base_url: self.test_messages_transport_base_url.clone(),
676 reasoning_stream_style: self.reasoning_stream_style.clone(),
677 stream_idle_timeout: self.stream_idle_timeout,
678 stream_open_timeout: self.stream_open_timeout,
679 force_http1: self.force_http1,
680 }
681 }
682 }
683
684 const MIN_EXACT_SECRET_CHARS: usize = 8;
685
686 pub(crate) use codewhale_config::apply_openrouter_vendor;
687
688 fn push_model_bound_secret(values: &mut Vec<String>, value: Option<&str>) {
689 let Some(value) = value
690 .map(str::trim)
691 .filter(|value| !value.is_empty() && value.chars().count() >= MIN_EXACT_SECRET_CHARS)
692 else {
693 return;
694 };
695 if !values.iter().any(|existing| existing == value) {
696 values.push(value.to_string());
697 }
698 }
699
700 fn model_bound_secret_store_slot(provider: ProviderKind) -> Option<&'static str> {
701 (provider != ProviderKind::Custom).then(|| provider.secret_store_slot())
702 }
703
704 fn push_file_backed_model_bound_secrets(values: &mut Vec<String>) {
705 // Unit tests must never inspect the developer's real credential store.
706 // The isolated regression below opts in with a temporary CODEWHALE_HOME,
707 // matching Config's existing secret-store test discipline.
708 #[cfg(test)]
709 if !codewhale_paths::codewhale_home_is_explicit()
710 || std::env::var_os("CODEWHALE_SECRET_BACKEND").is_none()
711 {
712 return;
713 }
714
715 // Redaction needs only a best-effort view of inactive file-backed
716 // credentials. It must not cause a legacy-store migration merely because a
717 // client is being constructed (notably for `doctor`'s live probe). Keep
718 // this file-only to avoid a burst of OS-keychain prompts for inactive
719 // providers; the active credential is already supplied by the route
720 // resolver.
721 let secrets = codewhale_secrets::Secrets::file_backed_read_only();
722 let mut slots = Vec::new();
723 for provider in codewhale_config::descriptors::provider_compatibility()
724 .iter()
725 .map(|row| row.kind)
726 {
727 let Some(slot) = model_bound_secret_store_slot(provider) else {
728 continue;
729 };
730 if !slots.contains(&slot) {
731 slots.push(slot);
732 }
733 }
734 // The legacy literal `provider = "custom"` route owns this durable slot.
735 slots.push("custom");
736
737 for slot in slots {
738 if let Ok(Some(secret)) = secrets.get(slot) {
739 push_model_bound_secret(values, Some(&secret));
740 }
741 }
742 }
743
744 pub(crate) fn configured_model_bound_secret_values(
745 config: &Config,
746 active_api_key: &str,
747 ) -> Vec<String> {
748 let mut values = Vec::new();
749 push_model_bound_secret(&mut values, Some(active_api_key));
750 push_model_bound_secret(&mut values, config.sandbox_api_key.as_deref());
751 push_model_bound_secret(
752 &mut values,
753 config
754 .search
755 .as_ref()
756 .and_then(|search| search.api_key.as_deref()),
757 );
758 push_model_bound_secret(
759 &mut values,
760 config
761 .vision_model
762 .as_ref()
763 .and_then(|vision| vision.api_key.as_deref()),
764 );
765
766 if let Some(headers) = config.http_headers.as_ref() {
767 for (name, value) in headers {
768 if is_upstream_auth_header(name) {
769 push_model_bound_secret(&mut values, Some(value));
770 }
771 }
772 }
773
774 for row in codewhale_config::descriptors::provider_compatibility()
775 .iter()
776 .filter(|row| row.kind != ProviderKind::Custom)
777 {
778 for env_name in row.kind.provider().env_vars() {
779 if let Ok(value) = std::env::var(env_name) {
780 push_model_bound_secret(&mut values, Some(&value));
781 }
782 }
783 let Some(provider_config) = config
784 .providers
785 .as_ref()
786 .and_then(|tables| codewhale_config::provider_config_table!(@read tables, row.id))
787 else {
788 continue;
789 };
790 push_model_bound_secret(&mut values, provider_config.api_key.as_deref());
791 if let Some(headers) = provider_config.http_headers.as_ref() {
792 for (name, value) in headers {
793 if is_upstream_auth_header(name) {
794 push_model_bound_secret(&mut values, Some(value));
795 }
796 }
797 }
798 }
799
800 if let Some(providers) = config.providers.as_ref() {
801 for provider_config in providers.custom.values() {
802 push_model_bound_secret(&mut values, provider_config.api_key.as_deref());
803 if let Some(env_name) = provider_config
804 .api_key_env
805 .as_deref()
806 .map(str::trim)
807 .filter(|name| !name.is_empty())
808 && let Ok(value) = std::env::var(env_name)
809 {
810 push_model_bound_secret(&mut values, Some(&value));
811 }
812 if let Some(headers) = provider_config.http_headers.as_ref() {
813 for (name, value) in headers {
814 if is_upstream_auth_header(name) {
815 push_model_bound_secret(&mut values, Some(value));
816 }
817 }
818 }
819 }
820 }
821
822 // The decision router's TypeSafe key is no chat provider's key; its env
823 // form must still never reach a model. (`for_decision_route` adds the key
824 // from every source to its own client.)
825 if let Ok(value) = std::env::var(system_one::TYPESAFE_API_KEY_ENV) {
826 push_model_bound_secret(&mut values, Some(&value));
827 }
828
829 push_file_backed_model_bound_secrets(&mut values);
830
831 // Replace longer values first in case one credential happens to contain
832 // another as a prefix.
833 values.sort_by_key(|value| std::cmp::Reverse(value.len()));
834 values
835 }
836
837 /// Everything a catalog probe could have sent that must not come back.
838 ///
839 /// Longest first, so a value that contains another is masked whole.
840 fn catalog_error_secret_values(
841 active_api_key: &str,
842 http_headers: &HashMap<String, String>,
843 model_bound: &[String],
844 ) -> Vec<String> {
845 let mut values: Vec<String> = Vec::new();
846 let mut push = |value: &str| {
847 let value = value.trim();
848 if !value.is_empty() && !values.iter().any(|existing| existing == value) {
849 values.push(value.to_string());
850 }
851 };
852 // No length floor here, unlike `push_model_bound_secret`: a short key
853 // echoed back by an untrusted endpoint is still a leaked key. The cost of
854 // being wrong is a suppressed message, which is what this path did before.
855 push(active_api_key);
856 for value in http_headers.values() {
857 push(value);
858 }
859 for value in model_bound {
860 push(value);
861 }
862 values.sort_by_key(|value| std::cmp::Reverse(value.len()));
863 values
864 }
865
866 /// The opaque values in a request's URL query — the pagination cursor.
867 ///
868 /// Collected in both forms: decoded, as a provider that parsed the cursor
869 /// would echo it, and raw, as one that quoted the URL back would.
870 fn request_query_secret_values(url: &reqwest::Url) -> Vec<String> {
871 let mut values: Vec<String> = Vec::new();
872 let mut push = |value: String| {
873 if !value.trim().is_empty() && !values.contains(&value) {
874 values.push(value);
875 }
876 };
877 for (_, value) in url.query_pairs() {
878 push(value.into_owned());
879 }
880 for pair in url.query().unwrap_or_default().split('&') {
881 if let Some((_, value)) = pair.split_once('=') {
882 push(value.to_string());
883 }
884 }
885 values.sort_by_key(|value| std::cmp::Reverse(value.len()));
886 values
887 }
888
889 pub(crate) fn redact_model_bound_text(text: &str, exact_secret_values: &[String]) -> String {
890 let mut redacted = text.to_string();
891 for secret in exact_secret_values {
892 redacted = redacted.replace(secret, codewhale_config::persistence::REDACTED);
893 }
894 // Tool results feed exact-match edits, so only credential-shaped values
895 // are masked here; key-only hits (`password: credentials?.password`) stay
896 // byte-exact. Logs and previews keep the broad key-based scrubber.
897 codewhale_config::persistence::redact_model_bound_secrets(&redacted)
898 }
899
900 pub(crate) fn redact_json_model_bound_text(
901 value: &serde_json::Value,
902 secrets: &[String],
903 ) -> serde_json::Value {
904 fn mask_exact_values(value: &mut serde_json::Value, secrets: &[String]) {
905 match value {
906 serde_json::Value::String(text) => *text = redact_model_bound_text(text, secrets),
907 serde_json::Value::Array(values) => {
908 for value in values {
909 mask_exact_values(value, secrets);
910 }
911 }
912 serde_json::Value::Object(values) => {
913 for value in values.values_mut() {
914 mask_exact_values(value, secrets);
915 }
916 }
917 _ => {}
918 }
919 }
920 // Preserve sensitive-key masking and the existing recursion-depth bound.
921 let mut redacted = codewhale_config::persistence::redact_json_model_bound_secrets(value);
922 mask_exact_values(&mut redacted, secrets);
923 redacted
924 }
925
926 // === Helpers ===
927
928 /// Maximum bytes to read from an error response body (64 KB).
929 pub(super) const ERROR_BODY_MAX_BYTES: usize = 64 * 1024;
930
931 /// How much of a provider's HTTP error body may be shown to the user.
932 ///
933 /// The catalog probe is the one request that contacts a `base_url` the user
934 /// typed during setup, and its URL carries an opaque pagination cursor, so a
935 /// provider — or anything answering at that URL — can echo a key, a custom
936 /// header value or the cursor back inside an error body. #3385 answered that
937 /// by discarding the body entirely, and the cost of the blunt version was
938 /// #6173: Gemini's real reason ("User location is not supported for the API
939 /// use") reached the user as `Invalid request (400): ` with nothing after the
940 /// colon, so a geo-block, a bad key and a wrong endpoint were indistinguishable.
941 pub(super) enum ErrorBodyDisclosure {
942 /// The endpoint is already established. Surface the provider's message.
943 Full,
944 /// Untrusted endpoint. Surface only what survives redaction — and nothing
945 /// at all if a known secret is still in the result.
946 Guarded { request_secrets: Vec<String> },
947 }
948
949 /// Read/overall timeout for the shared client's non-streaming requests
950 /// (`/models` listing, catalog refresh, health probes). Streaming requests
951 /// keep their own idle-timeout envelope; without this, a provider that accepts
952 /// the connection and never answers hangs model-list, catalog refresh, and
953 /// health checks forever (ops R4).
954 pub(super) const NON_STREAMING_HTTP_TIMEOUT: Duration = Duration::from_secs(30);
955 const PROVIDER_CATALOG_MAX_RESPONSE_BYTES: usize = 8 * 1024 * 1024;
956 const PROVIDER_CATALOG_MAX_ROWS: usize = 10_000;
957
958 #[derive(Clone, Copy)]
959 struct ModelsFetchLimits {
960 bytes: usize,
961 rows: usize,
962 pages: usize,
963 cursor_bytes: usize,
964 timeout: Duration,
965 }
966
967 const MODELS_FETCH_LIMITS: ModelsFetchLimits = ModelsFetchLimits {
968 bytes: PROVIDER_CATALOG_MAX_RESPONSE_BYTES,
969 rows: PROVIDER_CATALOG_MAX_ROWS,
970 pages: 1_000,
971 cursor_bytes: 4_096,
972 timeout: NON_STREAMING_HTTP_TIMEOUT,
973 };
974
975 #[derive(Clone, Copy)]
976 enum ModelsRequestMode {
977 Interactive,
978 Refresh,
979 }
980
981 #[derive(Debug)]
982 enum ModelsFetchError {
983 Catalog(CatalogRefreshError),
984 Interactive(anyhow::Error),
985 }
986
987 impl ModelsFetchError {
988 fn into_interactive(self) -> anyhow::Error {
989 match self {
990 Self::Catalog(reason) => anyhow::anyhow!("Failed to list models: {reason:?}"),
991 Self::Interactive(error) => error,
992 }
993 }
994
995 fn into_catalog(self) -> CatalogRefreshError {
996 match self {
997 Self::Catalog(reason) => reason,
998 Self::Interactive(_) => CatalogRefreshError::Network,
999 }
1000 }
1001 }
1002
1003 impl From<CatalogRefreshError> for ModelsFetchError {
1004 fn from(reason: CatalogRefreshError) -> Self {
1005 Self::Catalog(reason)
1006 }
1007 }
1008
1009 // Keep row bytes intact: a Value round trip would silently accept duplicate
1010 // fields that the existing typed provider parsers reject.
1011 #[derive(Deserialize)]
1012 struct ModelsPage<'a> {
1013 #[serde(borrow, alias = "models")]
1014 data: &'a serde_json::value::RawValue,
1015 #[serde(default)]
1016 has_more: bool,
1017 last_id: Option<String>,
1018 }
1019
1020 struct BoundedModelsRows(usize);
1021
1022 impl<'de> serde::de::Visitor<'de> for BoundedModelsRows {
1023 type Value = Vec<Box<serde_json::value::RawValue>>;
1024
1025 fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1026 formatter.write_str("a bounded model array")
1027 }
1028
1029 fn visit_seq<A: serde::de::SeqAccess<'de>>(self, mut seq: A) -> Result<Self::Value, A::Error> {
1030 let mut rows = Vec::new();
1031 while rows.len() < self.0 {
1032 let Some(row) = seq.next_element()? else {
1033 return Ok(rows);
1034 };
1035 rows.push(row);
1036 }
1037 if seq.next_element::<serde::de::IgnoredAny>()?.is_some() {
1038 return Err(serde::de::Error::custom("model row limit exceeded"));
1039 }
1040 Ok(rows)
1041 }
1042 }
1043
1044 impl<'de> serde::de::DeserializeSeed<'de> for BoundedModelsRows {
1045 type Value = Vec<Box<serde_json::value::RawValue>>;
1046
1047 fn deserialize<D: serde::Deserializer<'de>>(
1048 self,
1049 deserializer: D,
1050 ) -> Result<Self::Value, D::Error> {
1051 deserializer.deserialize_seq(self)
1052 }
1053 }
1054
1055 /// One complete traversal, independent of provider row parsing and transport
1056 /// retry policy. Only a verified endpoint contract supplies a cursor query key;
1057 /// generation wire format alone does not prove a model-list dialect.
1058 async fn collect_models_document<F, Fut>(
1059 endpoint: reqwest::Url,
1060 cursor_query: Option<&str>,
1061 limits: ModelsFetchLimits,
1062 mut fetch: F,
1063 ) -> Result<(String, tokio::time::Instant), ModelsFetchError>
1064 where
1065 F: FnMut(reqwest::Url) -> Fut,
1066 Fut: std::future::Future<Output = Result<reqwest::Response, ModelsFetchError>>,
1067 {
1068 use serde::de::DeserializeSeed;
1069 let deadline = tokio::time::Instant::now() + limits.timeout;
1070 let body = tokio::time::timeout_at(deadline, async {
1071 let mut rows = Vec::new();
1072 let mut bytes = 0usize;
1073 let mut cursor: Option<String> = None;
1074 let mut seen = std::collections::HashSet::new();
1075 for page_index in 0..limits.pages {
1076 let mut url = endpoint.clone();
1077 if let (Some(key), Some(value)) = (cursor_query, cursor.as_deref()) {
1078 let query: Vec<_> = url
1079 .query_pairs()
1080 .filter(|(name, _)| name != key)
1081 .map(|(name, value)| (name.into_owned(), value.into_owned()))
1082 .collect();
1083 url.query_pairs_mut()
1084 .clear()
1085 .extend_pairs(query)
1086 .append_pair(key, value);
1087 }
1088 let response = fetch(url).await?;
1089 let body =
1090 bounded_provider_catalog_text(response, limits.bytes.saturating_sub(bytes)).await?;
1091 bytes += body.len();
1092 let page: ModelsPage<'_> =
1093 serde_json::from_str(&body).map_err(|_| CatalogRefreshError::InvalidResponse)?;
1094 let page_rows = BoundedModelsRows(limits.rows.saturating_sub(rows.len()))
1095 .deserialize(&mut serde_json::Deserializer::from_str(page.data.get()))
1096 .map_err(|_| CatalogRefreshError::InvalidResponse)?;
1097 rows.extend(page_rows);
1098 if !page.has_more {
1099 if tokio::time::Instant::now() >= deadline {
1100 return Err(CatalogRefreshError::Network.into());
1101 }
1102 if page_index == 0 {
1103 return Ok(body);
1104 }
1105 #[derive(Serialize)]
1106 struct Document {
1107 data: Vec<Box<serde_json::value::RawValue>>,
1108 }
1109 return serde_json::to_string(&Document { data: rows })
1110 .map_err(|_| CatalogRefreshError::InvalidResponse.into());
1111 }
1112 let next = page
1113 .last_id
1114 .filter(|value| !value.is_empty() && value.len() <= limits.cursor_bytes)
1115 .ok_or(CatalogRefreshError::InvalidResponse)?;
1116 if cursor_query.is_none() || !seen.insert(next.clone()) {
1117 return Err(CatalogRefreshError::InvalidResponse.into());
1118 }
1119 cursor = Some(next);
1120 }
1121 Err(ModelsFetchError::Catalog(
1122 CatalogRefreshError::InvalidResponse,
1123 ))
1124 })
1125 .await
1126 .map_err(|_| CatalogRefreshError::Network)??;
1127 Ok((body, deadline))
1128 }
1129
1130 /// Read an error response body with a size limit to prevent unbounded allocation.
1131 pub(super) async fn bounded_error_text(response: reqwest::Response, max_bytes: usize) -> String {
1132 use futures_util::StreamExt;
1133 let mut stream = response.bytes_stream();
1134 let mut buf = Vec::with_capacity(max_bytes.min(8192));
1135 while let Some(chunk) = stream.next().await {
1136 let Ok(chunk) = chunk else { break };
1137 let remaining = max_bytes.saturating_sub(buf.len());
1138 if remaining == 0 {
1139 break;
1140 }
1141 buf.extend_from_slice(&chunk[..chunk.len().min(remaining)]);
1142 }
1143 String::from_utf8_lossy(&buf).into_owned()
1144 }
1145
1146 async fn bounded_provider_catalog_text(
1147 response: reqwest::Response,
1148 max_bytes: usize,
1149 ) -> Result<String, CatalogRefreshError> {
1150 if response
1151 .content_length()
1152 .is_some_and(|length| length > max_bytes as u64)
1153 {
1154 return Err(CatalogRefreshError::InvalidResponse);
1155 }
1156 let mut stream = response.bytes_stream();
1157 let mut body = Vec::new();
1158 while let Some(chunk) = stream.next().await {
1159 let chunk = chunk.map_err(|_| CatalogRefreshError::Network)?;
1160 if body.len().saturating_add(chunk.len()) > max_bytes {
1161 return Err(CatalogRefreshError::InvalidResponse);
1162 }
1163 body.extend_from_slice(&chunk);
1164 }
1165 String::from_utf8(body).map_err(|_| CatalogRefreshError::InvalidResponse)
1166 }
1167
1168 fn validate_base_url_security(base_url: &str, provider_allows_insecure_http: bool) -> Result<()> {
1169 let display_base_url = redact_url_for_display(base_url);
1170 let parsed = reqwest::Url::parse(base_url)
1171 .map_err(|_| anyhow::anyhow!("Refusing invalid base URL '{display_base_url}'"))?;
1172 let loopback = parsed.host_str().is_some_and(|host| {
1173 host.eq_ignore_ascii_case("localhost")
1174 || host
1175 .trim_matches(['[', ']'])
1176 .parse::<std::net::IpAddr>()
1177 .is_ok_and(|address| address.is_loopback())
1178 });
1179 if parsed.scheme() == "https" || (parsed.scheme() == "http" && loopback) {
1180 return Ok(());
1181 }
1182
1183 if parsed.scheme() == "http" && provider_allows_insecure_http {
1184 logging::warn(
1185 "Using insecure HTTP base URL because this provider sets allow_insecure_http = true in config.toml",
1186 );
1187 return Ok(());
1188 }
1189
1190 if parsed.scheme() == "http"
1191 && std::env::var(ALLOW_INSECURE_HTTP_ENV)
1192 .or_else(|_| std::env::var(LEGACY_ALLOW_INSECURE_HTTP_ENV))
1193 .ok()
1194 .as_deref()
1195 .is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
1196 {
1197 logging::warn(format!(
1198 "Using insecure HTTP base URL because {ALLOW_INSECURE_HTTP_ENV} is set"
1199 ));
1200 return Ok(());
1201 }
1202
1203 if parsed.scheme() == "http" {
1204 anyhow::bail!(
1205 "Refusing insecure base URL '{display_base_url}'.\n\
1206 \n\
1207 Loopback hosts (localhost, 127.0.0.1, [::1]) are auto-allowed.\n\
1208 For one trusted local provider (LAN, llama.cpp on a private IP, etc.) set\n\
1209 `allow_insecure_http = true` under its `[providers.<name>]` table in config.toml.\n\
1210 To allow it for every provider in this shell instead, set the env var\n\
1211 `{ALLOW_INSECURE_HTTP_ENV}=1` and re-run.",
1212 );
1213 }
1214
1215 anyhow::bail!(
1216 "Refusing base URL '{display_base_url}': only HTTPS (or explicitly allowed HTTP) URLs are supported.",
1217 )
1218 }
1219
1220 /// Mask credentials in a URL for display.
1221 ///
1222 /// Delegates to the single shared implementation in
1223 /// [`codewhale_secrets::sanitize`] (FEAT-025 D4) so command sanitization and
1224 /// client diagnostics cannot drift.
1225 pub(crate) fn redact_url_for_display(url: &str) -> String {
1226 codewhale_secrets::sanitize::redact_url_for_display(url)
1227 }
1228
1229 pub(super) fn versioned_base_url(base_url: &str) -> String {
1230 let trimmed = base_url.trim_end_matches('/');
1231 if base_url_has_version_suffix(trimmed) {
1232 trimmed.to_string()
1233 } else {
1234 format!("{trimmed}/v1")
1235 }
1236 }
1237
1238 fn unversioned_base_url(base_url: &str) -> String {
1239 let trimmed = base_url.trim_end_matches('/');
1240 trimmed
1241 .rsplit_once('/')
1242 .filter(|(_, segment)| is_version_segment(segment))
1243 .map(|(base, _)| base)
1244 .unwrap_or(trimmed)
1245 .to_string()
1246 }
1247
1248 fn base_url_has_version_suffix(trimmed: &str) -> bool {
1249 trimmed.rsplit('/').next().is_some_and(is_version_segment)
1250 }
1251
1252 fn is_version_segment(segment: &str) -> bool {
1253 segment.eq_ignore_ascii_case("beta")
1254 || segment
1255 .strip_prefix('v')
1256 .or_else(|| segment.strip_prefix('V'))
1257 .is_some_and(|rest| !rest.is_empty() && rest.chars().all(|ch| ch.is_ascii_digit()))
1258 }
1259
1260 pub(crate) fn api_url(base_url: &str, path: &str) -> String {
1261 api_url_with_suffix(base_url, path, None)
1262 }
1263
1264 fn responses_api_url(base_url: &str, provider: ProviderKind) -> String {
1265 let normalized = base_url.trim_end_matches('/').to_ascii_lowercase();
1266 let official_deepseek = matches!(provider, ProviderKind::Deepseek)
1267 && matches!(
1268 normalized.as_str(),
1269 "https://api.deepseek.com"
1270 | "https://api.deepseek.com/v1"
1271 | "https://api.deepseek.com/beta"
1272 | "https://api.deepseeki.com"
1273 | "https://api.deepseeki.com/v1"
1274 | "https://api.deepseeki.com/beta"
1275 );
1276 if official_deepseek {
1277 format!("{}/responses", unversioned_base_url(base_url))
1278 } else {
1279 api_url(base_url, "responses")
1280 }
1281 }
1282
1283 pub(super) fn api_url_with_suffix(base_url: &str, path: &str, path_suffix: Option<&str>) -> String {
1284 let path = path.trim_start_matches('/');
1285 if path.starts_with("beta/") {
1286 return format!("{}/{}", unversioned_base_url(base_url), path);
1287 }
1288 if let ("chat/completions", Some(suffix)) = (path, path_suffix) {
1289 return format!(
1290 "{}/{}",
1291 unversioned_base_url(base_url),
1292 suffix.trim_start_matches('/')
1293 );
1294 }
1295 let mut versioned = versioned_base_url(base_url);
1296 // The /beta suffix is not a real API version — it is an
1297 // opt-in surface for beta features. Only paths with an
1298 // explicit `beta/` prefix should hit the beta surface;
1299 // everything else (models, chat/completions, health, …)
1300 // must go to the standard /v1 surface.
1301 if versioned.ends_with("beta") {
1302 versioned = format!("{}/v1", unversioned_base_url(base_url));
1303 }
1304 format!("{}/{}", versioned.trim_end_matches('/'), path)
1305 }
1306
1307 /// Route strict DeepSeek tool requests through the beta Chat Completions
1308 /// surface while keeping every ordinary request on the canonical `/v1` path.
1309 ///
1310 /// DeepSeek requires its `/beta` base URL when a function opts into
1311 /// `strict: true`. The configured route URL remains semantic here because
1312 /// unit tests may replace only the transport origin with a local capture
1313 /// server.
1314 ///
1315 /// Source: <https://api-docs.deepseek.com/guides/tool_calls/> (verified 2026-07-22).
1316 fn chat_completions_url(
1317 transport_base_url: &str,
1318 route_base_url: &str,
1319 provider: ProviderKind,
1320 path_suffix: Option<&str>,
1321 body: &Value,
1322 ) -> String {
1323 let uses_deepseek_beta = matches!(provider, ProviderKind::Deepseek)
1324 && is_official_deepseek_beta_base_url(route_base_url)
1325 && body_uses_strict_tools(body)
1326 && path_suffix.is_none();
1327 let path = if uses_deepseek_beta {
1328 "beta/chat/completions"
1329 } else {
1330 "chat/completions"
1331 };
1332 api_url_with_suffix(transport_base_url, path, path_suffix)
1333 }
1334
1335 fn is_official_deepseek_beta_base_url(base_url: &str) -> bool {
1336 matches!(
1337 base_url.trim_end_matches('/').to_ascii_lowercase().as_str(),
1338 "https://api.deepseek.com/beta" | "https://api.deepseeki.com/beta"
1339 )
1340 }
1341
1342 fn body_uses_strict_tools(body: &Value) -> bool {
1343 body.get("tools")
1344 .and_then(Value::as_array)
1345 .is_some_and(|tools| {
1346 tools
1347 .iter()
1348 .any(|tool| tool.pointer("/function/strict").and_then(Value::as_bool) == Some(true))
1349 })
1350 }
1351
1352 fn normalize_audio_format(format: &str) -> String {
1353 let normalized = format.trim().to_ascii_lowercase();
1354 if normalized.is_empty() {
1355 "wav".to_string()
1356 } else {
1357 normalized
1358 }
1359 }
1360
1361 fn parse_speech_audio_response(payload: &Value) -> Result<(Vec<u8>, Option<String>)> {
1362 let audio = payload
1363 .get("choices")
1364 .and_then(Value::as_array)
1365 .and_then(|choices| choices.first())
1366 .and_then(|choice| {
1367 choice
1368 .get("message")
1369 .and_then(|message| message.get("audio"))
1370 .or_else(|| choice.get("delta").and_then(|delta| delta.get("audio")))
1371 })
1372 .or_else(|| payload.get("audio"))
1373 .context("Speech synthesis response did not include choices[0].message.audio")?;
1374
1375 let data = audio
1376 .get("data")
1377 .and_then(Value::as_str)
1378 .context("Speech synthesis response did not include audio.data")?
1379 .trim();
1380 let data = data
1381 .split_once(',')
1382 .map(|(_, base64)| base64.trim())
1383 .unwrap_or(data);
1384 let audio_bytes = general_purpose::STANDARD
1385 .decode(data)
1386 .context("Failed to decode speech audio base64 data")?;
1387 let transcript = audio
1388 .get("transcript")
1389 .and_then(Value::as_str)
1390 .map(str::to_string);
1391
1392 Ok((audio_bytes, transcript))
1393 }
1394
1395 fn build_speech_synthesis_body(
1396 model: &str,
1397 text: &str,
1398 instruction: Option<&str>,
1399 audio: Value,
1400 ) -> Value {
1401 let mut messages = Vec::new();
1402 if let Some(instruction) = instruction.map(str::trim).filter(|value| !value.is_empty()) {
1403 messages.push(json!({
1404 "role": "user",
1405 "content": instruction,
1406 }));
1407 }
1408 messages.push(json!({
1409 "role": "assistant",
1410 "content": text,
1411 }));
1412
1413 json!({
1414 "model": model,
1415 "messages": messages,
1416 "audio": audio,
1417 })
1418 }
1419
1420 // === CodewhaleClient ===
1421
1422 /// Returns true when CODEWHALE_FORCE_HTTP1 (legacy alias: DEEPSEEK_FORCE_HTTP1)
1423 /// is set to a truthy value (`1`, `true`, `yes`, `on`, case-insensitive). Read
1424 /// only by `Config::force_http1`, which ORs it with the selected stream config flag; every
1425 /// client builder and stream open takes that resolved value (#103, #6700). Anything else (unset, `0`,
1426 /// `false`, ...) leaves HTTP/2 on.
1427 pub(crate) fn force_http1_from_env() -> bool {
1428 std::env::var("CODEWHALE_FORCE_HTTP1")
1429 .or_else(|_| std::env::var("DEEPSEEK_FORCE_HTTP1"))
1430 .ok()
1431 .map(|v| v.trim().to_ascii_lowercase())
1432 .is_some_and(|v| matches!(v.as_str(), "1" | "true" | "yes" | "on"))
1433 }
1434
1435 /// Read `SSL_CERT_FILE` and add its contents as extra root
1436 /// certificates on the reqwest builder (#418). Tries the PEM-bundle
1437 /// parser first (covers single-cert files too), then falls back to
1438 /// DER. All failures log a warning and return the builder unchanged
1439 /// so a malformed env var degrades gracefully.
1440 fn add_extra_root_certs(
1441 mut builder: reqwest::ClientBuilder,
1442 cert_path: &str,
1443 ) -> reqwest::ClientBuilder {
1444 let bytes = match std::fs::read(cert_path) {
1445 Ok(b) => b,
1446 Err(err) => {
1447 logging::warn(format!(
1448 "SSL_CERT_FILE={cert_path} could not be read: {err}"
1449 ));
1450 return builder;
1451 }
1452 };
1453
1454 if let Ok(certs) = reqwest::Certificate::from_pem_bundle(&bytes) {
1455 let added = certs.len();
1456 for cert in certs {
1457 builder = builder.add_root_certificate(cert);
1458 }
1459 logging::info(format!(
1460 "SSL_CERT_FILE={cert_path} loaded ({added} cert(s))"
1461 ));
1462 return builder;
1463 }
1464
1465 match reqwest::Certificate::from_der(&bytes) {
1466 Ok(cert) => {
1467 builder = builder.add_root_certificate(cert);
1468 logging::info(format!("SSL_CERT_FILE={cert_path} loaded (1 DER cert)"));
1469 }
1470 Err(err) => {
1471 logging::warn(format!(
1472 "SSL_CERT_FILE={cert_path} could not be parsed as PEM bundle or DER: {err}"
1473 ));
1474 }
1475 }
1476 builder
1477 }
1478
1479 /// Admit a complete response's tool pairing ids before hooks, approvals or
1480 /// replayable history. Use the protocol frozen with the actual request.
1481 pub(crate) fn validate_tool_call_ids_for_protocol<'a>(
1482 protocol: WireFormat,
1483 ids: impl IntoIterator<Item = &'a str>,
1484 ) -> Result<()> {
1485 let mut seen = std::collections::HashSet::new();
1486 for id in ids {
1487 let pairing_id = if protocol == WireFormat::Responses {
1488 responses::parse_tool_use_id(id).0
1489 } else {
1490 id.to_string()
1491 };
1492 if pairing_id.trim().is_empty() {
1493 anyhow::bail!("Provider returned a tool call without a pairing id");
1494 }
1495 if !seen.insert(pairing_id) {
1496 anyhow::bail!("Provider returned duplicate tool call pairing ids in one response");
1497 }
1498 }
1499 Ok(())
1500 }
1501
1502 impl CodewhaleClient {
1503 pub(crate) fn wire_format(&self) -> WireFormat {
1504 self.wire_format
1505 }
1506
1507 #[cfg(test)]
1508 pub(crate) fn validate_tool_call_ids<'a>(
1509 &self,
1510 ids: impl IntoIterator<Item = &'a str>,
1511 ) -> Result<()> {
1512 validate_tool_call_ids_for_protocol(self.wire_format, ids)
1513 }
1514
1515 fn is_local_ds4_model(&self, model: &str) -> bool {
1516 self.api_provider == ProviderKind::Custom
1517 && self
1518 .admitted_identity
1519 .key
1520 .as_str()
1521 .eq_ignore_ascii_case("ds4")
1522 && crate::config::base_url_uses_local_host(&self.base_url)
1523 && matches!(
1524 model.trim().to_ascii_lowercase().as_str(),
1525 "deepseek-v4-flash" | "deepseek-v4-pro"
1526 )
1527 }
1528
1529 /// DS4 is configured as a named custom route so its endpoint and billing
1530 /// identity remain exact, but its chat payload deliberately speaks the
1531 /// first-party DeepSeek reasoning/tool dialect that DS4 implements.
1532 fn chat_shape_provider(&self, model: &str) -> ProviderKind {
1533 if self.is_local_ds4_model(model) {
1534 ProviderKind::Deepseek
1535 } else {
1536 self.api_provider
1537 }
1538 }
1539
1540 /// Create a DeepSeek client from CLI configuration.
1541 pub fn new(config: &Config) -> Result<Self> {
1542 let identity = config
1543 .active_provider_identity()
1544 .map_err(anyhow::Error::msg)?;
1545 let api_provider = identity.provider;
1546 let model_aware = api_provider.provider().wire_policy()
1547 == codewhale_config::provider::WirePolicy::ModelAware;
1548 let default_model = config.default_model();
1549 let unresolved_local_model = api_provider == ProviderKind::Ollama
1550 && matches!(default_model.trim(), "" | "auto" | "unknown");
1551 if model_aware || unresolved_local_model {
1552 let route =
1553 crate::route_runtime::resolve_runtime_route_for_identity(config, &identity, None)
1554 .map_err(anyhow::Error::msg)?;
1555 return Self::from_candidate(&route.config, &route.candidate);
1556 }
1557 let route_limits = crate::route_runtime::resolve_runtime_route_for_identity(
1558 config,
1559 &identity,
1560 Some(&default_model),
1561 )
1562 .ok()
1563 .and_then(|route| crate::route_budget::known_route_limits(route.candidate.limits()));
1564 Self::from_parts(
1565 config.active_route_base_url(),
1566 default_model,
1567 provider_wire_format_for_config(&identity, Some(config)),
1568 route_limits,
1569 config,
1570 )
1571 }
1572
1573 /// Construct only the model-list probe before a route has a concrete model.
1574 /// Catalog bootstrap must not depend on the catalog it is about to fetch.
1575 pub(crate) fn for_catalog_refresh(config: &Config) -> Result<Self> {
1576 let identity = config
1577 .active_provider_identity()
1578 .map_err(anyhow::Error::msg)?;
1579 Self::from_parts(
1580 config.active_route_base_url(),
1581 config.default_model(),
1582 provider_wire_format_for_config(&identity, Some(config)),
1583 None,
1584 config,
1585 )
1586 }
1587
1588 /// Create a DeepSeek client whose transport is bound to a runtime-resolved
1589 /// route (#3384).
1590 ///
1591 /// The base URL and default model come from the executable `candidate`, so
1592 /// the client talks to exactly the endpoint and wire model the resolver
1593 /// chose instead of re-deriving them from `Config`. Secrets stay in
1594 /// `Config`: `ReadyRouteCandidate` is secret-free by design (it carries only
1595 /// an auth-source *class*), so the API key and provider are still read from
1596 /// `config`.
1597 pub fn from_candidate(config: &Config, candidate: &ReadyRouteCandidate) -> Result<Self> {
1598 let identity = config
1599 .active_provider_identity()
1600 .map_err(anyhow::Error::msg)?;
1601 anyhow::ensure!(
1602 candidate.provider_kind() == identity.provider,
1603 "resolved candidate does not match admitted provider kind"
1604 );
1605 Self::from_parts(
1606 candidate.endpoint().base_url.clone(),
1607 candidate.wire_model_id().as_str().to_string(),
1608 candidate.protocol(),
1609 crate::route_budget::known_route_limits(candidate.limits()),
1610 config,
1611 )
1612 }
1613
1614 /// Shared constructor body for [`Self::new`] and [`Self::from_candidate`].
1615 ///
1616 /// `base_url` and `default_model` are the only inputs that differ between
1617 /// the two entry points; everything else (auth, provider, retry, headers,
1618 /// timeouts) is derived from `config` so the two paths cannot drift.
1619 fn from_parts(
1620 base_url: String,
1621 default_model: String,
1622 wire_format: WireFormat,
1623 route_limits: Option<RouteLimits>,
1624 config: &Config,
1625 ) -> Result<Self> {
1626 let admitted_identity = config
1627 .active_provider_identity()
1628 .map_err(anyhow::Error::msg)?;
1629 let api_provider = admitted_identity.provider;
1630 anyhow::ensure!(
1631 api_provider != ProviderKind::Antigravity,
1632 codewhale_config::LEGACY_ANTIGRAVITY_TOMBSTONE_MESSAGE
1633 );
1634 config
1635 .verify_provider_identity(&admitted_identity)
1636 .map_err(anyhow::Error::msg)?;
1637 let openrouter_vendor = config.openrouter_vendor()?;
1638 let billing_surface = crate::route_billing::billing_surface_for_dispatch(
1639 Some(config),
1640 &admitted_identity,
1641 Some(&base_url),
1642 )
1643 .map(str::to_string);
1644 let billing_mode =
1645 crate::route_billing::for_route_with_endpoint(config, &admitted_identity, &base_url)
1646 .into();
1647 if api_provider == ProviderKind::OpenaiCodex {
1648 anyhow::ensure!(
1649 !reqwest::Url::parse(&base_url)
1650 .ok()
1651 .is_some_and(|url| url.host_str() == Some("chatgpt.com")),
1652 "The legacy ChatGPT backend route is retired. Use https://api.openai.com/v1 and run codewhale auth chatgpt."
1653 );
1654 }
1655 if api_provider == ProviderKind::OpencodeGo {
1656 validate_route(api_provider, &default_model).map_err(anyhow::Error::msg)?;
1657 }
1658 anyhow::ensure!(
1659 api_provider != ProviderKind::OpenaiCodex
1660 || crate::pricing::is_official_chatgpt_api(&base_url)
1661 || config.provider_uses_custom_endpoint(&admitted_identity),
1662 "Refusing to send an official ChatGPT grant to a different endpoint"
1663 );
1664 // Guidance and opaque-reasoning scope come from the exact owned
1665 // credential snapshot this client sends, never a second store read.
1666 let (
1667 (api_key, api_key_source),
1668 codex_account_id,
1669 subscription_limit_guidance,
1670 chatgpt_reasoning_api,
1671 ) = if api_provider == ProviderKind::OpenaiCodex
1672 && crate::pricing::is_official_chatgpt_api(&base_url)
1673 {
1674 let credentials = config.codex_credentials()?;
1675 let guidance = crate::oauth::usage_limit_guidance(
1676 crate::oauth::OAuthProvider::Chatgpt,
1677 credentials.account_label.as_deref(),
1678 );
1679 let identity = serde_json::to_vec(&(
1680 credentials.issuer.as_str(),
1681 credentials.client_id.as_str(),
1682 credentials
1683 .account_id
1684 .as_deref()
1685 .context("Verified ChatGPT subject is missing")?,
1686 ))?;
1687 let reasoning_api = format!(
1688 "openai-responses-siwc-v1:{}",
1689 crate::hashing::sha256_hex(&identity),
1690 );
1691 (
1692 (credentials.access_token, "ChatGPT sign-in".to_string()),
1693 credentials.account_id,
1694 Some(guidance),
1695 Some(reasoning_api),
1696 )
1697 } else {
1698 // Only the resolver's xAI OAuth step carries a sign-in; an
1699 // OAuth mode that fell through to an API key names no account.
1700 let (resolved, xai_sign_in) = config.active_route_api_key_with_xai_sign_in()?;
1701 let guidance = xai_sign_in.map(|label| {
1702 crate::oauth::usage_limit_guidance(
1703 crate::oauth::OAuthProvider::Xai,
1704 label.as_deref(),
1705 )
1706 });
1707 (resolved, None, guidance, None)
1708 };
1709 let model_bound_secret_values =
1710 Arc::new(configured_model_bound_secret_values(config, &api_key));
1711 // The opt-out is effective only after an explicit startup confirmation;
1712 // every unconfirmed or absent request stays on the safe default.
1713 let model_bound_masking = !codewhale_config::redaction::effective_masking(
1714 config.model_bound_redaction(),
1715 config.loaded_config_path.as_deref(),
1716 )
1717 .is_disabled();
1718 validate_base_url_security(&base_url, config.allow_insecure_http())?;
1719 let retry = config.retry_policy();
1720 let stream_idle_timeout = Duration::from_secs(config.stream_chunk_timeout_secs());
1721 let stream_open_timeout = config.stream_open_timeout();
1722 let force_http1 = config.force_http1();
1723 if force_http1 {
1724 logging::info(
1725 "HTTP/1.1 pinned (stream configuration or environment) — HTTP/2 disabled",
1726 );
1727 }
1728 let http_headers = config.http_headers();
1729 let auth_disabled = auth_mode_disables_api_key(
1730 config.auth_mode_for_provider(&admitted_identity).as_deref(),
1731 );
1732 let insecure_skip_tls_verify = config.insecure_skip_tls_verify();
1733 let path_suffix = config
1734 .provider_config_for(&admitted_identity)
1735 .and_then(|p| p.path_suffix.clone());
1736 let reasoning_stream_style = config
1737 .provider_config_for(&admitted_identity)
1738 .and_then(|p| p.reasoning_stream_style.clone());
1739 let request_concurrency_limit = config.provider_max_concurrency(&admitted_identity);
1740
1741 logging::info(format!("API provider: {}", api_provider.as_str()));
1742 logging::info(format!(
1743 "API base URL: {}",
1744 redact_url_for_display(&base_url)
1745 ));
1746 if let Some(suffix) = &path_suffix {
1747 logging::info(format!("API path suffix override: {suffix}"));
1748 }
1749 if !http_headers.is_empty() {
1750 logging::info(format!(
1751 "{} custom HTTP header(s) configured",
1752 http_headers.len()
1753 ));
1754 }
1755 if insecure_skip_tls_verify {
1756 logging::warn(format!(
1757 "TLS certificate verification cannot be disabled for provider {}; use SSL_CERT_FILE with a trusted custom CA bundle instead",
1758 api_provider.as_str()
1759 ));
1760 bail!(
1761 "TLS certificate verification cannot be disabled for provider {}; configure SSL_CERT_FILE with a trusted custom CA bundle instead",
1762 api_provider.as_str()
1763 );
1764 }
1765 logging::info(format!(
1766 "Retry policy: enabled={}, max_retries={}, initial_delay={}s, max_delay={}s",
1767 retry.enabled, retry.max_retries, retry.initial_delay, retry.max_delay
1768 ));
1769 if let Some(limit) = request_concurrency_limit {
1770 logging::info(format!(
1771 "Provider request concurrency cap: {} in-flight request(s)",
1772 limit
1773 ));
1774 }
1775
1776 let http_client = Self::http_client_builder_with_auth_mode(
1777 &api_key,
1778 &http_headers,
1779 api_provider,
1780 &base_url,
1781 wire_format,
1782 auth_disabled,
1783 force_http1,
1784 config,
1785 )?
1786 .build()?;
1787 let models_http_client = Self::http_client_builder_with_auth_mode(
1788 &api_key,
1789 &http_headers,
1790 api_provider,
1791 &base_url,
1792 wire_format,
1793 auth_disabled,
1794 force_http1,
1795 config,
1796 )?
1797 .redirect(reqwest::redirect::Policy::none())
1798 .build()?;
1799 // Always keep an HTTP/1.1 twin for automatic stream-header fallback
1800 // when H2 stalls. When `force_http1` is pinned, both clients are
1801 // HTTP/1.1 and the fallback is a no-op retry path.
1802 let http1_client = Self::http_client_builder_with_auth_mode(
1803 &api_key,
1804 &http_headers,
1805 api_provider,
1806 &base_url,
1807 wire_format,
1808 auth_disabled,
1809 true,
1810 config,
1811 )?
1812 .build()?;
1813
1814 let catalog_error_secret_values = Arc::new(catalog_error_secret_values(
1815 &api_key,
1816 &http_headers,
1817 &model_bound_secret_values,
1818 ));
1819 Ok(Self {
1820 http_client,
1821 models_http_client,
1822 http1_client,
1823 api_key,
1824 api_key_source,
1825 subscription_limit_guidance,
1826 model_bound_secret_values,
1827 catalog_error_secret_values,
1828 model_bound_masking,
1829 base_url,
1830 api_provider,
1831 admitted_identity,
1832 openrouter_vendor,
1833 billing_surface,
1834 billing_mode,
1835 configured_models: Arc::new(config.custom_models.clone().unwrap_or_default()),
1836 route_limits,
1837 codex_account_id,
1838 chatgpt_reasoning_api,
1839 wire_format,
1840 retry,
1841 isolated_request_state: false,
1842 remote_control_inference_participant: !config.runtime_chat_isolated
1843 && !config.runtime_thread_inference_unrelated,
1844 default_model,
1845 connection_health: Arc::new(AsyncMutex::new(ConnectionHealth::default())),
1846 rate_limiter: Arc::new(AsyncMutex::new(TokenBucket::from_env())),
1847 request_concurrency: request_concurrency_limit.map(ProviderConcurrencyLimiter::new),
1848 path_suffix,
1849 #[cfg(test)]
1850 test_chat_transport_base_url: None,
1851 #[cfg(test)]
1852 test_messages_transport_base_url: None,
1853 reasoning_stream_style,
1854 stream_idle_timeout,
1855 stream_open_timeout,
1856 force_http1,
1857 })
1858 }
1859
1860 /// Map a failed HTTP response, naming the route, host and key source on
1861 /// authentication and unknown-model failures (#6528) so a rejected key
1862 /// explains which credential was sent where and how to replace it.
1863 fn http_error_with_route_context(
1864 &self,
1865 status: u16,
1866 body: &str,
1867 retry_after: Option<Duration>,
1868 ) -> LlmError {
1869 let route = self.admitted_identity.key.as_str();
1870 let error = match status {
1871 401 | 403 => LlmError::from_http_response_with_auth_context(
1872 status,
1873 body,
1874 Some(
1875 crate::llm_client::AuthenticationErrorContext::from_parts(
1876 Some(route),
1877 Some(&self.base_url),
1878 None,
1879 Some(&self.api_key_source),
1880 Some(&self.api_key),
1881 )
1882 .with_fix(auth_fix_hint(route, &self.api_key_source)),
1883 ),
1884 ),
1885 _ => LlmError::from_http_response_with_retry_after(status, body, retry_after),
1886 };
1887 match error {
1888 LlmError::ModelError(message) => LlmError::ModelError(format!(
1889 "{message} (provider route: {route}, host: {}; run `codewhale model resolve` or /model to pick a model this route serves)",
1890 crate::llm_client::base_url_authority(&self.base_url)
1891 .unwrap_or_else(|| redact_url_for_display(&self.base_url))
1892 )),
1893 LlmError::QuotaExhausted(error) => match self.subscription_limit_guidance.as_deref() {
1894 Some(guidance) => LlmError::QuotaExhausted(error.with_guidance(guidance)),
1895 None => LlmError::QuotaExhausted(error),
1896 },
1897 other => other,
1898 }
1899 }
1900
1901 /// Transport destination for Chat Completions requests.
1902 ///
1903 /// Production always uses the semantic route base URL. Unit tests may
1904 /// substitute a local capture server without changing the endpoint/model
1905 /// identity used by exact-route request shaping.
1906 pub(super) fn chat_transport_base_url(&self) -> &str {
1907 #[cfg(test)]
1908 if let Some(base_url) = self.test_chat_transport_base_url.as_deref() {
1909 return base_url;
1910 }
1911 &self.base_url
1912 }
1913
1914 /// Redirect Chat Completions *transport* to a local capture server while
1915 /// the semantic route (`base_url`, model, endpoint identity) stays exact.
1916 ///
1917 /// Test-only, and compiled out of release builds. Route shaping reads
1918 /// [`Self::base_url`], so an exact-route matrix can capture the real
1919 /// first-turn body for `api.z.ai`, `api.moonshot.ai`, `api.kimi.com`, or
1920 /// `api.minimax.io` without ever making a live provider call.
1921 #[cfg(test)]
1922 pub(crate) fn set_test_chat_transport_base_url(&mut self, base_url: String) {
1923 self.test_chat_transport_base_url = Some(base_url);
1924 }
1925
1926 /// Transport destination for a prepared Anthropic-compatible request.
1927 /// Production sends the exact prepared endpoint; tests may redirect the
1928 /// transport while preserving that immutable endpoint for route shaping.
1929 pub(super) fn messages_transport_url(&self, prepared_url: &str) -> String {
1930 #[cfg(test)]
1931 if let Some(base_url) = self.test_messages_transport_base_url.as_deref() {
1932 return anthropic::anthropic_messages_url(base_url);
1933 }
1934 prepared_url.to_string()
1935 }
1936
1937 /// Return a request whose tool results are safe to send to an upstream
1938 /// model provider.
1939 ///
1940 /// Tool output is untrusted model-bound data: it can contain a whole
1941 /// config file, a bare credential emitted by a shell command, or a
1942 /// spillover receipt whose backing content is later persisted by the chat
1943 /// adapter. Keep this boundary above all protocol adapters so Chat,
1944 /// Anthropic Messages, and OpenAI Responses — streaming and non-streaming
1945 /// alike — receive the same sanitized payload.
1946 fn prepare_model_bound_request(&self, mut request: MessageRequest) -> MessageRequest {
1947 let repair =
1948 crate::tool_history_repair::repair_tool_call_pairs_for_provider(&mut request.messages);
1949 if !repair.is_empty() {
1950 tracing::warn!(
1951 repaired_call_ids = ?repair.repaired_call_ids,
1952 duplicate_result_ids = ?repair.duplicate_result_ids,
1953 orphan_result_ids = ?repair.orphan_result_ids,
1954 "repaired tool call/result history before provider projection"
1955 );
1956 }
1957 for message in &mut request.messages {
1958 for block in &mut message.content {
1959 if let ContentBlock::ToolResult { content, .. } = block
1960 && self.model_bound_masking
1961 {
1962 *content = redact_model_bound_text(content, &self.model_bound_secret_values);
1963 }
1964 }
1965 }
1966 request
1967 }
1968
1969 /// Redact configured credentials from text that has been flattened into a
1970 /// normal model-bound text block. Unlike `prepare_model_bound_request`,
1971 /// this path always redacts: routing/classification prompts (which may
1972 /// summarize tool output) and durable goal-state text are not covered by
1973 /// the `[redaction] model_bound` opt-out, which exists so the model can
1974 /// quote file bytes back for exact edits — never to relax storage or
1975 /// routing summaries.
1976 pub(crate) fn redact_model_bound_text(&self, text: &str) -> String {
1977 redact_model_bound_text(text, &self.model_bound_secret_values)
1978 }
1979
1980 /// Redact tool output as it enters the transcript. Same masking as the
1981 /// request boundary, including the confirmed `[redaction] model_bound`
1982 /// opt-out: a user who chose to let the model see file bytes verbatim
1983 /// keeps that, and everyone else never stores a live credential.
1984 pub(crate) fn redact_tool_output_for_transcript(&self, text: &str) -> String {
1985 if self.model_bound_masking {
1986 redact_model_bound_text(text, &self.model_bound_secret_values)
1987 } else {
1988 text.to_string()
1989 }
1990 }
1991
1992 /// Alternate models share this client's frozen endpoint and declarations.
1993 /// Resolution still owns protocol admission, including closed rosters.
1994 pub(crate) fn resolve_model_route(&self, model: &str) -> Result<ReadyRouteCandidate> {
1995 static RESOLVER: OnceLock<RouteResolver> = OnceLock::new();
1996 let resolver = RESOLVER.get_or_init(RouteResolver::new);
1997 let request = RouteRequest {
1998 explicit_provider: Some(self.api_provider),
1999 model_selector: Some(LogicalModelRef::from(model)),
2000 saved_provider_model: None,
2001 base_url_override: Some(self.base_url.clone()),
2002 limit_overrides: Vec::new(),
2003 };
2004 if self.configured_models.is_empty() || self.api_provider == ProviderKind::OpenaiCodex {
2005 resolver.resolve(&request)
2006 } else {
2007 resolver
2008 .clone()
2009 .with_configured_models(
2010 &self.configured_models,
2011 self.admitted_identity.key.as_str(),
2012 self.api_provider,
2013 &self.base_url,
2014 )
2015 .resolve(&request)
2016 }
2017 .map_err(anyhow::Error::msg)
2018 }
2019
2020 fn declared_wire_model<'a>(&self, model: &'a str) -> Option<&'a str> {
2021 if self.api_provider == ProviderKind::OpenaiCodex
2022 || !self.configured_models.iter().any(|row| {
2023 row.id == model
2024 && row.matches_route(self.admitted_identity.key.as_str(), &self.base_url)
2025 })
2026 {
2027 return None;
2028 }
2029 // A declaration cannot admit a new model to a closed protocol roster,
2030 // or defeat a required protocol alias such as OpenCode Go's allowlist.
2031 self.resolve_model_route(model)
2032 .ok()
2033 .filter(|candidate| candidate.wire_model_id().as_str() == model)
2034 .map(|_| model)
2035 }
2036
2037 fn wire_model_for_route(&self, model: &str) -> String {
2038 self.declared_wire_model(model)
2039 .map(str::to_string)
2040 .unwrap_or_else(|| {
2041 wire_model_for_provider_route(self.api_provider, &self.base_url, model)
2042 })
2043 }
2044
2045 /// Resolve `model` through the central route resolver and rebuild this
2046 /// client whenever its exact wire identity, limits, or protocol differs
2047 /// from the route bound at construction (#5042). `Ok(None)` means the
2048 /// existing binding is already exact for `model`. This includes
2049 /// same-protocol switches: an OpenAI-compatible client for model A must not
2050 /// carry A's route limits into model B merely because both speak Chat.
2051 pub(crate) fn rebound_for_model_protocol(
2052 &self,
2053 config: Option<&Config>,
2054 model: &str,
2055 ) -> Result<Option<Self>> {
2056 // The bound model already owns its admitted protocol and limits,
2057 // including exact live-catalog facts absent from the offline catalog.
2058 if model == self.default_model {
2059 return Ok(None);
2060 }
2061 let candidate = self.resolve_model_route(model)?;
2062 let candidate_limits = crate::route_budget::known_route_limits(candidate.limits());
2063 if candidate.protocol() == self.wire_format
2064 && candidate.wire_model_id().as_str() == self.default_model
2065 && candidate_limits == self.route_limits
2066 {
2067 return Ok(None);
2068 }
2069 if candidate.protocol() == self.wire_format {
2070 // Same-protocol model switches keep the already-authenticated,
2071 // endpoint-bound transport but must freeze the alternate model's
2072 // exact wire identity and limits. This path is also what makes a
2073 // model-aware client usable in embedded/test runtimes that do not
2074 // retain the original Config after construction.
2075 let mut rebound = self.clone();
2076 rebound.default_model = candidate.wire_model_id().as_str().to_string();
2077 rebound.route_limits = candidate_limits;
2078 return Ok(Some(rebound));
2079 }
2080 let config = config.ok_or_else(|| {
2081 anyhow::anyhow!(
2082 "{} model {:?} uses {:?}, but this client is bound to {:?} and no configuration is available to rebuild it",
2083 self.api_provider.provider().display_name(),
2084 model,
2085 candidate.protocol(),
2086 self.wire_format
2087 )
2088 })?;
2089 let mut rebound = Self::from_candidate(config, &candidate)?;
2090 rebound.configured_models = Arc::clone(&self.configured_models);
2091 Ok(Some(rebound))
2092 }
2093
2094 fn bind_request_to_protocol(
2095 &self,
2096 mut request: MessageRequest,
2097 ) -> Result<(MessageRequest, Option<RouteLimits>)> {
2098 let model_aware = self.api_provider.provider().wire_policy()
2099 == codewhale_config::provider::WirePolicy::ModelAware;
2100 // Preserve the exact binding for model-aware providers too. Looking
2101 // up this same model in the offline catalog would discard protocol
2102 // and limit facts already admitted from an exact provider catalog.
2103 if request.model == self.default_model
2104 || (!model_aware && request.model.trim() == self.default_model)
2105 {
2106 return Ok((request, self.route_limits));
2107 }
2108
2109 let candidate = match self.resolve_model_route(&request.model) {
2110 Ok(candidate) => candidate,
2111 Err(error) if model_aware || self.api_provider == ProviderKind::OpencodeGo => {
2112 return Err(error);
2113 }
2114 Err(_) => {
2115 // A fixed-protocol gateway may legitimately accept an id that
2116 // is newer than our offline catalog. Preserve the caller's
2117 // model, but fail closed on limits instead of reusing the
2118 // bound default model's envelope.
2119 return Ok((request, None));
2120 }
2121 };
2122 if candidate.protocol() != self.wire_format {
2123 bail!(
2124 "{} model {:?} uses {:?}, but this client is bound to {:?}; resolve a new model route before sending",
2125 self.api_provider.provider().display_name(),
2126 request.model,
2127 candidate.protocol(),
2128 self.wire_format
2129 );
2130 }
2131 request.model = candidate.wire_model_id().as_str().to_string();
2132 let route_limits = crate::route_budget::known_route_limits(candidate.limits());
2133 Ok((request, route_limits))
2134 }
2135
2136 #[cfg(test)]
2137 fn build_http_client(
2138 api_key: &str,
2139 extra_headers: &HashMap<String, String>,
2140 api_provider: ProviderKind,
2141 base_url: &str,
2142 ) -> Result<reqwest::Client> {
2143 Self::http_client_builder_with_auth_mode(
2144 api_key,
2145 extra_headers,
2146 api_provider,
2147 base_url,
2148 provider_default_wire_format(api_provider),
2149 false,
2150 false,
2151 &Config::default(),
2152 )?
2153 .build()
2154 .map_err(Into::into)
2155 }
2156
2157 fn http_client_builder_with_auth_mode(
2158 api_key: &str,
2159 extra_headers: &HashMap<String, String>,
2160 api_provider: ProviderKind,
2161 base_url: &str,
2162 wire_format: WireFormat,
2163 auth_disabled: bool,
2164 force_http1: bool,
2165 config: &Config,
2166 ) -> Result<reqwest::ClientBuilder> {
2167 let headers = build_default_headers(
2168 api_key,
2169 extra_headers,
2170 api_provider,
2171 base_url,
2172 wire_format,
2173 auth_disabled,
2174 )?;
2175 let mut builder = crate::tls::reqwest_client_builder()
2176 .default_headers(headers)
2177 .user_agent(client_user_agent(api_provider))
2178 .connect_timeout(config.connect_timeout())
2179 .tcp_keepalive(config.tcp_keepalive())
2180 .http2_keep_alive_interval(config.http2_keep_alive_interval())
2181 .http2_keep_alive_timeout(config.http2_keep_alive_timeout())
2182 .min_tls_version(reqwest::tls::Version::TLS_1_2);
2183 if api_provider == ProviderKind::OpenaiCodex {
2184 builder = builder.redirect(reqwest::redirect::Policy::none());
2185 }
2186 if force_http1 {
2187 builder = builder.http1_only();
2188 }
2189 if let Ok(cert_path) = std::env::var("SSL_CERT_FILE")
2190 && !cert_path.is_empty()
2191 {
2192 builder = add_extra_root_certs(builder, &cert_path);
2193 }
2194 Ok(builder)
2195 }
2196
2197 /// HTTP/1.1 client for automatic stream-header fallback.
2198 #[must_use]
2199 pub(crate) fn http1_fallback_client(&self) -> &reqwest::Client {
2200 &self.http1_client
2201 }
2202
2203 #[cfg(test)]
2204 fn default_headers(
2205 api_key: &str,
2206 extra_headers: &HashMap<String, String>,
2207 ) -> Result<HeaderMap> {
2208 build_default_headers(
2209 api_key,
2210 extra_headers,
2211 ProviderKind::Deepseek,
2212 crate::config::DEFAULT_DEEPSEEK_BASE_URL,
2213 WireFormat::ChatCompletions,
2214 false,
2215 )
2216 }
2217
2218 #[cfg(test)]
2219 fn default_headers_for_provider(
2220 api_key: &str,
2221 extra_headers: &HashMap<String, String>,
2222 api_provider: ProviderKind,
2223 base_url: &str,
2224 ) -> Result<HeaderMap> {
2225 build_default_headers(
2226 api_key,
2227 extra_headers,
2228 api_provider,
2229 base_url,
2230 provider_default_wire_format(api_provider),
2231 false,
2232 )
2233 }
2234
2235 #[cfg(test)]
2236 fn default_headers_for_provider_with_auth_disabled(
2237 api_key: &str,
2238 extra_headers: &HashMap<String, String>,
2239 api_provider: ProviderKind,
2240 base_url: &str,
2241 ) -> Result<HeaderMap> {
2242 build_default_headers(
2243 api_key,
2244 extra_headers,
2245 api_provider,
2246 base_url,
2247 provider_default_wire_format(api_provider),
2248 true,
2249 )
2250 }
2251 }
2252
2253 fn build_default_headers(
2254 api_key: &str,
2255 extra_headers: &HashMap<String, String>,
2256 api_provider: ProviderKind,
2257 base_url: &str,
2258 wire_format: WireFormat,
2259 auth_disabled: bool,
2260 ) -> Result<HeaderMap> {
2261 let mut headers = HeaderMap::new();
2262 headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
2263 let api_key = api_key.trim();
2264 let uses_anthropic_messages = wire_format == WireFormat::AnthropicMessages;
2265 if uses_anthropic_messages {
2266 // #3014: most Messages API routes authenticate with `x-api-key`.
2267 // OpenModel also supports Bearer auth for Messages, and its `/models`
2268 // endpoint requires it, so the header chooser below keeps OpenModel on
2269 // Bearer while still pinning the Anthropic wire contract here. The
2270 // Codewhale API is the same shape: its Anthropic passthrough
2271 // authenticates the account key with `Authorization: Bearer` on every
2272 // protocol and does not accept `x-api-key`.
2273 headers.insert(
2274 HeaderName::from_static("anthropic-version"),
2275 HeaderValue::from_static("2023-06-01"),
2276 );
2277 }
2278 let auth_header_name = if auth_disabled {
2279 None
2280 } else if !api_key.is_empty()
2281 && uses_anthropic_messages
2282 && !matches!(
2283 api_provider,
2284 ProviderKind::Openmodel | ProviderKind::Codewhale
2285 )
2286 {
2287 Some(HeaderName::from_static("x-api-key"))
2288 } else if !api_key.is_empty()
2289 && api_provider == ProviderKind::XiaomiMimo
2290 && (xiaomi_mimo_base_url_uses_token_plan(base_url)
2291 || xiaomi_mimo_api_key_uses_token_plan(api_key))
2292 {
2293 Some(HeaderName::from_static("api-key"))
2294 } else if !api_key.is_empty() {
2295 Some(AUTHORIZATION)
2296 } else {
2297 None
2298 };
2299 if let Some(header_name) = auth_header_name.as_ref() {
2300 let header_value = if *header_name == AUTHORIZATION {
2301 HeaderValue::from_str(&format!("Bearer {api_key}"))?
2302 } else {
2303 HeaderValue::from_str(api_key)?
2304 };
2305 headers.insert(header_name.clone(), header_value);
2306 }
2307 // OpenRouter app attribution: these two headers are how apps appear on
2308 // openrouter.ai's app rankings — there is no manual submission. They
2309 // identify the app, never the user, and a user-configured header of the
2310 // same name below still wins (the extra-header loop overwrites).
2311 if api_provider == ProviderKind::Openrouter {
2312 // `HTTP-Referer` is the only required header: it is the app's unique
2313 // identifier, and without it OpenRouter creates no app page at all.
2314 headers.insert(
2315 HeaderName::from_static("http-referer"),
2316 HeaderValue::from_static("https://codewhale.net"),
2317 );
2318 // `X-OpenRouter-Title` is the current display-name header. `X-Title`
2319 // is only kept for backwards compatibility — OpenRouter still honours
2320 // it, and a base-URL override can point this provider at an
2321 // OpenRouter-compatible gateway that knows the old name and not the
2322 // new one. Two bytes of redundancy is cheaper than losing attribution.
2323 //
2324 // Both names carry one display title. A user who configures either
2325 // one (a fork, say, setting only the legacy `X-Title`) must control
2326 // the name on both, or the default left on the other header would
2327 // silently override theirs on OpenRouter.
2328 let user_title = |wanted: &str| {
2329 extra_headers.iter().find_map(|(name, value)| {
2330 let value = value.trim();
2331 (name.trim().eq_ignore_ascii_case(wanted) && !value.is_empty()).then_some(value)
2332 })
2333 };
2334 let title = HeaderValue::from_str(
2335 user_title("x-openrouter-title")
2336 .or_else(|| user_title("x-title"))
2337 .unwrap_or("Codewhale"),
2338 )?;
2339 headers.insert(HeaderName::from_static("x-openrouter-title"), title.clone());
2340 headers.insert(HeaderName::from_static("x-title"), title);
2341 // Marketplace categories, at most two per request. These place the app
2342 // under Coding → CLI Agents and Productivity → Personal Agents on
2343 // openrouter.ai/apps, which is how the rankings page groups entries.
2344 // Unrecognised values are dropped silently, so these must stay exactly
2345 // as OpenRouter spells them.
2346 headers.insert(
2347 HeaderName::from_static("x-openrouter-categories"),
2348 HeaderValue::from_static("cli-agent,personal-agent"),
2349 );
2350 }
2351 // OpenCode Go / OpenCode Zen gateways (https://opencode.ai/docs/go/)
2352 // ask clients to send a stable `x-opencode-session` header so the
2353 // service can optimize prompt caching and attribute traffic to a
2354 // conversation. Generate one UUID v4 per process so every request a
2355 // session makes to the gateway shares a single ID ("one stable ID per
2356 // conversation"). A user-configured header of the same name wins: the
2357 // extra-header loop below overwrites this value.
2358 if matches!(
2359 api_provider,
2360 ProviderKind::OpencodeGo | ProviderKind::OpencodeZen
2361 ) {
2362 static OPENCODE_SESSION: OnceLock<String> = OnceLock::new();
2363 let session = OPENCODE_SESSION.get_or_init(|| uuid::Uuid::new_v4().to_string());
2364 headers.insert(
2365 HeaderName::from_static("x-opencode-session"),
2366 HeaderValue::from_str(session)?,
2367 );
2368 }
2369 for (name, value) in extra_headers {
2370 let name = name.trim();
2371 let value = value.trim();
2372 if name.is_empty() || value.is_empty() {
2373 continue;
2374 }
2375 if auth_disabled && is_upstream_auth_header(name) {
2376 continue;
2377 }
2378 let header_name = HeaderName::from_bytes(name.as_bytes())?;
2379 if header_name == AUTHORIZATION
2380 || header_name == CONTENT_TYPE
2381 || auth_header_name.as_ref() == Some(&header_name)
2382 || (auth_header_name.is_some() && is_auth_dialect_header(&header_name))
2383 {
2384 continue;
2385 }
2386 headers.insert(header_name, HeaderValue::from_str(value)?);
2387 }
2388 Ok(headers)
2389 }
2390
2391 fn is_auth_dialect_header(header_name: &HeaderName) -> bool {
2392 header_name == AUTHORIZATION
2393 || header_name == HeaderName::from_static("api-key")
2394 || header_name == HeaderName::from_static("x-api-key")
2395 }
2396
2397 fn provider_default_wire_format(provider: ProviderKind) -> WireFormat {
2398 provider
2399 .provider()
2400 .wire_policy()
2401 .fixed()
2402 .unwrap_or_else(|| {
2403 if provider == ProviderKind::OpencodeZen {
2404 WireFormat::Responses
2405 } else {
2406 WireFormat::ChatCompletions
2407 }
2408 })
2409 }
2410
2411 /// Resolve the wire dialect for a dual-protocol vendor.
2412 ///
2413 /// Power-user toggle: `providers.<id>.wire = "openai" | "anthropic" | "responses"`.
2414 /// Legacy dialect kinds (`*Anthropic`) still force Messages. Custom providers
2415 /// honor `wire = "responses" | "anthropic" | "chat"` per-config (see
2416 /// `crates/config/src/provider.rs:Custom`). Everyone else keeps the descriptor's
2417 /// fixed policy (or Chat Completions).
2418 fn provider_wire_format_for_config(
2419 identity: &ProviderIdentity,
2420 config: Option<&crate::config::Config>,
2421 ) -> WireFormat {
2422 let api_provider = identity.provider;
2423 let wire = config.and_then(|cfg| cfg.provider_wire_dialect(identity));
2424 let prefers_anthropic = matches!(
2425 api_provider,
2426 ProviderKind::DeepseekAnthropic
2427 | ProviderKind::MinimaxAnthropic
2428 | ProviderKind::ModelstudioTokenPlanAnthropic
2429 | ProviderKind::ModelstudioCodingPlanAnthropic
2430 ) || wire_config_prefers_anthropic(wire);
2431
2432 if prefers_anthropic
2433 && matches!(
2434 api_provider,
2435 ProviderKind::Deepseek
2436 | ProviderKind::Minimax
2437 | ProviderKind::ModelstudioTokenPlan
2438 | ProviderKind::DeepseekAnthropic
2439 | ProviderKind::MinimaxAnthropic
2440 | ProviderKind::ModelstudioTokenPlanAnthropic
2441 | ProviderKind::ModelstudioCodingPlan
2442 | ProviderKind::ModelstudioCodingPlanAnthropic
2443 )
2444 {
2445 return WireFormat::AnthropicMessages;
2446 }
2447
2448 // Custom providers honor `wire = "anthropic"` / `wire = "responses"` explicitly.
2449 // The static `Custom::wire_policy()` remains `Chat` as a safe default; the
2450 // per-config override lives here (and in `provider_capability`) so existing
2451 // `[providers.<name>]` tables gain the three-way switch without changing the
2452 // provider registry trait. Supported aliases:
2453 // anthropic: "anthropic" | "messages" | "claude" | "anthropic-messages" | ...
2454 // responses: "responses" | "responses-api" | "openai-responses" | "openai_responses" | ...
2455 if api_provider == ProviderKind::Custom {
2456 if wire_config_prefers_anthropic(wire) {
2457 return WireFormat::AnthropicMessages;
2458 }
2459 if wire_config_prefers_responses(wire) {
2460 return WireFormat::Responses;
2461 }
2462 }
2463
2464 provider_default_wire_format(api_provider)
2465 }
2466
2467 fn wire_config_prefers_anthropic(wire: Option<&str>) -> bool {
2468 let Some(raw) = wire.map(str::trim).filter(|value| !value.is_empty()) else {
2469 return false;
2470 };
2471 let normalized = raw.to_ascii_lowercase().replace(['_', ' '], "-");
2472 matches!(
2473 normalized.as_str(),
2474 "anthropic"
2475 | "anthropic-messages"
2476 | "messages"
2477 | "claude"
2478 | "anthropic-compatible"
2479 | "anthropic-compat"
2480 )
2481 }
2482
2483 fn wire_config_prefers_responses(wire: Option<&str>) -> bool {
2484 let Some(raw) = wire.map(str::trim).filter(|value| !value.is_empty()) else {
2485 return false;
2486 };
2487 let normalized = raw.to_ascii_lowercase().replace(['_', ' '], "-");
2488 matches!(
2489 normalized.as_str(),
2490 "responses"
2491 | "responses-api"
2492 | "openai-responses"
2493 | "openai-responses-api"
2494 | "response"
2495 | "response-api"
2496 | "openai-responses-compat"
2497 | "responses-compat"
2498 )
2499 }
2500
2501 fn api_provider_skips_models_probe(api_provider: ProviderKind) -> bool {
2502 // Concentrate's `GET /v1/models` is explicitly unauthenticated
2503 // (docs/PROVIDERS.md): a 2xx proves nothing about the key, so guided
2504 // setup must not count it as verification — the key is observed on the
2505 // first authenticated call instead.
2506 matches!(
2507 api_provider,
2508 ProviderKind::DeepseekAnthropic | ProviderKind::Concentrate
2509 )
2510 }
2511
2512 #[must_use]
2513 pub(crate) fn provider_api_key_verification_is_observed(api_provider: ProviderKind) -> bool {
2514 !api_provider_skips_models_probe(api_provider)
2515 }
2516
2517 /// Verify a provider API key by hitting the `/models` endpoint
2518 /// (#3875). Builds a minimal HTTP client with the canonical auth
2519 /// headers for `provider`, issues a single GET, and returns
2520 /// `Ok(roster)` on a 2xx response or `Err(reason)` on any failure.
2521 ///
2522 /// `roster` is the listing that 2xx already carried, projected by the same
2523 /// rules as a catalog refresh and scoped to `provider`'s own id: the caller
2524 /// rescopes it to the exact route identity and publishes it, so guided setup
2525 /// offers the models this key can call today rather than bundled or
2526 /// Models.dev rows the provider may have retired. It is `None` when the body
2527 /// is not one complete, valid roster (a paginated first page, malformed or
2528 /// empty JSON); a valid 2xx still verifies the key, and the existing rows stay.
2529 ///
2530 /// This is intentionally a one-shot call — no retry, no rate-limit
2531 /// wait — so a bad key is surfaced immediately.
2532 pub async fn verify_provider_api_key(
2533 provider: ProviderKind,
2534 api_key: &str,
2535 base_url: &str,
2536 ) -> Result<Option<ProviderCatalogDelta>, String> {
2537 if api_provider_skips_models_probe(provider) {
2538 // Providers without a /models endpoint can't be verified this
2539 // way; accept the key optimistically (same as health_check).
2540 return Ok(None);
2541 }
2542 let headers = build_default_headers(
2543 api_key,
2544 &Default::default(),
2545 provider,
2546 base_url,
2547 provider_default_wire_format(provider),
2548 false,
2549 )
2550 .map_err(|err| format!("failed to build auth headers: {err:#}"))?;
2551 let client = crate::tls::reqwest_client_builder()
2552 .default_headers(headers)
2553 .redirect(reqwest::redirect::Policy::none())
2554 .user_agent(concat!(
2555 "Mozilla/5.0 (compatible; codewhale/",
2556 env!("CARGO_PKG_VERSION"),
2557 "; +https://github.com/codewhale-hq/CodeWhale)"
2558 ))
2559 .connect_timeout(Duration::from_secs(10))
2560 .timeout(Duration::from_secs(15))
2561 .build()
2562 .map_err(|err| format!("failed to build HTTP client: {err:#}"))?;
2563 let url = api_url(base_url, "models");
2564 let response = client
2565 .get(&url)
2566 .send()
2567 .await
2568 .map_err(|err| format!("request failed: {}", err.without_url()))?;
2569 let status = response.status();
2570 if status.is_success() {
2571 // A valid 2xx verifies the key even when the body is not a usable
2572 // roster; failure-preserving catalog semantics then keep the
2573 // existing rows.
2574 let body = bounded_provider_catalog_text(response, PROVIDER_CATALOG_MAX_RESPONSE_BYTES)
2575 .await
2576 .unwrap_or_default();
2577 let complete =
2578 serde_json::from_str::<ModelsPage<'_>>(&body).is_ok_and(|page| !page.has_more);
2579 let secrets = [api_key.trim().to_string()];
2580 Ok(complete
2581 .then(|| {
2582 catalog_delta_from_models_body(
2583 provider,
2584 provider.as_str().to_string(),
2585 base_url,
2586 &body,
2587 &secrets,
2588 )
2589 .ok()
2590 })
2591 .flatten())
2592 } else {
2593 let body = bounded_error_text(response, ERROR_BODY_MAX_BYTES).await;
2594 let body = if api_key.trim().is_empty() {
2595 body
2596 } else {
2597 redact_model_bound_text(&body, &[api_key.trim().to_string()])
2598 };
2599 let summary = if body.chars().count() > 200 {
2600 format!("{}...", body.chars().take(200).collect::<String>())
2601 } else {
2602 body
2603 };
2604 Err(format!("HTTP {status}: {summary}"))
2605 }
2606 }
2607
2608 fn translation_system_prompt(target_language: &str) -> String {
2609 format!(
2610 "You are a professional translator. Your ONLY task is to translate text to {target_language}. \
2611 Rules:\n\
2612 1. Output ONLY the translation, nothing else — no explanations, no notes, no quotes.\n\
2613 2. Preserve all code blocks (```...```), URLs, file paths, command names, \
2614 and technical terms like API names, function names, and library names untranslated.\n\
2615 3. Keep Markdown formatting (headings, lists, bold, italics, links) intact.\n\
2616 4. Translate all natural-language prose naturally and professionally.\n\
2617 5. Do NOT add any prefix, suffix, or commentary.\n\
2618 6. If the input is already in {target_language} or contains no prose to translate, \
2619 return it as-is."
2620 )
2621 }
2622
2623 fn translation_message_request(
2624 text: &str,
2625 model: String,
2626 target_language: &str,
2627 max_tokens: u32,
2628 ) -> MessageRequest {
2629 MessageRequest {
2630 model,
2631 messages: vec![Message {
2632 role: Role::User,
2633 content: vec![ContentBlock::Text {
2634 text: text.to_string(),
2635 cache_control: None,
2636 }],
2637 }],
2638 max_tokens,
2639 system: Some(SystemPrompt::Text(translation_system_prompt(
2640 target_language,
2641 ))),
2642 tools: None,
2643 tool_choice: None,
2644 metadata: None,
2645 thinking: None,
2646 reasoning_effort: Some("off".to_string()),
2647 stream: Some(false),
2648 temperature: None,
2649 top_p: None,
2650 }
2651 }
2652
2653 fn translation_text_from_response(response: &MessageResponse) -> Result<String> {
2654 let translated = response
2655 .content
2656 .iter()
2657 .filter_map(|block| match block {
2658 ContentBlock::Text { text, .. } => Some(text.as_str()),
2659 _ => None,
2660 })
2661 .collect::<Vec<_>>()
2662 .join("")
2663 .trim()
2664 .to_string();
2665 if translated.is_empty() {
2666 bail!("translate: Anthropic Messages response did not contain text content");
2667 }
2668 Ok(translated)
2669 }
2670
2671 fn xiaomi_mimo_base_url_uses_token_plan(base_url: &str) -> bool {
2672 let normalized = base_url.trim().to_ascii_lowercase();
2673 let without_scheme = normalized
2674 .strip_prefix("https://")
2675 .or_else(|| normalized.strip_prefix("http://"))
2676 .unwrap_or(&normalized);
2677 let host = without_scheme
2678 .split(['/', '?', '#'])
2679 .next()
2680 .unwrap_or_default();
2681 let host = host.split(':').next().unwrap_or(host);
2682 host.starts_with("token-plan-") && host.ends_with(".xiaomimimo.com")
2683 }
2684
2685 fn xiaomi_mimo_api_key_uses_token_plan(api_key: &str) -> bool {
2686 api_key.trim_start().starts_with("tp-")
2687 }
2688
2689 impl CodewhaleClient {
2690 /// Returns the API base URL used by this client.
2691 pub fn base_url(&self) -> &str {
2692 &self.base_url
2693 }
2694
2695 /// Prepare — but do not send — the exact outbound request for `request`.
2696 ///
2697 /// This is *the* outbound seam (#1004). Production dispatch
2698 /// (`create_message`, `create_message_stream`) and `/preview-request` both
2699 /// call it, so a preview cannot describe a request different from the one
2700 /// a turn would send.
2701 ///
2702 /// It runs, in production order:
2703 ///
2704 /// 1. tool-history repair and model-bound secret redaction
2705 /// ([`Self::prepare_model_bound_request`]);
2706 /// 2. protocol binding and route model re-resolution
2707 /// ([`Self::bind_request_to_protocol`]);
2708 /// 3. the dialect's own body builder — Chat Completions, Anthropic
2709 /// Messages, or OpenAI Responses — including every provider-specific
2710 /// sanitizer and reasoning shaper;
2711 /// 4. exact endpoint resolution for that dialect and route shape.
2712 ///
2713 /// It performs no I/O and mutates no client state.
2714 pub(crate) fn prepare_outbound_request(
2715 &self,
2716 request: MessageRequest,
2717 stream: bool,
2718 ) -> Result<PreparedOutboundRequest> {
2719 // Step 0: refuse role/dialect pairs this wire cannot represent, before
2720 // any dialect builds a body. Doing it here rather than inside each
2721 // adapter is what stops an unrepresentable role from being discovered
2722 // as an opaque provider 400 (Anthropic) or from vanishing silently
2723 // (the OpenAI-shaped dialects) depending on which adapter ran.
2724 let outbound_dialect = WireDialect::from_wire_format(self.wire_format);
2725 role_placement::reject_unsupported_roles(&request.messages, outbound_dialect)?;
2726 let clamp_output_cap = |mut request: MessageRequest, route_limits: Option<RouteLimits>| {
2727 let route_cap =
2728 self.effective_max_output_tokens_with_limits(&request.model, route_limits);
2729 if request.max_tokens > route_cap {
2730 tracing::debug!(
2731 requested_max_tokens = request.max_tokens,
2732 route_max_tokens = route_cap,
2733 model = %request.model,
2734 "clamped outbound max_tokens to the resolved route envelope"
2735 );
2736 request.max_tokens = route_cap;
2737 }
2738 request
2739 };
2740 let (request, request_route_limits) =
2741 self.bind_request_to_protocol(self.prepare_model_bound_request(request))?;
2742 let mut request = clamp_output_cap(request, request_route_limits);
2743 let declared_wire_model = self.declared_wire_model(&request.model).map(str::to_string);
2744 if self.is_local_ds4_model(&request.model)
2745 && let Some(tools) = request.tools.as_mut()
2746 {
2747 for tool in tools {
2748 tool.strict = None;
2749 }
2750 }
2751 let requested_effort = request.reasoning_effort.clone();
2752 // Same value computed for the seam above.
2753 let dialect = outbound_dialect;
2754 // `stream` is the caller's entry point, not a wire fact: each dialect
2755 // decides for itself what the body's `stream` field says.
2756 let entrypoint = CallerStreamMode::from_stream_flag(stream);
2757
2758 match self.wire_format {
2759 WireFormat::ChatCompletions => {
2760 let chat_shape_provider = self.chat_shape_provider(&request.model);
2761 let mut wire = chat::build_chat_wire_body(
2762 &request,
2763 chat_shape_provider,
2764 &self.base_url,
2765 stream,
2766 request_route_limits,
2767 )?;
2768 if let Some(model) = &declared_wire_model {
2769 wire.model.clone_from(model);
2770 wire.body["model"] = json!(model);
2771 }
2772 self.apply_provider_routing(&mut wire.body);
2773 let url = chat_completions_url(
2774 self.chat_transport_base_url(),
2775 &self.base_url,
2776 self.api_provider,
2777 self.path_suffix.as_deref(),
2778 &wire.body,
2779 );
2780 let shape = prepared::chat_route_shape(
2781 self.api_provider,
2782 &self.base_url,
2783 &wire.model,
2784 &url,
2785 );
2786 Ok(PreparedOutboundRequest::new(
2787 dialect,
2788 self.endpoint_identity(url, shape),
2789 wire.model,
2790 wire.body,
2791 requested_effort,
2792 wire.replay_input_tokens,
2793 entrypoint,
2794 )
2795 .with_omitted_tool_names(wire.omitted_tool_names))
2796 }
2797 WireFormat::AnthropicMessages => {
2798 let mut body = self.build_anthropic_body(&request, stream);
2799 if let Some(model) = &declared_wire_model {
2800 body["model"] = json!(model);
2801 }
2802 let url = anthropic::anthropic_messages_url(&self.base_url);
2803 let shape = if self.api_provider == ProviderKind::OpencodeZen {
2804 RouteShape::OpencodeZen
2805 } else if self.api_provider == ProviderKind::Custom {
2806 RouteShape::CustomCompatible
2807 } else {
2808 RouteShape::Standard
2809 };
2810 let wire_model = body
2811 .get("model")
2812 .and_then(Value::as_str)
2813 .unwrap_or(request.model.as_str())
2814 .to_string();
2815 Ok(PreparedOutboundRequest::new(
2816 dialect,
2817 self.endpoint_identity(url, shape),
2818 wire_model,
2819 body,
2820 requested_effort,
2821 None,
2822 entrypoint,
2823 ))
2824 }
2825 WireFormat::Responses => {
2826 let mut body = responses::build_responses_body_for_provider(
2827 &request,
2828 self.api_provider,
2829 self.chatgpt_reasoning_api.as_deref(),
2830 );
2831 if let Some(model) = &declared_wire_model {
2832 body["model"] = json!(model);
2833 }
2834 let is_codex = self.api_provider == ProviderKind::OpenaiCodex;
2835 let url = responses_api_url(&self.base_url, self.api_provider);
2836 let shape = if is_codex {
2837 RouteShape::CodexResponses
2838 } else if self.api_provider == ProviderKind::OpencodeZen {
2839 RouteShape::OpencodeZen
2840 } else if self.api_provider == ProviderKind::Custom {
2841 RouteShape::CustomCompatible
2842 } else {
2843 RouteShape::Standard
2844 };
2845 let wire_model = body
2846 .get("model")
2847 .and_then(Value::as_str)
2848 .unwrap_or(request.model.as_str())
2849 .to_string();
2850 Ok(PreparedOutboundRequest::new(
2851 dialect,
2852 self.endpoint_identity(url, shape),
2853 wire_model,
2854 body,
2855 requested_effort,
2856 None,
2857 entrypoint,
2858 ))
2859 }
2860 }
2861 }
2862
2863 pub(crate) fn apply_provider_routing(&self, body: &mut Value) {
2864 apply_openrouter_vendor(body, self.openrouter_vendor.as_deref());
2865 }
2866
2867 pub(crate) fn openrouter_vendor(&self) -> Option<&str> {
2868 self.openrouter_vendor.as_deref()
2869 }
2870
2871 /// Typed identity of the endpoint this client would POST to.
2872 ///
2873 /// `route_id` is left empty here on purpose: the client knows the provider
2874 /// and the URL, but only the caller's resolved turn plan knows whether the
2875 /// user reached this route through a named custom-provider entry. The
2876 /// engine attaches it with [`PreparedOutboundRequest::with_route_id`].
2877 fn endpoint_identity(&self, url: String, shape: RouteShape) -> EndpointIdentity {
2878 EndpointIdentity {
2879 provider_id: self.api_provider.as_str().to_string(),
2880 provider_display: self.api_provider.provider().display_name().to_string(),
2881 route_id: None,
2882 url,
2883 shape,
2884 }
2885 }
2886
2887 /// Returns the active API provider for this client.
2888 pub fn api_provider(&self) -> ProviderKind {
2889 self.api_provider
2890 }
2891
2892 pub(crate) fn admitted_provider_identity(&self) -> &ProviderIdentity {
2893 &self.admitted_identity
2894 }
2895
2896 /// Route limits frozen with this client at resolution time.
2897 #[must_use]
2898 pub fn route_limits(&self) -> Option<RouteLimits> {
2899 self.route_limits
2900 }
2901
2902 /// Output cap for a request dispatched by this exact client route.
2903 #[must_use]
2904 pub fn effective_max_output_tokens(&self, requested_model: &str) -> u32 {
2905 let route_limits = if requested_model.trim() == self.default_model {
2906 self.route_limits
2907 } else {
2908 self.resolve_model_route(requested_model)
2909 .ok()
2910 .and_then(|candidate| crate::route_budget::known_route_limits(candidate.limits()))
2911 };
2912 self.effective_max_output_tokens_with_limits(requested_model, route_limits)
2913 }
2914
2915 #[must_use]
2916 fn effective_max_output_tokens_with_limits(
2917 &self,
2918 requested_model: &str,
2919 route_limits: Option<RouteLimits>,
2920 ) -> u32 {
2921 let wire_model = self.wire_model_for_route(requested_model);
2922 crate::route_budget::effective_max_output_tokens_for_route(
2923 self.api_provider,
2924 &wire_model,
2925 route_limits,
2926 )
2927 }
2928
2929 /// Secret-free receipt for the exact base endpoint and credential
2930 /// generation this client was constructed with.
2931 ///
2932 /// This is the only way the API key leaves `client.rs`, and it leaves as a
2933 /// one-way digest. Minting the receipt here — rather than re-reading config
2934 /// at some later lifecycle point — is what makes it immutable proof of the
2935 /// route that was actually installed for the turn.
2936 #[must_use]
2937 pub fn turn_route_receipt(&self) -> crate::route_receipt::TurnRouteReceipt {
2938 crate::route_receipt::TurnRouteReceipt::from_admitted(
2939 &self.admitted_identity,
2940 &self.default_model,
2941 &self.base_url,
2942 &self.api_key,
2943 )
2944 .with_openrouter_vendor(self.openrouter_vendor.as_deref())
2945 }
2946
2947 /// The operator-declared `[[custom_models]]` rate for `model` on this
2948 /// client's exact endpoint, frozen at `dispatched_at`. Every dispatch
2949 /// boundary (background envelopes and main interactive turns alike) asks
2950 /// this before the provider lake, so a declared rate is honored the same
2951 /// way on each (#6690).
2952 #[must_use]
2953 pub(crate) fn configured_pricing_quote_at(
2954 &self,
2955 provider: ProviderKind,
2956 provider_identity: &str,
2957 model: &str,
2958 dispatched_at: u64,
2959 ) -> Option<crate::provider_catalog_live::ProviderLivePricingQuote> {
2960 crate::provider_catalog_live::configured_dispatch_pricing_quote_at(
2961 &self.configured_models,
2962 provider,
2963 provider_identity,
2964 model,
2965 &self.base_url,
2966 dispatched_at,
2967 )
2968 }
2969
2970 /// Capture the immutable, redacted route envelope at the caller's
2971 /// application-dispatch/admission time. This is not proof of network
2972 /// delivery or provider invoice-time pricing. The wire model is normalized
2973 /// exactly as the transport will normalize it; a provider-returned alias
2974 /// must never replace this billing identity later.
2975 #[must_use]
2976 pub fn effective_route_envelope(
2977 &self,
2978 requested_model: &str,
2979 dispatched_at: chrono::DateTime<chrono::Utc>,
2980 ) -> crate::cost_status::EffectiveRouteEnvelope {
2981 let model = self.wire_model_for_route(requested_model);
2982 let endpoint_fingerprint = crate::cost_status::endpoint_fingerprint(&self.base_url);
2983 let provider_live_pricing = u64::try_from(dispatched_at.timestamp())
2984 .ok()
2985 .and_then(|at| {
2986 crate::provider_catalog_live::declared_or_catalog_quote(
2987 self.configured_pricing_quote_at(
2988 self.api_provider,
2989 self.admitted_identity.key.as_str(),
2990 &model,
2991 at,
2992 ),
2993 || {
2994 crate::provider_catalog_live::fresh_dispatch_pricing_quote_at(
2995 self.api_provider,
2996 self.admitted_identity.key.as_str(),
2997 &model,
2998 &self.base_url,
2999 at,
3000 )
3001 },
3002 )
3003 });
3004 crate::cost_status::EffectiveRouteEnvelope {
3005 openrouter_vendor: self.openrouter_vendor.clone(),
3006 provider: self.api_provider,
3007 provider_identity: self.admitted_identity.key.to_string(),
3008 model,
3009 billing_surface: self.billing_surface.clone(),
3010 endpoint_fingerprint,
3011 provider_live_pricing,
3012 billing_mode: self.billing_mode,
3013 dispatched_at,
3014 }
3015 }
3016
3017 /// Resolved in-flight provider request cap, if one is active.
3018 #[must_use]
3019 pub fn provider_request_concurrency_limit(&self) -> Option<usize> {
3020 self.request_concurrency
3021 .as_ref()
3022 .map(ProviderConcurrencyLimiter::limit)
3023 }
3024
3025 /// Number of currently active requests held by this client's shared
3026 /// provider request limiter.
3027 #[must_use]
3028 pub fn active_provider_requests(&self) -> usize {
3029 self.request_concurrency
3030 .as_ref()
3031 .map_or(0, ProviderConcurrencyLimiter::active)
3032 }
3033
3034 async fn acquire_provider_request_permit(&self) -> Option<ProviderRequestPermit> {
3035 match self.request_concurrency.as_ref() {
3036 Some(limiter) => limiter.acquire().await,
3037 None => None,
3038 }
3039 }
3040
3041 pub(crate) async fn acquire_remote_control_inference_permit(
3042 &self,
3043 ) -> Option<RemoteControlInferencePermit> {
3044 if !self.remote_control_inference_participant {
3045 return None;
3046 }
3047 Some(acquire_remote_control_inference_participant().await)
3048 }
3049
3050 fn hold_provider_request_permit_for_stream(
3051 stream: crate::llm_client::StreamEventBox,
3052 permit: Option<ProviderRequestPermit>,
3053 ) -> crate::llm_client::StreamEventBox {
3054 Box::pin(async_stream::stream! {
3055 let _permit = permit;
3056 let mut stream = stream;
3057 while let Some(event) = stream.next().await {
3058 yield event;
3059 }
3060 })
3061 }
3062
3063 fn prepend_tool_projection_warning(
3064 stream: crate::llm_client::StreamEventBox,
3065 provider: String,
3066 omitted_tool_names: Vec<String>,
3067 omitted_tool_count: usize,
3068 ) -> crate::llm_client::StreamEventBox {
3069 Box::pin(async_stream::stream! {
3070 yield Ok(codewhale_models::StreamEvent::ToolProjectionWarning {
3071 provider,
3072 omitted_tool_names,
3073 omitted_tool_count,
3074 });
3075 let mut stream = stream;
3076 while let Some(event) = stream.next().await {
3077 yield event;
3078 }
3079 })
3080 }
3081
3082 fn hold_remote_control_inference_permit_for_stream(
3083 stream: crate::llm_client::StreamEventBox,
3084 permit: Option<RemoteControlInferencePermit>,
3085 ) -> crate::llm_client::StreamEventBox {
3086 Box::pin(async_stream::stream! {
3087 let _permit = permit;
3088 let mut stream = stream;
3089 while let Some(event) = stream.next().await {
3090 yield event;
3091 }
3092 })
3093 }
3094
3095 /// Translate text to the requested target language using a focused
3096 /// non-streaming chat completion call on the supplied model.
3097 ///
3098 /// This is a lightweight translation service — no tool calls, no
3099 /// streaming, no conversation history. The dedicated translation agent
3100 /// receives the source text and returns only the translated result.
3101 pub(crate) async fn translate_with_usage(
3102 &self,
3103 text: &str,
3104 model: &str,
3105 target_language: &str,
3106 ) -> Result<TranslationProviderResponse> {
3107 // Freeze pricing before either the remote-control gate or the provider
3108 // permit. A later live-catalog refresh must not reprice this request.
3109 let route = self.effective_route_envelope(model, chrono::Utc::now());
3110 let _inference = self.acquire_remote_control_inference_permit().await;
3111 let _permit = self.acquire_provider_request_permit().await;
3112 let model = self.wire_model_for_route(model);
3113 let max_tokens = self.effective_max_output_tokens(&model);
3114 if self.wire_format != WireFormat::ChatCompletions {
3115 // Non-Chat dialects reuse the prepared-request seam so translation
3116 // cannot drift from production shaping. Translation is still an
3117 // *auxiliary* call, not a primary agent turn: the Chat dialect
3118 // below builds its own small fixed body, and `/preview-request`
3119 // deliberately does not claim to describe either
3120 // (see `docs/PREVIEW_REQUEST.md`).
3121 let prepared = self.prepare_outbound_request(
3122 translation_message_request(text, model, target_language, max_tokens),
3123 false,
3124 )?;
3125 let response = match prepared.dialect {
3126 WireDialect::OpenAiResponses => self.handle_responses_message(&prepared).await?,
3127 WireDialect::AnthropicMessages => self.handle_anthropic_message(&prepared).await?,
3128 WireDialect::ChatCompletions => unreachable!(),
3129 };
3130 let usage = (response.usage != Usage::default()).then_some(response.usage.clone());
3131 let translated =
3132 if codewhale_models::is_incomplete_stop_reason(response.stop_reason.as_deref()) {
3133 Err(anyhow::anyhow!(
3134 "translate: provider response incomplete ({})",
3135 codewhale_models::stop_reason_detail(response.stop_reason.as_deref())
3136 ))
3137 } else {
3138 translation_text_from_response(&response)
3139 };
3140 return Ok(TranslationProviderResponse {
3141 translated,
3142 route,
3143 usage,
3144 });
3145 }
3146
3147 let url = api_url_with_suffix(
3148 self.chat_transport_base_url(),
3149 "chat/completions",
3150 self.path_suffix.as_deref(),
3151 );
3152 let mut body = serde_json::json!({
3153 "model": model,
3154 "messages": [
3155 {
3156 "role": "system",
3157 "content": translation_system_prompt(target_language)
3158 },
3159 {
3160 "role": "user",
3161 "content": text
3162 }
3163 ],
3164 "max_tokens": max_tokens,
3165 "stream": false
3166 });
3167 chat::apply_route_reasoning_controls(
3168 &mut body,
3169 self.api_provider,
3170 &self.base_url,
3171 &model,
3172 Some("off"),
3173 );
3174
3175 self.apply_provider_routing(&mut body);
3176 let response = self.send_json_with_retry(&url, &body).await?;
3177 let status = response.status();
3178 if !status.is_success() {
3179 let raw_error_text = bounded_error_text(response, ERROR_BODY_MAX_BYTES).await;
3180 let error_text = sanitize_http_error_body(
3181 Some(self.api_provider.provider().display_name()),
3182 status.as_u16(),
3183 &raw_error_text,
3184 );
3185 anyhow::bail!("translate: HTTP {status}: {error_text}");
3186 }
3187
3188 let value: serde_json::Value = response.json().await?;
3189 let usage_reported = value
3190 .get("usage")
3191 .and_then(Value::as_object)
3192 .is_some_and(|usage| {
3193 [
3194 "input_tokens",
3195 "prompt_tokens",
3196 "output_tokens",
3197 "completion_tokens",
3198 "total_tokens",
3199 ]
3200 .iter()
3201 .any(|field| usage.contains_key(*field))
3202 });
3203 let usage = parse_usage(value.get("usage"));
3204 let stop_reason = value["choices"][0]["finish_reason"].as_str();
3205 let usage = (usage_reported && usage != Usage::default()).then_some(usage);
3206 let translated = if codewhale_models::is_incomplete_stop_reason(stop_reason) {
3207 Err(anyhow::anyhow!(
3208 "translate: provider response incomplete ({})",
3209 codewhale_models::stop_reason_detail(stop_reason)
3210 ))
3211 } else {
3212 value["choices"][0]["message"]["content"]
3213 .as_str()
3214 .ok_or_else(|| anyhow::anyhow!("translate: unexpected API response shape"))
3215 .and_then(|translated| {
3216 let translated = translated.trim().to_string();
3217 if translated.is_empty() {
3218 bail!("translate: provider response did not contain text content");
3219 }
3220 Ok(translated)
3221 })
3222 };
3223
3224 Ok(TranslationProviderResponse {
3225 translated,
3226 route,
3227 usage,
3228 })
3229 }
3230
3231 /// Test adapter for asserting the translated text; production retains receipts.
3232 #[cfg(test)]
3233 async fn translate(&self, text: &str, model: &str, target_language: &str) -> Result<String> {
3234 self.translate_with_usage(text, model, target_language)
3235 .await?
3236 .translated
3237 }
3238
3239 /// List every available model under the endpoint's verified pagination contract.
3240 pub async fn list_models(&self) -> Result<Vec<AvailableModel>> {
3241 let (body, deadline) = self
3242 .models_document(ModelsRequestMode::Interactive)
3243 .await
3244 .map_err(ModelsFetchError::into_interactive)?;
3245 let models = parse_models_response_for_provider(&body, self.api_provider)
3246 .map(|models| apply_provider_model_cutline(self.api_provider, models))
3247 .map_err(|_| {
3248 ModelsFetchError::Catalog(CatalogRefreshError::InvalidResponse).into_interactive()
3249 })?;
3250 if tokio::time::Instant::now() >= deadline {
3251 return Err(ModelsFetchError::Catalog(CatalogRefreshError::Network).into_interactive());
3252 }
3253 if self.api_provider == ProviderKind::OpenaiCodex
3254 && models.iter().any(|model| {
3255 self.model_bound_secret_values.iter().any(|secret| {
3256 !secret.is_empty()
3257 && (model.id.contains(secret.as_str())
3258 || model
3259 .display_name
3260 .as_deref()
3261 .is_some_and(|name| name.contains(secret.as_str())))
3262 })
3263 })
3264 {
3265 return Err(
3266 ModelsFetchError::Catalog(CatalogRefreshError::InvalidResponse).into_interactive(),
3267 );
3268 }
3269 Ok(models)
3270 }
3271
3272 async fn models_document(
3273 &self,
3274 mode: ModelsRequestMode,
3275 ) -> Result<(String, tokio::time::Instant), ModelsFetchError> {
3276 let endpoint = reqwest::Url::parse(&api_url(&self.base_url, "models"))
3277 .map_err(|_| CatalogRefreshError::InvalidResponse)?;
3278 // https://platform.claude.com/docs/en/api/models/list specifies after_id.
3279 // Go is unpaginated. A Messages generation dialect or a custom identity
3280 // resembling a built-in provider does not establish this list contract.
3281 let cursor_query = (self.api_provider == ProviderKind::Anthropic).then_some("after_id");
3282 collect_models_document(
3283 endpoint,
3284 cursor_query,
3285 MODELS_FETCH_LIMITS,
3286 |url| async move {
3287 let build = || {
3288 self.models_http_client
3289 .get(url.clone())
3290 .timeout(NON_STREAMING_HTTP_TIMEOUT)
3291 };
3292 let response = match mode {
3293 // #6173: the provider's own words, minus this client's
3294 // secrets and this request's cursor. The endpoint is not
3295 // established here — it is whatever the user typed during
3296 // setup — so the body is guarded rather than trusted.
3297 ModelsRequestMode::Interactive => {
3298 let disclosure = ErrorBodyDisclosure::Guarded {
3299 request_secrets: request_query_secret_values(&url),
3300 };
3301 // The pinned 30s per-attempt total survives: the retry
3302 // loop's shared envelope is not allowed to overwrite a
3303 // caller's own budget.
3304 self.send_with_retry_total_error_body(
3305 NON_STREAMING_HTTP_TIMEOUT,
3306 build,
3307 &disclosure,
3308 )
3309 .await
3310 .map_err(ModelsFetchError::Interactive)?
3311 }
3312 ModelsRequestMode::Refresh => build()
3313 .send()
3314 .await
3315 .map_err(|_| CatalogRefreshError::Network)?,
3316 };
3317 if !response.status().is_success() {
3318 return Err(ModelsFetchError::Catalog(
3319 match response.status().as_u16() {
3320 401 => CatalogRefreshError::Unauthorized,
3321 403 => CatalogRefreshError::Forbidden,
3322 404 => CatalogRefreshError::NotFound,
3323 429 => CatalogRefreshError::RateLimited,
3324 _ => CatalogRefreshError::Network,
3325 },
3326 ));
3327 }
3328 Ok(response)
3329 },
3330 )
3331 .await
3332 }
3333
3334 /// The catalog provider id for this client (the `ProviderKind` slug, falling
3335 /// back to the `ProviderKind` slug for legacy variants without a kind). This
3336 /// is the id used as the cache scope and `CatalogOffering.provider`.
3337 fn catalog_provider_id(&self) -> String {
3338 self.admitted_identity.key.to_string()
3339 }
3340
3341 /// Whether this route speaks Baseten's `/models` dialect.
3342 ///
3343 /// Schema recognition is by endpoint, never by table name: whatever the
3344 /// user called the `[providers.<name>]` table, the wire shape is a fact
3345 /// about the host (#6289). Catalog ownership still uses the exact
3346 /// configured identity, so a renamed Baseten table keeps its own
3347 /// partition.
3348 #[cfg(test)]
3349 fn catalog_endpoint_is_baseten(&self) -> bool {
3350 catalog_endpoint_is_baseten(self.api_provider, &self.base_url)
3351 }
3352
3353 /// Fetch the provider's live `/models` listing as a secret-free
3354 /// [`ProviderCatalogDelta`] (#3385).
3355 ///
3356 /// Uses the same URL construction and auth client as [`Self::list_models`],
3357 /// but fetches pages without `send_with_retry` so a refresh
3358 /// failure stays typed and non-fatal — bundled / saved / static rows are
3359 /// untouched. The delta is scoped to the base-URL fingerprint and stamped
3360 /// with the fetch time; the API key authorizes the request but is **never**
3361 /// persisted into the delta or cache. Unknown live rows carry no canonical
3362 /// model, capabilities, or pricing, per the #3385 contract.
3363 pub async fn fetch_catalog_delta(&self) -> Result<ProviderCatalogDelta, CatalogRefreshError> {
3364 let (body, deadline) = self
3365 .models_document(ModelsRequestMode::Refresh)
3366 .await
3367 .map_err(ModelsFetchError::into_catalog)?;
3368
3369 let delta = catalog_delta_from_models_body(
3370 self.api_provider,
3371 self.catalog_provider_id(),
3372 &self.base_url,
3373 &body,
3374 &self.model_bound_secret_values,
3375 )?;
3376 if tokio::time::Instant::now() >= deadline {
3377 return Err(CatalogRefreshError::Network);
3378 }
3379 Ok(delta)
3380 }
3381
3382 /// Refresh `cache` for this client's provider + base URL, recording either a
3383 /// success or a typed failure (#3385). Returns the resulting status so the UI
3384 /// can surface a visible "fresh / failed(reason)" chip without inspecting the
3385 /// cache internals. A failed refresh preserves any previously cached rows.
3386 #[cfg(test)]
3387 pub async fn refresh_catalog_cache(
3388 &self,
3389 cache: &mut ProviderCatalogCache,
3390 ttl_secs: u64,
3391 ) -> CatalogStatus {
3392 match self.fetch_catalog_delta().await {
3393 Ok(delta) => {
3394 let provider = delta.provider.clone();
3395 let fingerprint = delta.base_url_fingerprint.clone();
3396 cache.record_success(delta, ttl_secs);
3397 publish_provider_lake_scope(cache, &provider, &fingerprint);
3398 CatalogStatus::Fresh
3399 }
3400 Err(reason) => {
3401 let provider = self.catalog_provider_id();
3402 let fingerprint = base_url_fingerprint(&self.base_url);
3403 cache.record_failure(&provider, &fingerprint, reason);
3404 publish_provider_lake_scope(cache, &provider, &fingerprint);
3405 CatalogStatus::Failed { reason }
3406 }
3407 }
3408 }
3409
3410 /// Best-effort background refresh of the active provider's own `/v1/models`
3411 /// catalog, replacing that provider's exact lake partition (#3385).
3412 ///
3413 /// Unlike `models_dev_live::spawn_background_refresh` (which fetches the
3414 /// cross-provider Models.dev catalog), this calls the provider's own
3415 /// `/v1/models` endpoint and merges the results into the existing live
3416 /// snapshot via `provider_catalog_live`, preserving other providers while
3417 /// allowing this provider's successful roster to retire removed ids.
3418 ///
3419 /// Activated for model-list authorities that are not satisfied by the
3420 /// cross-provider Models.dev snapshot: OpenRouter, named live gateways,
3421 /// and Baseten's account-scoped endpoint (no static snapshot can serve a
3422 /// per-credential roster). Custom OpenAI-compatible hosts are included
3423 /// too: a private relay is not in the Models.dev snapshot, so without a
3424 /// probe its `/model` picker stays empty even though the chat route
3425 /// already talks to the same endpoint (#6289 widened).
3426 /// The refresh is non-fatal: on failure, persisted prior rows and static
3427 /// seeds remain available with a typed failed receipt.
3428 pub fn spawn_active_provider_catalog_refresh(config: &Config) {
3429 // Unit tests use explicit fixture refreshes; never probe a developer's provider.
3430 #[cfg(test)]
3431 let _ = config;
3432 #[cfg(not(test))]
3433 {
3434 let Ok(identity) = config.active_provider_identity() else {
3435 return;
3436 };
3437 let provider = identity.provider;
3438 // Custom hosts include Baseten (its `/models` dialect is detected
3439 // at fetch time) and every other custom host. A private route is
3440 // the only place its roster exists, and a failed probe stays
3441 // non-fatal.
3442 if !crate::provider_catalog_live::provider_owns_live_catalog(provider) {
3443 return;
3444 }
3445
3446 // Invalidate older in-flight fetches before loading any reusable
3447 // scope. Baseten's ticket also clears account-scoped rows because the
3448 // same endpoint can expose a different workspace after a key change.
3449 let refresh_ticket = crate::provider_catalog_live::begin_refresh_for_identity(
3450 provider,
3451 identity.key.as_str(),
3452 &config.active_route_base_url(),
3453 );
3454
3455 // Publish the exact persisted scope immediately so opening `/model`
3456 // never waits on the network and another endpoint's rows cannot leak
3457 // into this route. Account-scoped Baseten rows deliberately do not
3458 // reload from disk until the current credential proves them again.
3459 crate::provider_catalog_live::maybe_load_persisted_cache_for_config(config);
3460
3461 let client = match CodewhaleClient::for_catalog_refresh(config) {
3462 Ok(client) => client,
3463 Err(err) => {
3464 tracing::debug!(
3465 target: "provider_catalog",
3466 error = %err,
3467 "skipping provider catalog refresh: client creation failed"
3468 );
3469 return;
3470 }
3471 };
3472
3473 tokio::spawn(async move {
3474 match client.fetch_catalog_delta().await {
3475 Ok(delta) => {
3476 let count = delta.offerings.len();
3477 if crate::provider_catalog_live::record_success_if_current(
3478 &refresh_ticket,
3479 delta,
3480 )
3481 .is_none()
3482 {
3483 tracing::debug!(
3484 target: "provider_catalog",
3485 "discarded provider catalog response superseded by a newer refresh"
3486 );
3487 return;
3488 }
3489 tracing::debug!(
3490 target: "provider_catalog",
3491 offering_count = count,
3492 "provider catalog refresh merged {count} offerings into provider lake"
3493 );
3494 }
3495 Err(err) => {
3496 if crate::provider_catalog_live::record_failure_if_current(
3497 &refresh_ticket,
3498 &client.catalog_provider_id(),
3499 &base_url_fingerprint(&client.base_url),
3500 err,
3501 )
3502 .is_none()
3503 {
3504 tracing::debug!(
3505 target: "provider_catalog",
3506 "discarded provider catalog failure superseded by a newer refresh"
3507 );
3508 return;
3509 }
3510 tracing::debug!(
3511 target: "provider_catalog",
3512 error = ?err,
3513 "provider catalog refresh failed; keeping existing rows"
3514 );
3515 }
3516 }
3517 });
3518 }
3519 }
3520
3521 /// Generate speech with Xiaomi MiMo TTS models.
3522 ///
3523 /// The spoken text is placed in an `assistant` message because Xiaomi
3524 /// MiMo's TTS chat-completions surface expects that shape. The optional
3525 /// `instruction` is a `user` message that controls style, voice design, or
3526 /// voice-clone performance and is not spoken verbatim.
3527 pub async fn synthesize_speech(
3528 &self,
3529 request: SpeechSynthesisRequest,
3530 ) -> Result<SpeechSynthesisResponse> {
3531 let _inference = self.acquire_remote_control_inference_permit().await;
3532 let _permit = self.acquire_provider_request_permit().await;
3533 if self.api_provider != crate::config::ProviderKind::XiaomiMimo {
3534 anyhow::bail!(
3535 "speech synthesis requires provider 'xiaomi-mimo' (current: {})",
3536 self.api_provider.as_str()
3537 );
3538 }
3539
3540 let model = request.model.trim().to_string();
3541 if model.is_empty() {
3542 anyhow::bail!("Speech model cannot be empty");
3543 }
3544 let text = request.text.trim().to_string();
3545 if text.is_empty() {
3546 anyhow::bail!("Speech text cannot be empty");
3547 }
3548
3549 let audio_format = normalize_audio_format(&request.audio_format);
3550 let model = self.wire_model_for_route(&model);
3551 let model_lower = model.to_ascii_lowercase();
3552 let instruction = request
3553 .instruction
3554 .as_deref()
3555 .map(str::trim)
3556 .filter(|value| !value.is_empty());
3557 let voice = request
3558 .voice
3559 .as_deref()
3560 .map(str::trim)
3561 .filter(|value| !value.is_empty())
3562 .map(str::to_string);
3563
3564 if model_lower.contains("voicedesign") && instruction.is_none() {
3565 anyhow::bail!(
3566 "Model '{model}' requires a voice design prompt. Pass --voice-prompt or --instruction."
3567 );
3568 }
3569 if model_lower.contains("voiceclone") && voice.is_none() {
3570 anyhow::bail!(
3571 "Model '{model}' requires cloned voice data. Pass --clone-voice <mp3|wav> or --voice <data-uri>."
3572 );
3573 }
3574
3575 let mut audio = json!({
3576 "format": audio_format.clone(),
3577 });
3578 if let Some(voice) = voice.as_deref() {
3579 audio["voice"] = json!(voice);
3580 }
3581
3582 let body = build_speech_synthesis_body(&model, &text, instruction, audio);
3583
3584 let url = api_url(&self.base_url, "chat/completions");
3585 let response = self.send_json_with_retry(&url, &body).await?;
3586 let status = response.status();
3587 if !status.is_success() {
3588 let raw_error_text = bounded_error_text(response, ERROR_BODY_MAX_BYTES).await;
3589 let error_text = sanitize_http_error_body(
3590 Some(self.api_provider.provider().display_name()),
3591 status.as_u16(),
3592 &raw_error_text,
3593 );
3594 anyhow::bail!("Speech synthesis failed: HTTP {status}: {error_text}");
3595 }
3596
3597 let response_text = response
3598 .text()
3599 .await
3600 .context("Failed to read speech synthesis response body")?;
3601 let payload: Value = serde_json::from_str(&response_text)
3602 .context("Failed to parse speech synthesis response JSON")?;
3603 let (audio_bytes, transcript) = parse_speech_audio_response(&payload)?;
3604
3605 Ok(SpeechSynthesisResponse {
3606 model,
3607 audio_format,
3608 audio_bytes,
3609 transcript,
3610 voice,
3611 })
3612 }
3613
3614 async fn wait_for_rate_limit(&self) {
3615 let maybe_delay = {
3616 let mut limiter = self.rate_limiter.lock().await;
3617 limiter.delay_until_available(1.0)
3618 };
3619 if let Some(delay) = maybe_delay {
3620 tokio::time::sleep(delay).await;
3621 }
3622 }
3623
3624 async fn mark_request_success(&self) {
3625 let mut health = self.connection_health.lock().await;
3626 if apply_request_success(&mut health, Instant::now()) {
3627 logging::info("Connection recovered");
3628 }
3629 }
3630
3631 async fn mark_request_failure(&self, reason: &str) {
3632 let mut health = self.connection_health.lock().await;
3633 apply_request_failure(&mut health, Instant::now());
3634 logging::warn(format!(
3635 "Connection degraded (failures={}): {}",
3636 health.consecutive_failures, reason
3637 ));
3638 }
3639
3640 async fn maybe_probe_recovery(&self) {
3641 let should_probe = {
3642 let mut health = self.connection_health.lock().await;
3643 mark_recovery_probe_if_due(&mut health, Instant::now())
3644 };
3645 if !should_probe {
3646 return;
3647 }
3648 if api_provider_skips_models_probe(self.api_provider) {
3649 self.mark_request_success().await;
3650 logging::info("Skipping /models recovery probe for provider without a models endpoint");
3651 return;
3652 }
3653 let health_url = api_url(&self.base_url, "models");
3654 let probe = self
3655 .models_http_client
3656 .get(health_url)
3657 .timeout(NON_STREAMING_HTTP_TIMEOUT)
3658 .send()
3659 .await;
3660 match probe {
3661 Ok(resp) if resp.status().is_success() => {
3662 // Consume the response body so the connection can be returned to the pool.
3663 let _ = bounded_error_text(resp, ERROR_BODY_MAX_BYTES).await;
3664 self.mark_request_success().await;
3665 logging::info("Recovery probe succeeded");
3666 }
3667 Ok(resp) => {
3668 self.mark_request_failure(&format!("probe status={}", resp.status()))
3669 .await;
3670 }
3671 Err(err) => {
3672 self.mark_request_failure(&format!("probe error={}", err.without_url()))
3673 .await;
3674 }
3675 }
3676 }
3677
3678 /// Apply `disclosure` to one provider error body.
3679 ///
3680 /// Redaction runs on the raw bytes, *before* `sanitize_http_error_body`
3681 /// truncates them, so a secret can never be split across the truncation
3682 /// boundary and survive as a fragment. The result is then checked against
3683 /// every value this client knows is secret; if one is still there the
3684 /// whole body is dropped, which is exactly the behaviour this path had
3685 /// before #6173. The fallback is the old contract, not a weaker one.
3686 ///
3687 /// Known limitation: this removes what the *client* knows is secret. A
3688 /// credential the user configured outside Codewhale — in a proxy, say —
3689 /// is not in that set and would pass through, the same limit
3690 /// `redact_model_bound_text` has.
3691 fn disclosed_http_error_body(
3692 &self,
3693 disclosure: &ErrorBodyDisclosure,
3694 status: u16,
3695 raw: &str,
3696 ) -> String {
3697 let provider = Some(self.api_provider.provider().display_name());
3698 let ErrorBodyDisclosure::Guarded { request_secrets } = disclosure else {
3699 return sanitize_http_error_body(provider, status, raw);
3700 };
3701 let mut redacted = raw.to_string();
3702 for secret in self
3703 .catalog_error_secret_values
3704 .iter()
3705 .chain(request_secrets.iter())
3706 {
3707 redacted = redacted.replace(secret.as_str(), codewhale_config::persistence::REDACTED);
3708 }
3709 let message = sanitize_http_error_body(provider, status, &redacted);
3710 let leaked = self
3711 .catalog_error_secret_values
3712 .iter()
3713 .chain(request_secrets.iter())
3714 .any(|secret| message.contains(secret.as_str()));
3715 if leaked { String::new() } else { message }
3716 }
3717
3718 pub(super) async fn send_with_retry<F>(&self, build: F) -> Result<reqwest::Response>
3719 where
3720 F: FnMut() -> reqwest::RequestBuilder,
3721 {
3722 self.send_with_retry_error_body(build, &ErrorBodyDisclosure::Full)
3723 .await
3724 }
3725
3726 /// Model-list errors can echo opaque cursors or credentials. Keep status
3727 /// and Retry-After classification either way; `disclosure` decides how
3728 /// much of the body reaches retry logs, state updates and the user.
3729 async fn send_with_retry_error_body<F>(
3730 &self,
3731 build: F,
3732 disclosure: &ErrorBodyDisclosure,
3733 ) -> Result<reqwest::Response>
3734 where
3735 F: FnMut() -> reqwest::RequestBuilder,
3736 {
3737 if self.isolated_request_state {
3738 return self.send_with_isolated_retry(build, disclosure).await;
3739 }
3740 self.send_retry_loop(build, Some(non_streaming_request_envelope()), disclosure)
3741 .await
3742 }
3743
3744 /// [`Self::send_with_retry_error_body`] with a caller-pinned per-attempt
3745 /// total (connect through body end). `list_models` pins its own 30s: the
3746 /// plain variant would otherwise stretch that pinned budget out to the
3747 /// shared envelope, because `.timeout()` on the builder is a pure
3748 /// overwrite.
3749 async fn send_with_retry_total_error_body<F>(
3750 &self,
3751 total: Duration,
3752 build: F,
3753 disclosure: &ErrorBodyDisclosure,
3754 ) -> Result<reqwest::Response>
3755 where
3756 F: FnMut() -> reqwest::RequestBuilder,
3757 {
3758 if self.isolated_request_state {
3759 return self.send_with_isolated_retry(build, disclosure).await;
3760 }
3761 self.send_retry_loop(build, Some(total), disclosure).await
3762 }
3763
3764 /// The streaming-open twin of [`Self::send_with_retry`]: the same retry
3765 /// and rate-limit handling with no total deadline anywhere. reqwest's
3766 /// per-request timeout wraps the response *body*, so a total set on the
3767 /// open would ride along inside the returned body and hard-cut a live
3768 /// stream mid-generation. Stream opens stay bounded by the caller's
3769 /// `stream_open_timeout` around the open and per-chunk idle checks on
3770 /// the returned body instead.
3771 pub(super) async fn send_stream_open_with_retry<F>(&self, build: F) -> Result<reqwest::Response>
3772 where
3773 F: FnMut() -> reqwest::RequestBuilder,
3774 {
3775 if self.isolated_request_state {
3776 return self
3777 .send_with_isolated_retry(build, &ErrorBodyDisclosure::Full)
3778 .await;
3779 }
3780 self.send_retry_loop(build, None, &ErrorBodyDisclosure::Full)
3781 .await
3782 }
3783
3784 async fn send_retry_loop<F>(
3785 &self,
3786 mut build: F,
3787 attempt_total: Option<Duration>,
3788 disclosure: &ErrorBodyDisclosure,
3789 ) -> Result<reqwest::Response>
3790 where
3791 F: FnMut() -> reqwest::RequestBuilder,
3792 {
3793 let retry_cfg: LlmRetryConfig = self.retry.clone().into();
3794 let pause_scope = self.rate_limit_scope();
3795 let callback_scope = pause_scope.clone();
3796 // Two bounded layers around a non-streaming completion
3797 // (`attempt_total` = `Some`): the per-attempt request total (connect
3798 // through body end) and an envelope around the whole retry loop (all
3799 // attempts + backoff + honored Retry-After). Streaming opens
3800 // (`None`) set no deadline at all: any total here would be
3801 // inherited by the returned body and truncate the stream, so those
3802 // calls keep only the caller's own open budget.
3803 let retry_future = with_retry(
3804 &retry_cfg,
3805 || {
3806 // Per-attempt total: unlike the loop envelope below,
3807 // reqwest's per-request timeout also covers the response
3808 // body, so a slow-drip body cannot outlive the budget.
3809 let request = match attempt_total {
3810 Some(total) => build().timeout(total),
3811 None => build(),
3812 };
3813 let pause_scope = pause_scope.as_str();
3814 async move {
3815 // Sleep in bounded slices rather than the full remaining
3816 // window: the pause is shared by every request to this
3817 // route, so a concurrent `clear_rate_limit()` (or a
3818 // shortened deadline) must release requests that are
3819 // already waiting instead of stranding them for the whole
3820 // original window.
3821 while let Some(delay) = crate::retry_status::rate_limit_remaining(pause_scope) {
3822 tokio::time::sleep(delay.min(RATE_LIMIT_PAUSE_RECHECK_INTERVAL)).await;
3823 }
3824 self.wait_for_rate_limit().await;
3825 let response = request
3826 .send()
3827 .await
3828 .map_err(|err| LlmError::from_reqwest(&err.without_url()))?;
3829 let status = response.status();
3830 if status.is_success() {
3831 return Ok(response);
3832 }
3833 let retry_after = extract_retry_after(response.headers());
3834 let raw = bounded_error_text(response, ERROR_BODY_MAX_BYTES).await;
3835 let body = self.disclosed_http_error_body(disclosure, status.as_u16(), &raw);
3836 Err(self.http_error_with_route_context(status.as_u16(), &body, retry_after))
3837 }
3838 },
3839 Some(Box::new(move |err, attempt, delay| {
3840 let (reason_label, human_reason) = retry_reason_label_and_human(err);
3841 logging::warn(format!(
3842 "HTTP retry reason={} attempt={} delay={:.2}s",
3843 reason_label,
3844 attempt + 1,
3845 delay.as_secs_f64(),
3846 ));
3847 if matches!(err, LlmError::RateLimited { .. }) {
3848 crate::retry_status::note_rate_limit(&callback_scope, delay);
3849 }
3850 crate::retry_status::start(attempt + 1, delay, human_reason);
3851 })),
3852 );
3853 let request_result = if let Some(total) = attempt_total {
3854 // The loop envelope must dominate the per-attempt total it
3855 // wraps: a caller-pinned budget (list_models' 30s) may exceed
3856 // the shared envelope, and the envelope must never strangle
3857 // its own attempts.
3858 let loop_envelope = total.max(non_streaming_request_envelope());
3859 match tokio::time::timeout(loop_envelope, retry_future).await {
3860 Ok(result) => result,
3861 Err(_elapsed) => {
3862 let last = LlmError::Timeout(loop_envelope);
3863 logging::warn(format!(
3864 "non-streaming request envelope exceeded ({loop_envelope:?}); retry loop aborted"
3865 ));
3866 crate::retry_status::failed(last.to_string());
3867 self.mark_request_failure("non-streaming request envelope exceeded")
3868 .await;
3869 self.maybe_probe_recovery().await;
3870 return Err(anyhow::Error::new(last));
3871 }
3872 }
3873 } else {
3874 retry_future.await
3875 };
3876
3877 match request_result {
3878 Ok(response) => {
3879 crate::retry_status::succeeded();
3880 self.mark_request_success().await;
3881 Ok(response)
3882 }
3883 Err(err) => {
3884 if let LlmError::RateLimited { retry_after, .. } = &err.last_error {
3885 crate::retry_status::note_rate_limit(
3886 &pause_scope,
3887 retry_after
3888 .unwrap_or_else(|| retry_cfg.delay_for_attempt(retry_cfg.max_retries)),
3889 );
3890 }
3891 let last = err.last_error.to_string();
3892 if err.attempts > 1 {
3893 crate::retry_status::failed(last.clone());
3894 } else {
3895 crate::retry_status::clear();
3896 }
3897 self.mark_request_failure(&last).await;
3898 self.maybe_probe_recovery().await;
3899 // Keep the structured `LlmError` downcastable so failure
3900 // surfaces can classify auth/rate-limit/invalid-request
3901 // instead of reporting an opaque string (#3884).
3902 Err(anyhow::Error::new(err.last_error))
3903 }
3904 }
3905 }
3906
3907 /// Key for this route's shared `Retry-After` pause: the configured route
3908 /// identity plus the host it reaches. A 429 from one provider pauses only
3909 /// requests that would hit the same limit, never another provider, a
3910 /// local runtime, or a sub-agent on a different route.
3911 pub(crate) fn rate_limit_scope(&self) -> String {
3912 let route = self.admitted_identity.key.as_str();
3913 let host = crate::llm_client::base_url_authority(&self.base_url)
3914 .unwrap_or_else(|| redact_url_for_display(&self.base_url));
3915 format!("{route}@{host}")
3916 }
3917
3918 /// The same bounded transport retry policy without process-global retry
3919 /// banners, provider-wide pause cells, or shared connection-health writes.
3920 /// Used only by the Auto classifier during read-only request inspection.
3921 async fn send_with_isolated_retry<F>(
3922 &self,
3923 mut build: F,
3924 disclosure: &ErrorBodyDisclosure,
3925 ) -> Result<reqwest::Response>
3926 where
3927 F: FnMut() -> reqwest::RequestBuilder,
3928 {
3929 let retry_cfg: LlmRetryConfig = self.retry.clone().into();
3930 let request_result = crate::llm_client::observe_request_retries(
3931 None,
3932 with_retry(
3933 &retry_cfg,
3934 || {
3935 let request = build();
3936 async move {
3937 self.wait_for_rate_limit().await;
3938 let response = request
3939 .send()
3940 .await
3941 .map_err(|err| LlmError::from_reqwest(&err.without_url()))?;
3942 let status = response.status();
3943 if status.is_success() {
3944 return Ok(response);
3945 }
3946 let retry_after = extract_retry_after(response.headers());
3947 let raw = bounded_error_text(response, ERROR_BODY_MAX_BYTES).await;
3948 let body =
3949 self.disclosed_http_error_body(disclosure, status.as_u16(), &raw);
3950 Err(self.http_error_with_route_context(status.as_u16(), &body, retry_after))
3951 }
3952 },
3953 Some(Box::new(|err, attempt, delay| {
3954 let (reason_label, _) = retry_reason_label_and_human(err);
3955 logging::warn(format!(
3956 "Isolated HTTP retry reason={} attempt={} delay={:.2}s",
3957 reason_label,
3958 attempt + 1,
3959 delay.as_secs_f64(),
3960 ));
3961 })),
3962 ),
3963 )
3964 .await;
3965
3966 request_result.map_err(|err| anyhow::Error::new(err.last_error))
3967 }
3968
3969 pub(super) async fn send_json_with_retry(
3970 &self,
3971 url: &str,
3972 body: &serde_json::Value,
3973 ) -> Result<reqwest::Response> {
3974 let request_body =
3975 serde_json::to_vec(body).context("Failed to serialize JSON request body")?;
3976 self.send_with_retry(|| {
3977 self.http_client
3978 .post(url)
3979 .header(CONTENT_TYPE, "application/json")
3980 .body(request_body.clone())
3981 })
3982 .await
3983 }
3984
3985 /// JSON POST through the streaming-open retry path: no total deadline,
3986 /// because the response body outlives the open (see
3987 /// [`Self::send_stream_open_with_retry`]).
3988 pub(super) async fn open_stream_json_with_retry(
3989 &self,
3990 url: &str,
3991 body: &serde_json::Value,
3992 ) -> Result<reqwest::Response> {
3993 let request_body =
3994 serde_json::to_vec(body).context("Failed to serialize JSON request body")?;
3995 self.send_stream_open_with_retry(|| {
3996 self.http_client
3997 .post(url)
3998 .header(CONTENT_TYPE, "application/json")
3999 .body(request_body.clone())
4000 })
4001 .await
4002 }
4003 }
4004
4005 /// Record that a request was routed to `provider` and came back with `status`.
4006 ///
4007 /// Called at every provider response site, **before** the error is built: an
4008 /// `LlmError` carries the raw provider body verbatim, so the status class has
4009 /// to be taken from the response itself.
4010 ///
4011 /// The provider is recorded as a `ProviderKind` by value. Every accessor that
4012 /// looks like the natural seam here — the persistence identity, the stream
4013 /// meta's `provider_id`, the planned route's effective label — returns the
4014 /// customer's own `[providers.<name>]` table key when the route is custom.
4015 /// `ProviderKind::Custom` yields the literal `"custom"` and nothing else, and
4016 /// no model id is sent for any provider.
4017 pub(crate) fn record_provider_response(provider: crate::config::ProviderKind, status: u16) {
4018 let counters = codewhale_telemetry::session_counters();
4019 counters.record_provider(provider);
4020 if let Some(counter) = codewhale_telemetry::counters::http_status_counter(status) {
4021 counters.bump_error(counter);
4022 }
4023 }
4024
4025 /// Translate the structured `LlmError` into both a categorical label
4026 /// (for structured logs / metrics) and a short human reason string
4027 /// (for the retry banner). Returning both from one match avoids the
4028 /// double-classification we had before.
4029 fn retry_reason_label_and_human(err: &LlmError) -> (&'static str, String) {
4030 // The variant, never the payload. Every `LlmError` variant carries the raw
4031 // provider HTTP body verbatim, and a 400 from a content filter routinely
4032 // echoes the prompt.
4033 if matches!(err, LlmError::NetworkError(_) | LlmError::Timeout(_)) {
4034 codewhale_telemetry::session_counters()
4035 .bump_error(codewhale_telemetry::ErrorCounter::NetworkError);
4036 }
4037 match err {
4038 LlmError::RateLimited { retry_after, .. } => {
4039 let human = if let Some(after) = retry_after {
4040 format!("rate limited (Retry-After {}s)", after.as_secs())
4041 } else {
4042 "rate limited".to_string()
4043 };
4044 ("rate_limited", human)
4045 }
4046 LlmError::ServerError { status, .. } => ("server_error", format!("upstream {status}")),
4047 LlmError::NetworkError(_) => ("network_error", "network error".to_string()),
4048 LlmError::Timeout(_) => ("timeout", "timeout".to_string()),
4049 _ => ("other", "other".to_string()),
4050 }
4051 }
4052
4053 impl CodewhaleClient {
4054 /// Execute a non-streaming request without consulting or updating the
4055 /// process-global response cache.
4056 ///
4057 /// Request previews use this only for Auto's auxiliary router classifier:
4058 /// the classifier may call its configured provider, but an inspection must
4059 /// not perturb later production routing through shared cache state.
4060 pub(crate) async fn create_message_without_response_cache(
4061 &self,
4062 request: MessageRequest,
4063 ) -> Result<MessageResponse> {
4064 let mut isolated = self.clone();
4065 isolated.isolated_request_state = true;
4066 // The ordinary clone shares its provider token bucket so concurrent
4067 // production calls observe one rate budget. Request inspection is an
4068 // auxiliary classifier call, however: it must neither consume nor
4069 // inherit that mutable foreground state.
4070 isolated.rate_limiter = Arc::new(AsyncMutex::new(TokenBucket::from_env()));
4071 isolated
4072 .create_message_with_cache_policy(request, false)
4073 .await
4074 }
4075
4076 async fn create_message_with_cache_policy(
4077 &self,
4078 request: MessageRequest,
4079 allow_response_cache: bool,
4080 ) -> Result<MessageResponse> {
4081 let _inference = self.acquire_remote_control_inference_permit().await;
4082 let _permit = self.acquire_provider_request_permit().await;
4083 let cacheable =
4084 allow_response_cache && crate::llm_response_cache::request_is_cacheable(&request);
4085 let prepared = self.prepare_outbound_request(request, false)?;
4086 match prepared.dialect {
4087 WireDialect::OpenAiResponses => self.handle_responses_message(&prepared).await,
4088 WireDialect::AnthropicMessages => self.handle_anthropic_message(&prepared).await,
4089 WireDialect::ChatCompletions => self.create_message_chat(&prepared, cacheable).await,
4090 }
4091 }
4092 }
4093
4094 impl LlmClient for CodewhaleClient {
4095 fn provider_name(&self) -> &'static str {
4096 self.api_provider.as_str()
4097 }
4098
4099 fn model(&self) -> &str {
4100 &self.default_model
4101 }
4102
4103 fn billing_base_url(&self) -> Option<&str> {
4104 Some(&self.base_url)
4105 }
4106
4107 fn route_limits(&self) -> Option<RouteLimits> {
4108 CodewhaleClient::route_limits(self)
4109 }
4110
4111 fn effective_max_output_tokens(&self, requested_model: &str) -> u32 {
4112 CodewhaleClient::effective_max_output_tokens(self, requested_model)
4113 }
4114
4115 fn effective_route_envelope(
4116 &self,
4117 requested_model: &str,
4118 dispatched_at: chrono::DateTime<chrono::Utc>,
4119 ) -> crate::cost_status::EffectiveRouteEnvelope {
4120 CodewhaleClient::effective_route_envelope(self, requested_model, dispatched_at)
4121 }
4122
4123 async fn health_check(&self) -> Result<bool> {
4124 if api_provider_skips_models_probe(self.api_provider) {
4125 self.mark_request_success().await;
4126 return Ok(true);
4127 }
4128 let health_url = api_url(&self.base_url, "models");
4129 self.wait_for_rate_limit().await;
4130 let response = self
4131 .models_http_client
4132 .get(health_url)
4133 .timeout(NON_STREAMING_HTTP_TIMEOUT)
4134 .send()
4135 .await;
4136 match response {
4137 Ok(resp) if resp.status().is_success() => {
4138 // Consume the response body so the connection can be returned to the pool.
4139 let _ = bounded_error_text(resp, ERROR_BODY_MAX_BYTES).await;
4140 self.mark_request_success().await;
4141 Ok(true)
4142 }
4143 Ok(resp) => {
4144 self.mark_request_failure(&format!("health status={}", resp.status()))
4145 .await;
4146 Ok(false)
4147 }
4148 Err(err) => {
4149 self.mark_request_failure(&format!("health error={}", err.without_url()))
4150 .await;
4151 Ok(false)
4152 }
4153 }
4154 }
4155
4156 async fn create_message(&self, request: MessageRequest) -> Result<MessageResponse> {
4157 self.create_message_with_cache_policy(request, true).await
4158 }
4159
4160 async fn create_message_uncached(&self, request: MessageRequest) -> Result<MessageResponse> {
4161 // Keep shared provider permits and rate limits. Only the response
4162 // cache is bypassed; a guardian is still real, metered inference.
4163 self.create_message_with_cache_policy(request, false).await
4164 }
4165
4166 async fn create_message_stream(
4167 &self,
4168 request: MessageRequest,
4169 ) -> Result<crate::llm_client::StreamEventBox> {
4170 let inference = self.acquire_remote_control_inference_permit().await;
4171 let permit = self.acquire_provider_request_permit().await;
4172 let prepared = self.prepare_outbound_request(request, true)?;
4173 let projection_warning = (!prepared.omitted_tool_names.is_empty()).then(|| {
4174 let omitted_tool_count = prepared.omitted_tool_names.len();
4175 (
4176 prepared.endpoint.provider_display.clone(),
4177 crate::core::events::bounded_tool_projection_warning_names(
4178 &prepared.omitted_tool_names,
4179 ),
4180 omitted_tool_count,
4181 )
4182 });
4183 let stream = match prepared.dialect {
4184 WireDialect::OpenAiResponses => self.handle_responses_stream(&prepared).await?,
4185 WireDialect::AnthropicMessages => self.handle_anthropic_stream(&prepared).await?,
4186 WireDialect::ChatCompletions => self.handle_chat_completion_stream(prepared).await?,
4187 };
4188 let stream = match projection_warning {
4189 Some((provider, omitted_tool_names, omitted_tool_count)) => {
4190 Self::prepend_tool_projection_warning(
4191 stream,
4192 provider,
4193 omitted_tool_names,
4194 omitted_tool_count,
4195 )
4196 }
4197 None => stream,
4198 };
4199 let stream = Self::hold_provider_request_permit_for_stream(stream, permit);
4200 Ok(Self::hold_remote_control_inference_permit_for_stream(
4201 stream, inference,
4202 ))
4203 }
4204 }
4205
4206 #[derive(Debug, Deserialize)]
4207 struct ModelsListResponse {
4208 data: Vec<ModelListItem>,
4209 }
4210
4211 /// The list envelope is validated as a whole; each row is decoded on its own
4212 /// so one malformed row cannot fail the entire roster (#6690). Rows stay raw
4213 /// text rather than `serde_json::Value`: a `Value` map keeps the last of a
4214 /// duplicated key, which would silently accept an ambiguous row (two `id`s,
4215 /// two `pricing.prompt`s) instead of skipping it as malformed.
4216 #[derive(Debug, Deserialize)]
4217 struct OpenRouterModelsResponse {
4218 data: Vec<Box<serde_json::value::RawValue>>,
4219 }
4220
4221 #[derive(Debug, Deserialize)]
4222 struct ModelListItem {
4223 id: String,
4224 #[serde(default)]
4225 owned_by: Option<String>,
4226 #[serde(default)]
4227 created: Option<u64>,
4228 }
4229
4230 /// OpenRouter `/models` response item with full capability metadata (#3385).
4231 #[derive(Debug, Deserialize)]
4232 struct OpenRouterModelItem {
4233 id: String,
4234 // Captured from OpenRouter for future display/deprecation surfaces. The
4235 // current CatalogOffering shape has no honest fields for these yet.
4236 #[serde(default)]
4237 #[expect(dead_code)]
4238 name: Option<String>,
4239 #[serde(default)]
4240 #[expect(dead_code)]
4241 created: Option<u64>,
4242 #[serde(default)]
4243 context_length: Option<u32>,
4244 #[serde(default)]
4245 pricing: Option<OpenRouterPricing>,
4246 #[serde(default)]
4247 top_provider: Option<OpenRouterTopProvider>,
4248 #[serde(default)]
4249 supported_parameters: Option<Vec<String>>,
4250 #[serde(default)]
4251 architecture: Option<OpenRouterArchitecture>,
4252 #[serde(default)]
4253 #[expect(dead_code)]
4254 expiration_date: Option<String>,
4255 }
4256
4257 #[derive(Debug, Deserialize)]
4258 struct OpenRouterPricing {
4259 #[serde(default)]
4260 prompt: Option<String>,
4261 #[serde(default)]
4262 completion: Option<String>,
4263 #[serde(default)]
4264 input_cache_read: Option<String>,
4265 /// Per-token cache-write (cache-creation) price. OpenRouter publishes this
4266 /// for the upstreams that charge a write premium (Anthropic, Qwen, …);
4267 /// dropping it undercounted every cache-creation turn on those routes.
4268 #[serde(default)]
4269 input_cache_write: Option<String>,
4270 }
4271
4272 #[derive(Debug, Deserialize)]
4273 struct OpenRouterTopProvider {
4274 #[serde(default)]
4275 context_length: Option<u32>,
4276 #[serde(default)]
4277 max_completion_tokens: Option<u32>,
4278 }
4279
4280 #[derive(Debug, Deserialize)]
4281 struct OpenRouterArchitecture {
4282 #[serde(default)]
4283 modality: Option<String>,
4284 #[serde(default)]
4285 input_modalities: Option<Vec<String>>,
4286 #[serde(default)]
4287 output_modalities: Option<Vec<String>>,
4288 }
4289
4290 /// Baseten Model APIs `/v1/models` item.
4291 ///
4292 /// Baseten publishes OpenAI-style ids and per-token prices, while current
4293 /// serving limits and feature fields are additive. Numeric fields accept JSON
4294 /// numbers or numeric strings because both appear in provider catalogs in the
4295 /// wild; malformed or negative known fields reject the refresh so the durable
4296 /// last-known-good snapshot remains authoritative.
4297 #[derive(Debug, Deserialize)]
4298 struct BasetenModelsResponse {
4299 data: Vec<BasetenModelItem>,
4300 }
4301
4302 #[derive(Debug, Deserialize)]
4303 struct BasetenModelItem {
4304 id: String,
4305 #[serde(default)]
4306 context_length: Option<CatalogNumber>,
4307 #[serde(default)]
4308 context_window: Option<CatalogNumber>,
4309 #[serde(default)]
4310 max_output_tokens: Option<CatalogNumber>,
4311 #[serde(default)]
4312 max_completion_tokens: Option<CatalogNumber>,
4313 #[serde(default)]
4314 limits: Option<BasetenLimits>,
4315 #[serde(default)]
4316 top_provider: Option<BasetenTopProvider>,
4317 #[serde(default)]
4318 pricing: Option<BasetenPricing>,
4319 #[serde(default)]
4320 supported_parameters: Option<Vec<String>>,
4321 #[serde(default)]
4322 supported_features: Option<Vec<String>>,
4323 #[serde(default)]
4324 features: Option<Vec<String>>,
4325 #[serde(default)]
4326 architecture: Option<BasetenArchitecture>,
4327 #[serde(default)]
4328 input_modalities: Option<Vec<String>>,
4329 #[serde(default)]
4330 output_modalities: Option<Vec<String>>,
4331 #[serde(default)]
4332 reasoning: Option<bool>,
4333 #[serde(default)]
4334 supports_reasoning: Option<bool>,
4335 #[serde(default)]
4336 supports_tools: Option<bool>,
4337 #[serde(default)]
4338 supports_structured_output: Option<bool>,
4339 #[serde(default)]
4340 reasoning_options: Vec<serde_json::Value>,
4341 }
4342
4343 #[derive(Debug, Deserialize)]
4344 #[serde(untagged)]
4345 enum CatalogNumber {
4346 Number(serde_json::Number),
4347 Text(String),
4348 }
4349
4350 #[derive(Debug, Default, Deserialize)]
4351 struct BasetenLimits {
4352 #[serde(default)]
4353 context: Option<CatalogNumber>,
4354 #[serde(default)]
4355 output: Option<CatalogNumber>,
4356 }
4357
4358 #[derive(Debug, Default, Deserialize)]
4359 struct BasetenTopProvider {
4360 #[serde(default)]
4361 context_length: Option<CatalogNumber>,
4362 #[serde(default)]
4363 max_completion_tokens: Option<CatalogNumber>,
4364 }
4365
4366 #[derive(Debug, Default, Deserialize)]
4367 struct BasetenPricing {
4368 #[serde(default)]
4369 prompt: Option<CatalogNumber>,
4370 #[serde(default)]
4371 completion: Option<CatalogNumber>,
4372 #[serde(default)]
4373 input_cache_read: Option<CatalogNumber>,
4374 #[serde(default)]
4375 input_cache_write: Option<CatalogNumber>,
4376 }
4377
4378 #[derive(Debug, Default, Deserialize)]
4379 struct BasetenArchitecture {
4380 #[serde(default)]
4381 input_modalities: Option<Vec<String>>,
4382 #[serde(default)]
4383 output_modalities: Option<Vec<String>>,
4384 }
4385
4386 fn parse_models_response_for_provider(
4387 payload: &str,
4388 provider: ProviderKind,
4389 ) -> Result<Vec<AvailableModel>> {
4390 if provider != ProviderKind::OpenaiCodex {
4391 return parse_models_response(payload);
4392 }
4393 #[derive(Deserialize)]
4394 struct Roster {
4395 models: Vec<Row>,
4396 }
4397 #[derive(Deserialize)]
4398 struct Row {
4399 slug: String,
4400 display_name: String,
4401 visibility: String,
4402 }
4403 let roster: Roster = serde_json::from_str(payload).context("Invalid ChatGPT model roster")?;
4404 anyhow::ensure!(
4405 roster.models.len() <= PROVIDER_CATALOG_MAX_ROWS,
4406 "ChatGPT model roster exceeds row limit"
4407 );
4408 let mut seen = std::collections::HashSet::new();
4409 let mut models = Vec::new();
4410 for row in roster.models {
4411 if row.visibility != "list" {
4412 continue;
4413 }
4414 anyhow::ensure!(
4415 crate::provider_lake::valid_catalog_model_id(&row.slug)
4416 && !row.display_name.trim().is_empty()
4417 && row.display_name.len() <= 256
4418 && !row.display_name.chars().any(char::is_control),
4419 "Invalid ChatGPT model row"
4420 );
4421 if seen.insert(row.slug.clone()) {
4422 models.push(AvailableModel {
4423 id: row.slug,
4424 display_name: Some(row.display_name),
4425 owned_by: None,
4426 created: None,
4427 });
4428 }
4429 }
4430 Ok(models)
4431 }
4432
4433 pub(crate) fn parse_models_response(payload: &str) -> Result<Vec<AvailableModel>> {
4434 let parsed: ModelsListResponse =
4435 serde_json::from_str(payload).context("Failed to parse model list JSON")?;
4436
4437 let mut models = parsed
4438 .data
4439 .into_iter()
4440 .map(|item| AvailableModel {
4441 id: item.id,
4442 owned_by: item.owned_by,
4443 created: item.created,
4444 display_name: None,
4445 })
4446 .collect::<Vec<_>>();
4447 models.sort_by(|a, b| a.id.cmp(&b.id));
4448 models.dedup_by(|a, b| a.id == b.id);
4449 Ok(models)
4450 }
4451
4452 /// Keep both live-model consumers on each provider's documented coding roster.
4453 /// Unknown models remain hidden until their wire contract is known.
4454 fn apply_provider_model_cutline(
4455 provider: ProviderKind,
4456 models: Vec<AvailableModel>,
4457 ) -> Vec<AvailableModel> {
4458 if provider == ProviderKind::Stepfun {
4459 // The shared /models endpoint also lists speech and image generators.
4460 // Only expose routes whose text/tool contract is in our catalog.
4461 let offerings = codewhale_config::catalog::bundled_catalog_offerings();
4462 return models
4463 .into_iter()
4464 .filter(|model| {
4465 offerings.iter().any(|row| {
4466 row.provider == "stepfun"
4467 && row.wire_model_id == model.id
4468 && row.tool_call == Some(true)
4469 })
4470 })
4471 .collect();
4472 }
4473 if provider != ProviderKind::OpencodeGo {
4474 return models;
4475 }
4476
4477 let mut models: Vec<_> = models
4478 .into_iter()
4479 .filter_map(|mut model| {
4480 let canonical = crate::config::opencode_go_model_id(&model.id)?;
4481 model.id = canonical.to_string();
4482 Some(model)
4483 })
4484 .collect();
4485 models.sort_by(|left, right| left.id.cmp(&right.id));
4486 models.dedup_by(|left, right| left.id == right.id);
4487 models
4488 }
4489
4490 /// Convert a named gateway's `/models` response into truthful provider-scoped
4491 /// catalog rows. Matching model ids on other providers prove no capabilities,
4492 /// limits, or prices; only an explicit same-provider bundled row may enrich a
4493 /// live offering.
4494 /// Parse the Codewhale API's authenticated `GET {base}/models` listing.
4495 ///
4496 /// The account control plane returns the OpenAI list shape, but every row
4497 /// carries a `codewhale` block naming the wire protocol Codewhale will use for
4498 /// that model (`chat-completions` → `{base}/chat/completions`,
4499 /// `anthropic-messages` → `{base}/messages`, `responses` →
4500 /// `{base}/responses`). That block is the whole reason this route is
4501 /// model-aware, so the protocol is read from the response rather
4502 /// than inferred — the namespace inference in
4503 /// [`codewhale_config::route::codewhale_endpoint_key_for_model`] is only the
4504 /// offline fallback for a row this listing did not describe.
4505 ///
4506 /// The listing is untrusted remote data: unusable rows are dropped rather than
4507 /// failing the refresh, and nothing here claims capabilities, limits, or
4508 /// pricing the account service did not state.
4509 fn codewhale_catalog_offerings_from_body(
4510 body: &str,
4511 provider: &str,
4512 fingerprint: &str,
4513 fetched_at: u64,
4514 ) -> Result<Vec<CatalogOffering>, CatalogRefreshError> {
4515 #[derive(serde::Deserialize)]
4516 struct Listing {
4517 #[serde(default)]
4518 data: Vec<Row>,
4519 }
4520 #[derive(serde::Deserialize)]
4521 struct Row {
4522 #[serde(default)]
4523 id: String,
4524 #[serde(default)]
4525 codewhale: Option<RowMeta>,
4526 }
4527 #[derive(serde::Deserialize)]
4528 struct RowMeta {
4529 #[serde(default)]
4530 protocol: Option<String>,
4531 #[serde(default)]
4532 default: bool,
4533 }
4534
4535 let listing: Listing =
4536 serde_json::from_str(body).map_err(|_| CatalogRefreshError::InvalidResponse)?;
4537 let mut offerings: Vec<CatalogOffering> = Vec::with_capacity(listing.data.len());
4538 let mut seen: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
4539 for row in listing.data {
4540 let id = row.id.trim().to_string();
4541 if id.is_empty() || !seen.insert(id.clone()) {
4542 continue;
4543 }
4544 // An unstated or unknown protocol falls back to the namespace rule
4545 // rather than being dropped: the account already proved it serves this
4546 // model by listing it, and a protocol label this build predates must
4547 // not make the row unreachable.
4548 let endpoint_key = match row
4549 .codewhale
4550 .as_ref()
4551 .and_then(|meta| meta.protocol.as_deref())
4552 .map(str::trim)
4553 {
4554 Some("anthropic-messages") => "messages",
4555 Some("chat-completions") => "chat",
4556 Some("responses") => "responses",
4557 _ => codewhale_config::route::codewhale_endpoint_key_for_model(&id),
4558 };
4559 offerings.push(CatalogOffering {
4560 cost_source: None,
4561 modalities_source: None,
4562 provider: provider.to_string(),
4563 wire_model_id: id,
4564 canonical_model: None,
4565 endpoint_key: endpoint_key.to_string(),
4566 default_for_provider: row.codewhale.is_some_and(|meta| meta.default),
4567 family: None,
4568 limit: None,
4569 cost: None,
4570 modalities: None,
4571 attachment: None,
4572 reasoning: None,
4573 tool_call: None,
4574 structured_output: None,
4575 reasoning_options: Vec::new(),
4576 source: CatalogSource::Live {
4577 base_url_fingerprint: fingerprint.to_string(),
4578 fetched_at,
4579 },
4580 });
4581 }
4582 if offerings.is_empty() {
4583 return Err(CatalogRefreshError::EmptyList);
4584 }
4585 Ok(offerings)
4586 }
4587
4588 /// Whether a route speaks Baseten's `/models` dialect: recognized by
4589 /// endpoint, never by table name (#6289).
4590 fn catalog_endpoint_is_baseten(api_provider: ProviderKind, base_url: &str) -> bool {
4591 api_provider == ProviderKind::Custom && codewhale_config::catalog::endpoint_is_baseten(base_url)
4592 }
4593
4594 /// Project one `/models` document onto a secret-free, endpoint-scoped roster
4595 /// delta (#3385). Shared by [`CodewhaleClient::fetch_catalog_delta`] and the
4596 /// guided-setup key probe ([`verify_provider_api_key`]), so the listing a key
4597 /// check already downloaded is read by exactly the rules a refresh uses.
4598 fn catalog_delta_from_models_body(
4599 api_provider: ProviderKind,
4600 provider: String,
4601 base_url: &str,
4602 body: &str,
4603 model_bound_secret_values: &[String],
4604 ) -> Result<ProviderCatalogDelta, CatalogRefreshError> {
4605 let fingerprint = base_url_fingerprint(base_url);
4606 let fetched_at = now_unix();
4607
4608 // OpenRouter returns extended capability metadata in its /models
4609 // response (#3385). Capture limits, pricing, reasoning, and modalities
4610 // from the live API instead of leaving them unknown.
4611 let offerings: Vec<CatalogOffering> = if api_provider == ProviderKind::Openrouter {
4612 let or_models = parse_openrouter_models_response(body)?;
4613 if or_models.is_empty() {
4614 return Err(CatalogRefreshError::EmptyList);
4615 }
4616 // A row with an unreadable or implausible price is skipped and
4617 // counted, not allowed to fail every other row (#6690).
4618 let offerings: Vec<_> = or_models
4619 .iter()
4620 .filter_map(|item| {
4621 openrouter_to_catalog_offering(item, &provider, &fingerprint, fetched_at).ok()
4622 })
4623 .collect();
4624 let skipped = or_models.len() - offerings.len();
4625 if skipped > 0 {
4626 tracing::warn!(
4627 skipped,
4628 listed = or_models.len(),
4629 "skipped OpenRouter model rows with invalid pricing"
4630 );
4631 }
4632 if offerings.is_empty() {
4633 return Err(CatalogRefreshError::InvalidResponse);
4634 }
4635 offerings
4636 } else if catalog_endpoint_is_baseten(api_provider, base_url) {
4637 let baseten_models = parse_baseten_models_response(body)?;
4638 if baseten_models.is_empty() {
4639 return Err(CatalogRefreshError::EmptyList);
4640 }
4641 baseten_models
4642 .iter()
4643 .map(|item| baseten_to_catalog_offering(item, &provider, &fingerprint, fetched_at))
4644 .collect::<Result<Vec<_>, _>>()?
4645 } else if api_provider == ProviderKind::Telecomjs {
4646 named_gateway_catalog_offerings_from_body(
4647 body,
4648 codewhale_config::ProviderKind::Telecomjs,
4649 &provider,
4650 &fingerprint,
4651 fetched_at,
4652 )?
4653 } else if api_provider == ProviderKind::Edenai {
4654 named_gateway_catalog_offerings_from_body(
4655 body,
4656 codewhale_config::ProviderKind::Edenai,
4657 &provider,
4658 &fingerprint,
4659 fetched_at,
4660 )?
4661 } else if api_provider == ProviderKind::Zenmux {
4662 named_gateway_catalog_offerings_from_body(
4663 body,
4664 codewhale_config::ProviderKind::Zenmux,
4665 &provider,
4666 &fingerprint,
4667 fetched_at,
4668 )?
4669 } else if provider == "codewhale" {
4670 // The Codewhale API's own listing states the wire protocol per
4671 // model, so it is the catalog authority for this route.
4672 codewhale_catalog_offerings_from_body(body, &provider, &fingerprint, fetched_at)?
4673 } else if provider == "concentrate" {
4674 // Concentrate's unauthenticated `GET /v1/models` is the same
4675 // OpenAI list shape (`{"object":"list","data":[{"id":..}]}`);
4676 // rows stay unclaimed unless a same-provider bundled row exists.
4677 named_gateway_catalog_offerings_from_body(
4678 body,
4679 codewhale_config::ProviderKind::Concentrate,
4680 &provider,
4681 &fingerprint,
4682 fetched_at,
4683 )?
4684 } else {
4685 let models = apply_provider_model_cutline(
4686 api_provider,
4687 parse_models_response_for_provider(body, api_provider)
4688 .map_err(|_| CatalogRefreshError::InvalidResponse)?,
4689 );
4690 if models.is_empty() {
4691 return Err(CatalogRefreshError::EmptyList);
4692 }
4693 models
4694 .into_iter()
4695 .map(|model| CatalogOffering {
4696 cost_source: None,
4697 modalities_source: None,
4698 provider: provider.clone(),
4699 endpoint_key: if api_provider == ProviderKind::OpencodeGo {
4700 codewhale_config::opencode_go_endpoint_key(&model.id)
4701 .expect("filtered Go roster")
4702 } else {
4703 "chat"
4704 }
4705 .to_string(),
4706 wire_model_id: model.id,
4707 canonical_model: None,
4708 default_for_provider: false,
4709 family: None,
4710 limit: None,
4711 cost: None,
4712 modalities: None,
4713 attachment: None,
4714 reasoning: None,
4715 tool_call: None,
4716 structured_output: None,
4717 reasoning_options: Vec::new(),
4718 source: CatalogSource::Live {
4719 base_url_fingerprint: fingerprint.clone(),
4720 fetched_at,
4721 },
4722 })
4723 .collect()
4724 };
4725
4726 if offerings.iter().any(|row| {
4727 !crate::provider_lake::valid_catalog_model_id(&row.wire_model_id)
4728 || model_bound_secret_values
4729 .iter()
4730 .any(|secret| !secret.is_empty() && row.wire_model_id.contains(secret))
4731 }) {
4732 return Err(CatalogRefreshError::InvalidResponse);
4733 }
4734 if offerings.len() > PROVIDER_CATALOG_MAX_ROWS {
4735 return Err(CatalogRefreshError::InvalidResponse);
4736 }
4737 Ok(ProviderCatalogDelta {
4738 provider,
4739 base_url_fingerprint: fingerprint,
4740 fetched_at,
4741 offerings,
4742 })
4743 }
4744
4745 fn named_gateway_catalog_offerings_from_body(
4746 body: &str,
4747 kind: codewhale_config::ProviderKind,
4748 provider: &str,
4749 fingerprint: &str,
4750 fetched_at: u64,
4751 ) -> Result<Vec<CatalogOffering>, CatalogRefreshError> {
4752 let models = parse_models_response(body).map_err(|_| CatalogRefreshError::InvalidResponse)?;
4753 if models.is_empty() {
4754 return Err(CatalogRefreshError::EmptyList);
4755 }
4756
4757 let bundled = codewhale_config::catalog::bundled_catalog_offerings();
4758 let default_model_id = kind.provider().default_model();
4759 Ok(models
4760 .into_iter()
4761 .map(|model| {
4762 let is_default = model.id.eq_ignore_ascii_case(default_model_id);
4763 let same_provider_match = bundled.iter().find(|offering| {
4764 offering.provider.eq_ignore_ascii_case(provider)
4765 && offering.wire_model_id.eq_ignore_ascii_case(&model.id)
4766 });
4767 if let Some(matched) = same_provider_match {
4768 CatalogOffering {
4769 cost_source: Some(matched.pricing_source().clone()),
4770 modalities_source: Some(matched.modalities_source().clone()),
4771 provider: provider.to_string(),
4772 wire_model_id: model.id,
4773 canonical_model: matched.canonical_model.clone(),
4774 endpoint_key: "chat".to_string(),
4775 default_for_provider: is_default,
4776 family: matched.family.clone(),
4777 limit: matched.limit.clone(),
4778 cost: matched.cost.clone(),
4779 modalities: matched.modalities.clone(),
4780 attachment: matched.attachment,
4781 reasoning: matched.reasoning,
4782 tool_call: matched.tool_call,
4783 structured_output: matched.structured_output,
4784 reasoning_options: matched.reasoning_options.clone(),
4785 source: CatalogSource::Live {
4786 base_url_fingerprint: fingerprint.to_string(),
4787 fetched_at,
4788 },
4789 }
4790 } else {
4791 CatalogOffering {
4792 cost_source: None,
4793 modalities_source: None,
4794 provider: provider.to_string(),
4795 wire_model_id: model.id,
4796 canonical_model: None,
4797 endpoint_key: "chat".to_string(),
4798 default_for_provider: is_default,
4799 family: None,
4800 limit: None,
4801 cost: None,
4802 modalities: None,
4803 attachment: None,
4804 reasoning: None,
4805 tool_call: None,
4806 structured_output: None,
4807 reasoning_options: Vec::new(),
4808 source: CatalogSource::Live {
4809 base_url_fingerprint: fingerprint.to_string(),
4810 fetched_at,
4811 },
4812 }
4813 }
4814 })
4815 .collect())
4816 }
4817
4818 /// Parse an OpenRouter `/models` response, preserving server-side ordering and
4819 /// capturing full capability metadata (#3385).
4820 fn parse_openrouter_models_response(
4821 payload: &str,
4822 ) -> Result<Vec<OpenRouterModelItem>, CatalogRefreshError> {
4823 let parsed: OpenRouterModelsResponse =
4824 serde_json::from_str(payload).map_err(|_| CatalogRefreshError::InvalidResponse)?;
4825 let listed = parsed.data.len();
4826 let mut seen = std::collections::HashSet::new();
4827 let mut malformed = 0usize;
4828 let mut models = Vec::with_capacity(listed);
4829 for row in parsed.data {
4830 // `~`-prefixed ids (`~deepseek/deepseek-pro-latest`) are OpenRouter's
4831 // moving "latest" aliases, not billing identities, so they are dropped
4832 // silently. Any other row that does not decode or carries an id the
4833 // catalog cannot hold is skipped and counted: one such row used to
4834 // fail the whole roster closed, so no OpenRouter route could ever be
4835 // priced from the lake (#6690).
4836 let Ok(item) = serde_json::from_str::<OpenRouterModelItem>(row.get()) else {
4837 malformed += 1;
4838 continue;
4839 };
4840 if item.id.starts_with('~') {
4841 continue;
4842 }
4843 if !crate::provider_lake::valid_catalog_model_id(&item.id) {
4844 malformed += 1;
4845 continue;
4846 }
4847 if seen.insert(item.id.clone()) {
4848 models.push(item);
4849 }
4850 }
4851 if malformed > 0 {
4852 tracing::warn!(
4853 malformed,
4854 listed,
4855 "skipped malformed OpenRouter model rows in the catalog refresh"
4856 );
4857 }
4858 // Rows were listed but not one survived: the response is not a roster.
4859 if models.is_empty() && malformed > 0 {
4860 return Err(CatalogRefreshError::InvalidResponse);
4861 }
4862 Ok(models)
4863 }
4864
4865 /// Parse Baseten's authenticated Model APIs catalog without inferring facts
4866 /// from an identically named model on another provider.
4867 fn parse_baseten_models_response(
4868 payload: &str,
4869 ) -> Result<Vec<BasetenModelItem>, CatalogRefreshError> {
4870 let parsed: BasetenModelsResponse =
4871 serde_json::from_str(payload).map_err(|_| CatalogRefreshError::InvalidResponse)?;
4872 let mut seen = std::collections::HashSet::new();
4873 let mut models = Vec::with_capacity(parsed.data.len());
4874 for mut item in parsed.data {
4875 item.id = item.id.trim().to_string();
4876 if item.id.is_empty() || !seen.insert(item.id.clone()) {
4877 return Err(CatalogRefreshError::InvalidResponse);
4878 }
4879 models.push(item);
4880 }
4881 Ok(models)
4882 }
4883
4884 fn catalog_number_f64(value: Option<&CatalogNumber>) -> Result<Option<f64>, CatalogRefreshError> {
4885 let Some(value) = value else {
4886 return Ok(None);
4887 };
4888 let parsed = match value {
4889 CatalogNumber::Number(number) => number.as_f64(),
4890 CatalogNumber::Text(text) => text.trim().parse::<f64>().ok(),
4891 }
4892 .filter(|number| number.is_finite() && *number >= 0.0)
4893 .ok_or(CatalogRefreshError::InvalidResponse)?;
4894 Ok(Some(parsed))
4895 }
4896
4897 fn catalog_number_u64(value: Option<&CatalogNumber>) -> Result<Option<u64>, CatalogRefreshError> {
4898 let Some(number) = catalog_number_f64(value)? else {
4899 return Ok(None);
4900 };
4901 if number.fract() != 0.0 || number > u64::MAX as f64 {
4902 return Err(CatalogRefreshError::InvalidResponse);
4903 }
4904 Ok(Some(number as u64))
4905 }
4906
4907 fn checked_per_token_to_per_million(value: f64) -> Result<f64, CatalogRefreshError> {
4908 let scaled = value * 1_000_000.0;
4909 (scaled.is_finite()
4910 && (0.0..=codewhale_config::pricing::MAX_PLAUSIBLE_PRICE_PER_MILLION).contains(&scaled))
4911 .then_some(scaled)
4912 .ok_or(CatalogRefreshError::InvalidResponse)
4913 }
4914
4915 fn catalog_price_per_million(
4916 value: Option<&CatalogNumber>,
4917 ) -> Result<Option<f64>, CatalogRefreshError> {
4918 catalog_number_f64(value)?
4919 .map(checked_per_token_to_per_million)
4920 .transpose()
4921 }
4922
4923 fn feature_matches_any(feature: &str, aliases: &[&str]) -> bool {
4924 let normalized = feature.replace('-', "_");
4925 aliases.iter().any(|alias| normalized == *alias)
4926 }
4927
4928 fn baseten_features(item: &BasetenModelItem) -> Option<Vec<String>> {
4929 let sources = [
4930 item.supported_parameters.as_ref(),
4931 item.supported_features.as_ref(),
4932 item.features.as_ref(),
4933 ];
4934 let mut features = Vec::new();
4935 let mut published = false;
4936 for source in sources.into_iter().flatten() {
4937 published = true;
4938 for feature in source {
4939 let normalized = feature.trim().to_ascii_lowercase();
4940 if !normalized.is_empty() && !features.contains(&normalized) {
4941 features.push(normalized);
4942 }
4943 }
4944 }
4945 published.then_some(features)
4946 }
4947
4948 fn baseten_to_catalog_offering(
4949 item: &BasetenModelItem,
4950 provider: &str,
4951 base_url_fingerprint: &str,
4952 fetched_at: u64,
4953 ) -> Result<CatalogOffering, CatalogRefreshError> {
4954 use codewhale_config::models_dev::{ModelsDevCost, ModelsDevLimit, ModelsDevModalities};
4955
4956 let context = catalog_number_u64(
4957 item.top_provider
4958 .as_ref()
4959 .and_then(|provider| provider.context_length.as_ref())
4960 .or(item.context_length.as_ref())
4961 .or(item.context_window.as_ref())
4962 .or_else(|| {
4963 item.limits
4964 .as_ref()
4965 .and_then(|limits| limits.context.as_ref())
4966 }),
4967 )?;
4968 let output = catalog_number_u64(
4969 item.top_provider
4970 .as_ref()
4971 .and_then(|provider| provider.max_completion_tokens.as_ref())
4972 .or(item.max_output_tokens.as_ref())
4973 .or(item.max_completion_tokens.as_ref())
4974 .or_else(|| {
4975 item.limits
4976 .as_ref()
4977 .and_then(|limits| limits.output.as_ref())
4978 }),
4979 )?;
4980 let limit = (context.is_some() || output.is_some()).then_some(ModelsDevLimit {
4981 context,
4982 input: context,
4983 output,
4984 });
4985
4986 let cost = if let Some(pricing) = item.pricing.as_ref() {
4987 let cost = ModelsDevCost {
4988 input: catalog_price_per_million(pricing.prompt.as_ref())?,
4989 output: catalog_price_per_million(pricing.completion.as_ref())?,
4990 cache_read: catalog_price_per_million(pricing.input_cache_read.as_ref())?,
4991 cache_write: catalog_price_per_million(pricing.input_cache_write.as_ref())?,
4992 };
4993 if !codewhale_config::pricing::catalog_cost_is_valid(&cost) {
4994 return Err(CatalogRefreshError::InvalidResponse);
4995 }
4996 (cost.input.is_some()
4997 || cost.output.is_some()
4998 || cost.cache_read.is_some()
4999 || cost.cache_write.is_some())
5000 .then_some(cost)
5001 } else {
5002 None
5003 };
5004
5005 let features = baseten_features(item);
5006 let has_feature = |needles: &[&str]| {
5007 features.as_ref().is_some_and(|features| {
5008 features
5009 .iter()
5010 .any(|feature| feature_matches_any(feature, needles))
5011 })
5012 };
5013 let mut input_modalities = item
5014 .architecture
5015 .as_ref()
5016 .and_then(|architecture| architecture.input_modalities.clone())
5017 .or_else(|| item.input_modalities.clone());
5018 let mut output_modalities = item
5019 .architecture
5020 .as_ref()
5021 .and_then(|architecture| architecture.output_modalities.clone())
5022 .or_else(|| item.output_modalities.clone());
5023 if input_modalities.is_none() {
5024 let supports_vision = has_feature(&["vision", "image", "image_input"]);
5025 let supports_audio = has_feature(&["audio", "audio_input"]);
5026 if supports_vision || supports_audio {
5027 let mut derived = vec!["text".to_string()];
5028 if supports_vision {
5029 derived.push("image".to_string());
5030 }
5031 if supports_audio {
5032 derived.push("audio".to_string());
5033 }
5034 input_modalities = Some(derived);
5035 output_modalities.get_or_insert_with(|| vec!["text".to_string()]);
5036 }
5037 }
5038 let modalities = if input_modalities.is_some() || output_modalities.is_some() {
5039 Some(ModelsDevModalities {
5040 input: input_modalities.unwrap_or_default(),
5041 output: output_modalities.unwrap_or_default(),
5042 })
5043 } else {
5044 None
5045 };
5046 let attachment = modalities.as_ref().map(|modalities| {
5047 modalities
5048 .input
5049 .iter()
5050 .any(|modality| !modality.eq_ignore_ascii_case("text") && !modality.trim().is_empty())
5051 });
5052
5053 let feature_support = |needles: &[&str]| {
5054 features.as_ref().map(|features| {
5055 features
5056 .iter()
5057 .any(|feature| feature_matches_any(feature, needles))
5058 })
5059 };
5060 let reasoning = item
5061 .reasoning
5062 .or(item.supports_reasoning)
5063 .or_else(|| feature_support(&["reasoning", "include_reasoning"]));
5064 // Baseten's current Model APIs contract states every catalog model supports
5065 // tool calling and structured outputs. Explicit upstream booleans still
5066 // win if the endpoint publishes a narrower model-specific fact.
5067 // Baseten's Model APIs contract applies these two capabilities to every
5068 // catalog model. `supported_features` is additive and may list only
5069 // model-variable facts such as `reasoning` or `vision`; absence from that
5070 // list is therefore not an explicit false. Only an upstream boolean may
5071 // narrow the universal contract for a specific row.
5072 let tool_call = item.supports_tools.or(Some(true));
5073 let structured_output = item.supports_structured_output.or(Some(true));
5074
5075 Ok(CatalogOffering {
5076 cost_source: None,
5077 modalities_source: None,
5078 provider: provider.to_string(),
5079 wire_model_id: item.id.clone(),
5080 canonical_model: None,
5081 endpoint_key: "chat".to_string(),
5082 default_for_provider: item
5083 .id
5084 .eq_ignore_ascii_case(codewhale_config::catalog::BASETEN_DEFAULT_MODEL),
5085 family: None,
5086 limit,
5087 cost,
5088 modalities,
5089 attachment,
5090 reasoning,
5091 tool_call,
5092 structured_output,
5093 reasoning_options: item.reasoning_options.clone(),
5094 source: CatalogSource::Live {
5095 base_url_fingerprint: base_url_fingerprint.to_string(),
5096 fetched_at,
5097 },
5098 })
5099 }
5100
5101 #[cfg(test)]
5102 fn publish_provider_lake_scope(cache: &ProviderCatalogCache, provider: &str, fingerprint: &str) {
5103 // Publish fresh *and* stale/prior rows so pickers keep live catalog coverage
5104 // after TTL expiry or a failed refresh (#4139). Exact replacement is
5105 // essential: a successful smaller roster must remove upstream-retired ids,
5106 // while a failure preserves the rows already stored in this cache scope.
5107 let offerings = cache
5108 .get(provider, fingerprint)
5109 .map(|entry| entry.offerings.clone())
5110 .unwrap_or_default();
5111 crate::provider_lake::replace_provider_live_snapshot(provider, CatalogSnapshot { offerings });
5112 }
5113
5114 /// Convert an OpenRouter model item into a [`CatalogOffering`] with live-sourced
5115 /// limits, pricing, reasoning, and modalities (#3385).
5116 fn openrouter_to_catalog_offering(
5117 item: &OpenRouterModelItem,
5118 provider: &str,
5119 base_url_fingerprint: &str,
5120 fetched_at: u64,
5121 ) -> Result<CatalogOffering, CatalogRefreshError> {
5122 use codewhale_config::models_dev::{ModelsDevCost, ModelsDevLimit, ModelsDevModalities};
5123
5124 let context_length = item
5125 .top_provider
5126 .as_ref()
5127 .and_then(|tp| tp.context_length)
5128 .or(item.context_length);
5129
5130 let max_output = item
5131 .top_provider
5132 .as_ref()
5133 .and_then(|tp| tp.max_completion_tokens);
5134
5135 let limit = if context_length.is_some() || max_output.is_some() {
5136 Some(ModelsDevLimit {
5137 context: context_length.map(u64::from),
5138 input: context_length.map(u64::from),
5139 output: max_output.map(u64::from),
5140 })
5141 } else {
5142 None
5143 };
5144
5145 let cost = if let Some(p) = item.pricing.as_ref() {
5146 // OpenRouter quotes per-token USD strings; ModelsDevCost is per million.
5147 // Its routers (`openrouter/auto`, `openrouter/fusion`, ...) publish
5148 // `"-1"`: the price depends on the model the router picks. A negative
5149 // price therefore means "no fixed rate" for the whole row, never a
5150 // partial one.
5151 let parse_price = |value: &Option<String>| -> Result<Option<f64>, CatalogRefreshError> {
5152 value
5153 .as_ref()
5154 .map(|value| {
5155 value
5156 .trim()
5157 .parse::<f64>()
5158 .ok()
5159 .filter(|value| value.is_finite())
5160 .ok_or(CatalogRefreshError::InvalidResponse)
5161 })
5162 .transpose()
5163 };
5164 let raw = [
5165 parse_price(&p.prompt)?,
5166 parse_price(&p.completion)?,
5167 parse_price(&p.input_cache_read)?,
5168 parse_price(&p.input_cache_write)?,
5169 ];
5170 if raw.iter().flatten().any(|price| *price < 0.0) {
5171 None
5172 } else {
5173 let per_million =
5174 |price: Option<f64>| price.map(checked_per_token_to_per_million).transpose();
5175 let cost = ModelsDevCost {
5176 input: per_million(raw[0])?,
5177 output: per_million(raw[1])?,
5178 cache_read: per_million(raw[2])?,
5179 cache_write: per_million(raw[3])?,
5180 };
5181 if !codewhale_config::pricing::catalog_cost_is_valid(&cost) {
5182 return Err(CatalogRefreshError::InvalidResponse);
5183 }
5184 Some(cost)
5185 }
5186 } else {
5187 None
5188 };
5189
5190 let reasoning = item.supported_parameters.as_ref().map(|params| {
5191 params
5192 .iter()
5193 .any(|p| p == "reasoning" || p == "include_reasoning" || p.contains("reasoning"))
5194 });
5195
5196 let tool_call = item.supported_parameters.as_ref().map(|params| {
5197 params
5198 .iter()
5199 .any(|p| p == "tools" || p == "tool_choice" || p == "functions" || p.contains("tool"))
5200 });
5201
5202 let modalities = item.architecture.as_ref().map(|arch| {
5203 let mut input = arch.input_modalities.clone().unwrap_or_default();
5204 let mut output = arch.output_modalities.clone().unwrap_or_default();
5205 if input.is_empty()
5206 && output.is_empty()
5207 && let Some((left, right)) = arch
5208 .modality
5209 .as_deref()
5210 .and_then(|value| value.split_once("->"))
5211 {
5212 input.extend(
5213 left.split('+')
5214 .map(str::trim)
5215 .filter(|value| !value.is_empty())
5216 .map(str::to_string),
5217 );
5218 output.extend(
5219 right
5220 .split('+')
5221 .map(str::trim)
5222 .filter(|value| !value.is_empty())
5223 .map(str::to_string),
5224 );
5225 }
5226 ModelsDevModalities { input, output }
5227 });
5228
5229 Ok(CatalogOffering {
5230 cost_source: None,
5231 modalities_source: None,
5232 provider: provider.to_string(),
5233 wire_model_id: item.id.clone(),
5234 canonical_model: None,
5235 endpoint_key: "chat".to_string(),
5236 default_for_provider: false,
5237 family: None,
5238 limit,
5239 cost,
5240 modalities,
5241 attachment: None,
5242 reasoning,
5243 tool_call,
5244 structured_output: None,
5245 reasoning_options: Vec::new(),
5246 source: CatalogSource::Live {
5247 base_url_fingerprint: base_url_fingerprint.to_string(),
5248 fetched_at,
5249 },
5250 })
5251 }
5252
5253 /// The rate a main interactive turn freezes at its dispatch boundary (#6690).
5254 /// An operator-declared `[[custom_models]]` rate wins, exactly as on the
5255 /// background envelope path; `client` must be installed on the endpoint this
5256 /// turn dispatches to, so a declaration never leaks onto another endpoint. A
5257 /// declaration with no rates yields to the catalog price for that endpoint.
5258 pub(crate) fn main_turn_pricing_quote_at(
5259 client: Option<&CodewhaleClient>,
5260 provider: ProviderKind,
5261 provider_identity: &str,
5262 model: &str,
5263 endpoint_fingerprint: &str,
5264 dispatched_at: u64,
5265 ) -> Option<crate::provider_catalog_live::ProviderLivePricingQuote> {
5266 let declared = client
5267 .filter(|client| {
5268 crate::cost_status::endpoint_fingerprint(client.base_url()).as_deref()
5269 == Some(endpoint_fingerprint)
5270 })
5271 .and_then(|client| {
5272 client.configured_pricing_quote_at(provider, provider_identity, model, dispatched_at)
5273 });
5274 crate::provider_catalog_live::declared_or_catalog_quote(declared, || {
5275 crate::provider_catalog_live::fresh_provider_live_pricing_quote_at(
5276 provider,
5277 provider_identity,
5278 model,
5279 endpoint_fingerprint,
5280 dispatched_at,
5281 )
5282 })
5283 }
5284
5285 pub(super) fn system_to_instructions(system: Option<SystemPrompt>) -> Option<String> {
5286 match system {
5287 Some(SystemPrompt::Text(text)) => Some(text),
5288 Some(SystemPrompt::Blocks(blocks)) => {
5289 let joined = blocks
5290 .into_iter()
5291 .map(|b| b.text)
5292 .collect::<Vec<_>>()
5293 .join("\n\n---\n\n");
5294 if joined.trim().is_empty() {
5295 None
5296 } else {
5297 Some(joined)
5298 }
5299 }
5300 None => None,
5301 }
5302 }
5303
5304 /// Write DeepSeek's Chat Completions thinking controls from the shared tier
5305 /// table.
5306 ///
5307 /// The table (`client::deepseek_effort`) is the only place the tier ladder is
5308 /// written down; this function is just the Chat wire's spelling of it. An
5309 /// effort string the table does not name is not a DeepSeek tier request, so
5310 /// nothing is written rather than guessing a field the user did not ask for.
5311 fn apply_deepseek_chat_reasoning_effort(body: &mut Value, normalized: &str) {
5312 let Some(tier) = deepseek_effort::deepseek_effort_tier(normalized) else {
5313 return;
5314 };
5315 if let Some(value) = tier.chat_reasoning_effort() {
5316 body["reasoning_effort"] = json!(value);
5317 }
5318 body["thinking"] = json!({
5319 "type": if tier.chat_thinking_enabled() { "enabled" } else { "disabled" },
5320 });
5321 }
5322
5323 pub(super) fn apply_reasoning_effort(
5324 body: &mut Value,
5325 effort: Option<&str>,
5326 provider: ProviderKind,
5327 ) {
5328 let Some(effort) = effort else {
5329 return;
5330 };
5331 let normalized = effort.trim().to_ascii_lowercase();
5332 // DeepSeek's first-party routes read their tier ladder from the one
5333 // annotated table (`client::deepseek_effort`), shared with the Responses
5334 // wire, so a documented mapping change is a single edit. Every other
5335 // provider keeps its own dialect below.
5336 if matches!(provider, ProviderKind::Deepseek) {
5337 apply_deepseek_chat_reasoning_effort(body, &normalized);
5338 return;
5339 }
5340 if provider == ProviderKind::Stepfun {
5341 // Step 3.5 has no documented effort selector. The other coding models
5342 // are always reasoning-capable; do not invent an Off wire value.
5343 let model = body.get("model").and_then(Value::as_str).unwrap_or("");
5344 let tier = match (model, normalized.as_str()) {
5345 ("step-5-preview" | "step-3.7-flash" | "step-3.5-flash-2603", "minimal" | "low") => {
5346 Some("low")
5347 }
5348 ("step-5-preview" | "step-3.7-flash", "medium") => Some("medium"),
5349 ("step-3.5-flash-2603", "medium") => Some("high"),
5350 (
5351 "step-5-preview" | "step-3.7-flash" | "step-3.5-flash-2603",
5352 "high" | "max" | "xhigh" | "ultra",
5353 ) => Some("high"),
5354 _ => None,
5355 };
5356 if let Some(tier) = tier {
5357 body["reasoning_effort"] = json!(tier);
5358 }
5359 return;
5360 }
5361 match normalized.as_str() {
5362 "off" | "disabled" | "none" | "false" => match provider {
5363 // Handled by the shared DeepSeek table above, before this match.
5364 ProviderKind::Deepseek => {}
5365 ProviderKind::Openrouter
5366 | ProviderKind::Orcarouter
5367 | ProviderKind::XiaomiMimo
5368 | ProviderKind::Novita
5369 | ProviderKind::Siliconflow
5370 | ProviderKind::SiliconflowCN
5371 | ProviderKind::Sglang
5372 | ProviderKind::Volcengine
5373 | ProviderKind::Deepinfra
5374 | ProviderKind::Together
5375 | ProviderKind::Atlascloud
5376 | ProviderKind::Zai => {
5377 body["thinking"] = json!({ "type": "disabled" });
5378 }
5379 // TelecomJS TokenHub: the gateway's OpenAI Chat Completions API
5380 // (POST /v1/chat/completions) does not document `reasoning_effort`
5381 // or `thinking` as supported parameters. The `thinking` field is
5382 // only available on the Anthropic Messages API (POST /v1/messages)
5383 // with a different shape ({"type":"enabled","budget_tokens":N}).
5384 // Since CodeWhale routes TelecomJS through the Chat Completions
5385 // path, we must NOT inject these fields — the gateway may silently
5386 // ignore them or reject the request, and not every gateway model
5387 // (qwen-max, deepseek-chat, gpt-4o, claude, etc.) accepts the same
5388 // reasoning dialect (#4188 review: verify against actual behavior).
5389 ProviderKind::Telecomjs => {}
5390 // Eden AI documents `thinking` only for Anthropic Claude models.
5391 // This gateway can route unrelated model families, so the generic
5392 // provider must not inject a model-specific reasoning dialect.
5393 ProviderKind::Edenai => {}
5394 ProviderKind::Zenmux => {}
5395 // CSDN 星图 is OpenAI-compatible but documents no provider-owned
5396 // reasoning dialect; do not invent one.
5397 ProviderKind::Csdn => {}
5398 // The Codewhale API is a passthrough to the account's own
5399 // connected provider; it documents no Codewhale-owned
5400 // reasoning-effort translation, so nothing is invented here.
5401 ProviderKind::Codewhale => {}
5402 // Concentrate rides the Responses wire (`reasoning.effort`), never
5403 // these Chat Completions controls.
5404 ProviderKind::Concentrate => {}
5405 // Model Studio (DashScope): its top-level controls are route- AND
5406 // model-specific, so the provider enum alone cannot decide them —
5407 // a custom `base_url` on the same identity is an arbitrary
5408 // gateway. `apply_modelstudio_route_reasoning_controls` in
5409 // client::chat is the sole writer; it strips these fields for all
5410 // four variants and re-adds them only on a verified Alibaba host.
5411 // Source: <https://www.alibabacloud.com/help/en/model-studio/deep-thinking>
5412 ProviderKind::ModelstudioTokenPlan
5413 | ProviderKind::ModelstudioTokenPlanAnthropic
5414 | ProviderKind::ModelstudioCodingPlan
5415 | ProviderKind::ModelstudioCodingPlanAnthropic => {}
5416 ProviderKind::OpenaiCodex => {
5417 // OpenAI Codex uses Responses API — thinking handled differently
5418 }
5419 ProviderKind::Fireworks => {}
5420 // vLLM is an OpenAI-protocol server, not an Anthropic-protocol one.
5421 // For Qwen3 / DeepSeek-R1 / other reasoning models hosted via vLLM,
5422 // the canonical OpenAI extension to disable thinking is
5423 // `chat_template_kwargs.enable_thinking`. The old
5424 // `thinking: {type: disabled}` field is Anthropic-native and
5425 // silently ignored by vLLM — the model still emits a full
5426 // reasoning trace into the `reasoning` field (which this client
5427 // doesn't surface), causing 10+ seconds of perceived "freeze"
5428 // before the first content token (PR #1480 by @h3c-hexin).
5429 ProviderKind::Vllm => {
5430 body["chat_template_kwargs"] = json!({
5431 "enable_thinking": false,
5432 });
5433 }
5434 ProviderKind::Openai
5435 | ProviderKind::WanjieArk
5436 | ProviderKind::Qianfan
5437 | ProviderKind::Arcee
5438 | ProviderKind::Huggingface
5439 | ProviderKind::Modelscope
5440 | ProviderKind::Custom => {}
5441 ProviderKind::Moonshot => {
5442 // #3024: Kimi models accept thinking enable/disable.
5443 body["thinking"] = json!({ "type": "disabled" });
5444 }
5445 ProviderKind::Ollama => {
5446 // #3024: Ollama OpenAI-compat endpoint accepts think param.
5447 body["think"] = json!(false);
5448 }
5449 ProviderKind::OllamaCloud => {
5450 // Ollama Cloud stays on the documented OpenAI-compatible
5451 // `/v1/chat/completions` wire. Native `/api/chat` uses
5452 // `think`; this wire uses `reasoning_effort`.
5453 body["reasoning_effort"] = json!("none");
5454 }
5455 ProviderKind::Anthropic
5456 | ProviderKind::DeepseekAnthropic
5457 | ProviderKind::MinimaxAnthropic
5458 | ProviderKind::Openmodel => {
5459 // Thinking shaping happens in the Messages adapter, which
5460 // applies each provider's supported control fields.
5461 }
5462 ProviderKind::NvidiaNim => {
5463 body["chat_template_kwargs"] = json!({
5464 "thinking": false,
5465 });
5466 }
5467 ProviderKind::Minimax => {}
5468 ProviderKind::Stepfun => {}
5469 ProviderKind::Sakana => {}
5470 ProviderKind::LongCat => {}
5471 ProviderKind::OpencodeGo | ProviderKind::OpencodeZen => {}
5472 ProviderKind::Meta => {}
5473 ProviderKind::Xai => {}
5474 ProviderKind::Mistral => {}
5475 ProviderKind::Google => {}
5476 ProviderKind::Antigravity => {}
5477 },
5478 "low" | "minimal" | "medium" | "mid" | "high" | "" => match provider {
5479 // Handled by the shared DeepSeek table above, before this match.
5480 ProviderKind::Deepseek => {}
5481 // DeepSeek-compatible hosted routes: low/medium both map to high.
5482 // Their own wire contracts are not verified here, so the historic
5483 // collapse stays rather than inventing unsupported wire values.
5484 ProviderKind::Siliconflow
5485 | ProviderKind::SiliconflowCN
5486 | ProviderKind::Sglang
5487 | ProviderKind::Volcengine
5488 | ProviderKind::Deepinfra
5489 | ProviderKind::Atlascloud => {
5490 body["reasoning_effort"] = json!("high");
5491 body["thinking"] = json!({ "type": "enabled" });
5492 }
5493 // TelecomJS: see comment in the "off" branch above — the gateway's
5494 // Chat Completions API does not support reasoning_effort or thinking.
5495 ProviderKind::Telecomjs => {}
5496 ProviderKind::Edenai => {}
5497 ProviderKind::Zenmux => {}
5498 // CSDN 星图 is OpenAI-compatible but documents no provider-owned
5499 // reasoning dialect; do not invent one.
5500 ProviderKind::Csdn => {}
5501 // The Codewhale API is a passthrough to the account's own
5502 // connected provider; it documents no Codewhale-owned
5503 // reasoning-effort translation, so nothing is invented here.
5504 ProviderKind::Codewhale => {}
5505 // Concentrate rides the Responses wire (`reasoning.effort`), never
5506 // these Chat Completions controls.
5507 ProviderKind::Concentrate => {}
5508 // Model Studio: see the "off" branch — the route- and model-aware
5509 // shaper in client::chat is the sole writer of these fields.
5510 ProviderKind::ModelstudioTokenPlan
5511 | ProviderKind::ModelstudioTokenPlanAnthropic
5512 | ProviderKind::ModelstudioCodingPlan
5513 | ProviderKind::ModelstudioCodingPlanAnthropic => {}
5514 // OpenRouter/OrcaRouter/Novita/Together: pass through the actual
5515 // user-chosen value. OpenRouter's unified scale is
5516 // none/minimal/low/medium/high/xhigh; DeepSeek models hosted there
5517 // accept those directly.
5518 ProviderKind::Openrouter
5519 | ProviderKind::Orcarouter
5520 | ProviderKind::Novita
5521 | ProviderKind::Together => {
5522 let value = match normalized.as_str() {
5523 "low" | "minimal" => "low",
5524 "medium" | "mid" => "medium",
5525 _ => "high",
5526 };
5527 body["reasoning_effort"] = json!(value);
5528 body["thinking"] = json!({ "type": "enabled" });
5529 }
5530 ProviderKind::XiaomiMimo => {
5531 body["thinking"] = json!({ "type": "enabled" });
5532 }
5533 ProviderKind::Arcee | ProviderKind::Huggingface | ProviderKind::Modelscope => {
5534 let value = match normalized.as_str() {
5535 "minimal" => "minimal",
5536 "low" => "low",
5537 "medium" | "mid" => "medium",
5538 _ => "high",
5539 };
5540 body["reasoning_effort"] = json!(value);
5541 }
5542 ProviderKind::Fireworks => {
5543 body["reasoning_effort"] = json!("high");
5544 }
5545 ProviderKind::Vllm => {
5546 body["chat_template_kwargs"] = json!({
5547 "enable_thinking": true,
5548 });
5549 // vLLM supports low/medium/high natively — pass through the
5550 // user-chosen value instead of hard-coding "high".
5551 let value = match normalized.as_str() {
5552 "low" | "minimal" => "low",
5553 "medium" | "mid" => "medium",
5554 _ => "high",
5555 };
5556 body["reasoning_effort"] = json!(value);
5557 }
5558 ProviderKind::Openai
5559 | ProviderKind::WanjieArk
5560 | ProviderKind::Qianfan
5561 | ProviderKind::OpenaiCodex
5562 | ProviderKind::Custom => {}
5563 ProviderKind::Moonshot => {
5564 // #3024: Kimi models accept thinking enable.
5565 body["thinking"] = json!({ "type": "enabled" });
5566 }
5567 ProviderKind::Ollama => {
5568 // #3024: Ollama think param.
5569 body["think"] = json!(true);
5570 }
5571 ProviderKind::OllamaCloud => {
5572 let value = match normalized.as_str() {
5573 "low" | "minimal" => "low",
5574 "medium" | "mid" => "medium",
5575 _ => "high",
5576 };
5577 body["reasoning_effort"] = json!(value);
5578 }
5579 ProviderKind::Anthropic
5580 | ProviderKind::DeepseekAnthropic
5581 | ProviderKind::MinimaxAnthropic
5582 | ProviderKind::Openmodel => {
5583 // Thinking shaping happens in the Messages adapter, which
5584 // applies each provider's supported control fields.
5585 }
5586 ProviderKind::NvidiaNim => {
5587 body["chat_template_kwargs"] = json!({
5588 "thinking": true,
5589 "reasoning_effort": "high",
5590 });
5591 }
5592 ProviderKind::Minimax => {}
5593 ProviderKind::Zai => {
5594 body["thinking"] = json!({
5595 "type": "enabled",
5596 "clear_thinking": false,
5597 });
5598 }
5599 ProviderKind::Stepfun => {}
5600 ProviderKind::Sakana => {}
5601 ProviderKind::LongCat => {}
5602 ProviderKind::OpencodeGo | ProviderKind::OpencodeZen => {}
5603 ProviderKind::Meta => {}
5604 ProviderKind::Xai => {}
5605 ProviderKind::Mistral => {}
5606 ProviderKind::Google => {}
5607 ProviderKind::Antigravity => {}
5608 },
5609 "xhigh" | "max" | "highest" | "ultra" | "ultracode" => match provider {
5610 // Handled by the shared DeepSeek table above, before this match.
5611 ProviderKind::Deepseek => {}
5612 ProviderKind::Siliconflow
5613 | ProviderKind::SiliconflowCN
5614 | ProviderKind::Sglang
5615 | ProviderKind::Volcengine
5616 | ProviderKind::Deepinfra
5617 | ProviderKind::Atlascloud => {
5618 body["reasoning_effort"] = json!("max");
5619 body["thinking"] = json!({ "type": "enabled" });
5620 }
5621 // TelecomJS: see comment in the "off" branch above — the gateway's
5622 // Chat Completions API does not support reasoning_effort or thinking.
5623 ProviderKind::Telecomjs => {}
5624 ProviderKind::Edenai => {}
5625 ProviderKind::Zenmux => {}
5626 // CSDN 星图 is OpenAI-compatible but documents no provider-owned
5627 // reasoning dialect; do not invent one.
5628 ProviderKind::Csdn => {}
5629 // The Codewhale API is a passthrough to the account's own
5630 // connected provider; it documents no Codewhale-owned
5631 // reasoning-effort translation, so nothing is invented here.
5632 ProviderKind::Codewhale => {}
5633 // Concentrate rides the Responses wire (`reasoning.effort`), never
5634 // these Chat Completions controls.
5635 ProviderKind::Concentrate => {}
5636 // Model Studio: see the "off" branch — the route- and model-aware
5637 // shaper in client::chat is the sole writer of these fields.
5638 ProviderKind::ModelstudioTokenPlan
5639 | ProviderKind::ModelstudioTokenPlanAnthropic
5640 | ProviderKind::ModelstudioCodingPlan
5641 | ProviderKind::ModelstudioCodingPlanAnthropic => {}
5642 ProviderKind::Openrouter
5643 | ProviderKind::Orcarouter
5644 | ProviderKind::Novita
5645 | ProviderKind::Together => {
5646 body["reasoning_effort"] = json!("xhigh");
5647 body["thinking"] = json!({ "type": "enabled" });
5648 }
5649 ProviderKind::XiaomiMimo => {
5650 body["thinking"] = json!({ "type": "enabled" });
5651 }
5652 ProviderKind::Arcee | ProviderKind::Huggingface | ProviderKind::Modelscope => {
5653 body["reasoning_effort"] = json!("high");
5654 }
5655 ProviderKind::Fireworks => {
5656 body["reasoning_effort"] = json!("max");
5657 }
5658 ProviderKind::Vllm => {
5659 body["chat_template_kwargs"] = json!({
5660 "enable_thinking": true,
5661 });
5662 // vLLM only supports none/low/medium/high — downgrade
5663 // "max" to "high" instead of sending an invalid value.
5664 body["reasoning_effort"] = json!("high");
5665 }
5666 ProviderKind::Openai
5667 | ProviderKind::WanjieArk
5668 | ProviderKind::Qianfan
5669 | ProviderKind::OpenaiCodex
5670 | ProviderKind::Custom => {}
5671 ProviderKind::Moonshot => {
5672 // #3024: Kimi models accept thinking enable.
5673 body["thinking"] = json!({ "type": "enabled" });
5674 }
5675 ProviderKind::Ollama => {
5676 // #3024: Ollama think param.
5677 body["think"] = json!(true);
5678 }
5679 ProviderKind::OllamaCloud => {
5680 body["reasoning_effort"] = json!("max");
5681 }
5682 ProviderKind::Anthropic
5683 | ProviderKind::DeepseekAnthropic
5684 | ProviderKind::MinimaxAnthropic
5685 | ProviderKind::Openmodel => {
5686 // Thinking shaping happens in the Messages adapter, which
5687 // applies each provider's supported control fields.
5688 }
5689 ProviderKind::NvidiaNim => {
5690 body["chat_template_kwargs"] = json!({
5691 "thinking": true,
5692 "reasoning_effort": "max",
5693 });
5694 }
5695 ProviderKind::Minimax => {}
5696 ProviderKind::Zai => {
5697 body["thinking"] = json!({
5698 "type": "enabled",
5699 "clear_thinking": false,
5700 });
5701 }
5702 ProviderKind::Stepfun => {}
5703 ProviderKind::Sakana => {}
5704 ProviderKind::LongCat => {}
5705 ProviderKind::OpencodeGo | ProviderKind::OpencodeZen => {}
5706 ProviderKind::Meta => {}
5707 ProviderKind::Xai => {}
5708 ProviderKind::Mistral => {}
5709 ProviderKind::Google => {}
5710 ProviderKind::Antigravity => {}
5711 },
5712 _ => {}
5713 }
5714 }
5715
5716 impl CodewhaleClient {
5717 /// Call the DeepSeek `/beta/completions` FIM endpoint.
5718 pub async fn fim_completion(
5719 &self,
5720 model: &str,
5721 prompt: &str,
5722 suffix: &str,
5723 max_tokens: u32,
5724 ) -> anyhow::Result<String> {
5725 let _inference = self.acquire_remote_control_inference_permit().await;
5726 let _permit = self.acquire_provider_request_permit().await;
5727 if self.api_provider == ProviderKind::OpencodeZen
5728 || self.wire_format != WireFormat::ChatCompletions
5729 {
5730 bail!(
5731 "FIM completion is not supported for {} because the route has no proven FIM wire contract ({:?})",
5732 self.api_provider.provider().display_name(),
5733 self.wire_format
5734 );
5735 }
5736 let url = api_url_with_suffix(&self.base_url, "beta/completions", None);
5737 let model = self.wire_model_for_route(model);
5738 let max_tokens = max_tokens.min(self.effective_max_output_tokens(&model));
5739 let body = json!({
5740 "model": model,
5741 "prompt": prompt,
5742 "suffix": suffix,
5743 "max_tokens": max_tokens,
5744 });
5745 let response = self.send_json_with_retry(&url, &body).await?;
5746 let status = response.status();
5747 if !status.is_success() {
5748 let raw_error_text = bounded_error_text(response, ERROR_BODY_MAX_BYTES).await;
5749 let error_text = sanitize_http_error_body(
5750 Some(self.api_provider.provider().display_name()),
5751 status.as_u16(),
5752 &raw_error_text,
5753 );
5754 anyhow::bail!("FIM API error: HTTP {status}: {error_text}");
5755 }
5756 let response_text = response
5757 .text()
5758 .await
5759 .context("Failed to read FIM API response body")?;
5760 let value: serde_json::Value =
5761 serde_json::from_str(&response_text).context("Failed to parse FIM API response")?;
5762 let text = value
5763 .pointer("/choices/0/text")
5764 .and_then(serde_json::Value::as_str)
5765 .ok_or_else(|| anyhow::anyhow!("FIM response missing choices[0].text"))?;
5766 Ok(text.to_string())
5767 }
5768 }
5769
5770 mod anthropic;
5771 mod chat;
5772 mod deepseek_effort;
5773 #[cfg(test)]
5774 mod ds4_tests;
5775 mod prepared;
5776 mod provider_native_search;
5777 mod responses;
5778 mod role_placement;
5779 mod stream_entry;
5780 pub(crate) mod system_one;
5781
5782 /// Longest a request may take to open its stream and deliver the first body
5783 /// byte before the client itself times out (#6184): the first-byte bound plus
5784 /// two header waits, because a failed or stalled open on the dual client
5785 /// retries once on the HTTP/1.1 twin under its own header wait. HTTP retries
5786 /// inside one open attempt run within that attempt's header wait. The engine
5787 /// heartbeat uses this as its awaiting-model bound, so it must cover the
5788 /// fallback or an in-progress recovery reads as a stall (#6711).
5789 #[must_use]
5790 pub(crate) fn stream_first_response_bound(open: Duration, idle: Duration) -> Duration {
5791 open.saturating_mul(2)
5792 .saturating_add(stream_entry::first_byte_timeout(idle))
5793 }
5794
5795 #[cfg(test)]
5796 pub(crate) use stream_entry::first_byte_timeout as stream_first_byte_timeout;
5797 pub(crate) use stream_entry::{is_stream_open_transport_failure, resolve_stream_open_timeout};
5798 mod wire;
5799
5800 // Retain the crate-visible accounting helpers at the existing client seam.
5801 pub(super) use wire::{parse_usage, saturating_u32};
5802
5803 #[cfg(test)]
5804 pub(crate) fn anthropic_tool_result_content_for_test(
5805 content: &str,
5806 content_blocks: Option<&[Value]>,
5807 ) -> Value {
5808 anthropic::anthropic_tool_result_content(content, content_blocks)
5809 }
5810
5811 #[cfg(test)]
5812 pub(crate) fn responses_tool_output_for_test(
5813 content: &str,
5814 content_blocks: Option<&[Value]>,
5815 ) -> Value {
5816 responses::responses_tool_output(content, content_blocks)
5817 }
5818
5819 #[cfg(test)]
5820 pub(crate) fn chat_messages_for_test(messages: &[codewhale_models::Message]) -> Vec<Value> {
5821 chat::build_chat_messages(None, messages, "gpt-4o")
5822 }
5823
5824 pub(crate) use chat::{
5825 CacheWarmupKey, PromptInspection, PromptLayerInspection, PromptLayerStability,
5826 ToolResultInspection, TurnMetaInspection,
5827 };
5828 pub(crate) use prepared::{
5829 CallerStreamMode, EndpointIdentity, PreparedOutboundRequest, RouteShape, WireBodyView,
5830 WireDialect, canonical_json,
5831 };
5832 pub(crate) use provider_native_search::{ProviderNativeSearchClient, ProviderNativeSearchRequest};
5833
5834 /// Whether a route speaks ordinary `/chat/completions` (not Messages or Responses).
5835 #[must_use]
5836 pub(crate) fn provider_speaks_chat_completions(api_provider: ProviderKind) -> bool {
5837 provider_default_wire_format(api_provider) == WireFormat::ChatCompletions
5838 }
5839
5840 pub(crate) fn inspect_prompt_for_request(request: &MessageRequest) -> PromptInspection {
5841 chat::inspect_prompt_for_request(request)
5842 }
5843
5844 pub(crate) fn build_cache_warmup_request(request: &MessageRequest) -> MessageRequest {
5845 chat::build_cache_warmup_request(request)
5846 }
5847
5848 pub(crate) use chat::CACHE_WARMUP_MAX_TOKENS;
5849 pub(crate) use chat::is_reasoning_replay_placeholder;
5850
5851 #[cfg(test)]
5852 mod tests {
5853 include!("client/test_cases_01.rs");
5854
5855 include!("client/test_cases_02.rs");
5856
5857 include!("client/test_cases_03.rs");
5858
5859 include!("client/test_cases_04.rs");
5860
5861 include!("client/test_cases_05.rs");
5862
5863 include!("client/test_cases_06.rs");
5864
5865 include!("client/test_cases_07.rs");
5866 }
5867
5868 #[cfg(test)]
5869 mod configured_model_client_tests {
5870 use super::*;
5871 use crate::config::{ProviderConfig, ProvidersConfig};
5872
5873 fn config() -> Config {
5874 Config {
5875 provider: Some("deepseek".into()),
5876 providers: Some(ProvidersConfig {
5877 deepseek: ProviderConfig {
5878 api_key: Some("configured-model-local-fixture".into()),
5879 base_url: Some("https://api.deepseek.com".into()),
5880 ..ProviderConfig::default()
5881 },
5882 ..ProvidersConfig::default()
5883 }),
5884 custom_models: Some(vec![
5885 toml::from_str(
5886 r#"
5887 provider = "deepseek"
5888 base_url = "https://api.deepseek.com"
5889 id = "deepseek-v4pro"
5890 limit = { context = 96000, output = 32 }
5891 cost = { input = 1.0, output = 2.0 }
5892 "#,
5893 )
5894 .unwrap(),
5895 ]),
5896 ..Config::default()
5897 }
5898 }
5899
5900 #[test]
5901 fn alternate_declared_model_freezes_wire_limits_and_price() {
5902 let _env = crate::test_support::lock_test_env();
5903 let mut config = config();
5904 let client = CodewhaleClient::from_parts(
5905 "https://api.deepseek.com".into(),
5906 "initial-model".into(),
5907 WireFormat::ChatCompletions,
5908 None,
5909 &config,
5910 )
5911 .unwrap();
5912 let model = "deepseek-v4pro";
5913 config.providers.as_mut().unwrap().deepseek.model = Some(model.into());
5914 crate::config::normalize_model_config_for_test(&mut config);
5915 assert_eq!(config.default_model(), model);
5916 assert_eq!(
5917 config.providers.as_ref().unwrap().deepseek.model.as_deref(),
5918 Some(model)
5919 );
5920 config.custom_models.as_mut().unwrap()[0]
5921 .limit
5922 .as_mut()
5923 .unwrap()
5924 .output = Some(512);
5925 config.custom_models.as_mut().unwrap()[0]
5926 .cost
5927 .as_mut()
5928 .unwrap()
5929 .output = Some(99.0);
5930
5931 assert_eq!(client.effective_max_output_tokens(model), 32);
5932 let prepared = client
5933 .prepare_outbound_request(
5934 translation_message_request("hello", model.into(), "English", 4096),
5935 false,
5936 )
5937 .unwrap();
5938 assert_eq!(prepared.wire_model, model);
5939 assert_eq!(prepared.body["model"], model);
5940 assert_eq!(prepared.body["max_tokens"], 32);
5941 let rebound = client
5942 .rebound_for_model_protocol(Some(&config), model)
5943 .unwrap()
5944 .unwrap();
5945 assert_eq!(rebound.default_model, model);
5946 assert_eq!(rebound.route_limits.unwrap().output_tokens, Some(32));
5947 let envelope = rebound.effective_route_envelope(model, chrono::Utc::now());
5948 assert_eq!(envelope.model, model);
5949 let quote = serde_json::to_value(envelope.provider_live_pricing.unwrap()).unwrap();
5950 assert_eq!(quote["output_per_million"], "2");
5951
5952 let mut wrong_endpoint = client.clone();
5953 wrong_endpoint.base_url = "https://other.example.test/v1".into();
5954 assert_eq!(wrong_endpoint.declared_wire_model(model), None);
5955 assert_ne!(wrong_endpoint.effective_max_output_tokens(model), 32);
5956 let mut wrong_identity = client;
5957 wrong_identity.admitted_identity.key = "other-provider".into();
5958 assert_eq!(wrong_identity.declared_wire_model(model), None);
5959 }
5960
5961 #[test]
5962 fn declaration_does_not_open_opencode_go_protocol_roster() {
5963 let _env = crate::test_support::lock_test_env();
5964 let mut config = config();
5965 config.provider = Some("opencode-go".into());
5966 config.providers.as_mut().unwrap().opencode_go = ProviderConfig {
5967 api_key: Some("configured-model-local-fixture".into()),
5968 ..ProviderConfig::default()
5969 };
5970 let base_url = ProviderKind::OpencodeGo.provider().default_base_url();
5971 let declaration = &mut config.custom_models.as_mut().unwrap()[0];
5972 declaration.provider = "opencode-go".into();
5973 declaration.base_url = base_url.into();
5974 declaration.id = "claude-sonnet-unproven".into();
5975 let client = CodewhaleClient::from_parts(
5976 base_url.into(),
5977 config.default_model(),
5978 WireFormat::ChatCompletions,
5979 None,
5980 &config,
5981 )
5982 .unwrap();
5983 assert_eq!(client.declared_wire_model("claude-sonnet-unproven"), None);
5984 assert!(
5985 client
5986 .rebound_for_model_protocol(None, "claude-sonnet-unproven")
5987 .is_err()
5988 );
5989 assert!(
5990 client
5991 .prepare_outbound_request(
5992 translation_message_request(
5993 "hello",
5994 "claude-sonnet-unproven".into(),
5995 "English",
5996 64
5997 ),
5998 false,
5999 )
6000 .is_err()
6001 );
6002 }
6003 }
6004
6005 #[cfg(test)]
6006 mod openrouter_vendor_tests {
6007 use super::*;
6008 use crate::config::{OPENROUTER_QWEN_3_6_FLASH_MODEL, ProviderConfig, ProvidersConfig};
6009 use futures_util::StreamExt;
6010 use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method};
6011
6012 fn config(base_url: &str, vendor: Option<&str>) -> Config {
6013 Config {
6014 provider: Some("openrouter".into()),
6015 providers: Some(ProvidersConfig {
6016 openrouter: ProviderConfig {
6017 api_key: Some("vendor-pin-local-fixture".into()),
6018 base_url: Some(base_url.into()),
6019 model: Some("deepseek/deepseek-v4-pro".into()),
6020 vendor: vendor.map(str::to_string),
6021 ..ProviderConfig::default()
6022 },
6023 ..ProvidersConfig::default()
6024 }),
6025 ..Config::default()
6026 }
6027 }
6028
6029 fn request() -> MessageRequest {
6030 translation_message_request("hello", "deepseek/deepseek-v4-pro".into(), "English", 64)
6031 }
6032
6033 #[tokio::test]
6034 async fn openrouter_vendor_is_serialized_on_stream_blocking_and_translation_requests() {
6035 let _env = crate::test_support::lock_test_env();
6036 for streaming in [false, true] {
6037 let server = MockServer::start().await;
6038 let response = if streaming {
6039 ResponseTemplate::new(200)
6040 .insert_header("content-type", "text/event-stream")
6041 .set_body_string("data: [DONE]\n\n")
6042 } else {
6043 ResponseTemplate::new(200).set_body_json(json!({
6044 "id": "chatcmpl-vendor-pin", "object": "chat.completion",
6045 "model": "deepseek/deepseek-v4-pro",
6046 "choices": [{"index": 0, "message": {"role": "assistant", "content": "ok"}, "finish_reason": "stop"}],
6047 "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}
6048 }))
6049 };
6050 Mock::given(method("POST"))
6051 .respond_with(response)
6052 .expect(if streaming { 1 } else { 2 })
6053 .mount(&server)
6054 .await;
6055 let client =
6056 CodewhaleClient::new(&config(&server.uri(), Some("deepinfra/turbo"))).unwrap();
6057 let mut input = request();
6058 input.stream = Some(streaming);
6059 let preview = client
6060 .prepare_outbound_request(input.clone(), streaming)
6061 .unwrap();
6062 if streaming {
6063 let mut stream = client.create_message_stream(input).await.unwrap();
6064 while let Some(event) = stream.next().await {
6065 event.unwrap();
6066 }
6067 } else {
6068 client.create_message(input).await.unwrap();
6069 assert_eq!(
6070 client
6071 .translate("hello", "deepseek/deepseek-v4-pro", "English")
6072 .await
6073 .unwrap(),
6074 "ok"
6075 );
6076 }
6077 let captured = server.received_requests().await.unwrap();
6078 for outbound in &captured {
6079 let body: Value = serde_json::from_slice(&outbound.body).unwrap();
6080 assert_eq!(
6081 body["provider"],
6082 json!({"order": ["deepinfra/turbo"], "allow_fallbacks": false})
6083 );
6084 assert_eq!(body["provider"], preview.body["provider"]);
6085 assert_eq!(body["model"], "deepseek/deepseek-v4-pro");
6086 }
6087 }
6088 }
6089
6090 #[test]
6091 fn openrouter_vendor_freezes_rebinds_and_partitions_cached_requests() {
6092 let _env = crate::test_support::lock_test_env();
6093 let initial = config("https://openrouter.ai/api/v1", Some("deepinfra/turbo"));
6094 let client = CodewhaleClient::new(&initial).unwrap();
6095 let mut updated = initial.clone();
6096 updated.providers.as_mut().unwrap().openrouter.vendor =
6097 Some("another-vendor/region".into());
6098 let fresh = CodewhaleClient::new(&updated).unwrap();
6099 let rebound = client
6100 .rebound_for_model_protocol(Some(&updated), OPENROUTER_QWEN_3_6_FLASH_MODEL)
6101 .unwrap()
6102 .unwrap();
6103 assert_eq!(rebound.openrouter_vendor(), Some("deepinfra/turbo"));
6104 assert_eq!(client.clone().openrouter_vendor(), Some("deepinfra/turbo"));
6105 let key = |client: &CodewhaleClient| {
6106 let body = client
6107 .prepare_outbound_request(request(), false)
6108 .unwrap()
6109 .body;
6110 crate::llm_response_cache::ResponseCache::make_key(
6111 "openrouter",
6112 &client.base_url,
6113 None,
6114 &client.api_key,
6115 &serde_json::to_vec(&body).unwrap(),
6116 )
6117 };
6118 assert_ne!(key(&client), key(&fresh));
6119 updated.providers.as_mut().unwrap().openrouter.vendor = Some(String::new());
6120 let cleared = CodewhaleClient::new(&updated).unwrap();
6121 assert!(
6122 cleared
6123 .prepare_outbound_request(request(), false)
6124 .unwrap()
6125 .body
6126 .get("provider")
6127 .is_none()
6128 );
6129 assert_ne!(key(&fresh), key(&cleared));
6130 assert_eq!(
6131 client.turn_route_receipt().openrouter_vendor(),
6132 Some("deepinfra/turbo")
6133 );
6134 assert!(
6135 CodewhaleClient::new(&config("https://openrouter.ai/api/v1", Some("bad vendor")))
6136 .is_err()
6137 );
6138 updated.provider = Some("openai".into());
6139 updated.providers.as_mut().unwrap().openai = ProviderConfig {
6140 api_key: Some("other-provider-fixture".into()),
6141 base_url: Some("https://openrouter.ai/api/v1".into()),
6142 model: Some("deepseek/deepseek-v4-pro".into()),
6143 ..ProviderConfig::default()
6144 };
6145 let other = CodewhaleClient::new(&updated).unwrap();
6146 assert!(
6147 other
6148 .prepare_outbound_request(request(), false)
6149 .unwrap()
6150 .body
6151 .get("provider")
6152 .is_none()
6153 );
6154 }
6155 }
6156
6156 lines RUST