feat: Make transport channel capacity configurable - #1040
Conversation
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## master #1040 +/- ##
==========================================
+ Coverage 73.81% 74.76% +0.95%
==========================================
Files 64 76 +12
Lines 7538 9488 +1950
==========================================
+ Hits 5564 7094 +1530
- Misses 1974 2394 +420 |
|
Hi @mvanhorn, thanks for the contribution! Before I proceed to a full review, please address the CI failures and the AI review agent comments. If any of the AI review comments are inaccurate, please comment on them and mark as resolved. Thanks! |
|
Addressed in 9c8f904:
|
szokeasaurusrex
left a comment
There was a problem hiding this comment.
Thanks again for the contribution. I think we should revise the public API shape to avoid public API breakages before merging this change.
The current implementation introduces public API breakage by:
- adding a new public field to
ClientOptions - changing the signatures of the public transport thread constructors
I would like to avoid those breakages here and instead expose this as additive API:
- add additive
with_capacity(...)constructors on the public transport thread types, i.e.StdTransportThread::with_capacity(send, channel_capacity)andTokioTransportThread::with_capacity(send, channel_capacity) - keep the existing
new(...)signatures unchanged, but make those constructors delegate towith_capacity(..., 30) - add transport-specific constructors such as
ReqwestHttpTransport::with_channel_capacity(options, channel_capacity)(and similarly forcurlandureq), which would use those newwith_capacity(...)thread constructors internally - document use through
ClientOptions.transport
That still gives users a way to override the transport queue capacity via the existing TransportFactory mechanism, without changing ClientOptions yet.
For example, initializing the SDK with a custom capacity could look like this:
let opts = ClientOptions {
transport: Some(Arc::new(move |opts| {
Arc::new(ReqwestHttpTransport::with_channel_capacity(opts, 256))
})),
..Default::default()
};If you are open to it, please refactor the PR in that direction. If not, let me know and I can take it over as a follow-up.
|
Thanks for the direction @szokeasaurusrex. Refactored to the additive shape you described in fbff3ea:
Your example works unchanged: let opts = ClientOptions {
transport: Some(Arc::new(move |opts| {
Arc::new(ReqwestHttpTransport::with_channel_capacity(opts, 256))
})),
..Default::default()
};
|
|
@szokeasaurusrex - the refactor in fbff3ea matches the additive-API shape you requested. I've also replied to and resolved the two stale bot comments from Apr 15, which were analyzing against the original design. Verified locally: |
|
@mvanhorn I have this PR on my list of things to review. I am busy with another project at the moment but will try to review this within the next week or so |
|
Thanks @szokeasaurusrex - no rush, whenever you get a chance. |
szokeasaurusrex
left a comment
There was a problem hiding this comment.
Thanks for your patience for my review. I added some small comments; this looks pretty good overall though!
| /// `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<SendFn, SendFuture>(mut send: SendFn, channel_capacity: usize) -> Self |
There was a problem hiding this comment.
m: Let's make this pub(crate) for now to limit public API surface. If folks want to have this as a public API in the future, we can expose it at that time.
| pub fn with_capacity<SendFn, SendFuture>(mut send: SendFn, channel_capacity: usize) -> Self | |
| pub(crate) fn with_capacity<SendFn, SendFuture>(mut send: SendFn, channel_capacity: usize) -> Self |
| /// `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<SendFn>(mut send: SendFn, channel_capacity: usize) -> Self |
There was a problem hiding this comment.
m: Let's make this pub(crate) for now to limit public API surface. If folks want to have this as a public API in the future, we can expose it at that time.
| pub fn with_capacity<SendFn>(mut send: SendFn, channel_capacity: usize) -> Self | |
| pub(crate) fn with_capacity<SendFn>(mut send: SendFn, channel_capacity: usize) -> Self |
| /// Creates a new Transport. | ||
| pub fn new(options: &ClientOptions) -> Self { | ||
| Self::new_internal(options, None) | ||
| Self::new_internal(options, None, 30) |
There was a problem hiding this comment.
m: As we use the number 30 as the default channel capacity in all the transports, we should extract it to a constant that we reuse in all of them.
|
Addressed in 6d1c9ff:
Verified |
szokeasaurusrex
left a comment
There was a problem hiding this comment.
Hey, thanks for addressing those! I just have one more thought about the clamping, then I think this is good to merge
| where | ||
| SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, | ||
| { | ||
| let (sender, receiver) = sync_channel(channel_capacity.max(1)); |
There was a problem hiding this comment.
I think we should honor 0 here instead of clamping it. This is an advanced, opt-in transport API, and sync_channel(0) has defined rendezvous/no-buffer semantics. A capacity of 1 can still drop most events under bursts; it is not a general safeguard, only a different buffering policy. If someone explicitly passes 0, treating that as “no queued buffering” seems reasonable.
I tested this end-to-end with the clamp removed and with_channel_capacity(..., 0): sending 10 rapid events accepted 1 and dropped 9. That matches the channel semantics: zero capacity does not necessarily drop everything; it accepts an envelope when try_send happens while the transport thread is already waiting on the receiver. If we keep support for 0, the doc comment should describe that behavior rather than saying it would silently drop envelopes generally.
| // NOTE: returning RateLimiter here, otherwise we are in borrow hell | ||
| SendFuture: std::future::Future<Output = RateLimiter>, | ||
| { | ||
| let (sender, receiver) = sync_channel(channel_capacity.max(1)); |
There was a problem hiding this comment.
Same concept applies here; I tested both transport thread variants with the clamp removed, and both accepted 1 of 10 rapid events with capacity 0 rather than dropping everything.
|
Done in c819ce9 - dropped the |
| pub(crate) fn with_capacity<SendFn>(mut send: SendFn, channel_capacity: usize) -> Self | ||
| where | ||
| SendFn: FnMut(Envelope, &mut RateLimiter) + Send + 'static, | ||
| { | ||
| let (sender, receiver) = sync_channel(channel_capacity); | ||
| let shutdown = Arc::new(AtomicBool::new(false)); | ||
| let shutdown_worker = shutdown.clone(); | ||
| let handle = thread::Builder::new() |
There was a problem hiding this comment.
Bug: With channel_capacity=0, flush() and drop() use a blocking send() for control tasks, which can cause the caller to hang if the worker thread is busy.
Severity: MEDIUM
Suggested Fix
Consider using try_send() for control tasks like Task::Flush and Task::Shutdown, similar to how envelopes are handled. This would align with the documented behavior and prevent blocking in flush() and drop() when the channel capacity is zero and the worker is busy.
Prompt for AI Agent
Review the code at the location below. A potential bug has been identified by an AI
agent. Verify if this is a real issue. If it is, propose a fix; if not, explain why it's
not valid.
Location: sentry/src/transports/thread.rs#L46-L53
Potential issue: When the transport thread is configured with `channel_capacity=0`, it
creates a rendezvous channel. However, the `flush()` and `drop()` methods use a blocking
`send()` to dispatch control tasks (`Task::Flush`, `Task::Shutdown`). If the worker
thread is occupied with a long-running operation, such as sending an envelope over HTTP,
it cannot receive new tasks. Consequently, the call to `flush()` will block until its
timeout is reached, and more critically, `drop()` will block the calling thread (e.g.,
during application shutdown) until the long-running operation completes. This can lead
to significant delays or hangs during shutdown for users who opt into this advanced
configuration. The documentation is also misleading, as it implies non-blocking behavior
for all sends.
Also affects:
sentry/src/transports/tokio_thread.rs:48~55
|
Rebased onto master and resolved the conflict with the new transport-options refactor. Two things:
Verified locally with |
|
Pushed 58951b1. Addressed the capacity-0 flush/drop hang the bot flagged: flush() and drop() now use try_send() for their control tasks (envelope enqueue was already try_send), and the sender is wrapped in Option so Drop can take it cleanly, so a rendezvous channel never blocks the caller on flush or shutdown. Added a regression test in both the std and tokio transports asserting flush() returns false instead of blocking on a busy capacity-0 channel. Also made with_channel_capacity pub(crate) in both TransportThreadOptions to keep the public surface limited as you asked. The shared DEFAULT_CHANNEL_CAPACITY constant and honoring 0 landed in the earlier round. |
|
Pushed 3cda078. Capacity 0 stays supported — flush and shutdown now go through a dedicated control channel rather than the envelope channel, so they can't deadlock at capacity 0, and flush drains pending envelopes instead of dropping them via |
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, enable autofix in the Cursor dashboard.
Want reviews to match your repository better? Bugbot Learning can learn team-specific rules from PR activity. A team admin can enable Learning in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 780d015. Configure here.
|
Pushed 780d015: Drop drains pending envelopes before shutting the worker down (no more silent loss when the transport is dropped mid-queue), and the control-channel check no longer delays idle flushes. Transport suite is green locally (6/6) along with fmt. |
|
Pushed the fix for the flagged transport channel-capacity handling — pending envelopes are now drained on |
|
This is ready for another look whenever you have time. The requested changes went in at 8c14fa8, which landed after the last review pass, and all 16 review threads are resolved. Re-verified locally on that commit: the timeout-preserving No rush, just flagging that the ball is back on your side. |
|
Addressed in afa4f65. The two HIGH Drop findings are the same defect from two angles: Drop sent a Flush and blocked on recv() with no bound. That hangs on a stuck worker, and it also made Client::close's timed shutdown pointless, since dropping right after blocked anyway. Drop now runs a shutdown bounded by DROP_FLUSH_TIMEOUT. I kept the join when the flush completed, so a healthy worker is still awaited and a worker panic still surfaces rather than being swallowed. Only a timed-out flush skips it, since join() has no timeout and would put the unbounded wait right back. Draining on a normal drop is unchanged, so 780d015 isn't regressed. On the capacity finding: I'd documented 0 as a deliberate rendezvous / no-buffer back-pressure mode, but the Bugbot point stands. With send using try_send, that reads as back-pressure and behaves as near-total silent loss, and nothing warned. Clamped to a minimum of 1. Happy to reinstate the rendezvous behaviour behind something more explicit if you'd rather have it available. Both the std and tokio transport threads. 9/9 transport tests pass locally, clippy and fmt clean. |
…ransports Add `channel_capacity` to the built-in background transports so callers can size the envelope queue, with `with_channel_capacity` keeping the change additive against the existing constructors. A capacity of 0 is honored as a rendezvous channel rather than silently clamped, and flush/shutdown are routed through a dedicated control channel so they cannot be starved by a full envelope queue or race the worker. Drop is bounded: it runs a shutdown with the configured timeout before the final drain instead of blocking indefinitely on recv(). Previously a stuck worker could hang Drop forever, which also defeated the timed shutdown in `Client::close`, since dropping right after would block anyway. Regression tests cover the rendezvous capacity, the control-channel path, and the bounded Drop for both the thread and tokio transports. Signed-off-by: Matt Van Horn <mvanhorn@gmail.com>
afa4f65 to
e428a56
Compare
|
Rebased onto master (0.49.1) and squashed to a single commit. The branch had accumulated 16 commits, most of them fork-sync duplicates of upstream work under different SHAs, which is also why it went conflicting. The transport code is byte-identical to what you last reviewed at afa4f65 - only the history and the CHANGELOG placement changed. My entry now sits under a new @szokeasaurusrex the clamping thought from your last pass is in, and both Bugbot
Full run: 9/9 on |
| /// 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. |
There was a problem hiding this comment.
Bug: The documentation for with_channel_capacity(0) incorrectly claims it creates a rendezvous channel. The implementation actually creates a buffered channel with a capacity of 1.
Severity: LOW
Suggested Fix
Update the documentation in reqwest.rs, ureq.rs, and curl.rs to accurately reflect the implementation. The documentation should state that a capacity of 0 is clamped to 1, resulting in a buffered channel.
Prompt for AI Agent
Review the code at the location below. A potential bug has been identified by an AI
agent. Verify if this is a real issue. If it is, propose a fix; if not, explain why it's
not valid.
Location: sentry/src/transports/reqwest.rs#L235-L241
Potential issue: The public API documentation for `with_channel_capacity` in
`ReqwestHttpTransportOptions`, `UreqHttpTransportOptions`, and
`CurlHttpTransportOptions` incorrectly states that passing a `channel_capacity` of `0`
will create a rendezvous channel. The implementation, however, uses a
`normalize_channel_capacity` function which clamps any value less than 1 to 1. This
means that `with_channel_capacity(0)` actually creates a buffered channel with a
capacity of 1, not a rendezvous channel. This discrepancy between the documented
behavior and the actual implementation can lead to unexpected queuing of one envelope,
violating the contract promised to the user.
Also affects:
sentry/src/transports/ureq.rs:257~263

Summary
Adds transport-level
with_channel_capacityconstructors so users can tune the bounded channel size used by the transport thread. The default remains 30, preserving existing behavior.The original design added a
transport_channel_capacityfield toClientOptions. Per @szokeasaurusrex's review, refactored to additive transport-level methods soClientOptionsstays minimal and the feature is opt-in at the transport boundary.Why this matters
In high-throughput scenarios (many transactions with single spans each), the hardcoded capacity of 30 can saturate quickly, leading to dropped envelopes. Identified in #923, tracked in #994. Making it configurable lets users trade memory for reliability based on their workload.
Changes
sentry/src/transports/thread.rs:TransportThread::new(send)keeps its original signature. NewTransportThread::with_capacity(send, channel_capacity)accepts a custom capacity, clamped to a minimum of 1 to avoid rendezvous channels.sentry/src/transports/tokio_thread.rs: Same pattern for the async transport variant.sentry/src/transports/curl.rs: AddedCurlHttpTransport::with_channel_capacity(options, channel_capacity).new(options)unchanged.sentry/src/transports/reqwest.rs: AddedReqwestHttpTransport::with_channel_capacity(options, channel_capacity).new(options)/with_client(options, client)unchanged.sentry/src/transports/ureq.rs: AddedUreqHttpTransport::with_channel_capacity(options, channel_capacity).new(options)/with_agent(options, agent)unchanged.Usage
Testing
cargo check --workspace --all-featurespassescargo fmt --allcleancargo clippy --workspace --all-featurescleanCloses #994
This contribution was developed with AI assistance (Claude Code).