diff --git a/src/dav/client.rs b/src/dav/client.rs index c6ebe81..2d764de 100644 --- a/src/dav/client.rs +++ b/src/dav/client.rs @@ -14,7 +14,6 @@ use ureq::Agent; use ureq::Body; use ureq::config::{Config, RedirectAuthHeaders}; use ureq::http::{Method, Request, Response}; -use ureq::tls::{RootCerts, TlsConfig}; use crate::dav::parse::{ControlStrippingReader, DavResponse, parse_multistatus}; use crate::dav::retry::{DavOutcome, classify}; @@ -22,6 +21,7 @@ use crate::jmap::error::JmapError; use crate::jmap::http::{Auth, RetryPolicy, retry_after_header}; use crate::jmap::retry::{self, RateLimitState}; use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger}; +use crate::net::{tls, with_timeouts}; const MAX_BODY: u64 = 512 * 1024 * 1024; const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10); @@ -64,21 +64,15 @@ pub struct DavClient { impl DavClient { pub fn new(auth: Auth, retry: RetryPolicy, allow_invalid_certs: bool) -> Self { - let config: Config = Config::builder() - .http_status_as_error(false) - .allow_non_standard_methods(true) - .max_redirects(0) - .redirect_auth_headers(RedirectAuthHeaders::SameHost) - .tls_config( - TlsConfig::builder() - .unversioned_rustls_crypto_provider(std::sync::Arc::new( - rustls::crypto::aws_lc_rs::default_provider(), - )) - .root_certs(RootCerts::PlatformVerifier) - .disable_verification(allow_invalid_certs) - .build(), - ) - .build(); + let config: Config = with_timeouts!( + Config::builder() + .http_status_as_error(false) + .allow_non_standard_methods(true) + .max_redirects(0) + .redirect_auth_headers(RedirectAuthHeaders::SameHost) + .tls_config(tls(allow_invalid_certs)) + ) + .build(); DavClient { inner: Arc::new(Inner { agent: config.new_agent(), @@ -942,6 +936,24 @@ fn truncate(body: &[u8]) -> String { mod tests { use super::*; + #[test] + fn every_timeout_is_a_retryable_transport_error() { + for t in [ + ureq::Timeout::Connect, + ureq::Timeout::SendRequest, + ureq::Timeout::SendBody, + ureq::Timeout::RecvResponse, + ureq::Timeout::RecvBody, + ] { + let err = map_ureq_error(ureq::Error::Timeout(t)); + assert!(matches!(err, JmapError::Transport(_)), "{t:?} -> {err:?}"); + assert!( + matches!(transport_disposition(&err), retry::Disposition::Retryable), + "{t:?} must be retried" + ); + } + } + #[test] fn client_constructs_cleanly() { let c = DavClient::new( diff --git a/src/exchange_ews/autodiscover.rs b/src/exchange_ews/autodiscover.rs index cb28b4f..4975cf0 100644 --- a/src/exchange_ews/autodiscover.rs +++ b/src/exchange_ews/autodiscover.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -9,10 +10,10 @@ use quick_xml::events::Event; use serde_json::Value; use ureq::Agent; use ureq::config::Config; -use ureq::tls::{RootCerts, TlsConfig}; use crate::exchange_ews::error::EwsError; use crate::exchange_ews::parse::entity_to_char; +use crate::net::{tls, with_timeouts}; const V2_HOST: &str = "https://outlook.office365.com"; const POX_REQ_NS: &str = @@ -131,18 +132,12 @@ pub fn discover( } fn build_agent(allow_invalid_certs: bool) -> Agent { - let config: Config = Config::builder() - .http_status_as_error(false) - .tls_config( - TlsConfig::builder() - .unversioned_rustls_crypto_provider(std::sync::Arc::new( - rustls::crypto::aws_lc_rs::default_provider(), - )) - .root_certs(RootCerts::PlatformVerifier) - .disable_verification(allow_invalid_certs) - .build(), - ) - .build(); + let config: Config = with_timeouts!( + Config::builder() + .http_status_as_error(false) + .tls_config(tls(allow_invalid_certs)) + ) + .build(); config.new_agent() } diff --git a/src/exchange_ews/client.rs b/src/exchange_ews/client.rs index d7ab073..e988075 100644 --- a/src/exchange_ews/client.rs +++ b/src/exchange_ews/client.rs @@ -12,7 +12,6 @@ use std::time::{Duration, Instant}; use ureq::Agent; use ureq::config::{Config, RedirectAuthHeaders}; -use ureq::tls::{RootCerts, TlsConfig}; use crate::exchange_ews::error::EwsError; use crate::exchange_ews::parse::{EnvelopeKind, SoapFault, read_envelope_summary}; @@ -22,6 +21,7 @@ use crate::exchange_ews::types::ServerVersion; use crate::jmap::http::{Auth, RetryPolicy, retry_after_header}; use crate::jmap::retry::{self, Disposition, RateLimitState}; use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger}; +use crate::net::{tls, with_timeouts}; const MAX_BODY: u64 = 2 * 1024 * 1024 * 1024; const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10); @@ -55,19 +55,13 @@ pub struct SoapResponse { impl EwsClient { pub fn new(auth: Auth, retry: RetryPolicy, allow_invalid_certs: bool) -> EwsClient { - let config: Config = Config::builder() - .http_status_as_error(false) - .redirect_auth_headers(RedirectAuthHeaders::SameHost) - .tls_config( - TlsConfig::builder() - .unversioned_rustls_crypto_provider(std::sync::Arc::new( - rustls::crypto::aws_lc_rs::default_provider(), - )) - .root_certs(RootCerts::PlatformVerifier) - .disable_verification(allow_invalid_certs) - .build(), - ) - .build(); + let config: Config = with_timeouts!( + Config::builder() + .http_status_as_error(false) + .redirect_auth_headers(RedirectAuthHeaders::SameHost) + .tls_config(tls(allow_invalid_certs)) + ) + .build(); EwsClient { inner: Arc::new(Inner { agent: config.new_agent(), @@ -536,6 +530,21 @@ fn truncate(body: &[u8]) -> String { mod tests { use super::*; + #[test] + fn every_timeout_is_a_transport_error_and_so_retried() { + // Every EwsError::Transport goes round the retry loop in `execute`. + for t in [ + ureq::Timeout::Connect, + ureq::Timeout::SendRequest, + ureq::Timeout::SendBody, + ureq::Timeout::RecvResponse, + ureq::Timeout::RecvBody, + ] { + let err = map_ureq_error(ureq::Error::Timeout(t)); + assert!(matches!(err, EwsError::Transport(_)), "{t:?} -> {err:?}"); + } + } + #[test] fn client_constructs_with_defaults() { let c = EwsClient::new( diff --git a/src/exchange_ews/oauth.rs b/src/exchange_ews/oauth.rs index e00e935..8e320e2 100644 --- a/src/exchange_ews/oauth.rs +++ b/src/exchange_ews/oauth.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -10,9 +11,9 @@ use std::time::Duration; use encodify::base64::{Base64, Padding, URL_SAFE}; use serde_json::Value; use ureq::config::Config; -use ureq::tls::{RootCerts, TlsConfig}; use crate::exchange_ews::error::EwsError; +use crate::net::{tls, with_timeouts}; pub const SCOPE_APP_ONLY: &str = "https://outlook.office365.com/.default"; pub const SCOPE_DELEGATED: &str = @@ -105,18 +106,12 @@ fn device_code_endpoint(tenant: &str) -> String { } fn build_agent(allow_invalid_certs: bool) -> ureq::Agent { - let config: Config = Config::builder() - .http_status_as_error(false) - .tls_config( - TlsConfig::builder() - .unversioned_rustls_crypto_provider(std::sync::Arc::new( - rustls::crypto::aws_lc_rs::default_provider(), - )) - .root_certs(RootCerts::PlatformVerifier) - .disable_verification(allow_invalid_certs) - .build(), - ) - .build(); + let config: Config = with_timeouts!( + Config::builder() + .http_status_as_error(false) + .tls_config(tls(allow_invalid_certs)) + ) + .build(); config.new_agent() } diff --git a/src/exchange_graph/client.rs b/src/exchange_graph/client.rs index 40945ec..e79d72b 100644 --- a/src/exchange_graph/client.rs +++ b/src/exchange_graph/client.rs @@ -13,7 +13,6 @@ use std::time::{Duration, Instant}; use serde_json::Value; use ureq::Agent; use ureq::config::{Config, RedirectAuthHeaders}; -use ureq::tls::{RootCerts, TlsConfig}; use ureq::{ResponseExt, http::Uri}; use crate::exchange_graph::error::GraphError; @@ -21,6 +20,7 @@ use crate::exchange_graph::retry::{HttpClass, classify_http_status, is_throttled use crate::jmap::http::{RetryPolicy, cross_host, retry_after_header}; use crate::jmap::retry::{self, RateLimitState}; use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger}; +use crate::net::{tls, with_timeouts}; const MAX_BODY: u64 = 256 * 1024 * 1024; const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10); @@ -90,19 +90,13 @@ enum Attempt { impl GraphClient { pub fn new(bearer: String, retry: RetryPolicy, allow_invalid_certs: bool) -> GraphClient { - let config: Config = Config::builder() - .http_status_as_error(false) - .redirect_auth_headers(RedirectAuthHeaders::SameHost) - .tls_config( - TlsConfig::builder() - .unversioned_rustls_crypto_provider(std::sync::Arc::new( - rustls::crypto::aws_lc_rs::default_provider(), - )) - .root_certs(RootCerts::PlatformVerifier) - .disable_verification(allow_invalid_certs) - .build(), - ) - .build(); + let config: Config = with_timeouts!( + Config::builder() + .http_status_as_error(false) + .redirect_auth_headers(RedirectAuthHeaders::SameHost) + .tls_config(tls(allow_invalid_certs)) + ) + .build(); GraphClient { inner: Arc::new(Inner { agent: config.new_agent(), @@ -475,6 +469,21 @@ fn format_retry_wait(d: Duration) -> String { mod tests { use super::*; + #[test] + fn every_timeout_is_a_transport_error_and_so_retried() { + // `execute` retries every GraphError::Transport; only Connect is fatal. + for t in [ + ureq::Timeout::Connect, + ureq::Timeout::SendRequest, + ureq::Timeout::SendBody, + ureq::Timeout::RecvResponse, + ureq::Timeout::RecvBody, + ] { + let err = map_ureq_error(ureq::Error::Timeout(t)); + assert!(matches!(err, GraphError::Transport(_)), "{t:?} -> {err:?}"); + } + } + #[test] fn defaults_construct() { let c = GraphClient::new("token".to_owned(), RetryPolicy::new(3), false); diff --git a/src/exchange_graph/oauth.rs b/src/exchange_graph/oauth.rs index b240111..a1388ac 100644 --- a/src/exchange_graph/oauth.rs +++ b/src/exchange_graph/oauth.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -10,9 +11,9 @@ use std::time::{Duration, Instant}; use encodify::base64::{Base64, Padding, URL_SAFE}; use serde_json::Value; use ureq::config::Config; -use ureq::tls::{RootCerts, TlsConfig}; use crate::exchange_graph::error::GraphError; +use crate::net::{tls, with_timeouts}; pub const SCOPES: &str = "offline_access User.Read Mail.Read MailboxSettings.Read Calendars.Read Contacts.Read"; @@ -79,18 +80,12 @@ pub struct AcquiredToken { } fn build_agent(allow_invalid_certs: bool) -> ureq::Agent { - let config: Config = Config::builder() - .http_status_as_error(false) - .tls_config( - TlsConfig::builder() - .unversioned_rustls_crypto_provider(std::sync::Arc::new( - rustls::crypto::aws_lc_rs::default_provider(), - )) - .root_certs(RootCerts::PlatformVerifier) - .disable_verification(allow_invalid_certs) - .build(), - ) - .build(); + let config: Config = with_timeouts!( + Config::builder() + .http_status_as_error(false) + .tls_config(tls(allow_invalid_certs)) + ) + .build(); config.new_agent() } diff --git a/src/jmap/http.rs b/src/jmap/http.rs index 77c8ed8..31f84fc 100644 --- a/src/jmap/http.rs +++ b/src/jmap/http.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -13,7 +14,6 @@ use encodify::base64::STANDARD; use serde_json::Value; use ureq::Agent; use ureq::config::{Config, RedirectAuthHeaders}; -use ureq::tls::{RootCerts, TlsConfig}; use ureq::{ResponseExt, http::Uri}; use crate::jmap::error::JmapError; @@ -21,6 +21,7 @@ use crate::jmap::inflight::{Permit, Semaphore}; use crate::jmap::retry::{self, Disposition, RateLimitState}; use crate::jmap::session::Limits; use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger}; +use crate::net::{send_body_budget, tls, with_timeouts}; const MAX_BODY: u64 = 512 * 1024 * 1024; @@ -104,19 +105,13 @@ enum Attempt { impl HttpClient { pub fn new(auth: Auth, retry: RetryPolicy, allow_invalid_certs: bool) -> Self { - let config: Config = Config::builder() - .http_status_as_error(false) - .redirect_auth_headers(RedirectAuthHeaders::SameHost) - .tls_config( - TlsConfig::builder() - .unversioned_rustls_crypto_provider(std::sync::Arc::new( - rustls::crypto::aws_lc_rs::default_provider(), - )) - .root_certs(RootCerts::PlatformVerifier) - .disable_verification(allow_invalid_certs) - .build(), - ) - .build(); + let config: Config = with_timeouts!( + Config::builder() + .http_status_as_error(false) + .redirect_auth_headers(RedirectAuthHeaders::SameHost) + .tls_config(tls(allow_invalid_certs)) + ) + .build(); HttpClient { inner: Arc::new(Inner { agent: config.new_agent(), @@ -393,7 +388,12 @@ impl HttpClient { if let Some(ct) = content_type { req = req.header("Content-Type", ct); } - req.send(payload) + // A blob upload can run to hundreds of megabytes, so its send + // budget grows with its size instead of the agent's flat default. + req.config() + .timeout_send_body(Some(send_body_budget(payload.len()))) + .build() + .send(payload) } else { self.inner .agent @@ -645,6 +645,24 @@ pub fn format_retry_wait(d: Duration) -> String { mod tests { use super::*; + #[test] + fn every_timeout_is_a_retryable_transport_error() { + for t in [ + ureq::Timeout::Connect, + ureq::Timeout::SendRequest, + ureq::Timeout::SendBody, + ureq::Timeout::RecvResponse, + ureq::Timeout::RecvBody, + ] { + let err = map_ureq_error(ureq::Error::Timeout(t)); + assert!(matches!(err, JmapError::Transport(_)), "{t:?} -> {err:?}"); + assert!( + matches!(transport_disposition(&err), Disposition::Retryable), + "{t:?} must be retried" + ); + } + } + #[test] fn basic_header_matches_rfc7617_example() { let auth = Auth::Basic { diff --git a/src/lib.rs b/src/lib.rs index d20a878..6c40edf 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -16,6 +17,7 @@ pub mod inspect; pub mod jmap; pub mod logging; pub mod managesieve; +pub mod net; pub mod secret; pub mod sync; pub mod types; diff --git a/src/net.rs b/src/net.rs new file mode 100644 index 0000000..41e16c9 --- /dev/null +++ b/src/net.rs @@ -0,0 +1,99 @@ +/* + * SPDX-FileCopyrightText: 2026 John Coffey + * + * SPDX-License-Identifier: Apache-2.0 OR MIT + */ + +//! Settings every HTTP agent shares: timeouts and TLS. + +use std::time::Duration; + +use ureq::tls::{RootCerts, TlsConfig}; + +/// Opening the socket and completing any TLS handshake. +pub const CONNECT: Duration = Duration::from_secs(30); + +/// Writing the request line and headers. +pub const SEND_REQUEST: Duration = Duration::from_secs(60); + +/// Waiting for the response headers once the request is sent. This is the +/// server's thinking time: a large `Email/import`, an EWS `FindItem` over a big +/// folder or a CalDAV REPORT can legitimately take a while before the first +/// byte comes back. +pub const RECV_RESPONSE: Duration = Duration::from_secs(5 * 60); + +/// Reading the whole response body. ureq counts this as one budget for the +/// entire body, not per read, so it has to cover the largest body a client +/// accepts (512 MiB) on a slow link: 30 minutes is about 300 KB/s. A stalled +/// transfer is abandoned and retried after at most this long. +pub const RECV_BODY: Duration = Duration::from_secs(30 * 60); + +/// Sending a request body when its size is not known in advance. Uploads know +/// their size and get [`send_body_budget`] instead. +pub const SEND_BODY: Duration = Duration::from_secs(30 * 60); + +/// The slowest upload rate a send budget allows for, in bytes per second. +const MIN_UPLOAD_RATE: u64 = 64 * 1024; + +/// The floor under every send budget, so small bodies still get a sensible +/// allowance on a slow or busy connection. +const SEND_BODY_FLOOR: Duration = Duration::from_secs(2 * 60); + +/// How long sending a body of `len` bytes may take: the floor plus the time it +/// takes at [`MIN_UPLOAD_RATE`]. +pub fn send_body_budget(len: usize) -> Duration { + SEND_BODY_FLOOR + Duration::from_secs(len as u64 / MIN_UPLOAD_RATE) +} + +/// Applies the shared timeouts to a ureq `ConfigBuilder`. A macro rather than +/// a function because ureq keeps the builder's scope types private, so a +/// function could not name them. +macro_rules! with_timeouts { + ($builder:expr) => { + $builder + .timeout_connect(Some($crate::net::CONNECT)) + .timeout_send_request(Some($crate::net::SEND_REQUEST)) + .timeout_send_body(Some($crate::net::SEND_BODY)) + .timeout_recv_response(Some($crate::net::RECV_RESPONSE)) + .timeout_recv_body(Some($crate::net::RECV_BODY)) + }; +} +pub(crate) use with_timeouts; + +/// TLS settings for an agent: the platform's roots, and certificate checks off +/// only when `accept_invalid` is set. +pub fn tls(accept_invalid: bool) -> TlsConfig { + TlsConfig::builder() + .unversioned_rustls_crypto_provider(std::sync::Arc::new( + rustls::crypto::aws_lc_rs::default_provider(), + )) + .root_certs(RootCerts::PlatformVerifier) + .disable_verification(accept_invalid) + .build() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn send_budget_grows_with_size() { + assert_eq!(send_body_budget(0), Duration::from_secs(120)); + assert_eq!( + send_body_budget(64 * 1024 * 600), + Duration::from_secs(120 + 600) + ); + assert!(send_body_budget(512 * 1024 * 1024) > Duration::from_secs(2 * 60 * 60)); + } + + #[test] + fn timeouts_are_applied_to_a_config() { + let config: ureq::config::Config = with_timeouts!(ureq::config::Config::builder()).build(); + let t = config.timeouts(); + assert_eq!(t.connect, Some(CONNECT)); + assert_eq!(t.send_request, Some(SEND_REQUEST)); + assert_eq!(t.send_body, Some(SEND_BODY)); + assert_eq!(t.recv_response, Some(RECV_RESPONSE)); + assert_eq!(t.recv_body, Some(RECV_BODY)); + } +}