From 748e6e74a36a05fc277e990e422e16d47358f4ad Mon Sep 17 00:00:00 2001 From: Filip Kujawa Date: Mon, 27 Jul 2026 10:43:54 -0700 Subject: [PATCH 1/6] fix(providers): stop killing streaming responses at the total request timeout reqwest's client-level timeout is a total-request deadline that includes the streamed body, so any model turn longer than 600s died mid-stream with a decode error. Apply the total deadline per request and exempt requests marked .streaming(true); streams are bounded by read_timeout (reset per chunk) and connect_timeout instead. Providers with hand-rolled clients (GCP Vertex, Kimi, Gemini OAuth, device flow, ChatGPT Codex, Bedrock Mantle) get the same connect/read bounds, plus per-request total deadlines on their non-streaming calls. Error-body reads on paths that serve streaming are bounded separately since those requests no longer carry a total deadline. --- crates/goose-provider-types/src/errors.rs | 8 + crates/goose-providers/Cargo.toml | 2 +- crates/goose-providers/src/anthropic.rs | 1 + crates/goose-providers/src/api_client.rs | 198 +++++++++++++++++- crates/goose-providers/src/databricks.rs | 11 +- crates/goose-providers/src/databricks_v2.rs | 3 + crates/goose-providers/src/google.rs | 1 + crates/goose-providers/src/http_status.rs | 13 +- crates/goose-providers/src/ollama.rs | 1 + crates/goose-providers/src/openai.rs | 2 + .../goose-providers/src/openai_compatible.rs | 1 + crates/goose-providers/src/snowflake.rs | 17 +- crates/goose/src/providers/anthropic_def.rs | 15 +- crates/goose/src/providers/base.rs | 1 + crates/goose/src/providers/bedrock.rs | 27 ++- crates/goose/src/providers/chatgpt_codex.rs | 31 ++- crates/goose/src/providers/gcpvertexai.rs | 30 ++- crates/goose/src/providers/gemini_oauth.rs | 32 ++- crates/goose/src/providers/githubcopilot.rs | 5 + crates/goose/src/providers/kimicode.rs | 24 ++- crates/goose/src/providers/nanogpt.rs | 1 + .../goose/src/providers/oauth_device_flow.rs | 7 +- crates/goose/src/providers/openrouter.rs | 1 + crates/goose/src/providers/tetrate.rs | 20 +- 24 files changed, 389 insertions(+), 63 deletions(-) diff --git a/crates/goose-provider-types/src/errors.rs b/crates/goose-provider-types/src/errors.rs index ed38de7e2542..f2a50562aa3c 100644 --- a/crates/goose-provider-types/src/errors.rs +++ b/crates/goose-provider-types/src/errors.rs @@ -134,6 +134,14 @@ impl From for ProviderError { if let Some(reqwest_err) = error.downcast_ref::() { return provider_error_from_reqwest(reqwest_err); } + if error + .downcast_ref::() + .is_some() + { + return ProviderError::NetworkError( + "Request timed out — check your network connection and try again.".to_string(), + ); + } ProviderError::ExecutionError(error.to_string()) } } diff --git a/crates/goose-providers/Cargo.toml b/crates/goose-providers/Cargo.toml index e17e0322f7ce..ab231235b947 100644 --- a/crates/goose-providers/Cargo.toml +++ b/crates/goose-providers/Cargo.toml @@ -58,7 +58,7 @@ include_dir = { workspace = true } [dev-dependencies] test-case = { workspace = true } tempfile = { workspace = true } -tokio = { workspace = true, features = ["rt-multi-thread"] } +tokio = { workspace = true, features = ["io-util", "macros", "net", "rt-multi-thread", "time"] } tokio-stream = { workspace = true } env-lock = { workspace = true } wiremock.workspace = true diff --git a/crates/goose-providers/src/anthropic.rs b/crates/goose-providers/src/anthropic.rs index 09391acdf4f7..8ec61519f4a3 100644 --- a/crates/goose-providers/src/anthropic.rs +++ b/crates/goose-providers/src/anthropic.rs @@ -188,6 +188,7 @@ impl AnthropicProvider { self.api_client .request("v1/messages") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?, ) diff --git a/crates/goose-providers/src/api_client.rs b/crates/goose-providers/src/api_client.rs index b41b35594330..f3a1dc1f637a 100644 --- a/crates/goose-providers/src/api_client.rs +++ b/crates/goose-providers/src/api_client.rs @@ -14,7 +14,8 @@ use std::path::PathBuf; use std::sync::Arc; use std::time::Duration; -const DEFAULT_PROVIDER_TIMEOUT_SECS: u64 = 600; +pub const DEFAULT_PROVIDER_TIMEOUT_SECS: u64 = 600; +pub const DEFAULT_CONNECT_TIMEOUT_SECS: u64 = 30; pub type RequestBuilderDecorator = Arc Result + Send + Sync>; @@ -233,6 +234,7 @@ pub struct ApiRequestBuilder<'a> { client: &'a ApiClient, path: &'a str, headers: HeaderMap, + streaming: bool, } impl ApiClient { @@ -255,7 +257,7 @@ impl ApiClient { timeout: Duration, tls_config: Option, ) -> Result { - let mut client_builder = Client::builder().timeout(timeout); + let mut client_builder = Self::client_builder(timeout); if let Some(ref config) = tls_config { client_builder = Self::configure_tls(client_builder, config)?; @@ -283,10 +285,19 @@ impl ApiClient { self.timeout } + /// The total-request deadline is applied per request in `send_request` so + /// streaming requests can opt out of it: an SSE body is the whole + /// generation, so streams are bounded by `read_timeout` (reset per chunk) + /// rather than end-to-end duration. + fn client_builder(timeout: Duration) -> reqwest::ClientBuilder { + Client::builder() + .connect_timeout(Duration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) + .read_timeout(timeout) + } + fn rebuild_client(&mut self) -> Result<()> { - let mut client_builder = Client::builder() - .timeout(self.timeout) - .default_headers(self.default_headers.clone()); + let mut client_builder = + Self::client_builder(self.timeout).default_headers(self.default_headers.clone()); // Configure TLS if needed if let Some(ref tls_config) = self.tls_config { @@ -361,6 +372,7 @@ impl ApiClient { client: self, path, headers: HeaderMap::new(), + streaming: false, } } @@ -434,14 +446,33 @@ impl<'a> ApiRequestBuilder<'a> { } } + /// Exempts this request from the total-request deadline; the stream is + /// bounded by the client's read timeout instead. + pub fn streaming(mut self, streaming: bool) -> Self { + self.streaming = streaming; + self + } + pub async fn api_post(self, payload: &Value) -> Result { let response = self.response_post(payload).await?; ApiResponse::from_response(response).await } + /// `send()` resolves when response headers arrive, so for streaming + /// requests (exempt from the per-request total deadline) this still bounds + /// connect, request upload, and time-to-first-byte; only the streamed + /// response body is unbounded. + async fn send_bounded(&self, request: reqwest::RequestBuilder) -> Result { + if self.streaming { + Ok(tokio::time::timeout(self.client.timeout, request.send()).await??) + } else { + Ok(request.send().await?) + } + } + pub async fn response_post(self, payload: &Value) -> Result { let request = self.send_request(|url, client| client.post(url)).await?; - Ok(request.json(payload).send().await?) + self.send_bounded(request.json(payload)).await } pub async fn multipart_post(self, form: reqwest::multipart::Form) -> Result { @@ -456,7 +487,7 @@ impl<'a> ApiRequestBuilder<'a> { pub async fn response_get(self) -> Result { let request = self.send_request(|url, client| client.get(url)).await?; - Ok(request.send().await?) + self.send_bounded(request).await } async fn send_request(&self, request_builder: F) -> Result @@ -468,6 +499,10 @@ impl<'a> ApiRequestBuilder<'a> { let mut request = request_builder(url, &self.client.client); request = request.headers(headers); + if !self.streaming { + request = request.timeout(self.client.timeout); + } + if let Some(decorator) = &self.client.request_builder { request = decorator(request)?; } @@ -623,6 +658,155 @@ ShGoCNbfNS+COlPMRAujyDlATZcLs9p4tA== #[cfg(test)] mod tests { use super::*; + use std::net::SocketAddr; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + async fn spawn_chunked_server(gap_ms: u64, chunks: usize) -> SocketAddr { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + loop { + let Ok((mut sock, _)) = listener.accept().await else { + break; + }; + tokio::spawn(async move { + let mut buf = [0u8; 8192]; + let _ = sock.read(&mut buf).await; + if sock + .write_all( + b"HTTP/1.1 200 OK\r\n\ + content-type: text/event-stream\r\n\ + transfer-encoding: chunked\r\n\r\n", + ) + .await + .is_err() + { + return; + } + for i in 0..chunks { + if i > 0 { + tokio::time::sleep(Duration::from_millis(gap_ms)).await; + } + let data = format!("data: {}\n\n", i); + let chunk = format!("{:x}\r\n{}\r\n", data.len(), data); + if sock.write_all(chunk.as_bytes()).await.is_err() { + return; + } + let _ = sock.flush().await; + } + let _ = sock.write_all(b"0\r\n\r\n").await; + }); + } + }); + addr + } + + fn client_with_timeout(addr: SocketAddr, timeout_ms: u64) -> ApiClient { + ApiClient::with_timeout_and_tls( + format!("http://{}", addr), + AuthMethod::NoAuth, + Duration::from_millis(timeout_ms), + None, + ) + .unwrap() + } + + async fn drain_counting_data_lines(mut response: Response) -> Result { + let mut count = 0; + while let Some(chunk) = response.chunk().await? { + count += String::from_utf8_lossy(&chunk).matches("data:").count(); + } + Ok(count) + } + + #[tokio::test] + async fn streaming_request_survives_beyond_total_timeout() { + let addr = spawn_chunked_server(50, 12).await; + let client = client_with_timeout(addr, 400); + + let response = client + .request("v1/messages") + .streaming(true) + .response_post(&serde_json::json!({})) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + + let count = drain_counting_data_lines(response).await.unwrap(); + assert_eq!(count, 12); + } + + #[tokio::test] + async fn streaming_request_fails_when_stream_stalls() { + let addr = spawn_chunked_server(5_000, 2).await; + let client = client_with_timeout(addr, 400); + + let response = client + .request("v1/messages") + .streaming(true) + .response_post(&serde_json::json!({})) + .await + .unwrap(); + + let err = drain_counting_data_lines(response) + .await + .expect_err("stalled stream should time out, not complete"); + assert!(err.is_timeout(), "expected a timeout error, got: {err}"); + } + + #[tokio::test] + async fn non_streaming_request_enforces_total_deadline() { + let addr = spawn_chunked_server(50, 12).await; + let client = client_with_timeout(addr, 400); + + let response = client + .request("v1/messages") + .response_post(&serde_json::json!({})) + .await + .unwrap(); + + let err = drain_counting_data_lines(response) + .await + .expect_err("total deadline should cut off the response body"); + assert!(err.is_timeout(), "expected a timeout error, got: {err}"); + } + + #[tokio::test] + async fn streaming_request_times_out_before_response_headers() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + loop { + let Ok((mut sock, _)) = listener.accept().await else { + break; + }; + tokio::spawn(async move { + let mut buf = [0u8; 8192]; + while sock.read(&mut buf).await.is_ok_and(|n| n > 0) {} + }); + } + }); + let client = client_with_timeout(addr, 400); + + let started = std::time::Instant::now(); + let err = client + .request("v1/messages") + .streaming(true) + .response_post(&serde_json::json!({})) + .await + .expect_err("the phase before the response body must stay bounded"); + assert!( + started.elapsed() < Duration::from_secs(5), + "should fail near the configured timeout, took {:?}", + started.elapsed() + ); + let timed_out = err + .downcast_ref::() + .is_some_and(reqwest::Error::is_timeout) + || err.downcast_ref::().is_some(); + assert!(timed_out, "expected a timeout error, got: {err}"); + } #[test] fn test_model_headers_applied_and_override_static_headers() { diff --git a/crates/goose-providers/src/databricks.rs b/crates/goose-providers/src/databricks.rs index ad95803b49ce..eb8a4380aa56 100644 --- a/crates/goose-providers/src/databricks.rs +++ b/crates/goose-providers/src/databricks.rs @@ -601,6 +601,7 @@ impl Provider for DatabricksProvider { .api_client .request(&path) .model_headers(model_config)? + .streaming(true) .response_post(&payload_clone) .await?; handle_status(resp).await @@ -666,12 +667,15 @@ impl Provider for DatabricksProvider { .api_client .request(&path) .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; if !resp.status().is_success() { let status = resp.status(); let url = sanitize_url(resp.url().as_str()); - let error_text = resp.text().await.unwrap_or_default(); + let error_text = crate::http_status::read_error_body(resp) + .await + .unwrap_or_default(); let json_payload = serde_json::from_str::(&error_text).ok(); return Err(map_http_error_to_provider_error(status, json_payload, &url)); @@ -688,12 +692,15 @@ impl Provider for DatabricksProvider { .api_client .request(&path) .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; if !resp.status().is_success() { let status = resp.status(); let url = sanitize_url(resp.url().as_str()); - let error_text = resp.text().await.unwrap_or_default(); + let error_text = crate::http_status::read_error_body(resp) + .await + .unwrap_or_default(); let json_payload = serde_json::from_str::(&error_text).ok(); return Err(map_http_error_to_provider_error( status, diff --git a/crates/goose-providers/src/databricks_v2.rs b/crates/goose-providers/src/databricks_v2.rs index 8a2c1f8c431f..0557676527a8 100644 --- a/crates/goose-providers/src/databricks_v2.rs +++ b/crates/goose-providers/src/databricks_v2.rs @@ -207,6 +207,7 @@ impl DatabricksV2Provider { .api_client .request("ai-gateway/openai/v1/responses") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; handle_status(resp).await @@ -245,6 +246,7 @@ impl DatabricksV2Provider { .api_client .request("ai-gateway/mlflow/v1/chat/completions") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; handle_status(resp).await @@ -281,6 +283,7 @@ impl DatabricksV2Provider { .api_client .request("ai-gateway/anthropic/v1/messages") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; handle_status(resp).await diff --git a/crates/goose-providers/src/google.rs b/crates/goose-providers/src/google.rs index 16599fd122b5..bc9e7eab7e1e 100644 --- a/crates/goose-providers/src/google.rs +++ b/crates/goose-providers/src/google.rs @@ -107,6 +107,7 @@ impl GoogleProvider { .api_client .request(&path) .model_headers(model_config)? + .streaming(true) .response_post(payload) .await?; handle_status(response).await diff --git a/crates/goose-providers/src/http_status.rs b/crates/goose-providers/src/http_status.rs index 13147ccadd4a..52fed17d927b 100644 --- a/crates/goose-providers/src/http_status.rs +++ b/crates/goose-providers/src/http_status.rs @@ -248,12 +248,23 @@ pub fn map_http_error_to_provider_error( error } +/// Streaming requests carry no total-request deadline, so error-body reads +/// must be bounded or a server that drips a non-2xx body stalls retries forever. +const ERROR_BODY_READ_TIMEOUT: Duration = Duration::from_secs(30); + +pub async fn read_error_body(response: Response) -> Option { + tokio::time::timeout(ERROR_BODY_READ_TIMEOUT, response.text()) + .await + .ok() + .and_then(Result::ok) +} + pub async fn handle_status(response: Response) -> Result { let status = response.status(); if !status.is_success() { let url = sanitize_url(response.url().as_str()); let headers = response.headers().clone(); - let body = response.text().await.unwrap_or_default(); + let body = read_error_body(response).await.unwrap_or_default(); let payload = serde_json::from_str::(&body).ok(); let mut err = map_http_error_to_provider_error(status, payload.clone(), &url); if let ProviderError::RateLimitExceeded { details, .. } = &err { diff --git a/crates/goose-providers/src/ollama.rs b/crates/goose-providers/src/ollama.rs index 2a4cdbab8486..c1de5a336370 100644 --- a/crates/goose-providers/src/ollama.rs +++ b/crates/goose-providers/src/ollama.rs @@ -428,6 +428,7 @@ impl Provider for OllamaProvider { .api_client .request("v1/chat/completions") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; handle_status(resp).await diff --git a/crates/goose-providers/src/openai.rs b/crates/goose-providers/src/openai.rs index 15a5c983295f..b256ca989da8 100644 --- a/crates/goose-providers/src/openai.rs +++ b/crates/goose-providers/src/openai.rs @@ -311,6 +311,7 @@ impl OpenAiProvider { OPEN_AI_DEFAULT_RESPONSES_PATH, )) .model_headers(model_config)? + .streaming(self.supports_streaming) .response_post(&payload) .await?, ) @@ -797,6 +798,7 @@ impl Provider for OpenAiProvider { .api_client .request(&self.base_path) .model_headers(model_config)? + .streaming(self.supports_streaming) .response_post(&payload) .await?; handle_status(resp).await diff --git a/crates/goose-providers/src/openai_compatible.rs b/crates/goose-providers/src/openai_compatible.rs index d7f752bde616..f40d756da26c 100644 --- a/crates/goose-providers/src/openai_compatible.rs +++ b/crates/goose-providers/src/openai_compatible.rs @@ -112,6 +112,7 @@ impl OpenAiCompatibleProvider { self.api_client .request(&path) .model_headers(model_config)? + .streaming(self.supports_streaming) .response_post(&payload) .await?, ) diff --git a/crates/goose-providers/src/snowflake.rs b/crates/goose-providers/src/snowflake.rs index d30cb26d9433..42d17b109f0d 100644 --- a/crates/goose-providers/src/snowflake.rs +++ b/crates/goose-providers/src/snowflake.rs @@ -100,12 +100,27 @@ impl SnowflakeProvider { .api_client .request("api/v2/cortex/inference:complete") .model_headers(model_config)? + .streaming(true) .response_post(payload) .await?; let status = response.status(); let url = sanitize_url(response.url().as_str()); - let payload_text: String = response.text().await.ok().unwrap_or_default(); + let is_json = response + .headers() + .get(reqwest::header::CONTENT_TYPE) + .and_then(|v| v.to_str().ok()) + .map(|v| v.to_ascii_lowercase()) + .is_some_and(|v| v.contains("json")); + // A 200 with a JSON body is an error payload, not the SSE generation + // stream; bound its read since this request has no total deadline. + let payload_text: String = if status.is_success() && !is_json { + response.text().await.ok().unwrap_or_default() + } else { + crate::http_status::read_error_body(response) + .await + .unwrap_or_default() + }; if status.is_success() { if let Ok(payload) = serde_json::from_str::(&payload_text) { diff --git a/crates/goose/src/providers/anthropic_def.rs b/crates/goose/src/providers/anthropic_def.rs index 6388511a49bc..2efffb277374 100644 --- a/crates/goose/src/providers/anthropic_def.rs +++ b/crates/goose/src/providers/anthropic_def.rs @@ -39,14 +39,23 @@ async fn from_env( .get_param("ANTHROPIC_HOST") .unwrap_or_else(|_| "https://api.anthropic.com".to_string()); + let timeout_secs: u64 = config + .get_param("ANTHROPIC_TIMEOUT") + .unwrap_or(crate::providers::base::DEFAULT_PROVIDER_TIMEOUT_SECS); + let auth = AuthMethod::ApiKey { header_name: "x-api-key".to_string(), key: api_key, }; - let api_client = ApiClient::new_with_tls(host, auth, tls_config)? - .with_request_builder(crate::session_context::session_id_request_builder()) - .with_header("anthropic-version", ANTHROPIC_API_VERSION)?; + let api_client = ApiClient::with_timeout_and_tls( + host, + auth, + std::time::Duration::from_secs(timeout_secs), + tls_config, + )? + .with_request_builder(crate::session_context::session_id_request_builder()) + .with_header("anthropic-version", ANTHROPIC_API_VERSION)?; Ok(AnthropicProviderBuilder::new(api_client).build()) } diff --git a/crates/goose/src/providers/base.rs b/crates/goose/src/providers/base.rs index 8baee4838874..c1807a6eab3a 100644 --- a/crates/goose/src/providers/base.rs +++ b/crates/goose/src/providers/base.rs @@ -7,6 +7,7 @@ pub use goose_providers::conversation::token_usage::{ use serde::{Deserialize, Serialize}; pub const DEFAULT_PROVIDER_TIMEOUT_SECS: u64 = 600; +pub use goose_providers::api_client::DEFAULT_CONNECT_TIMEOUT_SECS; use crate::config::ExtensionConfig; diff --git a/crates/goose/src/providers/bedrock.rs b/crates/goose/src/providers/bedrock.rs index d6fee4ceceee..bef936f7ad5d 100644 --- a/crates/goose/src/providers/bedrock.rs +++ b/crates/goose/src/providers/bedrock.rs @@ -1,6 +1,9 @@ use std::collections::HashMap; -use super::base::{ConfigKey, MessageStream, Provider, ProviderDef, ProviderMetadata}; +use super::base::{ + ConfigKey, MessageStream, Provider, ProviderDef, ProviderMetadata, + DEFAULT_CONNECT_TIMEOUT_SECS, DEFAULT_PROVIDER_TIMEOUT_SECS, +}; use super::openai_compatible::{handle_status, stream_responses_compat}; use super::retry::{ProviderRetry, RetryConfig}; use crate::conversation::message::Message; @@ -177,7 +180,12 @@ impl BedrockProvider { name: BEDROCK_PROVIDER_NAME.to_string(), region: resolved_region, bearer_token, - http_client: reqwest::Client::new(), + http_client: reqwest::Client::builder() + .connect_timeout(std::time::Duration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) + .read_timeout(std::time::Duration::from_secs( + DEFAULT_PROVIDER_TIMEOUT_SECS, + )) + .build()?, mantle_base_url: None, }) } @@ -270,10 +278,17 @@ impl BedrockProvider { } } - let response = req - .send() - .await - .map_err(|e| ProviderError::RequestFailed(format!("Mantle request failed: {}", e)))?; + let response = tokio::time::timeout( + std::time::Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), + req.send(), + ) + .await + .map_err(|_| { + ProviderError::NetworkError( + "Request timed out — check your network connection and try again.".to_string(), + ) + })? + .map_err(|e| ProviderError::RequestFailed(format!("Mantle request failed: {}", e)))?; handle_status(response).await } diff --git a/crates/goose/src/providers/chatgpt_codex.rs b/crates/goose/src/providers/chatgpt_codex.rs index d683a9eef40b..6c166860e0e7 100644 --- a/crates/goose/src/providers/chatgpt_codex.rs +++ b/crates/goose/src/providers/chatgpt_codex.rs @@ -1,7 +1,10 @@ use crate::config::paths::Paths; use crate::conversation::message::{Message, MessageContent}; use crate::providers::api_client::{AuthProvider, RequestBuilderDecorator}; -use crate::providers::base::{ConfigKey, MessageStream, Provider, ProviderDef, ProviderMetadata}; +use crate::providers::base::{ + ConfigKey, MessageStream, Provider, ProviderDef, ProviderMetadata, + DEFAULT_CONNECT_TIMEOUT_SECS, DEFAULT_PROVIDER_TIMEOUT_SECS, +}; use crate::providers::openai_compatible::handle_status; use crate::providers::private_file::write_private_file; use crate::providers::retry::ProviderRetry; @@ -921,7 +924,13 @@ impl ChatGptCodexProvider { ); } - let client = reqwest::Client::new(); + let client = reqwest::Client::builder() + .connect_timeout(std::time::Duration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) + .read_timeout(std::time::Duration::from_secs( + DEFAULT_PROVIDER_TIMEOUT_SECS, + )) + .build() + .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; let request = client .post(format!("{}/responses", CODEX_API_ENDPOINT)) .header( @@ -932,11 +941,19 @@ impl ChatGptCodexProvider { .headers(headers) .json(payload); - let response = (self.request_builder)(request) - .map_err(|e| ProviderError::ExecutionError(e.to_string()))? - .send() - .await - .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; + let request = (self.request_builder)(request) + .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; + let response = tokio::time::timeout( + std::time::Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), + request.send(), + ) + .await + .map_err(|_| { + ProviderError::NetworkError( + "Request timed out — check your network connection and try again.".to_string(), + ) + })? + .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; handle_status(response).await } diff --git a/crates/goose/src/providers/gcpvertexai.rs b/crates/goose/src/providers/gcpvertexai.rs index 6a2ea2ebc47f..545a4a5bb223 100644 --- a/crates/goose/src/providers/gcpvertexai.rs +++ b/crates/goose/src/providers/gcpvertexai.rs @@ -18,7 +18,7 @@ use crate::conversation::message::Message; use crate::providers::api_client::RequestBuilderDecorator; use crate::providers::base::{ ConfigKey, MessageStream, Provider, ProviderDef, ProviderMetadata, - DEFAULT_PROVIDER_TIMEOUT_SECS, + DEFAULT_CONNECT_TIMEOUT_SECS, DEFAULT_PROVIDER_TIMEOUT_SECS, }; use goose_providers::model::ModelConfig; @@ -30,6 +30,7 @@ use crate::providers::gcpauth::GcpAuth; use crate::providers::openai_compatible::{map_http_error_to_provider_error, sanitize_url}; use crate::providers::retry::RetryConfig; use goose_providers::errors::ProviderError; +use goose_providers::http_status::read_error_body; use goose_providers::request_log::{start_log, LoggerHandleExt}; use rmcp::model::Tool; @@ -174,7 +175,8 @@ impl GcpVertexAIProvider { let host = Self::build_host_url(&location); let client = Client::builder() - .timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) + .connect_timeout(Duration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) + .read_timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .build()?; let auth = GcpAuth::new().await?; @@ -327,11 +329,19 @@ impl GcpVertexAIProvider { } } - let response = (self.request_builder)(request) - .map_err(|e| ProviderError::ExecutionError(e.to_string()))? - .send() - .await - .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; + let request = (self.request_builder)(request) + .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; + let response = tokio::time::timeout( + Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), + request.send(), + ) + .await + .map_err(|_| { + ProviderError::NetworkError( + "Request timed out — check your network connection and try again.".to_string(), + ) + })? + .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; let status = response.status(); @@ -345,7 +355,8 @@ impl GcpVertexAIProvider { }), ); } - let msg = rate_limit_error_message(&response.text().await.unwrap_or_default()); + let msg = + rate_limit_error_message(&read_error_body(response).await.unwrap_or_default()); tracing::warn!("429 (attempt {rate_limit_attempts}/{max_retries}): {msg}"); last_error = Some(ProviderError::RateLimitExceeded { details: msg, @@ -389,7 +400,7 @@ impl GcpVertexAIProvider { ))); } else { let url = sanitize_url(response.url().as_str()); - let response_text = response.text().await.unwrap_or_default(); + let response_text = read_error_body(response).await.unwrap_or_default(); let payload = serde_json::from_str::(&response_text).ok(); return Err(map_http_error_to_provider_error(status, payload, &url)); } @@ -459,6 +470,7 @@ impl GcpVertexAIProvider { let response = match self .client .post(&url) + .timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .header("Authorization", &auth_header) .json(&payload) .send() diff --git a/crates/goose/src/providers/gemini_oauth.rs b/crates/goose/src/providers/gemini_oauth.rs index a954838e48a5..c3f8fca818b7 100644 --- a/crates/goose/src/providers/gemini_oauth.rs +++ b/crates/goose/src/providers/gemini_oauth.rs @@ -3,7 +3,7 @@ use crate::conversation::message::Message; use crate::providers::api_client::RequestBuilderDecorator; use crate::providers::base::{ ConfigKey, MessageStream, Provider, ProviderDef, ProviderMetadata, - DEFAULT_PROVIDER_TIMEOUT_SECS, + DEFAULT_CONNECT_TIMEOUT_SECS, DEFAULT_PROVIDER_TIMEOUT_SECS, }; use crate::providers::formats::google::{create_request, response_to_streaming_message}; use crate::providers::google::GOOGLE_DOC_URL; @@ -40,7 +40,8 @@ use tokio_util::io::StreamReader; static HTTP_CLIENT: LazyLock = LazyLock::new(|| { reqwest::Client::builder() - .timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) + .connect_timeout(Duration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) + .read_timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .build() .expect("failed to build HTTP client") }); @@ -254,6 +255,7 @@ async fn exchange_code_for_tokens( let resp = client .post(GOOGLE_TOKEN_ENDPOINT) + .timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .header("Content-Type", "application/x-www-form-urlencoded") .form(¶ms) .send() @@ -281,6 +283,7 @@ async fn refresh_access_token(refresh_token: &str) -> Result { let resp = client .post(GOOGLE_TOKEN_ENDPOINT) + .timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .header("Content-Type", "application/x-www-form-urlencoded") .form(¶ms) .send() @@ -343,6 +346,7 @@ async fn code_assist_request(access_token: &str, method: &str, body: &Value) -> let client = &*HTTP_CLIENT; let resp = client .post(&url) + .timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .header("Authorization", format!("Bearer {}", access_token)) .header("Content-Type", "application/json") .json(body) @@ -371,6 +375,7 @@ async fn code_assist_get(access_token: &str, path: &str) -> Result { let client = &*HTTP_CLIENT; let resp = client .get(&url) + .timeout(Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .header("Authorization", format!("Bearer {}", access_token)) .send() .await?; @@ -884,18 +889,25 @@ impl GeminiOAuthProvider { ) .header("Content-Type", "application/json"); - let response = (self.request_builder)(request.json(&wrapped)) - .map_err(|e| ProviderError::ExecutionError(e.to_string()))? - .send() - .await - .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; + let request = (self.request_builder)(request.json(&wrapped)) + .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; + let response = tokio::time::timeout( + Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), + request.send(), + ) + .await + .map_err(|_| { + ProviderError::NetworkError( + "Request timed out — check your network connection and try again.".to_string(), + ) + })? + .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; if !response.status().is_success() { let status = response.status(); - let text = response - .text() + let text = goose_providers::http_status::read_error_body(response) .await - .unwrap_or_else(|_| "unknown error".to_string()); + .unwrap_or_else(|| "unknown error".to_string()); if status == reqwest::StatusCode::TOO_MANY_REQUESTS { // Parse retry delay from the error message if available diff --git a/crates/goose/src/providers/githubcopilot.rs b/crates/goose/src/providers/githubcopilot.rs index e99b7b7034d8..d1550afcae9b 100644 --- a/crates/goose/src/providers/githubcopilot.rs +++ b/crates/goose/src/providers/githubcopilot.rs @@ -265,6 +265,7 @@ impl GithubCopilotProvider { is_user_initiated: bool, payload: &mut Value, has_images: bool, + streaming: bool, ) -> Result { let (endpoint, token) = self.get_api_info().await?; let auth = AuthMethod::BearerToken(token); @@ -281,6 +282,7 @@ impl GithubCopilotProvider { api_client .request(path) .model_headers(model_config)? + .streaming(streaming) .response_post(payload) .await .map_err(|e| e.into()) @@ -420,6 +422,7 @@ impl GithubCopilotProvider { is_user_initiated, &mut payload_clone, has_images, + true, ) .await?; handle_status(resp).await @@ -467,6 +470,7 @@ impl GithubCopilotProvider { is_user_initiated, &mut payload_clone, has_images, + true, ) .await?; handle_status(resp).await @@ -497,6 +501,7 @@ impl GithubCopilotProvider { is_user_initiated, &mut payload_clone, has_images, + false, ) .await }) diff --git a/crates/goose/src/providers/kimicode.rs b/crates/goose/src/providers/kimicode.rs index 5568ef9fd480..4c61193c6f99 100644 --- a/crates/goose/src/providers/kimicode.rs +++ b/crates/goose/src/providers/kimicode.rs @@ -19,7 +19,7 @@ use uuid::Uuid; use super::api_client::RequestBuilderDecorator; use super::base::{ ConfigKey, MessageStream, Provider, ProviderDef, ProviderMetadata, - DEFAULT_PROVIDER_TIMEOUT_SECS, + DEFAULT_CONNECT_TIMEOUT_SECS, DEFAULT_PROVIDER_TIMEOUT_SECS, }; use super::formats::anthropic::{create_request, response_to_streaming_message}; use super::oauth_device_flow::{ @@ -172,7 +172,8 @@ impl KimiCodeProvider { _tls_config: Option, ) -> Result { let client = Client::builder() - .timeout(StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) + .connect_timeout(StdDuration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) + .read_timeout(StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .build()?; let device_id = Self::get_or_create_device_id().await?; Ok(Self { @@ -326,11 +327,19 @@ impl KimiCodeProvider { .headers(self.kimi_headers()) .json(payload); - (self.request_builder)(builder) - .map_err(|e| ProviderError::ExecutionError(e.to_string()))? - .send() - .await - .map_err(|e| ProviderError::RequestFailed(e.to_string())) + let request = (self.request_builder)(builder) + .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; + tokio::time::timeout( + StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), + request.send(), + ) + .await + .map_err(|_| { + ProviderError::NetworkError( + "Request timed out — check your network connection and try again.".to_string(), + ) + })? + .map_err(|e| ProviderError::RequestFailed(e.to_string())) } } @@ -454,6 +463,7 @@ impl Provider for KimiCodeProvider { let resp = self .client .get(format!("{}/v1/models", self.api_base)) + .timeout(StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) .bearer_auth(access_token) .headers(self.kimi_headers()) .send() diff --git a/crates/goose/src/providers/nanogpt.rs b/crates/goose/src/providers/nanogpt.rs index df7fcb57fe57..2ced5d7b2893 100644 --- a/crates/goose/src/providers/nanogpt.rs +++ b/crates/goose/src/providers/nanogpt.rs @@ -198,6 +198,7 @@ impl Provider for NanoGptProvider { .api_client .request("chat/completions") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; handle_status(resp).await diff --git a/crates/goose/src/providers/oauth_device_flow.rs b/crates/goose/src/providers/oauth_device_flow.rs index d257d67b2df8..6baff9a54108 100644 --- a/crates/goose/src/providers/oauth_device_flow.rs +++ b/crates/goose/src/providers/oauth_device_flow.rs @@ -287,7 +287,12 @@ async fn send_request( url: &str, body: &T, ) -> reqwest::Result { - let builder = client.post(url).headers(cfg.extra_headers.clone()); + let builder = client + .post(url) + .timeout(std::time::Duration::from_secs( + super::base::DEFAULT_PROVIDER_TIMEOUT_SECS, + )) + .headers(cfg.extra_headers.clone()); let builder = match cfg.encoding { RequestEncoding::Form => builder.form(body), RequestEncoding::Json => builder.json(body), diff --git a/crates/goose/src/providers/openrouter.rs b/crates/goose/src/providers/openrouter.rs index 631745ccfe49..177a6356d2e8 100644 --- a/crates/goose/src/providers/openrouter.rs +++ b/crates/goose/src/providers/openrouter.rs @@ -350,6 +350,7 @@ impl Provider for OpenRouterProvider { .api_client .request("api/v1/chat/completions") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; handle_status(resp).await diff --git a/crates/goose/src/providers/tetrate.rs b/crates/goose/src/providers/tetrate.rs index 2bb64231d16b..a9e45c4f9c2f 100644 --- a/crates/goose/src/providers/tetrate.rs +++ b/crates/goose/src/providers/tetrate.rs @@ -154,6 +154,7 @@ impl Provider for TetrateProvider { .api_client .request("v1/chat/completions") .model_headers(model_config)? + .streaming(true) .response_post(&payload) .await?; let resp = handle_status(resp) @@ -169,15 +170,18 @@ impl Provider for TetrateProvider { if is_json { // Streaming responses should be SSE; when we get JSON instead, parse it to map - // explicit error payloads and otherwise fail as a protocol mismatch. - let body = handle_response_openai_compat(resp) + // explicit error payloads and otherwise fail as a protocol mismatch. The read + // is bounded: this request is exempt from the total deadline. + let body = goose_providers::http_status::read_error_body(resp) .await - .map_err(Self::enrich_credits_error)?; - if body.get("error").is_some() { - return Err(Self::error_from_tetrate_error_payload( - body, - "v1/chat/completions", - )); + .unwrap_or_default(); + if let Ok(payload) = serde_json::from_str::(&body) { + if payload.get("error").is_some() { + return Err(Self::error_from_tetrate_error_payload( + payload, + "v1/chat/completions", + )); + } } return Err(ProviderError::ExecutionError( From 357fd2b2a8a58136defca8375ac4688f206542b1 Mon Sep 17 00:00:00 2001 From: Douwe M Osinga Date: Tue, 4 Aug 2026 21:13:26 +0200 Subject: [PATCH 2/6] fix(providers): honor configured error body timeout --- crates/goose-providers/src/anthropic.rs | 3 ++- crates/goose-providers/src/api_client.rs | 4 ++++ crates/goose-providers/src/databricks.rs | 18 +++++++++------ crates/goose-providers/src/databricks_v2.rs | 6 ++--- crates/goose-providers/src/google.rs | 2 +- crates/goose-providers/src/http_status.rs | 22 ++++++++++--------- crates/goose-providers/src/ollama.rs | 2 +- crates/goose-providers/src/openai.rs | 5 +++-- .../goose-providers/src/openai_compatible.rs | 1 + crates/goose-providers/src/snowflake.rs | 2 +- crates/goose/src/providers/bedrock.rs | 2 +- crates/goose/src/providers/chatgpt_codex.rs | 2 +- crates/goose/src/providers/gcpvertexai.rs | 12 +++++++--- crates/goose/src/providers/gemini_oauth.rs | 9 +++++--- crates/goose/src/providers/githubcopilot.rs | 4 ++-- crates/goose/src/providers/kimicode.rs | 5 +++-- crates/goose/src/providers/nanogpt.rs | 2 +- crates/goose/src/providers/openrouter.rs | 2 +- crates/goose/src/providers/tetrate.rs | 11 ++++++---- 19 files changed, 70 insertions(+), 44 deletions(-) diff --git a/crates/goose-providers/src/anthropic.rs b/crates/goose-providers/src/anthropic.rs index 8ec61519f4a3..7a60b15bb216 100644 --- a/crates/goose-providers/src/anthropic.rs +++ b/crates/goose-providers/src/anthropic.rs @@ -191,6 +191,7 @@ impl AnthropicProvider { .streaming(true) .response_post(&payload) .await?, + self.api_client.timeout(), ) .await }) @@ -229,7 +230,7 @@ impl AnthropicProvider { return Err(ProviderError::EndpointNotFound(msg)); } - let response = handle_status(response).await?; + let response = handle_status(response, self.api_client.timeout()).await?; let body = response.bytes().await.map_err(|e| { ProviderError::NetworkError(format!("Failed to read response body: {}", e)) diff --git a/crates/goose-providers/src/api_client.rs b/crates/goose-providers/src/api_client.rs index f3a1dc1f637a..b4fb57619ddb 100644 --- a/crates/goose-providers/src/api_client.rs +++ b/crates/goose-providers/src/api_client.rs @@ -238,6 +238,10 @@ pub struct ApiRequestBuilder<'a> { } impl ApiClient { + pub fn timeout(&self) -> Duration { + self.timeout + } + pub fn new_with_tls( host: String, auth: AuthMethod, diff --git a/crates/goose-providers/src/databricks.rs b/crates/goose-providers/src/databricks.rs index eb8a4380aa56..bafa31cac957 100644 --- a/crates/goose-providers/src/databricks.rs +++ b/crates/goose-providers/src/databricks.rs @@ -604,7 +604,7 @@ impl Provider for DatabricksProvider { .streaming(true) .response_post(&payload_clone) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { @@ -673,9 +673,10 @@ impl Provider for DatabricksProvider { if !resp.status().is_success() { let status = resp.status(); let url = sanitize_url(resp.url().as_str()); - let error_text = crate::http_status::read_error_body(resp) - .await - .unwrap_or_default(); + let error_text = + crate::http_status::read_error_body(resp, self.api_client.timeout()) + .await + .unwrap_or_default(); let json_payload = serde_json::from_str::(&error_text).ok(); return Err(map_http_error_to_provider_error(status, json_payload, &url)); @@ -698,9 +699,12 @@ impl Provider for DatabricksProvider { if !resp.status().is_success() { let status = resp.status(); let url = sanitize_url(resp.url().as_str()); - let error_text = crate::http_status::read_error_body(resp) - .await - .unwrap_or_default(); + let error_text = crate::http_status::read_error_body( + resp, + self.api_client.timeout(), + ) + .await + .unwrap_or_default(); let json_payload = serde_json::from_str::(&error_text).ok(); return Err(map_http_error_to_provider_error( status, diff --git a/crates/goose-providers/src/databricks_v2.rs b/crates/goose-providers/src/databricks_v2.rs index 0557676527a8..24f10bce2ff6 100644 --- a/crates/goose-providers/src/databricks_v2.rs +++ b/crates/goose-providers/src/databricks_v2.rs @@ -210,7 +210,7 @@ impl DatabricksV2Provider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { @@ -249,7 +249,7 @@ impl DatabricksV2Provider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { @@ -286,7 +286,7 @@ impl DatabricksV2Provider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { diff --git a/crates/goose-providers/src/google.rs b/crates/goose-providers/src/google.rs index bc9e7eab7e1e..43e5e9a5727a 100644 --- a/crates/goose-providers/src/google.rs +++ b/crates/goose-providers/src/google.rs @@ -110,7 +110,7 @@ impl GoogleProvider { .streaming(true) .response_post(payload) .await?; - handle_status(response).await + handle_status(response, self.api_client.timeout()).await } } diff --git a/crates/goose-providers/src/http_status.rs b/crates/goose-providers/src/http_status.rs index 52fed17d927b..d91763e9aafd 100644 --- a/crates/goose-providers/src/http_status.rs +++ b/crates/goose-providers/src/http_status.rs @@ -248,23 +248,22 @@ pub fn map_http_error_to_provider_error( error } -/// Streaming requests carry no total-request deadline, so error-body reads -/// must be bounded or a server that drips a non-2xx body stalls retries forever. -const ERROR_BODY_READ_TIMEOUT: Duration = Duration::from_secs(30); - -pub async fn read_error_body(response: Response) -> Option { - tokio::time::timeout(ERROR_BODY_READ_TIMEOUT, response.text()) +pub async fn read_error_body(response: Response, timeout: Duration) -> Option { + tokio::time::timeout(timeout, response.text()) .await .ok() .and_then(Result::ok) } -pub async fn handle_status(response: Response) -> Result { +pub async fn handle_status( + response: Response, + timeout: Duration, +) -> Result { let status = response.status(); if !status.is_success() { let url = sanitize_url(response.url().as_str()); let headers = response.headers().clone(); - let body = read_error_body(response).await.unwrap_or_default(); + let body = read_error_body(response, timeout).await.unwrap_or_default(); let payload = serde_json::from_str::(&body).ok(); let mut err = map_http_error_to_provider_error(status, payload.clone(), &url); if let ProviderError::RateLimitExceeded { details, .. } = &err { @@ -278,8 +277,11 @@ pub async fn handle_status(response: Response) -> Result Result { - let response = handle_status(response).await?; +pub async fn handle_response( + response: Response, + timeout: Duration, +) -> Result { + let response = handle_status(response, timeout).await?; response.json::().await.map_err(|e| { ProviderError::RequestFailed(format!("Response body is not valid JSON: {}", e)) diff --git a/crates/goose-providers/src/ollama.rs b/crates/goose-providers/src/ollama.rs index c1de5a336370..59b9c1fc7cd7 100644 --- a/crates/goose-providers/src/ollama.rs +++ b/crates/goose-providers/src/ollama.rs @@ -431,7 +431,7 @@ impl Provider for OllamaProvider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { diff --git a/crates/goose-providers/src/openai.rs b/crates/goose-providers/src/openai.rs index b256ca989da8..746635d441db 100644 --- a/crates/goose-providers/src/openai.rs +++ b/crates/goose-providers/src/openai.rs @@ -314,6 +314,7 @@ impl OpenAiProvider { .streaming(self.supports_streaming) .response_post(&payload) .await?, + self.api_client.timeout(), ) .await }) @@ -544,7 +545,7 @@ impl OpenAiProvider { return Err(ProviderError::EndpointNotFound(body)); } - let response = handle_status(response).await?; + let response = handle_status(response, self.api_client.timeout()).await?; let body = response.bytes().await.map_err(|e| { ProviderError::NetworkError(format!("Failed to read response body: {}", e)) @@ -801,7 +802,7 @@ impl Provider for OpenAiProvider { .streaming(self.supports_streaming) .response_post(&payload) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { diff --git a/crates/goose-providers/src/openai_compatible.rs b/crates/goose-providers/src/openai_compatible.rs index f40d756da26c..63fbe28fd273 100644 --- a/crates/goose-providers/src/openai_compatible.rs +++ b/crates/goose-providers/src/openai_compatible.rs @@ -115,6 +115,7 @@ impl OpenAiCompatibleProvider { .streaming(self.supports_streaming) .response_post(&payload) .await?, + self.api_client.timeout(), ) .await }) diff --git a/crates/goose-providers/src/snowflake.rs b/crates/goose-providers/src/snowflake.rs index 42d17b109f0d..28aed2a5d5b9 100644 --- a/crates/goose-providers/src/snowflake.rs +++ b/crates/goose-providers/src/snowflake.rs @@ -117,7 +117,7 @@ impl SnowflakeProvider { let payload_text: String = if status.is_success() && !is_json { response.text().await.ok().unwrap_or_default() } else { - crate::http_status::read_error_body(response) + crate::http_status::read_error_body(response, self.api_client.timeout()) .await .unwrap_or_default() }; diff --git a/crates/goose/src/providers/bedrock.rs b/crates/goose/src/providers/bedrock.rs index bef936f7ad5d..f643cbe8a575 100644 --- a/crates/goose/src/providers/bedrock.rs +++ b/crates/goose/src/providers/bedrock.rs @@ -290,7 +290,7 @@ impl BedrockProvider { })? .map_err(|e| ProviderError::RequestFailed(format!("Mantle request failed: {}", e)))?; - handle_status(response).await + handle_status(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await } /// Build the request inputs shared by [`Self::converse`] and diff --git a/crates/goose/src/providers/chatgpt_codex.rs b/crates/goose/src/providers/chatgpt_codex.rs index 6c166860e0e7..fe3acfd71bbf 100644 --- a/crates/goose/src/providers/chatgpt_codex.rs +++ b/crates/goose/src/providers/chatgpt_codex.rs @@ -955,7 +955,7 @@ impl ChatGptCodexProvider { })? .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; - handle_status(response).await + handle_status(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await } } diff --git a/crates/goose/src/providers/gcpvertexai.rs b/crates/goose/src/providers/gcpvertexai.rs index 545a4a5bb223..8633256c2855 100644 --- a/crates/goose/src/providers/gcpvertexai.rs +++ b/crates/goose/src/providers/gcpvertexai.rs @@ -355,8 +355,11 @@ impl GcpVertexAIProvider { }), ); } - let msg = - rate_limit_error_message(&read_error_body(response).await.unwrap_or_default()); + let msg = rate_limit_error_message( + &read_error_body(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) + .await + .unwrap_or_default(), + ); tracing::warn!("429 (attempt {rate_limit_attempts}/{max_retries}): {msg}"); last_error = Some(ProviderError::RateLimitExceeded { details: msg, @@ -400,7 +403,10 @@ impl GcpVertexAIProvider { ))); } else { let url = sanitize_url(response.url().as_str()); - let response_text = read_error_body(response).await.unwrap_or_default(); + let response_text = + read_error_body(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) + .await + .unwrap_or_default(); let payload = serde_json::from_str::(&response_text).ok(); return Err(map_http_error_to_provider_error(status, payload, &url)); } diff --git a/crates/goose/src/providers/gemini_oauth.rs b/crates/goose/src/providers/gemini_oauth.rs index c3f8fca818b7..91cc8638b834 100644 --- a/crates/goose/src/providers/gemini_oauth.rs +++ b/crates/goose/src/providers/gemini_oauth.rs @@ -905,9 +905,12 @@ impl GeminiOAuthProvider { if !response.status().is_success() { let status = response.status(); - let text = goose_providers::http_status::read_error_body(response) - .await - .unwrap_or_else(|| "unknown error".to_string()); + let text = goose_providers::http_status::read_error_body( + response, + Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), + ) + .await + .unwrap_or_else(|| "unknown error".to_string()); if status == reqwest::StatusCode::TOO_MANY_REQUESTS { // Parse retry delay from the error message if available diff --git a/crates/goose/src/providers/githubcopilot.rs b/crates/goose/src/providers/githubcopilot.rs index d1550afcae9b..dd5cb9b08f7a 100644 --- a/crates/goose/src/providers/githubcopilot.rs +++ b/crates/goose/src/providers/githubcopilot.rs @@ -425,7 +425,7 @@ impl GithubCopilotProvider { true, ) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { @@ -473,7 +473,7 @@ impl GithubCopilotProvider { true, ) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { diff --git a/crates/goose/src/providers/kimicode.rs b/crates/goose/src/providers/kimicode.rs index 4c61193c6f99..09aabc8117d0 100644 --- a/crates/goose/src/providers/kimicode.rs +++ b/crates/goose/src/providers/kimicode.rs @@ -419,7 +419,7 @@ impl Provider for KimiCodeProvider { let response = self .with_retry(|| async { let resp = self.post(&payload).await?; - handle_status(resp).await + handle_status(resp, StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) }) .await .inspect_err(|e| { @@ -469,7 +469,8 @@ impl Provider for KimiCodeProvider { .send() .await .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; - let resp = handle_status(resp).await?; + let resp = + handle_status(resp, StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await?; let parsed: ModelsResp = resp.json().await.map_err(|e| { ProviderError::RequestFailed(format!("/v1/models body is not valid JSON: {}", e)) diff --git a/crates/goose/src/providers/nanogpt.rs b/crates/goose/src/providers/nanogpt.rs index 2ced5d7b2893..20a8112cbc35 100644 --- a/crates/goose/src/providers/nanogpt.rs +++ b/crates/goose/src/providers/nanogpt.rs @@ -201,7 +201,7 @@ impl Provider for NanoGptProvider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { diff --git a/crates/goose/src/providers/openrouter.rs b/crates/goose/src/providers/openrouter.rs index 177a6356d2e8..589fdb672e0d 100644 --- a/crates/goose/src/providers/openrouter.rs +++ b/crates/goose/src/providers/openrouter.rs @@ -353,7 +353,7 @@ impl Provider for OpenRouterProvider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp).await + handle_status(resp, self.api_client.timeout()).await }) .await .inspect_err(|e| { diff --git a/crates/goose/src/providers/tetrate.rs b/crates/goose/src/providers/tetrate.rs index a9e45c4f9c2f..24c32d2fdab5 100644 --- a/crates/goose/src/providers/tetrate.rs +++ b/crates/goose/src/providers/tetrate.rs @@ -157,7 +157,7 @@ impl Provider for TetrateProvider { .streaming(true) .response_post(&payload) .await?; - let resp = handle_status(resp) + let resp = handle_status(resp, self.api_client.timeout()) .await .map_err(Self::enrich_credits_error)?; @@ -172,9 +172,12 @@ impl Provider for TetrateProvider { // Streaming responses should be SSE; when we get JSON instead, parse it to map // explicit error payloads and otherwise fail as a protocol mismatch. The read // is bounded: this request is exempt from the total deadline. - let body = goose_providers::http_status::read_error_body(resp) - .await - .unwrap_or_default(); + let body = goose_providers::http_status::read_error_body( + resp, + self.api_client.timeout(), + ) + .await + .unwrap_or_default(); if let Ok(payload) = serde_json::from_str::(&body) { if payload.get("error").is_some() { return Err(Self::error_from_tetrate_error_payload( From 5ba84157d63ed024168e054e651454fb649146d9 Mon Sep 17 00:00:00 2001 From: Douwe M Osinga Date: Tue, 4 Aug 2026 21:14:59 +0200 Subject: [PATCH 3/6] refactor(providers): remove redundant timeout comments --- crates/goose-providers/src/api_client.rs | 10 ---------- crates/goose-providers/src/snowflake.rs | 2 -- crates/goose/src/providers/tetrate.rs | 3 --- 3 files changed, 15 deletions(-) diff --git a/crates/goose-providers/src/api_client.rs b/crates/goose-providers/src/api_client.rs index b4fb57619ddb..da1644916c65 100644 --- a/crates/goose-providers/src/api_client.rs +++ b/crates/goose-providers/src/api_client.rs @@ -289,10 +289,6 @@ impl ApiClient { self.timeout } - /// The total-request deadline is applied per request in `send_request` so - /// streaming requests can opt out of it: an SSE body is the whole - /// generation, so streams are bounded by `read_timeout` (reset per chunk) - /// rather than end-to-end duration. fn client_builder(timeout: Duration) -> reqwest::ClientBuilder { Client::builder() .connect_timeout(Duration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) @@ -450,8 +446,6 @@ impl<'a> ApiRequestBuilder<'a> { } } - /// Exempts this request from the total-request deadline; the stream is - /// bounded by the client's read timeout instead. pub fn streaming(mut self, streaming: bool) -> Self { self.streaming = streaming; self @@ -462,10 +456,6 @@ impl<'a> ApiRequestBuilder<'a> { ApiResponse::from_response(response).await } - /// `send()` resolves when response headers arrive, so for streaming - /// requests (exempt from the per-request total deadline) this still bounds - /// connect, request upload, and time-to-first-byte; only the streamed - /// response body is unbounded. async fn send_bounded(&self, request: reqwest::RequestBuilder) -> Result { if self.streaming { Ok(tokio::time::timeout(self.client.timeout, request.send()).await??) diff --git a/crates/goose-providers/src/snowflake.rs b/crates/goose-providers/src/snowflake.rs index 28aed2a5d5b9..9d85f3ffb337 100644 --- a/crates/goose-providers/src/snowflake.rs +++ b/crates/goose-providers/src/snowflake.rs @@ -112,8 +112,6 @@ impl SnowflakeProvider { .and_then(|v| v.to_str().ok()) .map(|v| v.to_ascii_lowercase()) .is_some_and(|v| v.contains("json")); - // A 200 with a JSON body is an error payload, not the SSE generation - // stream; bound its read since this request has no total deadline. let payload_text: String = if status.is_success() && !is_json { response.text().await.ok().unwrap_or_default() } else { diff --git a/crates/goose/src/providers/tetrate.rs b/crates/goose/src/providers/tetrate.rs index 24c32d2fdab5..07505a65993b 100644 --- a/crates/goose/src/providers/tetrate.rs +++ b/crates/goose/src/providers/tetrate.rs @@ -169,9 +169,6 @@ impl Provider for TetrateProvider { .is_some_and(|v| v.contains("json")); if is_json { - // Streaming responses should be SSE; when we get JSON instead, parse it to map - // explicit error payloads and otherwise fail as a protocol mismatch. The read - // is bounded: this request is exempt from the total deadline. let body = goose_providers::http_status::read_error_body( resp, self.api_client.timeout(), From 585460f563a8db9464054fad2d5ddbd155538467 Mon Sep 17 00:00:00 2001 From: Douwe M Osinga Date: Tue, 4 Aug 2026 21:26:41 +0200 Subject: [PATCH 4/6] fix(providers): pass timeout to response handlers --- crates/goose-providers/src/api_client.rs | 4 ---- crates/goose-providers/src/azure_foundry.rs | 3 ++- crates/goose-providers/src/openai.rs | 4 +++- crates/goose-providers/src/openai_compatible.rs | 2 +- crates/goose/src/providers/chatgpt_codex.rs | 1 + crates/goose/src/providers/githubcopilot.rs | 10 +++++++--- crates/goose/src/providers/kimicode.rs | 2 +- crates/goose/src/providers/litellm.rs | 2 +- crates/goose/src/providers/tetrate.rs | 2 +- 9 files changed, 17 insertions(+), 13 deletions(-) diff --git a/crates/goose-providers/src/api_client.rs b/crates/goose-providers/src/api_client.rs index da1644916c65..1f169931655a 100644 --- a/crates/goose-providers/src/api_client.rs +++ b/crates/goose-providers/src/api_client.rs @@ -238,10 +238,6 @@ pub struct ApiRequestBuilder<'a> { } impl ApiClient { - pub fn timeout(&self) -> Duration { - self.timeout - } - pub fn new_with_tls( host: String, auth: AuthMethod, diff --git a/crates/goose-providers/src/azure_foundry.rs b/crates/goose-providers/src/azure_foundry.rs index 39c260d70d69..8948b8319ceb 100644 --- a/crates/goose-providers/src/azure_foundry.rs +++ b/crates/goose-providers/src/azure_foundry.rs @@ -273,7 +273,8 @@ impl AzureFoundryProvider { .response_get(&path) .await .map_err(|error| ProviderError::NetworkError(error.to_string()))?; - let json = handle_response_openai_compat(response).await?; + let json = + handle_response_openai_compat(response, self.deployments_client.timeout()).await?; if let Some(items) = json.get("value").and_then(|value| value.as_array()) { for item in items { let Some(name) = item.get("name").and_then(|value| value.as_str()) else { diff --git a/crates/goose-providers/src/openai.rs b/crates/goose-providers/src/openai.rs index 746635d441db..bcbc4396f8e7 100644 --- a/crates/goose-providers/src/openai.rs +++ b/crates/goose-providers/src/openai.rs @@ -577,7 +577,9 @@ impl OpenAiProvider { .response_get() .await .ok()?; - let json = handle_response_openai_compat(response).await.ok()?; + let json = handle_response_openai_compat(response, self.api_client.timeout()) + .await + .ok()?; parse_n_ctx_from_models(&json, model_name) } } diff --git a/crates/goose-providers/src/openai_compatible.rs b/crates/goose-providers/src/openai_compatible.rs index 63fbe28fd273..c722795c3e6c 100644 --- a/crates/goose-providers/src/openai_compatible.rs +++ b/crates/goose-providers/src/openai_compatible.rs @@ -185,7 +185,7 @@ impl Provider for OpenAiCompatibleProvider { .response_get("models") .await .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; - let json = handle_response_openai_compat(response).await?; + let json = handle_response_openai_compat(response, self.api_client.timeout()).await?; if let Some(err_obj) = json.get("error") { let msg = err_obj diff --git a/crates/goose/src/providers/chatgpt_codex.rs b/crates/goose/src/providers/chatgpt_codex.rs index fe3acfd71bbf..4e085a69a3aa 100644 --- a/crates/goose/src/providers/chatgpt_codex.rs +++ b/crates/goose/src/providers/chatgpt_codex.rs @@ -29,6 +29,7 @@ use std::io; use std::net::SocketAddr; use std::path::PathBuf; use std::sync::{Arc, LazyLock}; +use std::time::Duration; use tokio::pin; use tokio::sync::{oneshot, Mutex as TokioMutex}; use tokio_util::codec::{FramedRead, LinesCodec}; diff --git a/crates/goose/src/providers/githubcopilot.rs b/crates/goose/src/providers/githubcopilot.rs index dd5cb9b08f7a..c5c82b391020 100644 --- a/crates/goose/src/providers/githubcopilot.rs +++ b/crates/goose/src/providers/githubcopilot.rs @@ -425,7 +425,7 @@ impl GithubCopilotProvider { true, ) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await }) .await .inspect_err(|e| { @@ -473,7 +473,7 @@ impl GithubCopilotProvider { true, ) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await }) .await .inspect_err(|e| { @@ -506,7 +506,11 @@ impl GithubCopilotProvider { .await }) .await?; - let response = handle_response_openai_compat(response).await?; + let response = handle_response_openai_compat( + response, + Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), + ) + .await?; let response = promote_tool_choice(response); diff --git a/crates/goose/src/providers/kimicode.rs b/crates/goose/src/providers/kimicode.rs index 09aabc8117d0..cacab1040965 100644 --- a/crates/goose/src/providers/kimicode.rs +++ b/crates/goose/src/providers/kimicode.rs @@ -419,7 +419,7 @@ impl Provider for KimiCodeProvider { let response = self .with_retry(|| async { let resp = self.post(&payload).await?; - handle_status(resp, StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) + handle_status(resp, StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await }) .await .inspect_err(|e| { diff --git a/crates/goose/src/providers/litellm.rs b/crates/goose/src/providers/litellm.rs index e9caf70f573a..dc85b4ee8903 100644 --- a/crates/goose/src/providers/litellm.rs +++ b/crates/goose/src/providers/litellm.rs @@ -147,7 +147,7 @@ impl LiteLLMProvider { .model_headers(model_config)? .response_post(payload) .await?; - handle_response_openai_compat(response).await + handle_response_openai_compat(response, self.api_client.timeout()).await } async fn supports_cache_control(&self, model: &ModelConfig) -> bool { diff --git a/crates/goose/src/providers/tetrate.rs b/crates/goose/src/providers/tetrate.rs index 07505a65993b..dfb6b653149c 100644 --- a/crates/goose/src/providers/tetrate.rs +++ b/crates/goose/src/providers/tetrate.rs @@ -207,7 +207,7 @@ impl Provider for TetrateProvider { .response_get("v1/models") .await .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; - let json = handle_response_openai_compat(response).await?; + let json = handle_response_openai_compat(response, self.api_client.timeout()).await?; // Tetrate can return errors in 200 OK responses, so check explicitly if json.get("error").is_some() { From 1c9c3242ccab9ac894baa0238a81ef1c7b299da7 Mon Sep 17 00:00:00 2001 From: Douwe M Osinga Date: Tue, 4 Aug 2026 21:59:51 +0200 Subject: [PATCH 5/6] fix(providers): share streaming request deadlines --- crates/goose-providers/src/anthropic.rs | 3 +- crates/goose-providers/src/api_client.rs | 65 +++++++++++++++---- crates/goose-providers/src/azure_foundry.rs | 3 +- crates/goose-providers/src/databricks.rs | 18 ++--- crates/goose-providers/src/databricks_v2.rs | 6 +- crates/goose-providers/src/google.rs | 2 +- crates/goose-providers/src/http_status.rs | 46 +++++++++---- crates/goose-providers/src/ollama.rs | 2 +- crates/goose-providers/src/openai.rs | 9 +-- .../goose-providers/src/openai_compatible.rs | 3 +- crates/goose-providers/src/snowflake.rs | 2 +- crates/goose/src/providers/base.rs | 5 +- crates/goose/src/providers/bedrock.rs | 14 ++-- crates/goose/src/providers/chatgpt_codex.rs | 15 ++--- crates/goose/src/providers/gcpvertexai.rs | 24 ++----- crates/goose/src/providers/gemini_oauth.rs | 21 ++---- crates/goose/src/providers/githubcopilot.rs | 10 +-- crates/goose/src/providers/kimicode.rs | 15 ++--- crates/goose/src/providers/litellm.rs | 2 +- crates/goose/src/providers/nanogpt.rs | 2 +- crates/goose/src/providers/openrouter.rs | 2 +- crates/goose/src/providers/tetrate.rs | 20 +++--- 22 files changed, 147 insertions(+), 142 deletions(-) diff --git a/crates/goose-providers/src/anthropic.rs b/crates/goose-providers/src/anthropic.rs index 7a60b15bb216..8ec61519f4a3 100644 --- a/crates/goose-providers/src/anthropic.rs +++ b/crates/goose-providers/src/anthropic.rs @@ -191,7 +191,6 @@ impl AnthropicProvider { .streaming(true) .response_post(&payload) .await?, - self.api_client.timeout(), ) .await }) @@ -230,7 +229,7 @@ impl AnthropicProvider { return Err(ProviderError::EndpointNotFound(msg)); } - let response = handle_status(response, self.api_client.timeout()).await?; + let response = handle_status(response).await?; let body = response.bytes().await.map_err(|e| { ProviderError::NetworkError(format!("Failed to read response body: {}", e)) diff --git a/crates/goose-providers/src/api_client.rs b/crates/goose-providers/src/api_client.rs index 1f169931655a..62213e62e7f5 100644 --- a/crates/goose-providers/src/api_client.rs +++ b/crates/goose-providers/src/api_client.rs @@ -454,7 +454,7 @@ impl<'a> ApiRequestBuilder<'a> { async fn send_bounded(&self, request: reqwest::RequestBuilder) -> Result { if self.streaming { - Ok(tokio::time::timeout(self.client.timeout, request.send()).await??) + Ok(crate::http_status::send_bounded(request, self.client.timeout).await?) } else { Ok(request.send().await?) } @@ -693,21 +693,28 @@ mod tests { } fn client_with_timeout(addr: SocketAddr, timeout_ms: u64) -> ApiClient { - ApiClient::with_timeout_and_tls( + let mut client = ApiClient::with_timeout_and_tls( format!("http://{}", addr), AuthMethod::NoAuth, Duration::from_millis(timeout_ms), None, ) - .unwrap() + .unwrap(); + client.client = Client::builder() + .no_proxy() + .connect_timeout(Duration::from_secs(DEFAULT_CONNECT_TIMEOUT_SECS)) + .read_timeout(client.timeout) + .build() + .unwrap(); + client } async fn drain_counting_data_lines(mut response: Response) -> Result { - let mut count = 0; + let mut body = Vec::new(); while let Some(chunk) = response.chunk().await? { - count += String::from_utf8_lossy(&chunk).matches("data:").count(); + body.extend_from_slice(&chunk); } - Ok(count) + Ok(String::from_utf8_lossy(&body).matches("data:").count()) } #[tokio::test] @@ -791,11 +798,47 @@ mod tests { "should fail near the configured timeout, took {:?}", started.elapsed() ); - let timed_out = err - .downcast_ref::() - .is_some_and(reqwest::Error::is_timeout) - || err.downcast_ref::().is_some(); - assert!(timed_out, "expected a timeout error, got: {err}"); + assert!(matches!( + err.downcast_ref::(), + Some(crate::errors::ProviderError::NetworkError(message)) + if message.starts_with("Request timed out") + )); + } + + #[tokio::test] + async fn streaming_error_body_shares_send_deadline() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.unwrap(); + let mut buf = [0u8; 8192]; + let _ = socket.read(&mut buf).await; + tokio::time::sleep(Duration::from_millis(300)).await; + socket + .write_all(b"HTTP/1.1 500 Internal Server Error\r\ncontent-length: 1\r\n\r\n") + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(300)).await; + let _ = socket.write_all(b"x").await; + }); + + let client = client_with_timeout(addr, 400); + let started = std::time::Instant::now(); + let response = client + .request("v1/messages") + .streaming(true) + .response_post(&serde_json::json!({})) + .await + .unwrap(); + crate::http_status::handle_status(response) + .await + .unwrap_err(); + + assert!( + started.elapsed() < Duration::from_millis(550), + "send and error body used separate deadlines: {:?}", + started.elapsed() + ); } #[test] diff --git a/crates/goose-providers/src/azure_foundry.rs b/crates/goose-providers/src/azure_foundry.rs index 8948b8319ceb..39c260d70d69 100644 --- a/crates/goose-providers/src/azure_foundry.rs +++ b/crates/goose-providers/src/azure_foundry.rs @@ -273,8 +273,7 @@ impl AzureFoundryProvider { .response_get(&path) .await .map_err(|error| ProviderError::NetworkError(error.to_string()))?; - let json = - handle_response_openai_compat(response, self.deployments_client.timeout()).await?; + let json = handle_response_openai_compat(response).await?; if let Some(items) = json.get("value").and_then(|value| value.as_array()) { for item in items { let Some(name) = item.get("name").and_then(|value| value.as_str()) else { diff --git a/crates/goose-providers/src/databricks.rs b/crates/goose-providers/src/databricks.rs index bafa31cac957..eb8a4380aa56 100644 --- a/crates/goose-providers/src/databricks.rs +++ b/crates/goose-providers/src/databricks.rs @@ -604,7 +604,7 @@ impl Provider for DatabricksProvider { .streaming(true) .response_post(&payload_clone) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { @@ -673,10 +673,9 @@ impl Provider for DatabricksProvider { if !resp.status().is_success() { let status = resp.status(); let url = sanitize_url(resp.url().as_str()); - let error_text = - crate::http_status::read_error_body(resp, self.api_client.timeout()) - .await - .unwrap_or_default(); + let error_text = crate::http_status::read_error_body(resp) + .await + .unwrap_or_default(); let json_payload = serde_json::from_str::(&error_text).ok(); return Err(map_http_error_to_provider_error(status, json_payload, &url)); @@ -699,12 +698,9 @@ impl Provider for DatabricksProvider { if !resp.status().is_success() { let status = resp.status(); let url = sanitize_url(resp.url().as_str()); - let error_text = crate::http_status::read_error_body( - resp, - self.api_client.timeout(), - ) - .await - .unwrap_or_default(); + let error_text = crate::http_status::read_error_body(resp) + .await + .unwrap_or_default(); let json_payload = serde_json::from_str::(&error_text).ok(); return Err(map_http_error_to_provider_error( status, diff --git a/crates/goose-providers/src/databricks_v2.rs b/crates/goose-providers/src/databricks_v2.rs index 24f10bce2ff6..0557676527a8 100644 --- a/crates/goose-providers/src/databricks_v2.rs +++ b/crates/goose-providers/src/databricks_v2.rs @@ -210,7 +210,7 @@ impl DatabricksV2Provider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { @@ -249,7 +249,7 @@ impl DatabricksV2Provider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { @@ -286,7 +286,7 @@ impl DatabricksV2Provider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { diff --git a/crates/goose-providers/src/google.rs b/crates/goose-providers/src/google.rs index 43e5e9a5727a..bc9e7eab7e1e 100644 --- a/crates/goose-providers/src/google.rs +++ b/crates/goose-providers/src/google.rs @@ -110,7 +110,7 @@ impl GoogleProvider { .streaming(true) .response_post(payload) .await?; - handle_status(response, self.api_client.timeout()).await + handle_status(response).await } } diff --git a/crates/goose-providers/src/http_status.rs b/crates/goose-providers/src/http_status.rs index d91763e9aafd..51f8e8e4ff0a 100644 --- a/crates/goose-providers/src/http_status.rs +++ b/crates/goose-providers/src/http_status.rs @@ -248,22 +248,45 @@ pub fn map_http_error_to_provider_error( error } -pub async fn read_error_body(response: Response, timeout: Duration) -> Option { - tokio::time::timeout(timeout, response.text()) - .await - .ok() - .and_then(Result::ok) +#[derive(Clone, Copy)] +pub struct ResponseDeadline(tokio::time::Instant); + +pub fn set_response_deadline(response: &mut Response, deadline: tokio::time::Instant) { + response.extensions_mut().insert(ResponseDeadline(deadline)); } -pub async fn handle_status( - response: Response, +pub async fn send_bounded( + request: reqwest::RequestBuilder, timeout: Duration, ) -> Result { + let deadline = tokio::time::Instant::now() + timeout; + let mut response = tokio::time::timeout_at(deadline, request.send()) + .await + .map_err(|_| { + ProviderError::NetworkError( + "Request timed out — check your network connection and try again.".to_string(), + ) + })??; + set_response_deadline(&mut response, deadline); + Ok(response) +} + +pub async fn read_error_body(response: Response) -> Option { + match response.extensions().get::().copied() { + Some(ResponseDeadline(deadline)) => tokio::time::timeout_at(deadline, response.text()) + .await + .ok() + .and_then(Result::ok), + None => response.text().await.ok(), + } +} + +pub async fn handle_status(response: Response) -> Result { let status = response.status(); if !status.is_success() { let url = sanitize_url(response.url().as_str()); let headers = response.headers().clone(); - let body = read_error_body(response, timeout).await.unwrap_or_default(); + let body = read_error_body(response).await.unwrap_or_default(); let payload = serde_json::from_str::(&body).ok(); let mut err = map_http_error_to_provider_error(status, payload.clone(), &url); if let ProviderError::RateLimitExceeded { details, .. } = &err { @@ -277,11 +300,8 @@ pub async fn handle_status( Ok(response) } -pub async fn handle_response( - response: Response, - timeout: Duration, -) -> Result { - let response = handle_status(response, timeout).await?; +pub async fn handle_response(response: Response) -> Result { + let response = handle_status(response).await?; response.json::().await.map_err(|e| { ProviderError::RequestFailed(format!("Response body is not valid JSON: {}", e)) diff --git a/crates/goose-providers/src/ollama.rs b/crates/goose-providers/src/ollama.rs index 59b9c1fc7cd7..c1de5a336370 100644 --- a/crates/goose-providers/src/ollama.rs +++ b/crates/goose-providers/src/ollama.rs @@ -431,7 +431,7 @@ impl Provider for OllamaProvider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { diff --git a/crates/goose-providers/src/openai.rs b/crates/goose-providers/src/openai.rs index bcbc4396f8e7..b256ca989da8 100644 --- a/crates/goose-providers/src/openai.rs +++ b/crates/goose-providers/src/openai.rs @@ -314,7 +314,6 @@ impl OpenAiProvider { .streaming(self.supports_streaming) .response_post(&payload) .await?, - self.api_client.timeout(), ) .await }) @@ -545,7 +544,7 @@ impl OpenAiProvider { return Err(ProviderError::EndpointNotFound(body)); } - let response = handle_status(response, self.api_client.timeout()).await?; + let response = handle_status(response).await?; let body = response.bytes().await.map_err(|e| { ProviderError::NetworkError(format!("Failed to read response body: {}", e)) @@ -577,9 +576,7 @@ impl OpenAiProvider { .response_get() .await .ok()?; - let json = handle_response_openai_compat(response, self.api_client.timeout()) - .await - .ok()?; + let json = handle_response_openai_compat(response).await.ok()?; parse_n_ctx_from_models(&json, model_name) } } @@ -804,7 +801,7 @@ impl Provider for OpenAiProvider { .streaming(self.supports_streaming) .response_post(&payload) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { diff --git a/crates/goose-providers/src/openai_compatible.rs b/crates/goose-providers/src/openai_compatible.rs index c722795c3e6c..f40d756da26c 100644 --- a/crates/goose-providers/src/openai_compatible.rs +++ b/crates/goose-providers/src/openai_compatible.rs @@ -115,7 +115,6 @@ impl OpenAiCompatibleProvider { .streaming(self.supports_streaming) .response_post(&payload) .await?, - self.api_client.timeout(), ) .await }) @@ -185,7 +184,7 @@ impl Provider for OpenAiCompatibleProvider { .response_get("models") .await .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; - let json = handle_response_openai_compat(response, self.api_client.timeout()).await?; + let json = handle_response_openai_compat(response).await?; if let Some(err_obj) = json.get("error") { let msg = err_obj diff --git a/crates/goose-providers/src/snowflake.rs b/crates/goose-providers/src/snowflake.rs index 9d85f3ffb337..d0fa2bc06566 100644 --- a/crates/goose-providers/src/snowflake.rs +++ b/crates/goose-providers/src/snowflake.rs @@ -115,7 +115,7 @@ impl SnowflakeProvider { let payload_text: String = if status.is_success() && !is_json { response.text().await.ok().unwrap_or_default() } else { - crate::http_status::read_error_body(response, self.api_client.timeout()) + crate::http_status::read_error_body(response) .await .unwrap_or_default() }; diff --git a/crates/goose/src/providers/base.rs b/crates/goose/src/providers/base.rs index c1807a6eab3a..5f6ce506fdfb 100644 --- a/crates/goose/src/providers/base.rs +++ b/crates/goose/src/providers/base.rs @@ -6,8 +6,9 @@ pub use goose_providers::conversation::token_usage::{ }; use serde::{Deserialize, Serialize}; -pub const DEFAULT_PROVIDER_TIMEOUT_SECS: u64 = 600; -pub use goose_providers::api_client::DEFAULT_CONNECT_TIMEOUT_SECS; +pub use goose_providers::api_client::{ + DEFAULT_CONNECT_TIMEOUT_SECS, DEFAULT_PROVIDER_TIMEOUT_SECS, +}; use crate::config::ExtensionConfig; diff --git a/crates/goose/src/providers/bedrock.rs b/crates/goose/src/providers/bedrock.rs index f643cbe8a575..7b48823baaa8 100644 --- a/crates/goose/src/providers/bedrock.rs +++ b/crates/goose/src/providers/bedrock.rs @@ -278,19 +278,13 @@ impl BedrockProvider { } } - let response = tokio::time::timeout( + let response = goose_providers::http_status::send_bounded( + req, std::time::Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), - req.send(), ) - .await - .map_err(|_| { - ProviderError::NetworkError( - "Request timed out — check your network connection and try again.".to_string(), - ) - })? - .map_err(|e| ProviderError::RequestFailed(format!("Mantle request failed: {}", e)))?; + .await?; - handle_status(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await + handle_status(response).await } /// Build the request inputs shared by [`Self::converse`] and diff --git a/crates/goose/src/providers/chatgpt_codex.rs b/crates/goose/src/providers/chatgpt_codex.rs index 4e085a69a3aa..3622898cd2f0 100644 --- a/crates/goose/src/providers/chatgpt_codex.rs +++ b/crates/goose/src/providers/chatgpt_codex.rs @@ -29,7 +29,6 @@ use std::io; use std::net::SocketAddr; use std::path::PathBuf; use std::sync::{Arc, LazyLock}; -use std::time::Duration; use tokio::pin; use tokio::sync::{oneshot, Mutex as TokioMutex}; use tokio_util::codec::{FramedRead, LinesCodec}; @@ -944,19 +943,13 @@ impl ChatGptCodexProvider { let request = (self.request_builder)(request) .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; - let response = tokio::time::timeout( + let response = goose_providers::http_status::send_bounded( + request, std::time::Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), - request.send(), ) - .await - .map_err(|_| { - ProviderError::NetworkError( - "Request timed out — check your network connection and try again.".to_string(), - ) - })? - .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; + .await?; - handle_status(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await + handle_status(response).await } } diff --git a/crates/goose/src/providers/gcpvertexai.rs b/crates/goose/src/providers/gcpvertexai.rs index 8633256c2855..07f1aa218396 100644 --- a/crates/goose/src/providers/gcpvertexai.rs +++ b/crates/goose/src/providers/gcpvertexai.rs @@ -331,17 +331,11 @@ impl GcpVertexAIProvider { let request = (self.request_builder)(request) .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; - let response = tokio::time::timeout( + let response = goose_providers::http_status::send_bounded( + request, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), - request.send(), ) - .await - .map_err(|_| { - ProviderError::NetworkError( - "Request timed out — check your network connection and try again.".to_string(), - ) - })? - .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; + .await?; let status = response.status(); @@ -355,11 +349,8 @@ impl GcpVertexAIProvider { }), ); } - let msg = rate_limit_error_message( - &read_error_body(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) - .await - .unwrap_or_default(), - ); + let msg = + rate_limit_error_message(&read_error_body(response).await.unwrap_or_default()); tracing::warn!("429 (attempt {rate_limit_attempts}/{max_retries}): {msg}"); last_error = Some(ProviderError::RateLimitExceeded { details: msg, @@ -403,10 +394,7 @@ impl GcpVertexAIProvider { ))); } else { let url = sanitize_url(response.url().as_str()); - let response_text = - read_error_body(response, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)) - .await - .unwrap_or_default(); + let response_text = read_error_body(response).await.unwrap_or_default(); let payload = serde_json::from_str::(&response_text).ok(); return Err(map_http_error_to_provider_error(status, payload, &url)); } diff --git a/crates/goose/src/providers/gemini_oauth.rs b/crates/goose/src/providers/gemini_oauth.rs index 91cc8638b834..39bdb7f437bb 100644 --- a/crates/goose/src/providers/gemini_oauth.rs +++ b/crates/goose/src/providers/gemini_oauth.rs @@ -891,26 +891,17 @@ impl GeminiOAuthProvider { let request = (self.request_builder)(request.json(&wrapped)) .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; - let response = tokio::time::timeout( + let response = goose_providers::http_status::send_bounded( + request, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), - request.send(), ) - .await - .map_err(|_| { - ProviderError::NetworkError( - "Request timed out — check your network connection and try again.".to_string(), - ) - })? - .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; + .await?; if !response.status().is_success() { let status = response.status(); - let text = goose_providers::http_status::read_error_body( - response, - Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), - ) - .await - .unwrap_or_else(|| "unknown error".to_string()); + let text = goose_providers::http_status::read_error_body(response) + .await + .unwrap_or_else(|| "unknown error".to_string()); if status == reqwest::StatusCode::TOO_MANY_REQUESTS { // Parse retry delay from the error message if available diff --git a/crates/goose/src/providers/githubcopilot.rs b/crates/goose/src/providers/githubcopilot.rs index c5c82b391020..d1550afcae9b 100644 --- a/crates/goose/src/providers/githubcopilot.rs +++ b/crates/goose/src/providers/githubcopilot.rs @@ -425,7 +425,7 @@ impl GithubCopilotProvider { true, ) .await?; - handle_status(resp, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await + handle_status(resp).await }) .await .inspect_err(|e| { @@ -473,7 +473,7 @@ impl GithubCopilotProvider { true, ) .await?; - handle_status(resp, Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await + handle_status(resp).await }) .await .inspect_err(|e| { @@ -506,11 +506,7 @@ impl GithubCopilotProvider { .await }) .await?; - let response = handle_response_openai_compat( - response, - Duration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), - ) - .await?; + let response = handle_response_openai_compat(response).await?; let response = promote_tool_choice(response); diff --git a/crates/goose/src/providers/kimicode.rs b/crates/goose/src/providers/kimicode.rs index cacab1040965..3e0a38e39194 100644 --- a/crates/goose/src/providers/kimicode.rs +++ b/crates/goose/src/providers/kimicode.rs @@ -329,17 +329,11 @@ impl KimiCodeProvider { let request = (self.request_builder)(builder) .map_err(|e| ProviderError::ExecutionError(e.to_string()))?; - tokio::time::timeout( + goose_providers::http_status::send_bounded( + request, StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS), - request.send(), ) .await - .map_err(|_| { - ProviderError::NetworkError( - "Request timed out — check your network connection and try again.".to_string(), - ) - })? - .map_err(|e| ProviderError::RequestFailed(e.to_string())) } } @@ -419,7 +413,7 @@ impl Provider for KimiCodeProvider { let response = self .with_retry(|| async { let resp = self.post(&payload).await?; - handle_status(resp, StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await + handle_status(resp).await }) .await .inspect_err(|e| { @@ -469,8 +463,7 @@ impl Provider for KimiCodeProvider { .send() .await .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; - let resp = - handle_status(resp, StdDuration::from_secs(DEFAULT_PROVIDER_TIMEOUT_SECS)).await?; + let resp = handle_status(resp).await?; let parsed: ModelsResp = resp.json().await.map_err(|e| { ProviderError::RequestFailed(format!("/v1/models body is not valid JSON: {}", e)) diff --git a/crates/goose/src/providers/litellm.rs b/crates/goose/src/providers/litellm.rs index dc85b4ee8903..e9caf70f573a 100644 --- a/crates/goose/src/providers/litellm.rs +++ b/crates/goose/src/providers/litellm.rs @@ -147,7 +147,7 @@ impl LiteLLMProvider { .model_headers(model_config)? .response_post(payload) .await?; - handle_response_openai_compat(response, self.api_client.timeout()).await + handle_response_openai_compat(response).await } async fn supports_cache_control(&self, model: &ModelConfig) -> bool { diff --git a/crates/goose/src/providers/nanogpt.rs b/crates/goose/src/providers/nanogpt.rs index 20a8112cbc35..2ced5d7b2893 100644 --- a/crates/goose/src/providers/nanogpt.rs +++ b/crates/goose/src/providers/nanogpt.rs @@ -201,7 +201,7 @@ impl Provider for NanoGptProvider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { diff --git a/crates/goose/src/providers/openrouter.rs b/crates/goose/src/providers/openrouter.rs index 589fdb672e0d..177a6356d2e8 100644 --- a/crates/goose/src/providers/openrouter.rs +++ b/crates/goose/src/providers/openrouter.rs @@ -353,7 +353,7 @@ impl Provider for OpenRouterProvider { .streaming(true) .response_post(&payload) .await?; - handle_status(resp, self.api_client.timeout()).await + handle_status(resp).await }) .await .inspect_err(|e| { diff --git a/crates/goose/src/providers/tetrate.rs b/crates/goose/src/providers/tetrate.rs index dfb6b653149c..0524a10e3822 100644 --- a/crates/goose/src/providers/tetrate.rs +++ b/crates/goose/src/providers/tetrate.rs @@ -157,7 +157,7 @@ impl Provider for TetrateProvider { .streaming(true) .response_post(&payload) .await?; - let resp = handle_status(resp, self.api_client.timeout()) + let resp = handle_status(resp) .await .map_err(Self::enrich_credits_error)?; @@ -169,12 +169,9 @@ impl Provider for TetrateProvider { .is_some_and(|v| v.contains("json")); if is_json { - let body = goose_providers::http_status::read_error_body( - resp, - self.api_client.timeout(), - ) - .await - .unwrap_or_default(); + let body = goose_providers::http_status::read_error_body(resp) + .await + .unwrap_or_default(); if let Ok(payload) = serde_json::from_str::(&body) { if payload.get("error").is_some() { return Err(Self::error_from_tetrate_error_payload( @@ -184,10 +181,9 @@ impl Provider for TetrateProvider { } } - return Err(ProviderError::ExecutionError( - "Expected streaming response but received non-streaming payload" - .to_string(), - )); + return Err(ProviderError::ExecutionError(format!( + "Expected streaming response but received non-streaming payload: {body}" + ))); } Ok(resp) @@ -207,7 +203,7 @@ impl Provider for TetrateProvider { .response_get("v1/models") .await .map_err(|e| ProviderError::RequestFailed(e.to_string()))?; - let json = handle_response_openai_compat(response, self.api_client.timeout()).await?; + let json = handle_response_openai_compat(response).await?; // Tetrate can return errors in 200 OK responses, so check explicitly if json.get("error").is_some() { From 0b064c19ecc50bc07d1c5ade08bd6b3f7acc4dff Mon Sep 17 00:00:00 2001 From: Douwe M Osinga Date: Tue, 4 Aug 2026 22:26:20 +0200 Subject: [PATCH 6/6] fix(providers): preserve errors across anyhow boundary --- crates/goose-provider-types/src/errors.rs | 3 +++ crates/goose-providers/src/api_client.rs | 4 ++-- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/crates/goose-provider-types/src/errors.rs b/crates/goose-provider-types/src/errors.rs index f2a50562aa3c..c8d099c29abc 100644 --- a/crates/goose-provider-types/src/errors.rs +++ b/crates/goose-provider-types/src/errors.rs @@ -131,6 +131,9 @@ fn provider_error_from_reqwest(error: &reqwest::Error) -> ProviderError { impl From for ProviderError { fn from(error: anyhow::Error) -> Self { + if let Some(provider_error) = error.downcast_ref::() { + return provider_error.clone(); + } if let Some(reqwest_err) = error.downcast_ref::() { return provider_error_from_reqwest(reqwest_err); } diff --git a/crates/goose-providers/src/api_client.rs b/crates/goose-providers/src/api_client.rs index 62213e62e7f5..ccb9946c0a91 100644 --- a/crates/goose-providers/src/api_client.rs +++ b/crates/goose-providers/src/api_client.rs @@ -799,8 +799,8 @@ mod tests { started.elapsed() ); assert!(matches!( - err.downcast_ref::(), - Some(crate::errors::ProviderError::NetworkError(message)) + crate::errors::ProviderError::from(err), + crate::errors::ProviderError::NetworkError(message) if message.starts_with("Request timed out") )); }