From 1484e3db27c71addc998ad94a430cdc687dea0a3 Mon Sep 17 00:00:00 2001 From: dittops Date: Tue, 29 Sep 2026 11:10:14 +0530 Subject: [PATCH 1/3] fix(gateway): resolve the Azure OpenAI speech SSRF check off the async workers `SelfHostedTTS::new_azure_openai` ran `validate_azure_openai_url`, whose resolve-then-validate step is a blocking `getaddrinfo`, inside the synchronous provider constructor -- on a tokio worker, once per /v1/audio/speech request. Under Kubernetes' `ndots:5` an external name takes ~70 ms to resolve, so the replica's two workers served ~25 speech requests a second at ~0.08 cores, and every other session on the pod waited with them: realtime first-audio p99 went 722 -> 1535 ms under 25 req/s of speech. Transcription already runs the same check through `spawn_blocking` (`transcribe.rs`). Speech now does the same: * `core::net::validate_url_for_ssrf_without_dns` -- every check except the resolution (scheme, blocked hostnames, IP literals). The constructor uses it, so a bad URL is still refused at construction with the same message. * `SelfHostedTTS::connect` runs the full check in `spawn_blocking` before the first dial, so a name that resolves to a private address is still refused. Tests: the DNS-free variant refuses everything the full check does except a name that only resolves privately; a speech endpoint on `ip6-localhost` constructs and is refused at `connect` (fails with the connect-time check removed). Co-Authored-By: Claude Opus 5.5 --- gateway/src/core/net.rs | 76 +++++++++++++++++++++++++---- gateway/src/core/tts/self_hosted.rs | 73 ++++++++++++++++++++++++--- 2 files changed, 132 insertions(+), 17 deletions(-) diff --git a/gateway/src/core/net.rs b/gateway/src/core/net.rs index 9da31bc6..4bbb4f4d 100644 --- a/gateway/src/core/net.rs +++ b/gateway/src/core/net.rs @@ -111,6 +111,21 @@ pub fn validate_url_for_ssrf(url: &str, allowed_schemes: &[&str]) -> Result<(), validate_url_for_ssrf_inner(url, allowed_schemes, loopback_endpoints_allowed()) } +/// [`validate_url_for_ssrf`] without the resolve-then-validate step: the scheme allowlist, the +/// blocked hostnames and every IP-literal spelling, and nothing that does I/O. +/// +/// For a synchronous constructor on the request path. `validate_url_for_ssrf` resolves a DNS +/// name with a blocking `getaddrinfo`, and a constructor that runs on a tokio worker holds that +/// worker for the whole lookup (~70 ms for an external name under Kubernetes' `ndots:5`), so +/// every other session on the replica waits with it. Pair this with the full check off the +/// workers (`tokio::task::spawn_blocking`) before the first dial. +pub fn validate_url_for_ssrf_without_dns( + url: &str, + allowed_schemes: &[&str], +) -> Result<(), String> { + ssrf_dns_host(url, allowed_schemes, loopback_endpoints_allowed()).map(|_| ()) +} + /// Build a reqwest redirect policy that validates every redirect target before /// following it. Use this for requests whose original URL passed /// [`validate_url_for_ssrf`]; reqwest follows redirects by default, and an @@ -158,6 +173,23 @@ fn validate_url_for_ssrf_inner( allowed_schemes: &[&str], loopback_allowed: bool, ) -> Result<(), String> { + match ssrf_dns_host(url, allowed_schemes, loopback_allowed)? { + // Resolve-then-validate: when the host is a DNS name (not an IP literal), + // resolve it and reject if ANY resolved address is private/internal. This + // closes DNS-rebinding / TOCTOU holes where a public-looking hostname + // resolves to a private/metadata address. + Some(host) => validate_resolved_host_for_ssrf(&host), + None => Ok(()), + } +} + +/// Every check [`validate_url_for_ssrf_inner`] makes that needs no I/O. `Ok(Some(host))` when +/// the host is a DNS name still to be resolved, `Ok(None)` when nothing is left to check. +fn ssrf_dns_host( + url: &str, + allowed_schemes: &[&str], + loopback_allowed: bool, +) -> Result, String> { let parsed = url::Url::parse(url).map_err(|e| format!("invalid URL '{}': {}", url, e))?; // Scheme allowlist — applies even when the loopback escape hatch is on. @@ -172,7 +204,7 @@ fn validate_url_for_ssrf_inner( // Test/local-mock escape hatch (opt-in, OFF by default). if loopback_allowed { - return Ok(()); + return Ok(None); } let host = parsed @@ -194,7 +226,7 @@ fn validate_url_for_ssrf_inner( ip )); } - return Ok(()); + return Ok(None); } // Bracketed IPv6 literal. @@ -208,7 +240,7 @@ fn validate_url_for_ssrf_inner( )); } // An IP literal — never DNS-resolved. - return Ok(()); + return Ok(None); } // DECIMAL/integer IPv4 literal (e.g. `http://3232235777` == 192.168.1.1). @@ -225,14 +257,10 @@ fn validate_url_for_ssrf_inner( host, ip )); } - return Ok(()); + return Ok(None); } - // Resolve-then-validate: when the host is a DNS name (not an IP literal), - // resolve it and reject if ANY resolved address is private/internal. This - // closes DNS-rebinding / TOCTOU holes where a public-looking hostname - // resolves to a private/metadata address. - validate_resolved_host_for_ssrf(host) + Ok(Some(host.to_string())) } /// Resolve a DNS hostname and reject if any resolved IP is private/internal. @@ -544,6 +572,36 @@ mod tests { ); } + /// The DNS-free variant makes every check but the resolution: a name that resolves to a + /// private address passes it (the full check refuses it), everything else is refused alike. + #[test] + fn without_dns_skips_only_the_resolution() { + let _guard = env_guard(); + for bad in [ + "ftp://example.com/x", + "https://127.0.0.1/x", + "https://10.0.0.5/x", + "https://[::1]/x", + "https://3232235777/x", + "https://localhost/x", + "https://169.254.169.254/x", + ] { + assert!( + validate_url_for_ssrf_without_dns(bad, HTTP_SCHEMES).is_err(), + "{bad}" + ); + assert!(validate_url_for_ssrf(bad, HTTP_SCHEMES).is_err(), "{bad}"); + } + assert!(validate_url_for_ssrf_without_dns("https://8.8.8.8/x", HTTP_SCHEMES).is_ok()); + // `localhost.` (trailing dot) is not on the blocked list but resolves to loopback. + if validate_resolved_host_for_ssrf("localhost.").is_err() { + assert!( + validate_url_for_ssrf_without_dns("https://localhost./x", HTTP_SCHEMES).is_ok() + ); + assert!(validate_url_for_ssrf("https://localhost./x", HTTP_SCHEMES).is_err()); + } + } + /// Loopback gate OFF (default): private targets rejected via the public, /// env-reading entry point (env var removed under the shared lock). #[test] diff --git a/gateway/src/core/tts/self_hosted.rs b/gateway/src/core/tts/self_hosted.rs index c7d2eb07..c9c65142 100644 --- a/gateway/src/core/tts/self_hosted.rs +++ b/gateway/src/core/tts/self_hosted.rs @@ -220,12 +220,25 @@ pub(crate) fn azure_openai_url_schemes() -> &'static [&'static str] { /// legitimate here. (An Azure Private Link endpoint resolves to a private address and is /// refused; that deployment shape would need an explicit allowance.) /// -/// Resolves DNS synchronously: call it at construction, as the crate's other SSRF checks are, -/// or off the async workers. +/// Resolves DNS synchronously, so run it off the async workers (`spawn_blocking`), as +/// transcription and [`SelfHostedTTS`]'s `connect` do. On a worker it holds that worker for the +/// whole lookup, and every other session on the replica waits with it. pub fn validate_azure_openai_url(url: &str) -> Result<(), String> { crate::core::net::validate_url_for_ssrf(url, azure_openai_url_schemes()) } +/// Everything [`validate_azure_openai_url`] checks that needs no DNS: https, blocked hostnames +/// and IP literals. Safe in a synchronous constructor on the request path. +pub fn validate_azure_openai_url_without_dns(url: &str) -> Result<(), String> { + crate::core::net::validate_url_for_ssrf_without_dns(url, azure_openai_url_schemes()) +} + +fn azure_openai_ssrf_rejection(msg: String) -> crate::core::tts::TTSError { + crate::core::tts::TTSError::InvalidConfiguration(format!( + "azure_openai api_base rejected (SSRF protection): {msg}" + )) +} + /// Attach an Azure OpenAI key as `api-key`, marked sensitive so it never renders in a `Debug`. /// /// A key with bytes a header cannot carry is left for reqwest to refuse as a builder error, @@ -338,7 +351,9 @@ impl SelfHostedTTS { /// [`AZURE_OPENAI_DEFAULT_API_VERSION`]. Everything that can be wrong with the target -- /// no base, no key, no deployment, not https, a private or loopback host -- is refused HERE, /// where the message can name the field, rather than surfacing later as a transport error - /// that reads like Azure is down. + /// that reads like Azure is down. The one exception is a DNS name that RESOLVES to a private + /// address: resolving blocks, and this runs on a tokio worker, so that check waits for + /// `connect`, which runs it off the workers before anything is dialled. pub fn new_azure_openai(config: TTSConfig, api_version: Option<&str>) -> TTSResult { use crate::core::tts::TTSError::InvalidConfiguration; @@ -363,11 +378,7 @@ impl SelfHostedTTS { } let url = azure_openai_audio_url(base, &config.model, AzureAudioRoute::Speech, api_version) .map_err(InvalidConfiguration)?; - validate_azure_openai_url(&url).map_err(|msg| { - InvalidConfiguration(format!( - "azure_openai api_base rejected (SSRF protection): {msg}" - )) - })?; + validate_azure_openai_url_without_dns(&url).map_err(azure_openai_ssrf_rejection)?; Ok(Self::with_upstream(config, Upstream::AzureOpenAi { url })) } @@ -419,6 +430,19 @@ impl BaseTTS for SelfHostedTTS { async fn connect(&mut self) -> TTSResult<()> { let url = self.request_builder.target_url(); + if let Upstream::AzureOpenAi { url: target } = &self.request_builder.upstream { + // The resolve-then-validate half of the SSRF gate, off the async workers + // (the constructor ran the rest). Transcription does the same (`transcribe.rs`). + let target = target.clone(); + tokio::task::spawn_blocking(move || validate_azure_openai_url(&target)) + .await + .map_err(|e| { + crate::core::tts::TTSError::InternalError(format!( + "the Azure OpenAI endpoint check did not complete: {e}" + )) + })? + .map_err(azure_openai_ssrf_rejection)?; + } self.provider .generic_connect_with_config(&url, &self.request_builder.config) .await @@ -842,6 +866,39 @@ mod tests { assert!(SelfHostedTTS::new_azure_openai(azure_config(AZ), None).is_ok()); } + /// A DNS name that resolves to a private address is still refused -- by `connect`, off the + /// async workers -- not by the constructor, which runs on a tokio worker where a blocking + /// `getaddrinfo` stalls every session on the replica (~25 speech req/s per pod at ~70 ms a + /// lookup, measured). `ip6-localhost` is in every Debian/Docker `/etc/hosts`, so no network. + #[test] + fn a_name_resolving_to_a_private_address_is_refused_at_connect_not_construction() { + use std::net::ToSocketAddrs; + let _env = crate::core::net::ssrf_env_lock(); + const HOST: &str = "ip6-localhost"; + let resolves_to_loopback = (HOST, 0u16) + .to_socket_addrs() + .map(|mut a| a.any(|s| s.ip().is_loopback())) + .unwrap_or(false); + if !resolves_to_loopback { + eprintln!("skipped: {HOST} does not resolve to loopback on this host"); + return; + } + + let mut tts = + SelfHostedTTS::new_azure_openai(azure_config(&format!("https://{HOST}")), None) + .expect("the constructor makes no DNS lookup, so a DNS name passes it"); + let err = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(tts.connect()) + .expect_err("connect must refuse a host that resolves to loopback"); + assert!( + err.to_string().contains("SSRF protection"), + "the refusal must name the SSRF guard, got: {err}" + ); + } + #[test] fn an_azure_endpoint_without_base_key_or_deployment_is_refused_by_name() { let refused = |cfg: TTSConfig| match SelfHostedTTS::new_azure_openai(cfg, None) { From 54850c59c7041d5628fe2287312b91a6a89a8515 Mon Sep 17 00:00:00 2001 From: dittops Date: Tue, 29 Sep 2026 12:38:52 +0530 Subject: [PATCH 2/3] perf(gateway): reuse vendor connections for one-shot speech and transcription Every /v1/audio/speech and /v1/audio/transcriptions call built a new HTTP client, so each paid DNS + TCP + TLS to the vendor, and speech providers without a shared per-vendor manager also ran ReqManager::warmup per call: 4 unauthenticated HEADs to the vendor before the real POST (measured at the mock: 10 speech calls -> 10 POST + 40 HEAD). At ~4 ms of CPU per request this made CPU the per-replica ceiling (~250-360 speech req/s on 2 CPUs). * Speech: `synthesize_once_standard` now hands the provider a pooled manager -- the shared per-vendor one when it exists, otherwise one per deployment (`CoreState::deployment_tts_req_manager`, keyed by vendor, api_base and timeouts; credentials stay per request). A provider given a manager does not warm up. Long-lived /ws session providers are unchanged. * Transcription: the Azure OpenAI, self-hosted and prerecorded paths use `core::net::shared_http_client`, one pooled client per configuration. * Shared clients are HTTP/1.1 only (`ReqManagerConfig::http1_only`): many 5-10 s requests on one shared HTTP/2 connection queue behind the vendor's concurrent-stream limit. * Permits: `WAAV_TTS_MAX_CONCURRENT_PER_DEPLOYMENT` (default 4096); the ReqManager bound rises from 1000 to 10000 (`MAX_CONCURRENT_REQUESTS`), since one manager now carries a replica's whole load for its deployment. Verified live: 0 HEADs per speech call, 3 pooled vendor connections and no TIME_WAIT after 20 calls. Tests: cache built once and shared per key, HTTP/1.1, permit ceiling; lib tests for the touched modules (2168) and the api_tests / transcribe_batch / voice_http_spans suites pass. Co-Authored-By: Claude Opus 5.5 --- gateway/src/core/net.rs | 47 ++++++++ gateway/src/core/state.rs | 69 +++++++++++- gateway/src/core/stt/prerecorded.rs | 14 ++- gateway/src/handlers/speak.rs | 9 +- gateway/src/handlers/transcribe.rs | 50 ++++++--- gateway/src/state/mod.rs | 8 ++ gateway/src/utils/req_manager.rs | 163 +++++++++++++++++++++++++--- 7 files changed, 317 insertions(+), 43 deletions(-) diff --git a/gateway/src/core/net.rs b/gateway/src/core/net.rs index 4bbb4f4d..3b32d7c1 100644 --- a/gateway/src/core/net.rs +++ b/gateway/src/core/net.rs @@ -102,6 +102,29 @@ fn loopback_flag_enabled(value: Option<&str>) -> Result { } } +/// A process-wide pooled HTTP client for vendor calls, one per `key`, built by `build` on first +/// use and reused after. +/// +/// A `reqwest::Client` IS its connection pool, so building one per request -- which the audio +/// paths did -- throws the pool away every time: each call paid DNS, TCP and TLS to the vendor. +/// Reusing it keeps keep-alive connections to each vendor host (one pool per host inside the +/// client), which is per-deployment reuse without keying on deployments. Key by whatever changes +/// the client's configuration (schemes, timeout). Build these HTTP/1.1-only: many long requests +/// on one shared HTTP/2 connection queue behind the vendor's concurrent-stream limit. +pub fn shared_http_client( + key: &str, + build: impl FnOnce() -> Result, +) -> Result { + static CLIENTS: std::sync::LazyLock> = + std::sync::LazyLock::new(dashmap::DashMap::new); + if let Some(client) = CLIENTS.get(key) { + return Ok(client.clone()); + } + let built = build()?; + // Two first calls may race to build; both get the one that landed in the map. + Ok(CLIENTS.entry(key.to_string()).or_insert(built).clone()) +} + /// Validate a URL for SSRF (Server-Side Request Forgery) protection. /// /// `allowed_schemes` must be lowercase (the URL's scheme is lowercased before @@ -411,6 +434,30 @@ pub(crate) fn ssrf_env_lock() -> std::sync::MutexGuard<'static, ()> { #[cfg(test)] mod tests { + + #[test] + fn a_shared_client_is_built_once_per_key() { + use std::sync::atomic::{AtomicUsize, Ordering}; + let builds = AtomicUsize::new(0); + let build = || { + builds.fetch_add(1, Ordering::SeqCst); + reqwest::Client::builder().http1_only().build() + }; + let key = "net-tests-shared-client-once"; + shared_http_client(key, build).unwrap(); + shared_http_client(key, build).unwrap(); + assert_eq!( + builds.load(Ordering::SeqCst), + 1, + "the second call reuses the first client" + ); + shared_http_client("net-tests-shared-client-other", build).unwrap(); + assert_eq!( + builds.load(Ordering::SeqCst), + 2, + "another key builds its own" + ); + } use super::*; const HTTP_SCHEMES: &[&str] = &["http", "https"]; diff --git a/gateway/src/core/state.rs b/gateway/src/core/state.rs index 29515d96..8f4e3728 100644 --- a/gateway/src/core/state.rs +++ b/gateway/src/core/state.rs @@ -22,7 +22,9 @@ use crate::core::turn_detect::TurnDetector; #[cfg(feature = "turn-detect")] use crate::core::turn_detect::{TurnDetector, TurnDetectorConfig}; use crate::state::SipHooksState; -use crate::utils::req_manager::ReqManager; +use crate::utils::req_manager::{ + DeploymentReqManagers, MAX_CONCURRENT_REQUESTS, ReqManager, ReqManagerConfig, +}; /// Core-specific shared state for the application. /// @@ -32,6 +34,10 @@ use crate::utils::req_manager::ReqManager; pub struct CoreState { /// HTTP request managers for TTS providers - key is provider name (e.g., "deepgram") pub tts_req_managers: Arc>>>, + /// One pooled request manager per deployment for one-shot synthesis (`/v1/audio/speech`) + /// with a vendor that has no shared per-vendor manager above -- see + /// [`CoreState::deployment_tts_req_manager`]. + pub deployment_tts_req_managers: Arc, /// Unified cache store (in-memory by default) pub cache: Arc, /// Turn detector for determining end of user speech turns @@ -160,6 +166,7 @@ impl CoreState { Ok(Arc::new(Self { tts_req_managers: Arc::new(RwLock::new(tts_req_managers)), + deployment_tts_req_managers: Arc::new(DeploymentReqManagers::new()), cache, turn_detector, sip_hooks_state, @@ -185,6 +192,53 @@ impl CoreState { self.tts_req_managers.read().await.get(provider).cloned() } + /// The pooled manager for one-shot synthesis against this deployment's endpoint, built on + /// first use. `None` only if the manager cannot be built, in which case the provider falls + /// back to building its own as before. + /// + /// Keyed by vendor, endpoint base and the transport timeouts -- what makes two deployments + /// need different connections. The credential is not part of it: each request carries its + /// own. The permit count is [`tts_max_concurrent_per_deployment`], and like the per-vendor + /// knob it bounds transport concurrency on this replica only; vendor-ACCOUNT concurrency is + /// the deployment's `max_concurrent`. + pub async fn deployment_tts_req_manager( + &self, + config: &crate::core::tts::TTSConfig, + ) -> Option> { + let key = format!( + "{}|{}|{:?}|{:?}", + config.provider.trim().to_ascii_lowercase(), + config + .api_base + .as_deref() + .map(str::trim) + .unwrap_or_default(), + config.connection_timeout, + config.request_timeout, + ); + let mut req_config = ReqManagerConfig { + max_concurrent_requests: tts_max_concurrent_per_deployment(), + ..Default::default() + }; + if let Some(secs) = config.connection_timeout { + req_config.connect_timeout = Duration::from_secs(secs); + } + if let Some(secs) = config.request_timeout { + req_config.request_timeout = Duration::from_secs(secs); + } + match self + .deployment_tts_req_managers + .get_or_create(&key, req_config) + .await + { + Ok(manager) => Some(manager), + Err(e) => { + tracing::warn!(error = %e, "could not build the per-deployment TTS request manager"); + None + } + } + } + #[cfg(feature = "turn-detect")] /// Initialize and warmup the Turn Detector model async fn initialize_turn_detector( @@ -348,6 +402,19 @@ fn parse_env_positive_usize(name: &str) -> Result, String> { /// `WAAV_TTS_MAX_CONCURRENT_PER_VENDOR` (default 64, 1–1000): concurrent vendor requests per TTS /// vendor per replica. It was a hard-coded 4 with an unbounded queue behind it. +/// Permits of each per-deployment one-shot TTS manager +/// ([`CoreState::deployment_tts_req_manager`]): `WAAV_TTS_MAX_CONCURRENT_PER_DEPLOYMENT`, +/// 1..=[`MAX_CONCURRENT_REQUESTS`], default 4096. It is shared by every speech request to the +/// deployment on this replica, so it has to hold the replica's whole load for that deployment: +/// at 5 s of vendor time per request, 4096 is ~800 requests a second before the bounded wait. +pub fn tts_max_concurrent_per_deployment() -> usize { + std::env::var("WAAV_TTS_MAX_CONCURRENT_PER_DEPLOYMENT") + .ok() + .and_then(|v| v.trim().parse::().ok()) + .filter(|n| (1..=MAX_CONCURRENT_REQUESTS).contains(n)) + .unwrap_or(4096) +} + pub fn tts_max_concurrent_per_vendor() -> usize { std::env::var("WAAV_TTS_MAX_CONCURRENT_PER_VENDOR") .ok() diff --git a/gateway/src/core/stt/prerecorded.rs b/gateway/src/core/stt/prerecorded.rs index 024eb18a..6fabaa14 100644 --- a/gateway/src/core/stt/prerecorded.rs +++ b/gateway/src/core/stt/prerecorded.rs @@ -172,12 +172,16 @@ type AsyncErrorCallback = Box< + Sync, >; +/// Pooled and shared by every prerecorded upload (`core::net::shared_http_client`): a client per +/// request threw its connections away, so each call paid DNS + TCP + TLS to the vendor. fn http_client() -> Result { - crate::core::net::ssrf_protected_client_builder(crate::core::net::HTTP_URL_SCHEMES) - .timeout(HTTP_TIMEOUT) - .pool_max_idle_per_host(4) - .pool_idle_timeout(Duration::from_secs(90)) - .build() + crate::core::net::shared_http_client("stt-prerecorded", || { + crate::core::net::ssrf_protected_client_builder(crate::core::net::HTTP_URL_SCHEMES) + .timeout(HTTP_TIMEOUT) + .pool_idle_timeout(Duration::from_secs(90)) + .http1_only() + .build() + }) } /// Buffers PCM and transcribes it against a vendor's prerecorded API on close. diff --git a/gateway/src/handlers/speak.rs b/gateway/src/handlers/speak.rs index 49649463..71862fa4 100644 --- a/gateway/src/handlers/speak.rs +++ b/gateway/src/handlers/speak.rs @@ -774,7 +774,14 @@ pub async fn synthesize_once_standard( // Connection pooling and per-provider metrics come from the shared manager; without this // the OpenAI route would open a fresh connection per request while `/speak` reuses them. - if let Some(req_manager) = state.get_tts_req_manager(&tts_config.provider).await { + // A pooled manager for the provider to use instead of building its own per request: the + // shared per-vendor one when the vendor has it, otherwise this deployment's. Without either + // every call paid DNS + TCP + TLS to the vendor plus a warm-up burst of HEADs. + let req_manager = match state.get_tts_req_manager(&tts_config.provider).await { + Some(shared) => Some(shared), + None => state.deployment_tts_req_manager(&tts_config).await, + }; + if let Some(req_manager) = req_manager { if let Some(p) = provider.get_provider() { p.set_req_manager(req_manager.clone()).await; } diff --git a/gateway/src/handlers/transcribe.rs b/gateway/src/handlers/transcribe.rs index 59c07f08..de6467de 100644 --- a/gateway/src/handlers/transcribe.rs +++ b/gateway/src/handlers/transcribe.rs @@ -783,15 +783,20 @@ pub async fn transcribe_self_hosted( transcription_url(api_base) }; - let client = reqwest::Client::builder() - .timeout(OVERALL_DEADLINE) - .build() - .map_err(|e| { - VoiceFailure::new( - VoiceErrorType::Internal, - format!("could not build the http client: {e}"), - ) - })?; + // One pooled client for every self-hosted backend: connections to each host are kept and + // reused instead of a new DNS + TCP + TLS per upload (`core::net::shared_http_client`). + let client = crate::core::net::shared_http_client("stt-self-hosted", || { + reqwest::Client::builder() + .timeout(OVERALL_DEADLINE) + .http1_only() + .build() + }) + .map_err(|e| { + VoiceFailure::new( + VoiceErrorType::Internal, + format!("could not build the http client: {e}"), + ) + })?; let fields = openai_transcription_fields(model, settings); let call = upload_call( @@ -871,15 +876,24 @@ pub(crate) async fn transcribe_azure_openai( )) })?; - let client = crate::core::net::ssrf_protected_client_builder(azure_openai_url_schemes()) - .timeout(OVERALL_DEADLINE) - .build() - .map_err(|e| { - VoiceFailure::new( - VoiceErrorType::Internal, - format!("could not build the http client: {e}"), - ) - })?; + // Pooled and reused across uploads (`core::net::shared_http_client`); keyed by the scheme + // set because the redirect policy is built from it. + let schemes = azure_openai_url_schemes(); + let client = crate::core::net::shared_http_client( + &format!("stt-azure-openai|{}", schemes.join(",")), + || { + crate::core::net::ssrf_protected_client_builder(schemes) + .timeout(OVERALL_DEADLINE) + .http1_only() + .build() + }, + ) + .map_err(|e| { + VoiceFailure::new( + VoiceErrorType::Internal, + format!("could not build the http client: {e}"), + ) + })?; let model = deployment.trim(); let fields = openai_transcription_fields(model, settings); diff --git a/gateway/src/state/mod.rs b/gateway/src/state/mod.rs index 2b06ca5f..639bc458 100644 --- a/gateway/src/state/mod.rs +++ b/gateway/src/state/mod.rs @@ -473,6 +473,14 @@ impl AppState { self.core_state.get_tts_req_manager(provider).await } + /// See [`crate::core::state::CoreState::deployment_tts_req_manager`]. + pub async fn deployment_tts_req_manager( + &self, + config: &crate::core::tts::TTSConfig, + ) -> Option> { + self.core_state.deployment_tts_req_manager(config).await + } + /// Get a handle to the application's cache store pub fn cache(&self) -> Arc { self.core_state.cache.clone() diff --git a/gateway/src/utils/req_manager.rs b/gateway/src/utils/req_manager.rs index bbfbe19e..8f39dada 100644 --- a/gateway/src/utils/req_manager.rs +++ b/gateway/src/utils/req_manager.rs @@ -283,8 +283,18 @@ pub struct ReqManagerConfig { /// Longest a caller waits for a free slot before [`Saturated`] (FRD-022 §6.6). `None` waits /// forever — the unbounded queue this replaced. pub acquire_timeout: Option, + /// Speak HTTP/1.1 only: every in-flight request gets its own keep-alive connection from the + /// pool. For a manager SHARED by many concurrent long requests (one per deployment), where + /// HTTP/2 would put them all on one connection and queue whatever exceeds the vendor's + /// `SETTINGS_MAX_CONCURRENT_STREAMS` -- a hidden ceiling at 5-10 s per request. + pub http1_only: bool, } +/// Upper bound on [`ReqManagerConfig::max_concurrent_requests`]. A per-deployment manager is +/// shared by every one-shot request to that deployment on the replica, so at 5-10 s of vendor +/// time per request a four-digit bound is ordinary load, not a runaway. +pub const MAX_CONCURRENT_REQUESTS: usize = 10_000; + /// The pool stayed full for the whole bounded wait. Local back-pressure: the caller is told /// 503 + `Retry-After: 1` rather than joining an unbounded queue. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -329,6 +339,7 @@ impl Default for ReqManagerConfig { retry_max_delay: Duration::from_millis(500), per_request_timeout: Duration::from_secs(30), // 30s per request for TTS acquire_timeout: Some(Duration::from_secs(2)), + http1_only: false, } } } @@ -351,6 +362,7 @@ impl ReqManagerConfig { retry_max_delay: Duration::from_millis(300), per_request_timeout: Duration::from_secs(2), acquire_timeout: Some(Duration::from_secs(2)), + http1_only: false, } } @@ -371,6 +383,7 @@ impl ReqManagerConfig { retry_max_delay: Duration::from_secs(1), per_request_timeout: Duration::from_secs(5), acquire_timeout: Some(Duration::from_secs(2)), + http1_only: false, } } } @@ -379,7 +392,8 @@ impl ReqManager { /// Create a new request manager with the specified maximum concurrent requests /// /// # Arguments - /// * `max_concurrent_requests` - Maximum number of concurrent requests allowed (1-1000) + /// * `max_concurrent_requests` - Maximum number of concurrent requests allowed + /// (1..=[`MAX_CONCURRENT_REQUESTS`]) /// /// # Returns /// A new `ReqManager` instance with default configuration @@ -418,8 +432,11 @@ impl ReqManager { if config.max_concurrent_requests == 0 { return Err("max_concurrent_requests must be greater than 0".into()); } - if config.max_concurrent_requests > 1000 { - return Err("max_concurrent_requests must not exceed 1000".into()); + if config.max_concurrent_requests > MAX_CONCURRENT_REQUESTS { + return Err(format!( + "max_concurrent_requests must not exceed {MAX_CONCURRENT_REQUESTS}" + ) + .into()); } if config.max_retries == u32::MAX { return Err("max_retries must be less than u32::MAX".into()); @@ -438,21 +455,26 @@ impl ReqManager { /// Create an optimized HTTP/2 client with advanced connection pooling fn create_optimized_client(config: &ReqManagerConfig) -> Result { - crate::core::net::ssrf_protected_client_builder(crate::core::net::HTTP_URL_SCHEMES) - .http2_initial_stream_window_size(config.http2_stream_window_size) - .http2_initial_connection_window_size(config.http2_connection_window_size) - .http2_keep_alive_interval(Some(config.http2_keep_alive_interval)) - .http2_keep_alive_timeout(config.http2_keep_alive_timeout) - .http2_keep_alive_while_idle(true) - .http2_adaptive_window(true) - .pool_idle_timeout(None) - .pool_max_idle_per_host(config.pool_max_idle_per_host) - .tcp_keepalive(config.tcp_keepalive) - .tcp_nodelay(true) - .connect_timeout(config.connect_timeout) - .timeout(config.request_timeout) - .user_agent("waav-gateway-req-manager/2.0") - .build() + let builder = + crate::core::net::ssrf_protected_client_builder(crate::core::net::HTTP_URL_SCHEMES) + .http2_initial_stream_window_size(config.http2_stream_window_size) + .http2_initial_connection_window_size(config.http2_connection_window_size) + .http2_keep_alive_interval(Some(config.http2_keep_alive_interval)) + .http2_keep_alive_timeout(config.http2_keep_alive_timeout) + .http2_keep_alive_while_idle(true) + .http2_adaptive_window(true) + .pool_idle_timeout(None) + .pool_max_idle_per_host(config.pool_max_idle_per_host) + .tcp_keepalive(config.tcp_keepalive) + .tcp_nodelay(true) + .connect_timeout(config.connect_timeout) + .timeout(config.request_timeout) + .user_agent("waav-gateway-req-manager/2.0"); + if config.http1_only { + builder.http1_only().build() + } else { + builder.build() + } } /// Acquire a client from the pool with automatic metrics tracking @@ -823,8 +845,113 @@ impl ReqManager { } } +/// One [`ReqManager`] per deployment, shared by every one-shot request to it on this replica. +/// +/// A one-shot synthesis builds its provider per request, and a provider with no manager builds +/// its own: a new HTTP client (DNS, TCP and TLS to the vendor on every call) plus a warm-up that +/// sends the vendor a burst of unauthenticated `HEAD`s before the real request -- measured at 4 +/// per speech call, and ~4 ms of CPU per request, which made CPU the per-pod ceiling. Handing the +/// provider a cached manager skips both: connections are pooled and reused across requests, and +/// `generic_connect_with_config` does not warm up a manager it was given. +/// +/// Keyed by whatever makes two deployments need different transports (the caller's key: vendor, +/// endpoint base, timeouts). Credentials are NOT part of the transport -- every request carries +/// its own -- so two deployments on the same endpoint may share a pool. +#[derive(Default)] +pub struct DeploymentReqManagers { + managers: dashmap::DashMap>, +} + +impl DeploymentReqManagers { + pub fn new() -> Self { + Self::default() + } + + /// The manager for `key`, built from `config` the first time. HTTP/1.1 only (see + /// [`ReqManagerConfig::http1_only`]), with an idle pool as deep as the permit count so a + /// steady load keeps its connections. + pub async fn get_or_create( + &self, + key: &str, + config: ReqManagerConfig, + ) -> Result, Box> { + if let Some(existing) = self.managers.get(key) { + return Ok(existing.clone()); + } + let config = ReqManagerConfig { + http1_only: true, + pool_max_idle_per_host: config.max_concurrent_requests, + ..config + }; + let built = Arc::new(ReqManager::with_config(config).await?); + // Two first requests may race to build; both get the one that landed in the map. + Ok(self + .managers + .entry(key.to_string()) + .or_insert(built) + .clone()) + } + + /// Deployments with a live manager (for tests and diagnostics). + pub fn len(&self) -> usize { + self.managers.len() + } + + pub fn is_empty(&self) -> bool { + self.managers.is_empty() + } +} + #[cfg(test)] mod tests { + + #[tokio::test] + async fn a_deployment_manager_is_built_once_and_shared() { + let cache = DeploymentReqManagers::new(); + let cfg = || ReqManagerConfig { + max_concurrent_requests: 2048, + ..Default::default() + }; + let a = cache + .get_or_create("azure_openai|https://a", cfg()) + .await + .unwrap(); + let again = cache + .get_or_create("azure_openai|https://a", cfg()) + .await + .unwrap(); + let b = cache + .get_or_create("azure_openai|https://b", cfg()) + .await + .unwrap(); + assert!( + Arc::ptr_eq(&a, &again), + "the same deployment reuses its manager" + ); + assert!(!Arc::ptr_eq(&a, &b), "another endpoint gets its own"); + assert_eq!(cache.len(), 2); + assert_eq!(a.max_concurrent_requests, 2048, "the permit count survives"); + assert!(a.config.http1_only, "a shared manager is HTTP/1.1 only"); + assert_eq!(a.config.pool_max_idle_per_host, 2048); + } + + #[tokio::test] + async fn the_permit_ceiling_is_max_concurrent_requests() { + let at = |n| ReqManagerConfig { + max_concurrent_requests: n, + ..Default::default() + }; + assert!( + ReqManager::with_config(at(MAX_CONCURRENT_REQUESTS)) + .await + .is_ok() + ); + assert!( + ReqManager::with_config(at(MAX_CONCURRENT_REQUESTS + 1)) + .await + .is_err() + ); + } use super::*; use std::io::ErrorKind; use std::sync::atomic::{AtomicUsize, Ordering}; From dbeb9acdf27747d79ba7ffcf2839914e91e8f310 Mon Sep 17 00:00:00 2001 From: dittops Date: Tue, 29 Sep 2026 15:49:29 +0530 Subject: [PATCH 3/3] fix(gateway): build without turn-detect -- qualify Duration in the per-deployment TTS manager `core/state.rs` imports `std::time::Duration` only under `turn-detect`, so the default, dag-routing and noise-filter builds failed on the timeouts `deployment_tts_req_manager` copies from TTSConfig (E0433). Fully qualified. Checked: `cargo check --lib --bins --tests` with no features, dag-routing, noise-filter and dag-routing,turn-ensemble,noise-filter,openapi; changed-module lib tests with default features (117 passed). Co-Authored-By: Claude Opus 5.5 --- gateway/src/core/state.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/gateway/src/core/state.rs b/gateway/src/core/state.rs index 8f4e3728..54682645 100644 --- a/gateway/src/core/state.rs +++ b/gateway/src/core/state.rs @@ -221,10 +221,10 @@ impl CoreState { ..Default::default() }; if let Some(secs) = config.connection_timeout { - req_config.connect_timeout = Duration::from_secs(secs); + req_config.connect_timeout = std::time::Duration::from_secs(secs); } if let Some(secs) = config.request_timeout { - req_config.request_timeout = Duration::from_secs(secs); + req_config.request_timeout = std::time::Duration::from_secs(secs); } match self .deployment_tts_req_managers