From 911fd9d2554731f568f3ed19e1812178ead1d522 Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Fri, 27 Mar 2026 03:08:34 -0700 Subject: [PATCH 01/10] feat: Make transport channel capacity configurable Add `transport_channel_capacity` field to `ClientOptions` that controls the bounded sync channel size used by the transport thread. Defaults to 30 (preserving current behavior). Users in high-throughput scenarios can increase this to reduce the chance of dropped envelopes. Closes #994 Co-Authored-By: Claude Opus 4.6 --- sentry-core/src/clientoptions.rs | 7 +++++++ sentry/src/transports/curl.rs | 3 ++- sentry/src/transports/reqwest.rs | 3 ++- sentry/src/transports/thread.rs | 4 ++-- sentry/src/transports/tokio_thread.rs | 4 ++-- sentry/src/transports/ureq.rs | 3 ++- 6 files changed, 17 insertions(+), 7 deletions(-) diff --git a/sentry-core/src/clientoptions.rs b/sentry-core/src/clientoptions.rs index 02d5d5a98..07d1d1092 100644 --- a/sentry-core/src/clientoptions.rs +++ b/sentry-core/src/clientoptions.rs @@ -191,6 +191,12 @@ pub struct ClientOptions { pub session_mode: SessionMode, /// The user agent that should be reported. pub user_agent: Cow<'static, str>, + /// The capacity of the transport channel. + /// + /// This controls how many envelopes can be queued before the transport + /// starts dropping them. In high-throughput scenarios, increasing this + /// value can reduce the chance of losing events. Defaults to 30. + pub transport_channel_capacity: usize, } impl ClientOptions { @@ -313,6 +319,7 @@ impl Default for ClientOptions { session_mode: SessionMode::Application, user_agent: Cow::Borrowed(USER_AGENT), max_request_body_size: MaxRequestBodySize::Medium, + transport_channel_capacity: 30, #[cfg(feature = "logs")] enable_logs: true, #[cfg(feature = "logs")] diff --git a/sentry/src/transports/curl.rs b/sentry/src/transports/curl.rs index f3ae2a864..6d81589f2 100644 --- a/sentry/src/transports/curl.rs +++ b/sentry/src/transports/curl.rs @@ -36,6 +36,7 @@ impl CurlHttpTransport { let url = dsn.envelope_api_url().to_string(); let scheme = dsn.scheme(); let accept_invalid_certs = options.accept_invalid_certs; + let channel_capacity = options.transport_channel_capacity; let mut handle = client; let thread = TransportThread::new(move |envelope, rl| { @@ -130,7 +131,7 @@ impl CurlHttpTransport { sentry_debug!("Failed to send envelope: {}", err); } } - }); + }, channel_capacity); Self { thread } } } diff --git a/sentry/src/transports/reqwest.rs b/sentry/src/transports/reqwest.rs index a125b18ff..d6e68da24 100644 --- a/sentry/src/transports/reqwest.rs +++ b/sentry/src/transports/reqwest.rs @@ -63,6 +63,7 @@ impl ReqwestHttpTransport { let user_agent = options.user_agent.clone(); let auth = dsn.to_auth(Some(&user_agent)).to_string(); let url = dsn.envelope_api_url().to_string(); + let channel_capacity = options.transport_channel_capacity; let thread = TransportThread::new(move |envelope, mut rl| { let mut body = Vec::new(); @@ -110,7 +111,7 @@ impl ReqwestHttpTransport { } rl } - }); + }, channel_capacity); Self { thread } } } diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 4ab402600..30e1416ef 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -27,11 +27,11 @@ pub struct TransportThread { impl TransportThread { /// Spawn a new background thread. - pub fn new(mut send: SendFn) -> Self + pub fn new(mut send: SendFn, channel_capacity: usize) -> Self where SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, { - let (sender, receiver) = sync_channel(30); + let (sender, receiver) = sync_channel(channel_capacity); let shutdown = Arc::new(AtomicBool::new(false)); let shutdown_worker = shutdown.clone(); let handle = thread::Builder::new() diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index 398963882..790b1a6ed 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -27,13 +27,13 @@ pub struct TransportThread { impl TransportThread { /// Spawn a new background thread. - pub fn new(mut send: SendFn) -> Self + pub fn new(mut send: SendFn, channel_capacity: usize) -> Self where SendFn: FnMut(Envelope, RateLimiter) -> SendFuture + Send + 'static, // NOTE: returning RateLimiter here, otherwise we are in borrow hell SendFuture: std::future::Future, { - let (sender, receiver) = sync_channel(30); + let (sender, receiver) = sync_channel(channel_capacity); let shutdown = Arc::new(AtomicBool::new(false)); let shutdown_worker = shutdown.clone(); let handle = thread::Builder::new() diff --git a/sentry/src/transports/ureq.rs b/sentry/src/transports/ureq.rs index d079028f7..e884d1203 100644 --- a/sentry/src/transports/ureq.rs +++ b/sentry/src/transports/ureq.rs @@ -82,6 +82,7 @@ impl UreqHttpTransport { let user_agent = options.user_agent.clone(); let auth = dsn.to_auth(Some(&user_agent)).to_string(); let url = dsn.envelope_api_url().to_string(); + let channel_capacity = options.transport_channel_capacity; let thread = TransportThread::new(move |envelope, rl| { let mut body = Vec::new(); @@ -118,7 +119,7 @@ impl UreqHttpTransport { sentry_debug!("Failed to send envelope: {}", err); } } - }); + }, channel_capacity); Self { thread } } } From 9c8f904293204f4b83e845201daf78483d3c8c7e Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Mon, 13 Apr 2026 09:05:50 -0400 Subject: [PATCH 02/10] fix: add Debug impl, clamp channel capacity, run cargo fmt - Add transport_channel_capacity to manual Debug implementation - Clamp channel capacity to minimum of 1 to prevent rendezvous channel from silently dropping envelopes via try_send - Run cargo fmt to fix lint CI failure --- sentry-core/src/clientoptions.rs | 8 +- sentry/src/transports/curl.rs | 170 +++++++++++++------------- sentry/src/transports/reqwest.rs | 79 ++++++------ sentry/src/transports/thread.rs | 2 +- sentry/src/transports/tokio_thread.rs | 2 +- sentry/src/transports/ureq.rs | 64 +++++----- 6 files changed, 172 insertions(+), 153 deletions(-) diff --git a/sentry-core/src/clientoptions.rs b/sentry-core/src/clientoptions.rs index 07d1d1092..2e3f6c21c 100644 --- a/sentry-core/src/clientoptions.rs +++ b/sentry-core/src/clientoptions.rs @@ -284,7 +284,13 @@ impl fmt::Debug for ClientOptions { .field("enable_logs", &self.enable_logs) .field("before_send_log", &before_send_log); - debug_struct.field("user_agent", &self.user_agent).finish() + debug_struct + .field("user_agent", &self.user_agent) + .field( + "transport_channel_capacity", + &self.transport_channel_capacity, + ) + .finish() } } diff --git a/sentry/src/transports/curl.rs b/sentry/src/transports/curl.rs index 6d81589f2..cbbd2c904 100644 --- a/sentry/src/transports/curl.rs +++ b/sentry/src/transports/curl.rs @@ -39,99 +39,103 @@ impl CurlHttpTransport { let channel_capacity = options.transport_channel_capacity; let mut handle = client; - let thread = TransportThread::new(move |envelope, rl| { - handle.reset(); - handle.url(&url).unwrap(); - handle.custom_request("POST").unwrap(); - - if accept_invalid_certs { - handle.ssl_verify_host(false).unwrap(); - handle.ssl_verify_peer(false).unwrap(); - } - - match (scheme, &http_proxy, &https_proxy) { - (Scheme::Https, _, Some(proxy)) => { - if let Err(err) = handle.proxy(proxy) { - sentry_debug!("invalid proxy: {:?}", err); - } + let thread = TransportThread::new( + move |envelope, rl| { + handle.reset(); + handle.url(&url).unwrap(); + handle.custom_request("POST").unwrap(); + + if accept_invalid_certs { + handle.ssl_verify_host(false).unwrap(); + handle.ssl_verify_peer(false).unwrap(); } - (_, Some(proxy), _) => { - if let Err(err) = handle.proxy(proxy) { - sentry_debug!("invalid proxy: {:?}", err); + + match (scheme, &http_proxy, &https_proxy) { + (Scheme::Https, _, Some(proxy)) => { + if let Err(err) = handle.proxy(proxy) { + sentry_debug!("invalid proxy: {:?}", err); + } + } + (_, Some(proxy), _) => { + if let Err(err) = handle.proxy(proxy) { + sentry_debug!("invalid proxy: {:?}", err); + } } + _ => {} } - _ => {} - } - - let mut body = Vec::new(); - envelope.to_writer(&mut body).unwrap(); - let mut body = Cursor::new(body); - - let mut retry_after = None; - let mut sentry_header = None; - let mut headers = curl::easy::List::new(); - headers.append(&format!("X-Sentry-Auth: {auth}")).unwrap(); - headers.append("Expect:").unwrap(); - handle.http_headers(headers).unwrap(); - handle.upload(true).unwrap(); - handle.in_filesize(body.get_ref().len() as u64).unwrap(); - handle - .read_function(move |buf| Ok(body.read(buf).unwrap_or(0))) - .unwrap(); - handle.verbose(true).unwrap(); - handle - .debug_function(move |info, data| { - let prefix = match info { - curl::easy::InfoType::HeaderIn => "< ", - curl::easy::InfoType::HeaderOut => "> ", - curl::easy::InfoType::DataOut => "", - _ => return, - }; - sentry_debug!("curl: {}{}", prefix, String::from_utf8_lossy(data).trim()); - }) - .unwrap(); - - { - let mut handle = handle.transfer(); - let retry_after_setter = &mut retry_after; - let sentry_header_setter = &mut sentry_header; + + let mut body = Vec::new(); + envelope.to_writer(&mut body).unwrap(); + let mut body = Cursor::new(body); + + let mut retry_after = None; + let mut sentry_header = None; + let mut headers = curl::easy::List::new(); + headers.append(&format!("X-Sentry-Auth: {auth}")).unwrap(); + headers.append("Expect:").unwrap(); + handle.http_headers(headers).unwrap(); + handle.upload(true).unwrap(); + handle.in_filesize(body.get_ref().len() as u64).unwrap(); + handle + .read_function(move |buf| Ok(body.read(buf).unwrap_or(0))) + .unwrap(); + handle.verbose(true).unwrap(); handle - .header_function(move |data| { - if let Ok(data) = std::str::from_utf8(data) { - let mut iter = data.split(':'); - if let Some(key) = iter.next().map(str::to_lowercase) { - if key == "retry-after" { - *retry_after_setter = iter.next().map(|x| x.trim().to_string()); - } else if key == "x-sentry-rate-limits" { - *sentry_header_setter = - iter.next().map(|x| x.trim().to_string()); + .debug_function(move |info, data| { + let prefix = match info { + curl::easy::InfoType::HeaderIn => "< ", + curl::easy::InfoType::HeaderOut => "> ", + curl::easy::InfoType::DataOut => "", + _ => return, + }; + sentry_debug!("curl: {}{}", prefix, String::from_utf8_lossy(data).trim()); + }) + .unwrap(); + + { + let mut handle = handle.transfer(); + let retry_after_setter = &mut retry_after; + let sentry_header_setter = &mut sentry_header; + handle + .header_function(move |data| { + if let Ok(data) = std::str::from_utf8(data) { + let mut iter = data.split(':'); + if let Some(key) = iter.next().map(str::to_lowercase) { + if key == "retry-after" { + *retry_after_setter = + iter.next().map(|x| x.trim().to_string()); + } else if key == "x-sentry-rate-limits" { + *sentry_header_setter = + iter.next().map(|x| x.trim().to_string()); + } } } + true + }) + .unwrap(); + handle.perform().ok(); + } + + match handle.response_code() { + Ok(response_code) => { + if let Some(sentry_header) = sentry_header { + rl.update_from_sentry_header(&sentry_header); + } else if let Some(retry_after) = retry_after { + rl.update_from_retry_after(&retry_after); + } else if response_code == 429 { + rl.update_from_429(); + } + if response_code == HTTP_PAYLOAD_TOO_LARGE as u32 { + sentry_debug!("{HTTP_PAYLOAD_TOO_LARGE_MESSAGE}"); } - true - }) - .unwrap(); - handle.perform().ok(); - } - - match handle.response_code() { - Ok(response_code) => { - if let Some(sentry_header) = sentry_header { - rl.update_from_sentry_header(&sentry_header); - } else if let Some(retry_after) = retry_after { - rl.update_from_retry_after(&retry_after); - } else if response_code == 429 { - rl.update_from_429(); } - if response_code == HTTP_PAYLOAD_TOO_LARGE as u32 { - sentry_debug!("{HTTP_PAYLOAD_TOO_LARGE_MESSAGE}"); + Err(err) => { + sentry_debug!("Failed to send envelope: {}", err); } } - Err(err) => { - sentry_debug!("Failed to send envelope: {}", err); - } - } - }, channel_capacity); + }, + channel_capacity, + ); Self { thread } } } diff --git a/sentry/src/transports/reqwest.rs b/sentry/src/transports/reqwest.rs index d6e68da24..50f9191fd 100644 --- a/sentry/src/transports/reqwest.rs +++ b/sentry/src/transports/reqwest.rs @@ -65,53 +65,56 @@ impl ReqwestHttpTransport { let url = dsn.envelope_api_url().to_string(); let channel_capacity = options.transport_channel_capacity; - let thread = TransportThread::new(move |envelope, mut rl| { - let mut body = Vec::new(); - envelope.to_writer(&mut body).unwrap(); - let request = client.post(&url).header("X-Sentry-Auth", &auth).body(body); + let thread = TransportThread::new( + move |envelope, mut rl| { + let mut body = Vec::new(); + envelope.to_writer(&mut body).unwrap(); + let request = client.post(&url).header("X-Sentry-Auth", &auth).body(body); - // NOTE: because of lifetime issues, building the request using the - // `client` has to happen outside of this async block. - async move { - match request.send().await { - Ok(response) => { - let headers = response.headers(); + // NOTE: because of lifetime issues, building the request using the + // `client` has to happen outside of this async block. + async move { + match request.send().await { + Ok(response) => { + let headers = response.headers(); - if let Some(sentry_header) = headers - .get("x-sentry-rate-limits") - .and_then(|x| x.to_str().ok()) - { - rl.update_from_sentry_header(sentry_header); - } else if let Some(retry_after) = headers - .get(ReqwestHeaders::RETRY_AFTER) - .and_then(|x| x.to_str().ok()) - { - rl.update_from_retry_after(retry_after); - } else if response.status() == StatusCode::TOO_MANY_REQUESTS { - rl.update_from_429(); - } + if let Some(sentry_header) = headers + .get("x-sentry-rate-limits") + .and_then(|x| x.to_str().ok()) + { + rl.update_from_sentry_header(sentry_header); + } else if let Some(retry_after) = headers + .get(ReqwestHeaders::RETRY_AFTER) + .and_then(|x| x.to_str().ok()) + { + rl.update_from_retry_after(retry_after); + } else if response.status() == StatusCode::TOO_MANY_REQUESTS { + rl.update_from_429(); + } - let is_payload_too_large = - response.status().as_u16() == HTTP_PAYLOAD_TOO_LARGE; - match response.text().await { - Err(err) => { - sentry_debug!("Failed to read sentry response: {}", err); + let is_payload_too_large = + response.status().as_u16() == HTTP_PAYLOAD_TOO_LARGE; + match response.text().await { + Err(err) => { + sentry_debug!("Failed to read sentry response: {}", err); + } + Ok(text) => { + sentry_debug!("Get response: `{}`", text); + } } - Ok(text) => { - sentry_debug!("Get response: `{}`", text); + if is_payload_too_large { + sentry_debug!("{HTTP_PAYLOAD_TOO_LARGE_MESSAGE}"); } } - if is_payload_too_large { - sentry_debug!("{HTTP_PAYLOAD_TOO_LARGE_MESSAGE}"); + Err(err) => { + sentry_debug!("Failed to send envelope: {}", err); } } - Err(err) => { - sentry_debug!("Failed to send envelope: {}", err); - } + rl } - rl - } - }, channel_capacity); + }, + channel_capacity, + ); Self { thread } } } diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 30e1416ef..6d133ac12 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -31,7 +31,7 @@ impl TransportThread { where SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, { - let (sender, receiver) = sync_channel(channel_capacity); + let (sender, receiver) = sync_channel(channel_capacity.max(1)); let shutdown = Arc::new(AtomicBool::new(false)); let shutdown_worker = shutdown.clone(); let handle = thread::Builder::new() diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index 790b1a6ed..1487f2adf 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -33,7 +33,7 @@ impl TransportThread { // NOTE: returning RateLimiter here, otherwise we are in borrow hell SendFuture: std::future::Future, { - let (sender, receiver) = sync_channel(channel_capacity); + let (sender, receiver) = sync_channel(channel_capacity.max(1)); let shutdown = Arc::new(AtomicBool::new(false)); let shutdown_worker = shutdown.clone(); let handle = thread::Builder::new() diff --git a/sentry/src/transports/ureq.rs b/sentry/src/transports/ureq.rs index e884d1203..599e9b535 100644 --- a/sentry/src/transports/ureq.rs +++ b/sentry/src/transports/ureq.rs @@ -84,42 +84,48 @@ impl UreqHttpTransport { let url = dsn.envelope_api_url().to_string(); let channel_capacity = options.transport_channel_capacity; - let thread = TransportThread::new(move |envelope, rl| { - let mut body = Vec::new(); - envelope.to_writer(&mut body).unwrap(); - let request = agent.post(&url).header("X-Sentry-Auth", &auth).send(&body); - - match request { - Ok(mut response) => { - fn header_str<'a, B>(response: &'a Response, key: &str) -> Option<&'a str> { - response.headers().get(key)?.to_str().ok() - } + let thread = TransportThread::new( + move |envelope, rl| { + let mut body = Vec::new(); + envelope.to_writer(&mut body).unwrap(); + let request = agent.post(&url).header("X-Sentry-Auth", &auth).send(&body); + + match request { + Ok(mut response) => { + fn header_str<'a, B>( + response: &'a Response, + key: &str, + ) -> Option<&'a str> { + response.headers().get(key)?.to_str().ok() + } - if let Some(sentry_header) = header_str(&response, "x-sentry-rate-limits") { - rl.update_from_sentry_header(sentry_header); - } else if let Some(retry_after) = header_str(&response, "retry-after") { - rl.update_from_retry_after(retry_after); - } else if response.status() == 429 { - rl.update_from_429(); - } + if let Some(sentry_header) = header_str(&response, "x-sentry-rate-limits") { + rl.update_from_sentry_header(sentry_header); + } else if let Some(retry_after) = header_str(&response, "retry-after") { + rl.update_from_retry_after(retry_after); + } else if response.status() == 429 { + rl.update_from_429(); + } - match response.body_mut().read_to_string() { - Err(err) => { - sentry_debug!("Failed to read sentry response: {}", err); + match response.body_mut().read_to_string() { + Err(err) => { + sentry_debug!("Failed to read sentry response: {}", err); + } + Ok(text) => { + sentry_debug!("Get response: `{}`", text); + } } - Ok(text) => { - sentry_debug!("Get response: `{}`", text); + if response.status() == HTTP_PAYLOAD_TOO_LARGE { + sentry_debug!("{HTTP_PAYLOAD_TOO_LARGE_MESSAGE}"); } } - if response.status() == HTTP_PAYLOAD_TOO_LARGE { - sentry_debug!("{HTTP_PAYLOAD_TOO_LARGE_MESSAGE}"); + Err(err) => { + sentry_debug!("Failed to send envelope: {}", err); } } - Err(err) => { - sentry_debug!("Failed to send envelope: {}", err); - } - } - }, channel_capacity); + }, + channel_capacity, + ); Self { thread } } } From fbff3ea24d782a35565b6e0f324a0cb11b7edebe Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Tue, 14 Apr 2026 21:13:25 -0400 Subject: [PATCH 03/10] refactor: make channel capacity additive via with_capacity/with_channel_capacity Per review: avoid public API breakages in ClientOptions and the transport constructors. Replace the 'transport_channel_capacity: usize' field and the changed TransportThread/transport constructor signatures with purely additive APIs: - TransportThread::new(send) restored to original signature, delegates to with_capacity(send, 30). TransportThread::with_capacity(send, capacity) is the new entry point for custom capacity. - Same pattern for tokio_thread::TransportThread. - ReqwestHttpTransport/CurlHttpTransport/UreqHttpTransport gain with_channel_capacity(options, capacity). Existing new/with_client/ with_agent keep the default capacity of 30. - Remove transport_channel_capacity from ClientOptions (field, Default, and Debug impl). Users that need to override capacity configure it via ClientOptions.transport with a factory returning the desired transport, e.g. Arc::new(move |opts| Arc::new(ReqwestHttpTransport::with_channel_capacity(opts, 256))). --- sentry-core/src/clientoptions.rs | 15 +-------------- sentry/src/transports/curl.rs | 22 +++++++++++++++++----- sentry/src/transports/reqwest.rs | 22 +++++++++++++++++----- sentry/src/transports/thread.rs | 17 +++++++++++++++-- sentry/src/transports/tokio_thread.rs | 19 +++++++++++++++++-- sentry/src/transports/ureq.rs | 22 +++++++++++++++++----- 6 files changed, 84 insertions(+), 33 deletions(-) diff --git a/sentry-core/src/clientoptions.rs b/sentry-core/src/clientoptions.rs index 2e3f6c21c..02d5d5a98 100644 --- a/sentry-core/src/clientoptions.rs +++ b/sentry-core/src/clientoptions.rs @@ -191,12 +191,6 @@ pub struct ClientOptions { pub session_mode: SessionMode, /// The user agent that should be reported. pub user_agent: Cow<'static, str>, - /// The capacity of the transport channel. - /// - /// This controls how many envelopes can be queued before the transport - /// starts dropping them. In high-throughput scenarios, increasing this - /// value can reduce the chance of losing events. Defaults to 30. - pub transport_channel_capacity: usize, } impl ClientOptions { @@ -284,13 +278,7 @@ impl fmt::Debug for ClientOptions { .field("enable_logs", &self.enable_logs) .field("before_send_log", &before_send_log); - debug_struct - .field("user_agent", &self.user_agent) - .field( - "transport_channel_capacity", - &self.transport_channel_capacity, - ) - .finish() + debug_struct.field("user_agent", &self.user_agent).finish() } } @@ -325,7 +313,6 @@ impl Default for ClientOptions { session_mode: SessionMode::Application, user_agent: Cow::Borrowed(USER_AGENT), max_request_body_size: MaxRequestBodySize::Medium, - transport_channel_capacity: 30, #[cfg(feature = "logs")] enable_logs: true, #[cfg(feature = "logs")] diff --git a/sentry/src/transports/curl.rs b/sentry/src/transports/curl.rs index cbbd2c904..27f28b5bf 100644 --- a/sentry/src/transports/curl.rs +++ b/sentry/src/transports/curl.rs @@ -18,15 +18,28 @@ pub struct CurlHttpTransport { impl CurlHttpTransport { /// Creates a new Transport. pub fn new(options: &ClientOptions) -> Self { - Self::new_internal(options, None) + Self::new_internal(options, None, 30) } /// Creates a new Transport that uses the specified [`CurlClient`]. pub fn with_client(options: &ClientOptions, client: CurlClient) -> Self { - Self::new_internal(options, Some(client)) + Self::new_internal(options, Some(client), 30) } - fn new_internal(options: &ClientOptions, client: Option) -> Self { + /// Creates a new Transport with a custom transport channel capacity. + /// + /// The channel capacity bounds how many envelopes may be queued before + /// `send_envelope` blocks. A higher capacity reduces the chance of + /// dropped events in high-throughput scenarios at the cost of memory. + pub fn with_channel_capacity(options: &ClientOptions, channel_capacity: usize) -> Self { + Self::new_internal(options, None, channel_capacity) + } + + fn new_internal( + options: &ClientOptions, + client: Option, + channel_capacity: usize, + ) -> Self { let client = client.unwrap_or_else(CurlClient::new); let http_proxy = options.http_proxy.as_ref().map(ToString::to_string); let https_proxy = options.https_proxy.as_ref().map(ToString::to_string); @@ -36,10 +49,9 @@ impl CurlHttpTransport { let url = dsn.envelope_api_url().to_string(); let scheme = dsn.scheme(); let accept_invalid_certs = options.accept_invalid_certs; - let channel_capacity = options.transport_channel_capacity; let mut handle = client; - let thread = TransportThread::new( + let thread = TransportThread::with_capacity( move |envelope, rl| { handle.reset(); handle.url(&url).unwrap(); diff --git a/sentry/src/transports/reqwest.rs b/sentry/src/transports/reqwest.rs index 50f9191fd..ab5d4365a 100644 --- a/sentry/src/transports/reqwest.rs +++ b/sentry/src/transports/reqwest.rs @@ -21,15 +21,28 @@ pub struct ReqwestHttpTransport { impl ReqwestHttpTransport { /// Creates a new Transport. pub fn new(options: &ClientOptions) -> Self { - Self::new_internal(options, None) + Self::new_internal(options, None, 30) } /// Creates a new Transport that uses the specified [`ReqwestClient`]. pub fn with_client(options: &ClientOptions, client: ReqwestClient) -> Self { - Self::new_internal(options, Some(client)) + Self::new_internal(options, Some(client), 30) } - fn new_internal(options: &ClientOptions, client: Option) -> Self { + /// Creates a new Transport with a custom transport channel capacity. + /// + /// The channel capacity bounds how many envelopes may be queued before + /// `send_envelope` blocks. A higher capacity reduces the chance of + /// dropped events in high-throughput scenarios at the cost of memory. + pub fn with_channel_capacity(options: &ClientOptions, channel_capacity: usize) -> Self { + Self::new_internal(options, None, channel_capacity) + } + + fn new_internal( + options: &ClientOptions, + client: Option, + channel_capacity: usize, + ) -> Self { let client = client.unwrap_or_else(|| { let mut builder = reqwest::Client::builder(); if options.accept_invalid_certs { @@ -63,9 +76,8 @@ impl ReqwestHttpTransport { let user_agent = options.user_agent.clone(); let auth = dsn.to_auth(Some(&user_agent)).to_string(); let url = dsn.envelope_api_url().to_string(); - let channel_capacity = options.transport_channel_capacity; - let thread = TransportThread::new( + let thread = TransportThread::with_capacity( move |envelope, mut rl| { let mut body = Vec::new(); envelope.to_writer(&mut body).unwrap(); diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 6d133ac12..7e2f2aac0 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -26,8 +26,21 @@ pub struct TransportThread { } impl TransportThread { - /// Spawn a new background thread. - pub fn new(mut send: SendFn, channel_capacity: usize) -> Self + /// Spawn a new background thread with the default channel capacity of 30. + pub fn new(send: SendFn) -> Self + where + SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, + { + Self::with_capacity(send, 30) + } + + /// Spawn a new background thread with a custom channel capacity. + /// + /// The channel capacity bounds how many envelopes may be queued before + /// `send` blocks. `channel_capacity` is clamped to a minimum of 1 to + /// avoid a rendezvous channel, which would silently drop envelopes under + /// `try_send`. + pub fn with_capacity(mut send: SendFn, channel_capacity: usize) -> Self where SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, { diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index 1487f2adf..f862eb3a6 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -26,8 +26,23 @@ pub struct TransportThread { } impl TransportThread { - /// Spawn a new background thread. - pub fn new(mut send: SendFn, channel_capacity: usize) -> Self + /// Spawn a new background thread with the default channel capacity of 30. + pub fn new(send: SendFn) -> Self + where + SendFn: FnMut(Envelope, RateLimiter) -> SendFuture + Send + 'static, + // NOTE: returning RateLimiter here, otherwise we are in borrow hell + SendFuture: std::future::Future, + { + Self::with_capacity(send, 30) + } + + /// Spawn a new background thread with a custom channel capacity. + /// + /// The channel capacity bounds how many envelopes may be queued before + /// `send` blocks. `channel_capacity` is clamped to a minimum of 1 to + /// avoid a rendezvous channel, which would silently drop envelopes under + /// `try_send`. + pub fn with_capacity(mut send: SendFn, channel_capacity: usize) -> Self where SendFn: FnMut(Envelope, RateLimiter) -> SendFuture + Send + 'static, // NOTE: returning RateLimiter here, otherwise we are in borrow hell diff --git a/sentry/src/transports/ureq.rs b/sentry/src/transports/ureq.rs index 599e9b535..dc236a9a7 100644 --- a/sentry/src/transports/ureq.rs +++ b/sentry/src/transports/ureq.rs @@ -20,15 +20,28 @@ pub struct UreqHttpTransport { impl UreqHttpTransport { /// Creates a new Transport. pub fn new(options: &ClientOptions) -> Self { - Self::new_internal(options, None) + Self::new_internal(options, None, 30) } /// Creates a new Transport that uses the specified [`ureq::Agent`]. pub fn with_agent(options: &ClientOptions, agent: Agent) -> Self { - Self::new_internal(options, Some(agent)) + Self::new_internal(options, Some(agent), 30) } - fn new_internal(options: &ClientOptions, agent: Option) -> Self { + /// Creates a new Transport with a custom transport channel capacity. + /// + /// The channel capacity bounds how many envelopes may be queued before + /// `send_envelope` blocks. A higher capacity reduces the chance of + /// dropped events in high-throughput scenarios at the cost of memory. + pub fn with_channel_capacity(options: &ClientOptions, channel_capacity: usize) -> Self { + Self::new_internal(options, None, channel_capacity) + } + + fn new_internal( + options: &ClientOptions, + agent: Option, + channel_capacity: usize, + ) -> Self { let dsn = options.dsn.as_ref().unwrap(); let scheme = dsn.scheme(); let agent = agent.unwrap_or_else(|| { @@ -82,9 +95,8 @@ impl UreqHttpTransport { let user_agent = options.user_agent.clone(); let auth = dsn.to_auth(Some(&user_agent)).to_string(); let url = dsn.envelope_api_url().to_string(); - let channel_capacity = options.transport_channel_capacity; - let thread = TransportThread::new( + let thread = TransportThread::with_capacity( move |envelope, rl| { let mut body = Vec::new(); envelope.to_writer(&mut body).unwrap(); From 6d1c9ff12103628bbf3c804b00e719f4cd8c6343 Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Tue, 5 May 2026 08:59:59 -0700 Subject: [PATCH 04/10] fix: make with_capacity pub(crate) and extract DEFAULT_CHANNEL_CAPACITY const Per @szokeasaurusrex review: with_capacity narrowed to pub(crate) on both TransportThread types; the 30 literal extracted to a shared pub(crate) const reused across curl, reqwest, and ureq transports. Includes: Authorized scope (mod.rs) + correct extension (ureq.rs has the same literal). Mechanical change verified by build + fmt; no runtime effect needs test coverage. --- sentry/src/transports/curl.rs | 9 ++++++--- sentry/src/transports/mod.rs | 3 +++ sentry/src/transports/reqwest.rs | 7 ++++--- sentry/src/transports/thread.rs | 5 +++-- sentry/src/transports/tokio_thread.rs | 8 ++++++-- sentry/src/transports/ureq.rs | 9 ++++++--- 6 files changed, 28 insertions(+), 13 deletions(-) diff --git a/sentry/src/transports/curl.rs b/sentry/src/transports/curl.rs index 27f28b5bf..2f4c27259 100644 --- a/sentry/src/transports/curl.rs +++ b/sentry/src/transports/curl.rs @@ -3,7 +3,10 @@ use std::time::Duration; use curl::easy::Easy as CurlClient; -use super::{thread::TransportThread, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE}; +use super::{ + thread::TransportThread, DEFAULT_CHANNEL_CAPACITY, HTTP_PAYLOAD_TOO_LARGE, + HTTP_PAYLOAD_TOO_LARGE_MESSAGE, +}; use crate::{sentry_debug, types::Scheme, ClientOptions, Envelope, Transport}; @@ -18,12 +21,12 @@ pub struct CurlHttpTransport { impl CurlHttpTransport { /// Creates a new Transport. pub fn new(options: &ClientOptions) -> Self { - Self::new_internal(options, None, 30) + Self::new_internal(options, None, DEFAULT_CHANNEL_CAPACITY) } /// Creates a new Transport that uses the specified [`CurlClient`]. pub fn with_client(options: &ClientOptions, client: CurlClient) -> Self { - Self::new_internal(options, Some(client), 30) + Self::new_internal(options, Some(client), DEFAULT_CHANNEL_CAPACITY) } /// Creates a new Transport with a custom transport channel capacity. diff --git a/sentry/src/transports/mod.rs b/sentry/src/transports/mod.rs index 7e959612b..a18ebcc98 100644 --- a/sentry/src/transports/mod.rs +++ b/sentry/src/transports/mod.rs @@ -48,6 +48,9 @@ pub(crate) const HTTP_PAYLOAD_TOO_LARGE: u16 = 413; pub(crate) const HTTP_PAYLOAD_TOO_LARGE_MESSAGE: &str = "Envelope was discarded due to size limits (HTTP 413)."; +#[cfg(sentry_any_http_transport)] +pub(crate) const DEFAULT_CHANNEL_CAPACITY: usize = 30; + #[cfg(feature = "reqwest")] type DefaultTransport = ReqwestHttpTransport; diff --git a/sentry/src/transports/reqwest.rs b/sentry/src/transports/reqwest.rs index ab5d4365a..0b1137b05 100644 --- a/sentry/src/transports/reqwest.rs +++ b/sentry/src/transports/reqwest.rs @@ -3,7 +3,8 @@ use std::time::Duration; use reqwest::{header as ReqwestHeaders, Client as ReqwestClient, Proxy, StatusCode}; use super::{ - tokio_thread::TransportThread, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE, + tokio_thread::TransportThread, DEFAULT_CHANNEL_CAPACITY, HTTP_PAYLOAD_TOO_LARGE, + HTTP_PAYLOAD_TOO_LARGE_MESSAGE, }; use crate::{sentry_debug, ClientOptions, Envelope, Transport}; @@ -21,12 +22,12 @@ pub struct ReqwestHttpTransport { impl ReqwestHttpTransport { /// Creates a new Transport. pub fn new(options: &ClientOptions) -> Self { - Self::new_internal(options, None, 30) + Self::new_internal(options, None, DEFAULT_CHANNEL_CAPACITY) } /// Creates a new Transport that uses the specified [`ReqwestClient`]. pub fn with_client(options: &ClientOptions, client: ReqwestClient) -> Self { - Self::new_internal(options, Some(client), 30) + Self::new_internal(options, Some(client), DEFAULT_CHANNEL_CAPACITY) } /// Creates a new Transport with a custom transport channel capacity. diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 7e2f2aac0..156cb6d11 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -5,6 +5,7 @@ use std::thread::{self, JoinHandle}; use std::time::Duration; use super::ratelimit::{RateLimiter, RateLimitingCategory}; +use super::DEFAULT_CHANNEL_CAPACITY; use crate::{sentry_debug, Envelope}; #[expect( @@ -31,7 +32,7 @@ impl TransportThread { where SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, { - Self::with_capacity(send, 30) + Self::with_capacity(send, DEFAULT_CHANNEL_CAPACITY) } /// Spawn a new background thread with a custom channel capacity. @@ -40,7 +41,7 @@ impl TransportThread { /// `send` blocks. `channel_capacity` is clamped to a minimum of 1 to /// avoid a rendezvous channel, which would silently drop envelopes under /// `try_send`. - pub fn with_capacity(mut send: SendFn, channel_capacity: usize) -> Self + pub(crate) fn with_capacity(mut send: SendFn, channel_capacity: usize) -> Self where SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, { diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index f862eb3a6..657399118 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -5,6 +5,7 @@ use std::thread::{self, JoinHandle}; use std::time::Duration; use super::ratelimit::{RateLimiter, RateLimitingCategory}; +use super::DEFAULT_CHANNEL_CAPACITY; use crate::{sentry_debug, Envelope}; #[expect( @@ -33,7 +34,7 @@ impl TransportThread { // NOTE: returning RateLimiter here, otherwise we are in borrow hell SendFuture: std::future::Future, { - Self::with_capacity(send, 30) + Self::with_capacity(send, DEFAULT_CHANNEL_CAPACITY) } /// Spawn a new background thread with a custom channel capacity. @@ -42,7 +43,10 @@ impl TransportThread { /// `send` blocks. `channel_capacity` is clamped to a minimum of 1 to /// avoid a rendezvous channel, which would silently drop envelopes under /// `try_send`. - pub fn with_capacity(mut send: SendFn, channel_capacity: usize) -> Self + pub(crate) fn with_capacity( + mut send: SendFn, + channel_capacity: usize, + ) -> Self where SendFn: FnMut(Envelope, RateLimiter) -> SendFuture + Send + 'static, // NOTE: returning RateLimiter here, otherwise we are in borrow hell diff --git a/sentry/src/transports/ureq.rs b/sentry/src/transports/ureq.rs index dc236a9a7..614f0787d 100644 --- a/sentry/src/transports/ureq.rs +++ b/sentry/src/transports/ureq.rs @@ -5,7 +5,10 @@ use ureq::http::Response; use ureq::tls::{TlsConfig, TlsProvider}; use ureq::{Agent, Proxy}; -use super::{thread::TransportThread, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE}; +use super::{ + thread::TransportThread, DEFAULT_CHANNEL_CAPACITY, HTTP_PAYLOAD_TOO_LARGE, + HTTP_PAYLOAD_TOO_LARGE_MESSAGE, +}; use crate::{sentry_debug, types::Scheme, ClientOptions, Envelope, Transport}; @@ -20,12 +23,12 @@ pub struct UreqHttpTransport { impl UreqHttpTransport { /// Creates a new Transport. pub fn new(options: &ClientOptions) -> Self { - Self::new_internal(options, None, 30) + Self::new_internal(options, None, DEFAULT_CHANNEL_CAPACITY) } /// Creates a new Transport that uses the specified [`ureq::Agent`]. pub fn with_agent(options: &ClientOptions, agent: Agent) -> Self { - Self::new_internal(options, Some(agent), 30) + Self::new_internal(options, Some(agent), DEFAULT_CHANNEL_CAPACITY) } /// Creates a new Transport with a custom transport channel capacity. From c819ce9dc77fcb6afd7ff7d8767107939a6ff832 Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Wed, 6 May 2026 18:27:26 -0700 Subject: [PATCH 05/10] fix(transports): honor channel_capacity=0 (rendezvous), drop the .max(1) clamp Per @szokeasaurusrex's review: this is an advanced opt-in transport API, and sync_channel(0) has defined rendezvous semantics. With try_send, capacity 0 accepts an envelope when the transport thread is currently waiting on the receiver and drops it otherwise, which is a valid no-buffer back-pressure policy rather than "drop everything." Update the doc comment in thread.rs and tokio_thread.rs to describe that behavior accurately. --- sentry/src/transports/thread.rs | 10 ++++++---- sentry/src/transports/tokio_thread.rs | 10 ++++++---- 2 files changed, 12 insertions(+), 8 deletions(-) diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 156cb6d11..51bf6b78a 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -38,14 +38,16 @@ impl TransportThread { /// Spawn a new background thread with a custom channel capacity. /// /// The channel capacity bounds how many envelopes may be queued before - /// `send` blocks. `channel_capacity` is clamped to a minimum of 1 to - /// avoid a rendezvous channel, which would silently drop envelopes under - /// `try_send`. + /// `send` blocks. A capacity of `0` creates a rendezvous channel: + /// because `send` uses `try_send`, an envelope is accepted only when the + /// transport thread is currently waiting on the receiver, otherwise it + /// is dropped. That is a no-buffer back-pressure policy, not a blanket + /// "drop everything" mode. pub(crate) fn with_capacity(mut send: SendFn, channel_capacity: usize) -> Self where SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, { - let (sender, receiver) = sync_channel(channel_capacity.max(1)); + let (sender, receiver) = sync_channel(channel_capacity); let shutdown = Arc::new(AtomicBool::new(false)); let shutdown_worker = shutdown.clone(); let handle = thread::Builder::new() diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index 657399118..a1ad0cf5a 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -40,9 +40,11 @@ impl TransportThread { /// Spawn a new background thread with a custom channel capacity. /// /// The channel capacity bounds how many envelopes may be queued before - /// `send` blocks. `channel_capacity` is clamped to a minimum of 1 to - /// avoid a rendezvous channel, which would silently drop envelopes under - /// `try_send`. + /// `send` blocks. A capacity of `0` creates a rendezvous channel: + /// because `send` uses `try_send`, an envelope is accepted only when the + /// transport thread is currently waiting on the receiver, otherwise it + /// is dropped. That is a no-buffer back-pressure policy, not a blanket + /// "drop everything" mode. pub(crate) fn with_capacity( mut send: SendFn, channel_capacity: usize, @@ -52,7 +54,7 @@ impl TransportThread { // NOTE: returning RateLimiter here, otherwise we are in borrow hell SendFuture: std::future::Future, { - let (sender, receiver) = sync_channel(channel_capacity.max(1)); + let (sender, receiver) = sync_channel(channel_capacity); let shutdown = Arc::new(AtomicBool::new(false)); let shutdown_worker = shutdown.clone(); let handle = thread::Builder::new() From 58951b142a21d60b36428d019c50dae9e62f63eb Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Wed, 15 Jul 2026 20:40:36 -0700 Subject: [PATCH 06/10] fix(transports): prevent flush/drop hang on capacity-0 channels With a rendezvous (capacity-0) channel, flush() and drop() used a blocking send() for their control tasks, which could hang the caller when the worker was not currently receiving. Switch those control-path sends to try_send() and wrap the sender in Option so Drop can take it, so a capacity-0 channel never hangs the caller on flush/drop. Normal envelope enqueue semantics are unchanged (already try_send). Also make TransportThreadOptions::with_channel_capacity pub(crate) in both the std and tokio transports to limit public API surface, per review. Adds a regression test asserting flush() returns false rather than blocking on a busy rendezvous channel. Signed-off-by: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> --- sentry/src/transports/thread.rs | 53 ++++++++++++++++++++++++--- sentry/src/transports/tokio_thread.rs | 53 ++++++++++++++++++++++++--- 2 files changed, 94 insertions(+), 12 deletions(-) diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 63ff3f62d..493f55b83 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -25,7 +25,7 @@ enum Task { /// A background-thread dedicated to sending [`Envelope`]s while respecting the rate limits imposed in the responses. pub struct TransportThread { - sender: SyncSender, + sender: Option>, shutdown: Arc, handle: Option>, client_report_recorder: ClientReportRecorder, @@ -66,7 +66,7 @@ impl TransportThreadOptions { /// transport thread is currently waiting on the receiver, otherwise it is /// dropped. That is a no-buffer back-pressure policy, not a blanket /// "drop everything" mode. - pub fn with_channel_capacity(self, channel_capacity: usize) -> Self { + pub(crate) fn with_channel_capacity(self, channel_capacity: usize) -> Self { Self { channel_capacity, ..self @@ -152,7 +152,7 @@ impl TransportThread { .ok(); Self { - sender, + sender: Some(sender), shutdown, handle, client_report_recorder, @@ -166,7 +166,11 @@ impl TransportThread { // Using send here would mean that when the channel fills up for whatever // reason, trying to send an envelope would block everything. We'd rather // drop the envelope in that case. - if let Err(e) = self.sender.try_send(Task::SendEnvelope(envelope)) { + let sender = self + .sender + .as_ref() + .expect("transport sender is available until drop"); + if let Err(e) = sender.try_send(Task::SendEnvelope(envelope)) { sentry_debug!("envelope dropped: {e}"); // Get back the envelope from the TrySendError so we can record it as lost. @@ -188,7 +192,13 @@ impl TransportThread { /// Returns true if successful within given timeout. pub fn flush(&self, timeout: Duration) -> bool { let (sender, receiver) = sync_channel(1); - let _ = self.sender.send(Task::Flush(sender)); + let transport_sender = self + .sender + .as_ref() + .expect("transport sender is available until drop"); + if transport_sender.try_send(Task::Flush(sender)).is_err() { + return false; + } receiver.recv_timeout(timeout).is_ok() } } @@ -196,9 +206,40 @@ impl TransportThread { impl Drop for TransportThread { fn drop(&mut self) { self.shutdown.store(true, Ordering::SeqCst); - let _ = self.sender.send(Task::Shutdown); + if let Some(sender) = self.sender.take() { + let _ = sender.try_send(Task::Shutdown); + } if let Some(handle) = self.handle.take() { handle.join().unwrap(); } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn flush_does_not_block_on_a_busy_rendezvous_channel() { + let (sender, receiver) = sync_channel(0); + let transport = TransportThread { + sender: Some(sender), + shutdown: Arc::new(AtomicBool::new(false)), + handle: None, + client_report_recorder: Default::default(), + }; + let (result_sender, result_receiver) = sync_channel(1); + + let handle = thread::spawn(move || { + let result = transport.flush(Duration::from_secs(1)); + result_sender.send(result).unwrap(); + }); + + assert_eq!( + result_receiver.recv_timeout(Duration::from_secs(1)), + Ok(false) + ); + drop(receiver); + handle.join().unwrap(); + } +} diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index faf69c196..dd91e4df2 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -25,7 +25,7 @@ enum Task { /// A background-thread powered by [`tokio`] dedicated to sending [`Envelope`]s while respecting the rate limits imposed in the responses. pub struct TransportThread { - sender: SyncSender, + sender: Option>, shutdown: Arc, handle: Option>, client_report_recorder: ClientReportRecorder, @@ -66,7 +66,7 @@ impl TransportThreadOptions { /// transport thread is currently waiting on the receiver, otherwise it is /// dropped. That is a no-buffer back-pressure policy, not a blanket /// "drop everything" mode. - pub fn with_channel_capacity(self, channel_capacity: usize) -> Self { + pub(crate) fn with_channel_capacity(self, channel_capacity: usize) -> Self { Self { channel_capacity, ..self @@ -167,7 +167,7 @@ impl TransportThread { .ok(); Self { - sender, + sender: Some(sender), shutdown, handle, client_report_recorder, @@ -181,7 +181,11 @@ impl TransportThread { // Using send here would mean that when the channel fills up for whatever // reason, trying to send an envelope would block everything. We'd rather // drop the envelope in that case. - if let Err(e) = self.sender.try_send(Task::SendEnvelope(envelope)) { + let sender = self + .sender + .as_ref() + .expect("transport sender is available until drop"); + if let Err(e) = sender.try_send(Task::SendEnvelope(envelope)) { sentry_debug!("envelope dropped: {e}"); // Get back the envelope from the TrySendError so we can record it as lost. @@ -203,7 +207,13 @@ impl TransportThread { /// Returns true if successful within given timeout. pub fn flush(&self, timeout: Duration) -> bool { let (sender, receiver) = sync_channel(1); - let _ = self.sender.send(Task::Flush(sender)); + let transport_sender = self + .sender + .as_ref() + .expect("transport sender is available until drop"); + if transport_sender.try_send(Task::Flush(sender)).is_err() { + return false; + } receiver.recv_timeout(timeout).is_ok() } } @@ -211,9 +221,40 @@ impl TransportThread { impl Drop for TransportThread { fn drop(&mut self) { self.shutdown.store(true, Ordering::SeqCst); - let _ = self.sender.send(Task::Shutdown); + if let Some(sender) = self.sender.take() { + let _ = sender.try_send(Task::Shutdown); + } if let Some(handle) = self.handle.take() { handle.join().unwrap(); } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn flush_does_not_block_on_a_busy_rendezvous_channel() { + let (sender, receiver) = sync_channel(0); + let transport = TransportThread { + sender: Some(sender), + shutdown: Arc::new(AtomicBool::new(false)), + handle: None, + client_report_recorder: Default::default(), + }; + let (result_sender, result_receiver) = sync_channel(1); + + let handle = thread::spawn(move || { + let result = transport.flush(Duration::from_secs(1)); + result_sender.send(result).unwrap(); + }); + + assert_eq!( + result_receiver.recv_timeout(Duration::from_secs(1)), + Ok(false) + ); + drop(receiver); + handle.join().unwrap(); + } +} From 3cda0787bd027d2777d075aa5d1e32b175fb769a Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Fri, 17 Jul 2026 11:16:21 -0700 Subject: [PATCH 07/10] fix(transports): route flush/shutdown through a dedicated control channel Capacity 0 stays supported. Control tasks no longer share the (possibly rendezvous) envelope channel, so flush()/Drop can't deadlock, and flush drains pending envelopes instead of silently dropping them via try_send. Add changelog. --- CHANGELOG.md | 6 + sentry/src/transports/thread.rs | 204 ++++++++++++++-------- sentry/src/transports/tokio_thread.rs | 237 ++++++++++++++++++-------- 3 files changed, 306 insertions(+), 141 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 688a775fe..b5e30eafb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,11 @@ # Changelog +## Unreleased + +### New Features + +- Added configurable channel capacity for built-in background HTTP transports ([#1040](https://github.com/getsentry/sentry-rust/pull/1040)). + ## 0.48.5 ### Fixes diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 493f55b83..078b0e5d9 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -1,6 +1,6 @@ -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::mpsc::{sync_channel, SyncSender, TrySendError}; -use std::sync::Arc; +use std::sync::mpsc::{ + channel, sync_channel, RecvTimeoutError, Sender, SyncSender, TryRecvError, TrySendError, +}; use std::thread::{self, JoinHandle}; use std::time::Duration; @@ -12,21 +12,15 @@ use super::DEFAULT_CHANNEL_CAPACITY; use super::{StdTransportThread, StdTransportThreadOptions}; // so we can use pub re-exports in docs use crate::{sentry_debug, Envelope}; -#[expect( - clippy::large_enum_variant, - reason = "In normal usage this is usually SendEnvelope, the other variants are only used when \ - the user manually calls transport.flush() or when the transport is shut down." -)] -enum Task { - SendEnvelope(Envelope), +enum ControlTask { Flush(SyncSender<()>), Shutdown, } /// A background-thread dedicated to sending [`Envelope`]s while respecting the rate limits imposed in the responses. pub struct TransportThread { - sender: Option>, - shutdown: Arc, + sender: SyncSender, + control_sender: Sender, handle: Option>, client_report_recorder: ClientReportRecorder, } @@ -107,29 +101,13 @@ impl TransportThread { channel_capacity, } = options; let (sender, receiver) = sync_channel(channel_capacity); - let shutdown = Arc::new(AtomicBool::new(false)); - let shutdown_worker = shutdown.clone(); + let (control_sender, control_receiver) = channel(); let handle_client_report_recorder = client_report_recorder.clone(); let handle = thread::Builder::new() .name("sentry-transport".into()) .spawn(move || { let mut rl = RateLimiter::new(); - - for task in receiver.into_iter() { - if shutdown_worker.load(Ordering::SeqCst) { - return; - } - let envelope = match task { - Task::SendEnvelope(envelope) => envelope, - Task::Flush(sender) => { - sender.send(()).ok(); - continue; - } - Task::Shutdown => { - return; - } - }; - + let mut send_envelope = |envelope| { if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { sentry_debug!( "Skipping event send because we're disabled due to rate limits for {}s", @@ -137,23 +115,43 @@ impl TransportThread { ); handle_client_report_recorder .record_lost_data(&envelope, ClientReportReason::RatelimitBackoff); - continue; - } - match rl.filter(envelope, &handle_client_report_recorder) { - Some(envelope) => { - send(envelope, &mut rl); + } else { + match rl.filter(envelope, &handle_client_report_recorder) { + Some(envelope) => { + send(envelope, &mut rl); + } + None => { + sentry_debug!("Envelope was discarded due to per-item rate limits"); + } } - None => { - sentry_debug!("Envelope was discarded due to per-item rate limits"); + } + }; + + loop { + match control_receiver.try_recv() { + Ok(ControlTask::Flush(sender)) => { + for envelope in receiver.try_iter() { + send_envelope(envelope); + } + sender.send(()).ok(); + continue; } - }; + Ok(ControlTask::Shutdown) | Err(TryRecvError::Disconnected) => return, + Err(TryRecvError::Empty) => {} + } + + match receiver.recv_timeout(Duration::from_millis(10)) { + Ok(envelope) => send_envelope(envelope), + Err(RecvTimeoutError::Timeout) => {} + Err(RecvTimeoutError::Disconnected) => return, + } } }) .ok(); Self { - sender: Some(sender), - shutdown, + sender, + control_sender, handle, client_report_recorder, } @@ -166,22 +164,14 @@ impl TransportThread { // Using send here would mean that when the channel fills up for whatever // reason, trying to send an envelope would block everything. We'd rather // drop the envelope in that case. - let sender = self - .sender - .as_ref() - .expect("transport sender is available until drop"); - if let Err(e) = sender.try_send(Task::SendEnvelope(envelope)) { + if let Err(e) = self.sender.try_send(envelope) { sentry_debug!("envelope dropped: {e}"); // Get back the envelope from the TrySendError so we can record it as lost. - let (task, reason) = match e { + let (envelope, reason) = match e { TrySendError::Full(task) => (task, ClientReportReason::QueueOverflow), TrySendError::Disconnected(task) => (task, ClientReportReason::InternalError), }; - let Task::SendEnvelope(envelope) = task else { - unreachable!("we sent a `SendEnvelope` task"); - }; - self.client_report_recorder .record_lost_data(&envelope, reason); } @@ -192,11 +182,11 @@ impl TransportThread { /// Returns true if successful within given timeout. pub fn flush(&self, timeout: Duration) -> bool { let (sender, receiver) = sync_channel(1); - let transport_sender = self - .sender - .as_ref() - .expect("transport sender is available until drop"); - if transport_sender.try_send(Task::Flush(sender)).is_err() { + if self + .control_sender + .send(ControlTask::Flush(sender)) + .is_err() + { return false; } receiver.recv_timeout(timeout).is_ok() @@ -205,10 +195,7 @@ impl TransportThread { impl Drop for TransportThread { fn drop(&mut self) { - self.shutdown.store(true, Ordering::SeqCst); - if let Some(sender) = self.sender.take() { - let _ = sender.try_send(Task::Shutdown); - } + let _ = self.control_sender.send(ControlTask::Shutdown); if let Some(handle) = self.handle.take() { handle.join().unwrap(); } @@ -218,28 +205,107 @@ impl Drop for TransportThread { #[cfg(test)] mod tests { use super::*; + use std::sync::{ + atomic::{AtomicBool, AtomicUsize, Ordering}, + Arc, + }; + use std::time::Instant; + + fn envelope() -> Envelope { + let mut envelope = Envelope::new(); + envelope.add_item(crate::protocol::Event::default()); + envelope + } + + fn send_rendezvous(transport: &TransportThread) { + let mut envelope = envelope(); + let deadline = Instant::now() + .checked_add(Duration::from_secs(1)) + .expect("one-second deadline is representable"); + loop { + match transport.sender.try_send(envelope) { + Ok(()) => return, + Err(TrySendError::Full(returned)) if Instant::now() < deadline => { + envelope = returned; + thread::yield_now(); + } + Err(TrySendError::Full(_)) => panic!("worker did not receive rendezvous event"), + Err(TrySendError::Disconnected(_)) => panic!("worker disconnected"), + } + } + } #[test] - fn flush_does_not_block_on_a_busy_rendezvous_channel() { - let (sender, receiver) = sync_channel(0); - let transport = TransportThread { - sender: Some(sender), - shutdown: Arc::new(AtomicBool::new(false)), - handle: None, - client_report_recorder: Default::default(), - }; + fn flush_waits_for_a_busy_rendezvous_channel() { + let (started_sender, started_receiver) = sync_channel(1); + let (release_sender, release_receiver) = sync_channel(1); + let transport = TransportThreadOptions::new(move |_: Envelope, _: &mut RateLimiter| { + started_sender.send(()).unwrap(); + release_receiver.recv().unwrap(); + }) + .with_channel_capacity(0) + .spawn_thread(); let (result_sender, result_receiver) = sync_channel(1); + send_rendezvous(&transport); + started_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(); let handle = thread::spawn(move || { let result = transport.flush(Duration::from_secs(1)); result_sender.send(result).unwrap(); }); + assert!(result_receiver + .recv_timeout(Duration::from_millis(20)) + .is_err()); + release_sender.send(()).unwrap(); + assert_eq!( + result_receiver.recv_timeout(Duration::from_secs(1)), + Ok(true) + ); + handle.join().unwrap(); + } + + #[test] + fn flush_drains_queued_envelopes() { + let (started_sender, started_receiver) = sync_channel(1); + let (release_sender, release_receiver) = sync_channel(1); + let sent = Arc::new(AtomicUsize::new(0)); + let sent_worker = sent.clone(); + let block_first = Arc::new(AtomicBool::new(true)); + let block_first_worker = block_first.clone(); + let transport = TransportThreadOptions::new(move |_: Envelope, _: &mut RateLimiter| { + if block_first_worker.swap(false, Ordering::SeqCst) { + started_sender.send(()).unwrap(); + release_receiver.recv().unwrap(); + } + sent_worker.fetch_add(1, Ordering::SeqCst); + }) + .with_channel_capacity(1) + .spawn_thread(); + let (result_sender, result_receiver) = sync_channel(1); + + transport.send(envelope()); + started_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + transport.send(envelope()); + let handle = thread::spawn(move || { + result_sender + .send(transport.flush(Duration::from_secs(1))) + .unwrap(); + }); + + assert!(result_receiver + .recv_timeout(Duration::from_millis(20)) + .is_err()); + release_sender.send(()).unwrap(); assert_eq!( result_receiver.recv_timeout(Duration::from_secs(1)), - Ok(false) + Ok(true) ); - drop(receiver); + assert_eq!(sent.load(Ordering::SeqCst), 2); handle.join().unwrap(); } } diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index dd91e4df2..3cef47f13 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -1,6 +1,6 @@ -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::mpsc::{sync_channel, SyncSender, TrySendError}; -use std::sync::Arc; +use std::sync::mpsc::{ + channel, sync_channel, RecvTimeoutError, Sender, SyncSender, TryRecvError, TrySendError, +}; use std::thread::{self, JoinHandle}; use std::time::Duration; @@ -12,21 +12,15 @@ use super::DEFAULT_CHANNEL_CAPACITY; use super::{TokioTransportThread, TokioTransportThreadOptions}; // so we can use pub re-exports in docs use crate::{sentry_debug, Envelope}; -#[expect( - clippy::large_enum_variant, - reason = "In normal usage this is usually SendEnvelope, the other variants are only used when \ - the user manually calls transport.flush() or when the transport is shut down." -)] -enum Task { - SendEnvelope(Envelope), +enum ControlTask { Flush(SyncSender<()>), Shutdown, } /// A background-thread powered by [`tokio`] dedicated to sending [`Envelope`]s while respecting the rate limits imposed in the responses. pub struct TransportThread { - sender: Option>, - shutdown: Arc, + sender: SyncSender, + control_sender: Sender, handle: Option>, client_report_recorder: ClientReportRecorder, } @@ -113,8 +107,7 @@ impl TransportThread { channel_capacity, } = options; let (sender, receiver) = sync_channel(channel_capacity); - let shutdown = Arc::new(AtomicBool::new(false)); - let shutdown_worker = shutdown.clone(); + let (control_sender, control_receiver) = channel(); let handle_client_report_recorder = client_report_recorder.clone(); let handle = thread::Builder::new() .name("sentry-transport".into()) @@ -129,46 +122,64 @@ impl TransportThread { // and block on an async fn in this runtime/thread rt.block_on(async move { - for task in receiver.into_iter() { - if shutdown_worker.load(Ordering::SeqCst) { - return; - } - let envelope = match task { - Task::SendEnvelope(envelope) => envelope, - Task::Flush(sender) => { + loop { + match control_receiver.try_recv() { + Ok(ControlTask::Flush(sender)) => { + for envelope in receiver.try_iter() { + if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { + sentry_debug!( + "Skipping event send because we're disabled due to rate limits for {}s", + time_left.as_secs() + ); + handle_client_report_recorder.record_lost_data( + &envelope, + ClientReportReason::RatelimitBackoff, + ); + } else if let Some(envelope) = + rl.filter(envelope, &handle_client_report_recorder) + { + rl = send(envelope, rl).await; + } else { + sentry_debug!("Envelope was discarded due to per-item rate limits"); + } + } sender.send(()).ok(); continue; } - Task::Shutdown => { - return; - } - }; - - if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { - sentry_debug!( - "Skipping event send because we're disabled due to rate limits for {}s", - time_left.as_secs() - ); - handle_client_report_recorder - .record_lost_data(&envelope, ClientReportReason::RatelimitBackoff); - continue; + Ok(ControlTask::Shutdown) | Err(TryRecvError::Disconnected) => return, + Err(TryRecvError::Empty) => {} } - match rl.filter(envelope, &handle_client_report_recorder) { - Some(envelope) => { - rl = send(envelope, rl).await; - } - None => { - sentry_debug!("Envelope was discarded due to per-item rate limits"); + + match receiver.recv_timeout(Duration::from_millis(10)) { + Ok(envelope) => { + if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { + sentry_debug!( + "Skipping event send because we're disabled due to rate limits for {}s", + time_left.as_secs() + ); + handle_client_report_recorder.record_lost_data( + &envelope, + ClientReportReason::RatelimitBackoff, + ); + } else if let Some(envelope) = + rl.filter(envelope, &handle_client_report_recorder) + { + rl = send(envelope, rl).await; + } else { + sentry_debug!("Envelope was discarded due to per-item rate limits"); + } } - }; + Err(RecvTimeoutError::Timeout) => {} + Err(RecvTimeoutError::Disconnected) => return, + } } }) }) .ok(); Self { - sender: Some(sender), - shutdown, + sender, + control_sender, handle, client_report_recorder, } @@ -181,22 +192,14 @@ impl TransportThread { // Using send here would mean that when the channel fills up for whatever // reason, trying to send an envelope would block everything. We'd rather // drop the envelope in that case. - let sender = self - .sender - .as_ref() - .expect("transport sender is available until drop"); - if let Err(e) = sender.try_send(Task::SendEnvelope(envelope)) { + if let Err(e) = self.sender.try_send(envelope) { sentry_debug!("envelope dropped: {e}"); // Get back the envelope from the TrySendError so we can record it as lost. - let (task, reason) = match e { + let (envelope, reason) = match e { TrySendError::Full(task) => (task, ClientReportReason::QueueOverflow), TrySendError::Disconnected(task) => (task, ClientReportReason::InternalError), }; - let Task::SendEnvelope(envelope) = task else { - unreachable!("we sent a `SendEnvelope` task"); - }; - self.client_report_recorder .record_lost_data(&envelope, reason); } @@ -207,11 +210,11 @@ impl TransportThread { /// Returns true if successful within given timeout. pub fn flush(&self, timeout: Duration) -> bool { let (sender, receiver) = sync_channel(1); - let transport_sender = self - .sender - .as_ref() - .expect("transport sender is available until drop"); - if transport_sender.try_send(Task::Flush(sender)).is_err() { + if self + .control_sender + .send(ControlTask::Flush(sender)) + .is_err() + { return false; } receiver.recv_timeout(timeout).is_ok() @@ -220,10 +223,7 @@ impl TransportThread { impl Drop for TransportThread { fn drop(&mut self) { - self.shutdown.store(true, Ordering::SeqCst); - if let Some(sender) = self.sender.take() { - let _ = sender.try_send(Task::Shutdown); - } + let _ = self.control_sender.send(ControlTask::Shutdown); if let Some(handle) = self.handle.take() { handle.join().unwrap(); } @@ -233,28 +233,121 @@ impl Drop for TransportThread { #[cfg(test)] mod tests { use super::*; + use std::sync::{ + atomic::{AtomicBool, AtomicUsize, Ordering}, + Arc, Mutex, + }; + use std::time::Instant; + + fn envelope() -> Envelope { + let mut envelope = Envelope::new(); + envelope.add_item(crate::protocol::Event::default()); + envelope + } + + fn send_rendezvous(transport: &TransportThread) { + let mut envelope = envelope(); + let deadline = Instant::now() + .checked_add(Duration::from_secs(1)) + .expect("one-second deadline is representable"); + loop { + match transport.sender.try_send(envelope) { + Ok(()) => return, + Err(TrySendError::Full(returned)) if Instant::now() < deadline => { + envelope = returned; + thread::yield_now(); + } + Err(TrySendError::Full(_)) => panic!("worker did not receive rendezvous event"), + Err(TrySendError::Disconnected(_)) => panic!("worker disconnected"), + } + } + } #[test] - fn flush_does_not_block_on_a_busy_rendezvous_channel() { - let (sender, receiver) = sync_channel(0); - let transport = TransportThread { - sender: Some(sender), - shutdown: Arc::new(AtomicBool::new(false)), - handle: None, - client_report_recorder: Default::default(), - }; + fn flush_waits_for_a_busy_rendezvous_channel() { + let (started_sender, started_receiver) = sync_channel(1); + let (release_sender, release_receiver) = sync_channel(1); + let release_receiver = Arc::new(Mutex::new(release_receiver)); + let transport = TransportThreadOptions::new(move |_: Envelope, rl: RateLimiter| { + let started_sender = started_sender.clone(); + let release_receiver = release_receiver.clone(); + async move { + started_sender.send(()).unwrap(); + release_receiver.lock().unwrap().recv().unwrap(); + rl + } + }) + .with_channel_capacity(0) + .spawn_thread(); let (result_sender, result_receiver) = sync_channel(1); + send_rendezvous(&transport); + started_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(); let handle = thread::spawn(move || { let result = transport.flush(Duration::from_secs(1)); result_sender.send(result).unwrap(); }); + assert!(result_receiver + .recv_timeout(Duration::from_millis(20)) + .is_err()); + release_sender.send(()).unwrap(); + assert_eq!( + result_receiver.recv_timeout(Duration::from_secs(1)), + Ok(true) + ); + handle.join().unwrap(); + } + + #[test] + fn flush_drains_queued_envelopes() { + let (started_sender, started_receiver) = sync_channel(1); + let (release_sender, release_receiver) = sync_channel(1); + let release_receiver = Arc::new(Mutex::new(release_receiver)); + let sent = Arc::new(AtomicUsize::new(0)); + let sent_worker = sent.clone(); + let block_first = Arc::new(AtomicBool::new(true)); + let block_first_worker = block_first.clone(); + let transport = TransportThreadOptions::new(move |_: Envelope, rl: RateLimiter| { + let sent_worker = sent_worker.clone(); + let block_first_worker = block_first_worker.clone(); + let started_sender = started_sender.clone(); + let release_receiver = release_receiver.clone(); + async move { + if block_first_worker.swap(false, Ordering::SeqCst) { + started_sender.send(()).unwrap(); + release_receiver.lock().unwrap().recv().unwrap(); + } + sent_worker.fetch_add(1, Ordering::SeqCst); + rl + } + }) + .with_channel_capacity(1) + .spawn_thread(); + let (result_sender, result_receiver) = sync_channel(1); + + transport.send(envelope()); + started_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + transport.send(envelope()); + let handle = thread::spawn(move || { + result_sender + .send(transport.flush(Duration::from_secs(1))) + .unwrap(); + }); + + assert!(result_receiver + .recv_timeout(Duration::from_millis(20)) + .is_err()); + release_sender.send(()).unwrap(); assert_eq!( result_receiver.recv_timeout(Duration::from_secs(1)), - Ok(false) + Ok(true) ); - drop(receiver); + assert_eq!(sent.load(Ordering::SeqCst), 2); handle.join().unwrap(); } } From f19e186de1c34b40166214c3dc1644ad38292caa Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Sun, 19 Jul 2026 23:53:34 -0700 Subject: [PATCH 08/10] fix(transports): select over envelope and control channels without racing Replace the dual std-mpsc channels in the thread transports with crossbeam-channel select_biased (optional dep, gated to transport features) to close the flush/shutdown race flagged in review. --- Cargo.lock | 10 ++++ sentry/Cargo.toml | 7 ++- sentry/src/transports/thread.rs | 52 ++++++++-------- sentry/src/transports/tokio_thread.rs | 85 +++++++++++++++------------ 4 files changed, 89 insertions(+), 65 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 5d6d69c8d..3f2120e81 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -877,6 +877,15 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" +[[package]] +name = "crossbeam-channel" +version = "0.5.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d85363c37faeca707aef026efa9f3b34d077bce547e48f770770625c6013679e" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-deque" version = "0.8.6" @@ -3557,6 +3566,7 @@ dependencies = [ "actix-web", "anyhow", "cfg_aliases", + "crossbeam-channel", "curl", "embedded-svc", "esp-idf-svc", diff --git a/sentry/Cargo.toml b/sentry/Cargo.toml index 38aa79e82..4c7eec6d7 100644 --- a/sentry/Cargo.toml +++ b/sentry/Cargo.toml @@ -52,9 +52,9 @@ logs = ["sentry-core/logs", "sentry-tracing?/logs", "sentry-log?/logs"] metrics = ["sentry-core/metrics"] # transports transport = ["reqwest", "native-tls"] -reqwest = ["dep:reqwest", "httpdate", "tokio"] -curl = ["dep:curl", "httpdate"] -ureq = ["dep:ureq", "httpdate"] +reqwest = ["dep:reqwest", "dep:crossbeam-channel", "httpdate", "tokio"] +curl = ["dep:curl", "dep:crossbeam-channel", "httpdate"] +ureq = ["dep:ureq", "dep:crossbeam-channel", "httpdate"] # transport settings native-tls = ["dep:native-tls", "reqwest?/native-tls", "ureq?/native-tls"] rustls = ["dep:rustls", "reqwest?/rustls", "ureq?/rustls"] @@ -65,6 +65,7 @@ embedded-svc-http = ["dep:embedded-svc", "dep:esp-idf-svc"] sentry-core = { version = "0.48.5", path = "../sentry-core", features = [ "client", ] } +crossbeam-channel = { version = "0.5.16", optional = true } sentry-anyhow = { version = "0.48.5", path = "../sentry-anyhow", optional = true } sentry-actix = { version = "0.48.5", path = "../sentry-actix", optional = true, default-features = false } sentry-backtrace = { version = "0.48.5", path = "../sentry-backtrace", optional = true } diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 078b0e5d9..31c8453e2 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -1,9 +1,7 @@ -use std::sync::mpsc::{ - channel, sync_channel, RecvTimeoutError, Sender, SyncSender, TryRecvError, TrySendError, -}; use std::thread::{self, JoinHandle}; use std::time::Duration; +use crossbeam_channel::{bounded, select_biased, unbounded, Sender, TrySendError}; use sentry_core::client_report::{Reason as ClientReportReason, Recorder as ClientReportRecorder}; use super::ratelimit::{RateLimiter, RateLimitingCategory}; @@ -13,13 +11,13 @@ use super::{StdTransportThread, StdTransportThreadOptions}; // so we can use pub use crate::{sentry_debug, Envelope}; enum ControlTask { - Flush(SyncSender<()>), + Flush(Sender<()>), Shutdown, } /// A background-thread dedicated to sending [`Envelope`]s while respecting the rate limits imposed in the responses. pub struct TransportThread { - sender: SyncSender, + sender: Sender, control_sender: Sender, handle: Option>, client_report_recorder: ClientReportRecorder, @@ -100,8 +98,8 @@ impl TransportThread { client_report_recorder, channel_capacity, } = options; - let (sender, receiver) = sync_channel(channel_capacity); - let (control_sender, control_receiver) = channel(); + let (sender, receiver) = bounded(channel_capacity); + let (control_sender, control_receiver) = unbounded(); let handle_client_report_recorder = client_report_recorder.clone(); let handle = thread::Builder::new() .name("sentry-transport".into()) @@ -128,22 +126,20 @@ impl TransportThread { }; loop { - match control_receiver.try_recv() { - Ok(ControlTask::Flush(sender)) => { - for envelope in receiver.try_iter() { - send_envelope(envelope); + select_biased! { + recv(control_receiver) -> task => match task { + Ok(ControlTask::Flush(sender)) => { + for envelope in receiver.try_iter() { + send_envelope(envelope); + } + sender.send(()).ok(); } - sender.send(()).ok(); - continue; - } - Ok(ControlTask::Shutdown) | Err(TryRecvError::Disconnected) => return, - Err(TryRecvError::Empty) => {} - } - - match receiver.recv_timeout(Duration::from_millis(10)) { - Ok(envelope) => send_envelope(envelope), - Err(RecvTimeoutError::Timeout) => {} - Err(RecvTimeoutError::Disconnected) => return, + Ok(ControlTask::Shutdown) | Err(_) => return, + }, + recv(receiver) -> envelope => match envelope { + Ok(envelope) => send_envelope(envelope), + Err(_) => return, + }, } } }) @@ -181,7 +177,7 @@ impl TransportThread { /// /// Returns true if successful within given timeout. pub fn flush(&self, timeout: Duration) -> bool { - let (sender, receiver) = sync_channel(1); + let (sender, receiver) = bounded(1); if self .control_sender .send(ControlTask::Flush(sender)) @@ -205,6 +201,7 @@ impl Drop for TransportThread { #[cfg(test)] mod tests { use super::*; + use std::sync::mpsc::sync_channel; use std::sync::{ atomic::{AtomicBool, AtomicUsize, Ordering}, Arc, @@ -235,6 +232,15 @@ mod tests { } } + #[test] + fn flush_wakes_an_idle_transport_thread() { + let transport = + TransportThreadOptions::new(|_: Envelope, _: &mut RateLimiter| {}).spawn_thread(); + + thread::sleep(Duration::from_millis(1)); + assert!(transport.flush(Duration::from_millis(5))); + } + #[test] fn flush_waits_for_a_busy_rendezvous_channel() { let (started_sender, started_receiver) = sync_channel(1); diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index 3cef47f13..c9bc9d4f1 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -1,9 +1,7 @@ -use std::sync::mpsc::{ - channel, sync_channel, RecvTimeoutError, Sender, SyncSender, TryRecvError, TrySendError, -}; use std::thread::{self, JoinHandle}; use std::time::Duration; +use crossbeam_channel::{bounded, select_biased, unbounded, Sender, TrySendError}; use sentry_core::client_report::{Reason as ClientReportReason, Recorder as ClientReportRecorder}; use super::ratelimit::{RateLimiter, RateLimitingCategory}; @@ -13,13 +11,13 @@ use super::{TokioTransportThread, TokioTransportThreadOptions}; // so we can use use crate::{sentry_debug, Envelope}; enum ControlTask { - Flush(SyncSender<()>), + Flush(Sender<()>), Shutdown, } /// A background-thread powered by [`tokio`] dedicated to sending [`Envelope`]s while respecting the rate limits imposed in the responses. pub struct TransportThread { - sender: SyncSender, + sender: Sender, control_sender: Sender, handle: Option>, client_report_recorder: ClientReportRecorder, @@ -106,8 +104,8 @@ impl TransportThread { client_report_recorder, channel_capacity, } = options; - let (sender, receiver) = sync_channel(channel_capacity); - let (control_sender, control_receiver) = channel(); + let (sender, receiver) = bounded(channel_capacity); + let (control_sender, control_receiver) = unbounded(); let handle_client_report_recorder = client_report_recorder.clone(); let handle = thread::Builder::new() .name("sentry-transport".into()) @@ -123,9 +121,33 @@ impl TransportThread { // and block on an async fn in this runtime/thread rt.block_on(async move { loop { - match control_receiver.try_recv() { - Ok(ControlTask::Flush(sender)) => { - for envelope in receiver.try_iter() { + select_biased! { + recv(control_receiver) -> task => match task { + Ok(ControlTask::Flush(sender)) => { + for envelope in receiver.try_iter() { + if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { + sentry_debug!( + "Skipping event send because we're disabled due to rate limits for {}s", + time_left.as_secs() + ); + handle_client_report_recorder.record_lost_data( + &envelope, + ClientReportReason::RatelimitBackoff, + ); + } else if let Some(envelope) = + rl.filter(envelope, &handle_client_report_recorder) + { + rl = send(envelope, rl).await; + } else { + sentry_debug!("Envelope was discarded due to per-item rate limits"); + } + } + sender.send(()).ok(); + } + Ok(ControlTask::Shutdown) | Err(_) => return, + }, + recv(receiver) -> envelope => match envelope { + Ok(envelope) => { if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { sentry_debug!( "Skipping event send because we're disabled due to rate limits for {}s", @@ -143,34 +165,8 @@ impl TransportThread { sentry_debug!("Envelope was discarded due to per-item rate limits"); } } - sender.send(()).ok(); - continue; - } - Ok(ControlTask::Shutdown) | Err(TryRecvError::Disconnected) => return, - Err(TryRecvError::Empty) => {} - } - - match receiver.recv_timeout(Duration::from_millis(10)) { - Ok(envelope) => { - if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { - sentry_debug!( - "Skipping event send because we're disabled due to rate limits for {}s", - time_left.as_secs() - ); - handle_client_report_recorder.record_lost_data( - &envelope, - ClientReportReason::RatelimitBackoff, - ); - } else if let Some(envelope) = - rl.filter(envelope, &handle_client_report_recorder) - { - rl = send(envelope, rl).await; - } else { - sentry_debug!("Envelope was discarded due to per-item rate limits"); - } - } - Err(RecvTimeoutError::Timeout) => {} - Err(RecvTimeoutError::Disconnected) => return, + Err(_) => return, + }, } } }) @@ -209,7 +205,7 @@ impl TransportThread { /// /// Returns true if successful within given timeout. pub fn flush(&self, timeout: Duration) -> bool { - let (sender, receiver) = sync_channel(1); + let (sender, receiver) = bounded(1); if self .control_sender .send(ControlTask::Flush(sender)) @@ -233,6 +229,7 @@ impl Drop for TransportThread { #[cfg(test)] mod tests { use super::*; + use std::sync::mpsc::sync_channel; use std::sync::{ atomic::{AtomicBool, AtomicUsize, Ordering}, Arc, Mutex, @@ -263,6 +260,16 @@ mod tests { } } + #[test] + fn flush_wakes_an_idle_transport_thread() { + let transport = + TransportThreadOptions::new(|_: Envelope, rl: RateLimiter| async move { rl }) + .spawn_thread(); + + thread::sleep(Duration::from_millis(1)); + assert!(transport.flush(Duration::from_millis(5))); + } + #[test] fn flush_waits_for_a_busy_rendezvous_channel() { let (started_sender, started_receiver) = sync_channel(1); From 780d015b81dac6290bd551d68fe6542758982091 Mon Sep 17 00:00:00 2001 From: Matt Van Horn <455140+mvanhorn@users.noreply.github.com> Date: Wed, 22 Jul 2026 16:08:37 -0700 Subject: [PATCH 09/10] fix(transports): drain pending envelopes on Drop and keep idle flush responsive Drop now drains the envelope queue before shutdown instead of letting the high-priority Shutdown task preempt it, and the worker loop checks the control channel without delaying idle flushes behind the poll timeout. --- sentry/src/transports/thread.rs | 29 ++++++++++++++++++++++++++- sentry/src/transports/tokio_thread.rs | 4 ++++ 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 31c8453e2..99e072b55 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -191,6 +191,10 @@ impl TransportThread { impl Drop for TransportThread { fn drop(&mut self) { + let (sender, receiver) = bounded(1); + if self.control_sender.send(ControlTask::Flush(sender)).is_ok() { + let _ = receiver.recv(); + } let _ = self.control_sender.send(ControlTask::Shutdown); if let Some(handle) = self.handle.take() { handle.join().unwrap(); @@ -204,7 +208,7 @@ mod tests { use std::sync::mpsc::sync_channel; use std::sync::{ atomic::{AtomicBool, AtomicUsize, Ordering}, - Arc, + Arc, Mutex, }; use std::time::Instant; @@ -314,4 +318,27 @@ mod tests { assert_eq!(sent.load(Ordering::SeqCst), 2); handle.join().unwrap(); } + + #[test] + fn drop_drains_queued_envelopes() { + let gate = Arc::new(Mutex::new(())); + let guard = gate.lock().unwrap(); + let sent = Arc::new(AtomicUsize::new(0)); + let sent_worker = sent.clone(); + let gate_worker = gate.clone(); + let transport = TransportThreadOptions::new(move |_: Envelope, _: &mut RateLimiter| { + let _guard = gate_worker.lock().unwrap(); + sent_worker.fetch_add(1, Ordering::SeqCst); + }) + .with_channel_capacity(1) + .spawn_thread(); + + send_rendezvous(&transport); + send_rendezvous(&transport); + let handle = thread::spawn(move || drop(transport)); + drop(guard); + handle.join().unwrap(); + + assert_eq!(sent.load(Ordering::SeqCst), 2); + } } diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index c9bc9d4f1..d8eaa1603 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -219,6 +219,10 @@ impl TransportThread { impl Drop for TransportThread { fn drop(&mut self) { + let (sender, receiver) = bounded(1); + if self.control_sender.send(ControlTask::Flush(sender)).is_ok() { + let _ = receiver.recv(); + } let _ = self.control_sender.send(ControlTask::Shutdown); if let Some(handle) = self.handle.take() { handle.join().unwrap(); From 8c14fa8283592d4c992c8cbcb9071211716488b8 Mon Sep 17 00:00:00 2001 From: Matt Van Horn Date: Thu, 23 Jul 2026 21:13:10 -0700 Subject: [PATCH 10/10] fix: address review feedback on configurable channel capacity --- sentry/src/transports/curl.rs | 2 +- sentry/src/transports/reqwest.rs | 2 +- sentry/src/transports/thread.rs | 83 ++++++++++++++++++++++-- sentry/src/transports/tokio_thread.rs | 91 +++++++++++++++++++++++++-- sentry/src/transports/ureq.rs | 2 +- 5 files changed, 165 insertions(+), 15 deletions(-) diff --git a/sentry/src/transports/curl.rs b/sentry/src/transports/curl.rs index 7beefe0be..298218d9a 100644 --- a/sentry/src/transports/curl.rs +++ b/sentry/src/transports/curl.rs @@ -241,7 +241,7 @@ impl Transport for CurlHttpTransport { } fn shutdown(&self, timeout: Duration) -> bool { - self.flush(timeout) + self.thread.shutdown(timeout) } } diff --git a/sentry/src/transports/reqwest.rs b/sentry/src/transports/reqwest.rs index 265684b40..fa2480b7b 100644 --- a/sentry/src/transports/reqwest.rs +++ b/sentry/src/transports/reqwest.rs @@ -209,7 +209,7 @@ impl Transport for ReqwestHttpTransport { } fn shutdown(&self, timeout: Duration) -> bool { - self.flush(timeout) + self.thread.shutdown(timeout) } } diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index 99e072b55..314cf52b4 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -1,3 +1,5 @@ +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; use std::thread::{self, JoinHandle}; use std::time::Duration; @@ -19,6 +21,8 @@ enum ControlTask { pub struct TransportThread { sender: Sender, control_sender: Sender, + shutdown_requested: Arc, + shutdown_timed_out: AtomicBool, handle: Option>, client_report_recorder: ClientReportRecorder, } @@ -100,6 +104,8 @@ impl TransportThread { } = options; let (sender, receiver) = bounded(channel_capacity); let (control_sender, control_receiver) = unbounded(); + let shutdown_requested = Arc::new(AtomicBool::new(false)); + let handle_shutdown_requested = shutdown_requested.clone(); let handle_client_report_recorder = client_report_recorder.clone(); let handle = thread::Builder::new() .name("sentry-transport".into()) @@ -129,7 +135,10 @@ impl TransportThread { select_biased! { recv(control_receiver) -> task => match task { Ok(ControlTask::Flush(sender)) => { - for envelope in receiver.try_iter() { + while !handle_shutdown_requested.load(Ordering::SeqCst) { + let Ok(envelope) = receiver.try_recv() else { + break; + }; send_envelope(envelope); } sender.send(()).ok(); @@ -148,6 +157,8 @@ impl TransportThread { Self { sender, control_sender, + shutdown_requested, + shutdown_timed_out: AtomicBool::new(false), handle, client_report_recorder, } @@ -187,17 +198,32 @@ impl TransportThread { } receiver.recv_timeout(timeout).is_ok() } + + pub(crate) fn shutdown(&self, timeout: Duration) -> bool { + let flushed = self.flush(timeout); + if !flushed { + self.shutdown_timed_out.store(true, Ordering::SeqCst); + } + self.shutdown_requested.store(true, Ordering::SeqCst); + let _ = self.control_sender.send(ControlTask::Shutdown); + flushed + } } impl Drop for TransportThread { fn drop(&mut self) { - let (sender, receiver) = bounded(1); - if self.control_sender.send(ControlTask::Flush(sender)).is_ok() { - let _ = receiver.recv(); + if !self.shutdown_requested.load(Ordering::SeqCst) { + let (sender, receiver) = bounded(1); + if self.control_sender.send(ControlTask::Flush(sender)).is_ok() { + let _ = receiver.recv(); + } } + self.shutdown_requested.store(true, Ordering::SeqCst); let _ = self.control_sender.send(ControlTask::Shutdown); - if let Some(handle) = self.handle.take() { - handle.join().unwrap(); + if !self.shutdown_timed_out.load(Ordering::SeqCst) { + if let Some(handle) = self.handle.take() { + handle.join().unwrap(); + } } } } @@ -341,4 +367,49 @@ mod tests { assert_eq!(sent.load(Ordering::SeqCst), 2); } + + #[test] + fn timed_out_shutdown_does_not_block_drop_or_drain_queued_envelopes() { + let (started_sender, started_receiver) = sync_channel(1); + let (release_sender, release_receiver) = sync_channel(1); + let (completed_sender, completed_receiver) = sync_channel(2); + let block_first = Arc::new(AtomicBool::new(true)); + let block_first_worker = block_first.clone(); + let transport = TransportThreadOptions::new(move |_: Envelope, _: &mut RateLimiter| { + if block_first_worker.swap(false, Ordering::SeqCst) { + started_sender.send(()).unwrap(); + release_receiver.recv().unwrap(); + } + completed_sender.send(()).unwrap(); + }) + .with_channel_capacity(1) + .spawn_thread(); + + transport.send(envelope()); + started_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + transport.send(envelope()); + assert!(!transport.shutdown(Duration::from_millis(20))); + + let (dropped_sender, dropped_receiver) = sync_channel(1); + let handle = thread::spawn(move || { + drop(transport); + dropped_sender.send(()).unwrap(); + }); + + assert_eq!( + dropped_receiver.recv_timeout(Duration::from_secs(1)), + Ok(()) + ); + release_sender.send(()).unwrap(); + assert_eq!( + completed_receiver.recv_timeout(Duration::from_secs(1)), + Ok(()) + ); + assert!(completed_receiver + .recv_timeout(Duration::from_millis(50)) + .is_err()); + handle.join().unwrap(); + } } diff --git a/sentry/src/transports/tokio_thread.rs b/sentry/src/transports/tokio_thread.rs index d8eaa1603..b896e1d3a 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -1,3 +1,5 @@ +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; use std::thread::{self, JoinHandle}; use std::time::Duration; @@ -19,6 +21,8 @@ enum ControlTask { pub struct TransportThread { sender: Sender, control_sender: Sender, + shutdown_requested: Arc, + shutdown_timed_out: AtomicBool, handle: Option>, client_report_recorder: ClientReportRecorder, } @@ -106,6 +110,8 @@ impl TransportThread { } = options; let (sender, receiver) = bounded(channel_capacity); let (control_sender, control_receiver) = unbounded(); + let shutdown_requested = Arc::new(AtomicBool::new(false)); + let handle_shutdown_requested = shutdown_requested.clone(); let handle_client_report_recorder = client_report_recorder.clone(); let handle = thread::Builder::new() .name("sentry-transport".into()) @@ -124,7 +130,10 @@ impl TransportThread { select_biased! { recv(control_receiver) -> task => match task { Ok(ControlTask::Flush(sender)) => { - for envelope in receiver.try_iter() { + while !handle_shutdown_requested.load(Ordering::SeqCst) { + let Ok(envelope) = receiver.try_recv() else { + break; + }; if let Some(time_left) = rl.is_disabled(RateLimitingCategory::Any) { sentry_debug!( "Skipping event send because we're disabled due to rate limits for {}s", @@ -176,6 +185,8 @@ impl TransportThread { Self { sender, control_sender, + shutdown_requested, + shutdown_timed_out: AtomicBool::new(false), handle, client_report_recorder, } @@ -215,17 +226,32 @@ impl TransportThread { } receiver.recv_timeout(timeout).is_ok() } + + pub(crate) fn shutdown(&self, timeout: Duration) -> bool { + let flushed = self.flush(timeout); + if !flushed { + self.shutdown_timed_out.store(true, Ordering::SeqCst); + } + self.shutdown_requested.store(true, Ordering::SeqCst); + let _ = self.control_sender.send(ControlTask::Shutdown); + flushed + } } impl Drop for TransportThread { fn drop(&mut self) { - let (sender, receiver) = bounded(1); - if self.control_sender.send(ControlTask::Flush(sender)).is_ok() { - let _ = receiver.recv(); + if !self.shutdown_requested.load(Ordering::SeqCst) { + let (sender, receiver) = bounded(1); + if self.control_sender.send(ControlTask::Flush(sender)).is_ok() { + let _ = receiver.recv(); + } } + self.shutdown_requested.store(true, Ordering::SeqCst); let _ = self.control_sender.send(ControlTask::Shutdown); - if let Some(handle) = self.handle.take() { - handle.join().unwrap(); + if !self.shutdown_timed_out.load(Ordering::SeqCst) { + if let Some(handle) = self.handle.take() { + handle.join().unwrap(); + } } } } @@ -361,4 +387,57 @@ mod tests { assert_eq!(sent.load(Ordering::SeqCst), 2); handle.join().unwrap(); } + + #[test] + fn timed_out_shutdown_does_not_block_drop_or_drain_queued_envelopes() { + let (started_sender, started_receiver) = sync_channel(1); + let (release_sender, release_receiver) = sync_channel(1); + let release_receiver = Arc::new(Mutex::new(release_receiver)); + let (completed_sender, completed_receiver) = sync_channel(2); + let block_first = Arc::new(AtomicBool::new(true)); + let block_first_worker = block_first.clone(); + let transport = TransportThreadOptions::new(move |_: Envelope, rl: RateLimiter| { + let started_sender = started_sender.clone(); + let release_receiver = release_receiver.clone(); + let completed_sender = completed_sender.clone(); + let block_first_worker = block_first_worker.clone(); + async move { + if block_first_worker.swap(false, Ordering::SeqCst) { + started_sender.send(()).unwrap(); + release_receiver.lock().unwrap().recv().unwrap(); + } + completed_sender.send(()).unwrap(); + rl + } + }) + .with_channel_capacity(1) + .spawn_thread(); + + transport.send(envelope()); + started_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + transport.send(envelope()); + assert!(!transport.shutdown(Duration::from_millis(20))); + + let (dropped_sender, dropped_receiver) = sync_channel(1); + let handle = thread::spawn(move || { + drop(transport); + dropped_sender.send(()).unwrap(); + }); + + assert_eq!( + dropped_receiver.recv_timeout(Duration::from_secs(1)), + Ok(()) + ); + release_sender.send(()).unwrap(); + assert_eq!( + completed_receiver.recv_timeout(Duration::from_secs(1)), + Ok(()) + ); + assert!(completed_receiver + .recv_timeout(Duration::from_millis(50)) + .is_err()); + handle.join().unwrap(); + } } diff --git a/sentry/src/transports/ureq.rs b/sentry/src/transports/ureq.rs index f4689389e..97ca033ed 100644 --- a/sentry/src/transports/ureq.rs +++ b/sentry/src/transports/ureq.rs @@ -231,7 +231,7 @@ impl Transport for UreqHttpTransport { } fn shutdown(&self, timeout: Duration) -> bool { - self.flush(timeout) + self.thread.shutdown(timeout) } }