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/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/curl.rs b/sentry/src/transports/curl.rs index 871eff584..298218d9a 100644 --- a/sentry/src/transports/curl.rs +++ b/sentry/src/transports/curl.rs @@ -7,7 +7,7 @@ use sentry_core::TransportOptions; use super::{ thread::{TransportThread, TransportThreadOptions}, - RateLimiter, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE, + RateLimiter, DEFAULT_CHANNEL_CAPACITY, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE, }; use crate::{sentry_debug, types::Scheme, ClientOptions, Envelope, Transport}; @@ -33,6 +33,7 @@ pub struct CurlHttpTransport { pub struct CurlHttpTransportOptions { general_options: TransportOptions, client: Option, + channel_capacity: usize, } impl CurlHttpTransport { @@ -86,6 +87,7 @@ impl CurlHttpTransport { .. }, client, + channel_capacity, } = options; let client = client.unwrap_or_else(CurlClient::new); @@ -224,6 +226,7 @@ impl CurlHttpTransport { let thread = TransportThreadOptions::new(send_fn) .with_client_report_recorder(client_report_recorder) + .with_channel_capacity(channel_capacity) .spawn_thread(); Self { thread } } @@ -238,7 +241,7 @@ impl Transport for CurlHttpTransport { } fn shutdown(&self, timeout: Duration) -> bool { - self.flush(timeout) + self.thread.shutdown(timeout) } } @@ -248,6 +251,7 @@ impl From for CurlHttpTransportOptions { Self { general_options: value, client: None, + channel_capacity: DEFAULT_CHANNEL_CAPACITY, } } } @@ -260,6 +264,21 @@ impl CurlHttpTransportOptions { Self { client, ..self } } + /// Set the capacity of the channel that queues envelopes for the background + /// transport thread (default: 30). + /// + /// A capacity of `0` creates a rendezvous channel: an envelope is accepted + /// only when the transport thread is currently waiting on the receiver, + /// otherwise it is dropped. A higher capacity reduces the chance of dropped + /// events in high-throughput scenarios at the cost of memory. + #[inline] + pub fn with_channel_capacity(self, channel_capacity: usize) -> Self { + Self { + channel_capacity, + ..self + } + } + /// Create a [`CurlHttpTransport`] using these options. #[inline] pub fn build(self) -> CurlHttpTransport { diff --git a/sentry/src/transports/mod.rs b/sentry/src/transports/mod.rs index f2bf72e25..5a60aa3dd 100644 --- a/sentry/src/transports/mod.rs +++ b/sentry/src/transports/mod.rs @@ -54,6 +54,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 175904361..fa2480b7b 100644 --- a/sentry/src/transports/reqwest.rs +++ b/sentry/src/transports/reqwest.rs @@ -6,7 +6,7 @@ use sentry_core::TransportOptions; use super::{ tokio_thread::{TransportThread, TransportThreadOptions}, - RateLimiter, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE, + RateLimiter, DEFAULT_CHANNEL_CAPACITY, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE, }; use crate::{sentry_debug, ClientOptions, Envelope, Transport}; @@ -34,6 +34,7 @@ pub struct ReqwestHttpTransport { pub struct ReqwestHttpTransportOptions { general_options: TransportOptions, client: Option, + channel_capacity: usize, } impl ReqwestHttpTransport { @@ -87,6 +88,7 @@ impl ReqwestHttpTransport { .. }, client, + channel_capacity, } = options; let client = client.unwrap_or_else(|| { @@ -192,6 +194,7 @@ impl ReqwestHttpTransport { let thread = TransportThreadOptions::new(send_fn) .with_client_report_recorder(client_report_recorder) + .with_channel_capacity(channel_capacity) .spawn_thread(); Self { thread } } @@ -206,7 +209,7 @@ impl Transport for ReqwestHttpTransport { } fn shutdown(&self, timeout: Duration) -> bool { - self.flush(timeout) + self.thread.shutdown(timeout) } } @@ -216,6 +219,7 @@ impl From for ReqwestHttpTransportOptions { Self { general_options: value, client: None, + channel_capacity: DEFAULT_CHANNEL_CAPACITY, } } } @@ -228,6 +232,21 @@ impl ReqwestHttpTransportOptions { Self { client, ..self } } + /// Set the capacity of the channel that queues envelopes for the background + /// transport thread (default: 30). + /// + /// A capacity of `0` creates a rendezvous channel: an envelope is accepted + /// only when the transport thread is currently waiting on the receiver, + /// otherwise it is dropped. A higher capacity reduces the chance of dropped + /// events in high-throughput scenarios at the cost of memory. + #[inline] + pub fn with_channel_capacity(self, channel_capacity: usize) -> Self { + Self { + channel_capacity, + ..self + } + } + /// Create a [`ReqwestHttpTransport`] using these options. #[inline] pub fn build(self) -> ReqwestHttpTransport { diff --git a/sentry/src/transports/thread.rs b/sentry/src/transports/thread.rs index c08dac483..314cf52b4 100644 --- a/sentry/src/transports/thread.rs +++ b/sentry/src/transports/thread.rs @@ -1,31 +1,28 @@ use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::mpsc::{sync_channel, SyncSender, TrySendError}; use std::sync::Arc; 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}; +use super::DEFAULT_CHANNEL_CAPACITY; #[cfg(doc)] 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), - Flush(SyncSender<()>), +enum ControlTask { + 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, - shutdown: Arc, + sender: Sender, + control_sender: Sender, + shutdown_requested: Arc, + shutdown_timed_out: AtomicBool, handle: Option>, client_report_recorder: ClientReportRecorder, } @@ -35,6 +32,7 @@ pub struct TransportThread { pub struct TransportThreadOptions { send_fn: F, client_report_recorder: ClientReportRecorder, + channel_capacity: usize, } impl TransportThreadOptions { @@ -43,6 +41,7 @@ impl TransportThreadOptions { Self { send_fn, client_report_recorder: Default::default(), + channel_capacity: DEFAULT_CHANNEL_CAPACITY, } } @@ -53,6 +52,22 @@ impl TransportThreadOptions { ..self } } + + /// Set the capacity of the channel that queues envelopes for the background + /// thread. + /// + /// The capacity bounds how many envelopes may be queued before `send` + /// starts dropping them. 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_channel_capacity(self, channel_capacity: usize) -> Self { + Self { + channel_capacity, + ..self + } + } } impl TransportThreadOptions @@ -85,31 +100,18 @@ impl TransportThread { let TransportThreadOptions { send_fn: mut send, client_report_recorder, + channel_capacity, } = options; - let (sender, receiver) = sync_channel(30); - let shutdown = Arc::new(AtomicBool::new(false)); - let shutdown_worker = shutdown.clone(); + 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()) .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", @@ -117,23 +119,46 @@ 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); - } - None => { - sentry_debug!("Envelope was discarded due to per-item rate limits"); + } 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"); + } } - }; + } + }; + + loop { + select_biased! { + recv(control_receiver) -> task => match task { + Ok(ControlTask::Flush(sender)) => { + while !handle_shutdown_requested.load(Ordering::SeqCst) { + let Ok(envelope) = receiver.try_recv() else { + break; + }; + send_envelope(envelope); + } + sender.send(()).ok(); + } + Ok(ControlTask::Shutdown) | Err(_) => return, + }, + recv(receiver) -> envelope => match envelope { + Ok(envelope) => send_envelope(envelope), + Err(_) => return, + }, + } } }) .ok(); Self { sender, - shutdown, + control_sender, + shutdown_requested, + shutdown_timed_out: AtomicBool::new(false), handle, client_report_recorder, } @@ -146,18 +171,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. - if let Err(e) = self.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); } @@ -167,18 +188,228 @@ 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 (sender, receiver) = bounded(1); + if self + .control_sender + .send(ControlTask::Flush(sender)) + .is_err() + { + return false; + } 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) { - self.shutdown.store(true, Ordering::SeqCst); - let _ = self.sender.send(Task::Shutdown); - if let Some(handle) = self.handle.take() { - handle.join().unwrap(); + 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 !self.shutdown_timed_out.load(Ordering::SeqCst) { + if let Some(handle) = self.handle.take() { + handle.join().unwrap(); + } } } } + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::mpsc::sync_channel; + 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_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); + 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(true) + ); + 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); + } + + #[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 69cd22a12..b896e1d3a 100644 --- a/sentry/src/transports/tokio_thread.rs +++ b/sentry/src/transports/tokio_thread.rs @@ -1,31 +1,28 @@ use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::mpsc::{sync_channel, SyncSender, TrySendError}; use std::sync::Arc; 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}; +use super::DEFAULT_CHANNEL_CAPACITY; #[cfg(doc)] 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), - Flush(SyncSender<()>), +enum ControlTask { + 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, - shutdown: Arc, + sender: Sender, + control_sender: Sender, + shutdown_requested: Arc, + shutdown_timed_out: AtomicBool, handle: Option>, client_report_recorder: ClientReportRecorder, } @@ -35,6 +32,7 @@ pub struct TransportThread { pub struct TransportThreadOptions { send_fn: F, client_report_recorder: ClientReportRecorder, + channel_capacity: usize, } impl TransportThreadOptions { @@ -43,6 +41,7 @@ impl TransportThreadOptions { Self { send_fn, client_report_recorder: Default::default(), + channel_capacity: DEFAULT_CHANNEL_CAPACITY, } } @@ -53,6 +52,22 @@ impl TransportThreadOptions { ..self } } + + /// Set the capacity of the channel that queues envelopes for the background + /// thread. + /// + /// The capacity bounds how many envelopes may be queued before `send` + /// starts dropping them. 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_channel_capacity(self, channel_capacity: usize) -> Self { + Self { + channel_capacity, + ..self + } + } } impl TransportThreadOptions @@ -91,10 +106,12 @@ impl TransportThread { let TransportThreadOptions { send_fn: mut send, client_report_recorder, + channel_capacity, } = options; - let (sender, receiver) = sync_channel(30); - let shutdown = Arc::new(AtomicBool::new(false)); - let shutdown_worker = shutdown.clone(); + 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()) @@ -109,38 +126,57 @@ 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) => { - 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; + loop { + select_biased! { + recv(control_receiver) -> task => match task { + Ok(ControlTask::Flush(sender)) => { + 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", + 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", + 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(_) => return, + }, } - 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"); - } - }; } }) }) @@ -148,7 +184,9 @@ impl TransportThread { Self { sender, - shutdown, + control_sender, + shutdown_requested, + shutdown_timed_out: AtomicBool::new(false), handle, client_report_recorder, } @@ -161,18 +199,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. - if let Err(e) = self.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); } @@ -182,18 +216,228 @@ 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 (sender, receiver) = bounded(1); + if self + .control_sender + .send(ControlTask::Flush(sender)) + .is_err() + { + return false; + } 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) { - self.shutdown.store(true, Ordering::SeqCst); - let _ = self.sender.send(Task::Shutdown); - if let Some(handle) = self.handle.take() { - handle.join().unwrap(); + 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 !self.shutdown_timed_out.load(Ordering::SeqCst) { + if let Some(handle) = self.handle.take() { + handle.join().unwrap(); + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::mpsc::sync_channel; + 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_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); + 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(true) + ); + 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 521281d04..97ca033ed 100644 --- a/sentry/src/transports/ureq.rs +++ b/sentry/src/transports/ureq.rs @@ -13,7 +13,7 @@ use ureq::{Agent, Proxy}; use super::{ thread::{TransportThread, TransportThreadOptions}, - RateLimiter, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE, + RateLimiter, DEFAULT_CHANNEL_CAPACITY, HTTP_PAYLOAD_TOO_LARGE, HTTP_PAYLOAD_TOO_LARGE_MESSAGE, }; use crate::{sentry_debug, types::Scheme, ClientOptions, Envelope, Transport}; @@ -39,6 +39,7 @@ pub struct UreqHttpTransport { pub struct UreqHttpTransportOptions { general_options: TransportOptions, agent: Option, + channel_capacity: usize, } impl UreqHttpTransport { @@ -97,6 +98,7 @@ impl UreqHttpTransport { .. }, agent, + channel_capacity, } = options; let scheme = dsn.scheme(); let agent = agent.unwrap_or_else(|| { @@ -214,6 +216,7 @@ impl UreqHttpTransport { let thread = TransportThreadOptions::new(send_fn) .with_client_report_recorder(client_report_recorder) + .with_channel_capacity(channel_capacity) .spawn_thread(); Self { thread } } @@ -228,7 +231,7 @@ impl Transport for UreqHttpTransport { } fn shutdown(&self, timeout: Duration) -> bool { - self.flush(timeout) + self.thread.shutdown(timeout) } } @@ -238,6 +241,7 @@ impl From for UreqHttpTransportOptions { Self { general_options: value, agent: None, + channel_capacity: DEFAULT_CHANNEL_CAPACITY, } } } @@ -250,6 +254,21 @@ impl UreqHttpTransportOptions { Self { agent, ..self } } + /// Set the capacity of the channel that queues envelopes for the background + /// transport thread (default: 30). + /// + /// A capacity of `0` creates a rendezvous channel: an envelope is accepted + /// only when the transport thread is currently waiting on the receiver, + /// otherwise it is dropped. A higher capacity reduces the chance of dropped + /// events in high-throughput scenarios at the cost of memory. + #[inline] + pub fn with_channel_capacity(self, channel_capacity: usize) -> Self { + Self { + channel_capacity, + ..self + } + } + /// Create a [`UreqHttpTransport`] using these options. #[inline] pub fn build(self) -> UreqHttpTransport {