Skip to main content

relay_server/utils/
rate_limits.rs

1use std::fmt::{self, Write};
2use std::future::Future;
3use std::marker::PhantomData;
4
5use relay_profiling::ProfileType;
6use relay_quotas::{
7    DataCategory, ItemScoping, QuotaScope, RateLimit, RateLimitScope, RateLimits, ReasonCode,
8    Scoping,
9};
10use smallvec::SmallVec;
11
12use crate::envelope::{AttachmentParentType, AttachmentType, Envelope, Item, ItemType};
13use crate::integrations::Integration;
14use crate::managed::Managed;
15use crate::services::outcome::Outcome;
16
17/// Name of the rate limits header.
18pub const RATE_LIMITS_HEADER: &str = "X-Sentry-Rate-Limits";
19
20/// Formats the `X-Sentry-Rate-Limits` header.
21pub fn format_rate_limits(rate_limits: &RateLimits) -> String {
22    let mut header = String::new();
23
24    for rate_limit in rate_limits {
25        if !header.is_empty() {
26            header.push_str(", ");
27        }
28
29        write!(header, "{}:", rate_limit.retry_after.remaining_seconds()).ok();
30
31        for (index, category) in rate_limit.categories.iter().enumerate() {
32            if index > 0 {
33                header.push(';');
34            }
35            write!(header, "{category}").ok();
36        }
37
38        write!(header, ":{}", rate_limit.scope.name()).ok();
39
40        if let Some(ref reason_code) = rate_limit.reason_code {
41            write!(header, ":{reason_code}").ok();
42        } else if !rate_limit.namespaces.is_empty() {
43            write!(header, ":").ok(); // delimits the empty reason code for namespaces
44        }
45
46        for (index, namespace) in rate_limit.namespaces.iter().enumerate() {
47            header.push(if index == 0 { ':' } else { ';' });
48            write!(header, "{namespace}").ok();
49        }
50    }
51
52    header
53}
54
55/// Parses the `X-Sentry-Rate-Limits` header.
56pub fn parse_rate_limits(scoping: &Scoping, string: &str) -> RateLimits {
57    let mut rate_limits = RateLimits::new();
58
59    for limit in string.split(',') {
60        let limit = limit.trim();
61        if limit.is_empty() {
62            continue;
63        }
64
65        let mut components = limit.split(':');
66
67        let retry_after = match components.next().and_then(|s| s.parse().ok()) {
68            Some(retry_after) => retry_after,
69            None => continue,
70        };
71
72        let categories = components
73            .next()
74            .unwrap_or("")
75            .split(';')
76            .filter(|category| !category.is_empty())
77            .map(DataCategory::from_name)
78            .collect();
79
80        let quota_scope = QuotaScope::from_name(components.next().unwrap_or(""));
81        let scope = RateLimitScope::for_quota(scoping, quota_scope);
82
83        let reason_code = components
84            .next()
85            .filter(|s| !s.is_empty())
86            .map(ReasonCode::new);
87
88        let namespace = components
89            .next()
90            .unwrap_or("")
91            .split(';')
92            .filter(|s| !s.is_empty())
93            .filter_map(|s| s.parse().ok())
94            .collect();
95
96        rate_limits.add(RateLimit {
97            categories,
98            scope,
99            reason_code,
100            retry_after,
101            namespaces: namespace,
102        });
103    }
104
105    rate_limits
106}
107
108/// Infer the data category from an item.
109///
110/// Categories depend mostly on the item type, with a few special cases:
111/// - `Event`: the category is inferred from the event type. This requires the `event_type` header
112///   to be set on the event item.
113/// - `Attachment`: If the attachment creates an event (e.g. for minidumps), the category is assumed
114///   to be `Error`.
115fn infer_event_category(item: &Item) -> Option<DataCategory> {
116    match item.ty() {
117        ItemType::Event => Some(DataCategory::Error),
118        ItemType::Transaction => Some(DataCategory::Transaction),
119        ItemType::Security | ItemType::RawSecurity => Some(DataCategory::Security),
120        ItemType::UnrealReport => Some(DataCategory::Error),
121        ItemType::UserReportV2 => Some(DataCategory::UserReportV2),
122        ItemType::Attachment if item.creates_event() => Some(DataCategory::Error),
123        ItemType::Attachment => None,
124        ItemType::Session => None,
125        ItemType::Sessions => None,
126        ItemType::MetricBuckets => None,
127        ItemType::FormData => None,
128        ItemType::UserReport => None,
129        ItemType::Profile => None,
130        ItemType::ReplayEvent => None,
131        ItemType::ReplayRecording => None,
132        ItemType::ReplayVideo => None,
133        ItemType::ClientReport => None,
134        ItemType::CheckIn => None,
135        ItemType::Log => None,
136        ItemType::TraceMetric => None,
137        ItemType::Span => None,
138        ItemType::ProfileChunk => None,
139        ItemType::Integration => None,
140        ItemType::Unknown(_) => None,
141    }
142}
143
144/// Quantity metrics for a single category of attachments.
145///
146/// Tracks both the count of attachments and size in bytes.
147#[derive(Clone, Copy, Debug, Default)]
148pub struct AttachmentQuantity {
149    /// Number of attachment items.
150    pub count: usize,
151    /// Total size of attachments in bytes.
152    pub bytes: usize,
153}
154
155impl AttachmentQuantity {
156    pub fn is_empty(&self) -> bool {
157        self.count == 0 && self.bytes == 0
158    }
159}
160
161/// Aggregated attachment quantities grouped by [`AttachmentParentType`].
162///
163/// This separation is necessary since rate limiting logic varies by [`AttachmentParentType`].
164#[derive(Clone, Copy, Debug, Default)]
165pub struct AttachmentQuantities {
166    /// Quantities of Event Attachments.
167    ///
168    /// See also: [`AttachmentParentType::Event`].
169    pub event: AttachmentQuantity,
170    /// Quantities of trace V2 Attachments.
171    pub trace: AttachmentQuantity,
172    /// Quantities of span V2 Attachments.
173    pub span: AttachmentQuantity,
174}
175
176impl AttachmentQuantities {
177    /// Returns the total count of all attachments across all parent types.
178    pub fn count(&self) -> usize {
179        let AttachmentQuantities { event, trace, span } = self;
180        event.count + trace.count + span.count
181    }
182
183    /// Returns the total size in bytes of all attachments across all parent types.
184    pub fn bytes(&self) -> usize {
185        let AttachmentQuantities { event, trace, span } = self;
186        event.bytes + trace.bytes + span.bytes
187    }
188}
189
190/// Collection of all transaction profile quantities.
191#[derive(Clone, Copy, Debug, Default)]
192pub struct ProfileQuantities {
193    /// All transaction profiles in the backend category.
194    pub backend: usize,
195    /// All transaction profiles in the ui category.
196    pub ui: usize,
197    /// All transaction profiles, includes profiles in the backend and ui categories as well as
198    /// profiles which are in neither category.
199    pub total: usize,
200}
201
202/// A summary of `Envelope` contents.
203///
204/// Summarizes the contained event, size of attachments, session updates, and whether there are
205/// plain attachments. This is used for efficient rate limiting or outcome handling.
206#[non_exhaustive]
207#[derive(Clone, Copy, Debug, Default)]
208pub struct EnvelopeSummary {
209    /// The data category of the event in the envelope. `None` if there is no event.
210    pub event_category: Option<DataCategory>,
211
212    /// The quantities of all attachments combined.
213    pub attachment_quantities: AttachmentQuantities,
214
215    /// The number of all session updates.
216    pub session_quantity: usize,
217
218    /// The number of profiles.
219    pub profile_quantity: ProfileQuantities,
220
221    /// The number of replays.
222    pub replay_quantity: usize,
223
224    /// The number of user reports (legacy item type for user feedback).
225    pub user_report_quantity: usize,
226
227    /// The number of monitor check-ins.
228    pub monitor_quantity: usize,
229
230    /// The number of log for the log product sent.
231    pub log_item_quantity: usize,
232
233    /// The number of log bytes for the log product sent, in bytes
234    pub log_byte_quantity: usize,
235
236    /// Secondary number of transactions.
237    ///
238    /// This is 0 for envelopes which contain a transaction,
239    /// only secondary transaction quantity should be tracked here,
240    /// these are for example transaction counts extracted from metrics.
241    ///
242    /// A "primary" transaction is contained within the envelope,
243    /// marking the envelope data category a [`DataCategory::Transaction`].
244    pub secondary_transaction_quantity: usize,
245
246    /// See `secondary_transaction_quantity`.
247    pub secondary_span_quantity: usize,
248
249    /// The number of standalone spans.
250    pub span_quantity: usize,
251
252    /// Indicates that the envelope contains regular attachments that do not create event payloads.
253    pub has_plain_attachments: bool,
254
255    /// The payload size of this envelope.
256    pub payload_size: usize,
257
258    /// The number of profile chunks in this envelope.
259    pub profile_chunk_quantity: usize,
260    /// The number of UI profile chunks in this envelope.
261    pub profile_chunk_ui_quantity: usize,
262
263    /// The number of trace metrics in this envelope.
264    pub trace_metric_quantity: usize,
265
266    /// The number of trace metric bytes in this envelope.
267    pub trace_metric_byte_quantity: usize,
268}
269
270impl EnvelopeSummary {
271    /// Creates an empty summary.
272    pub fn empty() -> Self {
273        Self::default()
274    }
275
276    /// Creates an envelope summary and aggregates the given envelope.
277    pub fn compute(envelope: &Envelope) -> Self {
278        Self::compute_items(envelope.items())
279    }
280
281    pub fn compute_items<'a>(items: impl IntoIterator<Item = &'a Item>) -> Self {
282        let mut summary = Self::empty();
283
284        for item in items {
285            if item.creates_event() {
286                summary.infer_category(item);
287            } else if item.ty() == &ItemType::Attachment {
288                // Plain attachments do not create events.
289                summary.has_plain_attachments = true;
290            }
291
292            // If the item has been rate limited before, the quota has been consumed and outcomes
293            // emitted. We can skip it here.
294            if item.rate_limited() {
295                continue;
296            }
297
298            if let Some(source_quantities) = item.source_quantities() {
299                summary.secondary_transaction_quantity += source_quantities.transactions;
300                summary.secondary_span_quantity += source_quantities.spans;
301            }
302
303            summary.payload_size += item.len();
304
305            summary.add_quantities(item);
306
307            // Special case since v1 and v2 share a data category.
308            // Adding this in add_quantity would include v2 in the count.
309            if item.ty() == &ItemType::UserReport {
310                summary.user_report_quantity += 1;
311            }
312        }
313
314        summary
315    }
316
317    fn add_quantities(&mut self, item: &Item) {
318        // The Nintendo switch item is a special case which should've been modelled like the
319        // `Unreal4Context` as potentially a separate item type which does not have its own data
320        // category.
321        //
322        // Currently there is no outcome category for this item, as it will be dissolved into
323        // multiple different items once processed.
324        if item.attachment_type() == Some(AttachmentType::NintendoSwitchDyingMessage) {
325            return;
326        }
327
328        for (category, quantity) in item.quantities() {
329            let target_quantity = match category {
330                DataCategory::Attachment => match item.attachment_parent_type() {
331                    AttachmentParentType::Span => &mut self.attachment_quantities.span.bytes,
332                    AttachmentParentType::Trace => &mut self.attachment_quantities.trace.bytes,
333                    AttachmentParentType::Event => &mut self.attachment_quantities.event.bytes,
334                },
335                DataCategory::AttachmentItem => match item.attachment_parent_type() {
336                    AttachmentParentType::Span => &mut self.attachment_quantities.span.count,
337                    AttachmentParentType::Trace => &mut self.attachment_quantities.trace.count,
338                    AttachmentParentType::Event => &mut self.attachment_quantities.event.count,
339                },
340                DataCategory::Session => &mut self.session_quantity,
341                DataCategory::Profile => &mut self.profile_quantity.total,
342                DataCategory::ProfileBackend => &mut self.profile_quantity.backend,
343                DataCategory::ProfileUi => &mut self.profile_quantity.ui,
344                DataCategory::Replay => &mut self.replay_quantity,
345                DataCategory::DoNotUseReplayVideo => &mut self.replay_quantity,
346                DataCategory::Monitor => &mut self.monitor_quantity,
347                DataCategory::Span => &mut self.span_quantity,
348                DataCategory::TraceMetric => &mut self.trace_metric_quantity,
349                DataCategory::TraceMetricByte => &mut self.trace_metric_byte_quantity,
350                DataCategory::LogItem => &mut self.log_item_quantity,
351                DataCategory::LogByte => &mut self.log_byte_quantity,
352                DataCategory::ProfileChunk => &mut self.profile_chunk_quantity,
353                DataCategory::ProfileChunkUi => &mut self.profile_chunk_ui_quantity,
354                // TODO: This catch-all looks dangerous
355                _ => continue,
356            };
357            *target_quantity += quantity;
358        }
359    }
360
361    /// Infers the appropriate [`DataCategory`] for the envelope [`Item`].
362    ///
363    /// The inferred category is only applied to the [`EnvelopeSummary`] if there is not yet
364    /// a category set.
365    fn infer_category(&mut self, item: &Item) {
366        if matches!(self.event_category, None | Some(DataCategory::Default))
367            && let Some(category) = infer_event_category(item)
368        {
369            self.event_category = Some(category);
370        }
371    }
372
373    /// Returns `true` if the envelope contains items that depend on spans.
374    ///
375    /// This is used to determined if we should be checking span quota, as the quota should be
376    /// checked both if there are spans or if there are span dependent items (e.g. span attachments).
377    pub fn has_span_dependent_items(&self) -> bool {
378        !self.attachment_quantities.span.is_empty()
379    }
380}
381
382/// Rate limiting information for a data category.
383#[derive(Debug, Default, PartialEq)]
384#[cfg_attr(test, derive(Clone))]
385pub struct CategoryLimit {
386    /// The limited data category.
387    category: Option<DataCategory>,
388    /// Additional and optional data categories in which outcomes will be produced.
389    extra_outcome_categories: SmallVec<[DataCategory; 1]>,
390    /// The total rate limited quantity across all items.
391    ///
392    /// This will be `0` if nothing was rate limited.
393    quantity: usize,
394    /// The reason code of the applied rate limit.
395    ///
396    /// Defaults to `None` if the quota does not declare a reason code.
397    reason_code: Option<ReasonCode>,
398}
399
400impl CategoryLimit {
401    /// Creates a new `CategoryLimit`.
402    ///
403    /// Returns an inactive limit if `rate_limit` is `None`.
404    fn new(category: DataCategory, quantity: usize, rate_limit: Option<&RateLimit>) -> Self {
405        match rate_limit {
406            Some(limit) => Self {
407                category: Some(category),
408                quantity,
409                extra_outcome_categories: Default::default(),
410                reason_code: limit.reason_code.clone(),
411            },
412            None => Self::default(),
413        }
414    }
415
416    /// Adds an additional outcome in the specified category to the limit.
417    pub fn add_outcome_category(mut self, category: DataCategory) -> Self {
418        self.extra_outcome_categories.push(category);
419        self
420    }
421
422    /// Recreates the category limit, if active, for a new category with the same reason.
423    pub fn clone_for(&self, category: DataCategory, quantity: usize) -> CategoryLimit {
424        if !self.is_active() {
425            return Self::default();
426        }
427
428        Self {
429            category: Some(category),
430            extra_outcome_categories: Default::default(),
431            quantity,
432            reason_code: self.reason_code.clone(),
433        }
434    }
435
436    /// Returns `true` if this is an active limit.
437    ///
438    /// Inactive limits are placeholders with no category set.
439    pub fn is_active(&self) -> bool {
440        self.category.is_some()
441    }
442
443    fn outcomes(self) -> impl Iterator<Item = (Outcome, DataCategory, usize)> {
444        let Self {
445            category,
446            extra_outcome_categories,
447            quantity,
448            reason_code,
449        } = self;
450
451        if category.is_none() || quantity == 0 {
452            return either::Either::Left(std::iter::empty());
453        }
454
455        let outcomes = std::iter::chain(category, extra_outcome_categories).map(move |category| {
456            (
457                Outcome::RateLimited(reason_code.clone()),
458                category,
459                quantity,
460            )
461        });
462
463        either::Either::Right(outcomes)
464    }
465}
466
467/// Rate limiting information for a single category of attachments.
468#[derive(Default, Debug)]
469#[cfg_attr(test, derive(Clone))]
470pub struct AttachmentLimits {
471    /// Rate limit applied to attachment bytes ([`DataCategory::Attachment`]).
472    pub bytes: CategoryLimit,
473    /// Rate limit applied to attachment item count ([`DataCategory::AttachmentItem`]).
474    pub count: CategoryLimit,
475}
476
477impl AttachmentLimits {
478    fn is_active(&self) -> bool {
479        self.bytes.is_active() || self.count.is_active()
480    }
481}
482
483/// Rate limiting information for attachments grouped by [`AttachmentParentType`].
484///
485/// See [`AttachmentQuantities`] for the corresponding quantity tracking.
486#[derive(Default, Debug)]
487#[cfg_attr(test, derive(Clone))]
488pub struct AttachmentsLimits {
489    /// Limits for V1 Attachments.
490    pub event: AttachmentLimits,
491    /// Limits for trace V2 Attachments.
492    pub trace: AttachmentLimits,
493    /// Limits for span V2 Attachments.
494    pub span: AttachmentLimits,
495}
496
497/// Information on the limited quantities returned by [`EnvelopeLimiter::compute`].
498#[derive(Default, Debug)]
499#[cfg_attr(test, derive(Clone))]
500pub struct Enforcement {
501    /// The event item rate limit.
502    pub event: CategoryLimit,
503    /// The rate limit for the indexed category of the event.
504    pub event_indexed: CategoryLimit,
505    /// The attachments limits
506    pub attachments_limits: AttachmentsLimits,
507    /// The combined session item rate limit.
508    pub sessions: CategoryLimit,
509    /// The combined transaction profile item rate limits, for all transaction profiles.
510    ///
511    /// This is at least the sum of [`Self::profiles_backend`] and [`Self::profiles_ui`],
512    /// potentially more if there are profiles without a known platform.
513    pub profiles: CategoryLimit,
514    /// The combined backend transaction profile item rate limit.
515    pub profiles_backend: CategoryLimit,
516    /// The combined ui transaction profile item rate limit.
517    pub profiles_ui: CategoryLimit,
518    /// The rate limit for the indexed profiles category.
519    pub profiles_indexed: CategoryLimit,
520    /// The combined replay item rate limit.
521    pub replays: CategoryLimit,
522    /// The combined check-in item rate limit.
523    pub check_ins: CategoryLimit,
524    /// The combined logs (our product logs) rate limit.
525    pub log_items: CategoryLimit,
526    /// The combined logs (our product logs) rate limit.
527    pub log_bytes: CategoryLimit,
528    /// The combined spans rate limit.
529    pub spans: CategoryLimit,
530    /// The rate limit for the indexed span category.
531    pub spans_indexed: CategoryLimit,
532    /// The rate limit for user report v1.
533    pub user_reports: CategoryLimit,
534    /// The combined profile chunk item rate limit.
535    pub profile_chunks: CategoryLimit,
536    /// The combined profile chunk ui item rate limit.
537    pub profile_chunks_ui: CategoryLimit,
538    /// The combined trace metric item rate limit.
539    pub trace_metrics: CategoryLimit,
540    /// The combined trace metric byte rate limit.
541    pub trace_metrics_bytes: CategoryLimit,
542}
543
544impl Enforcement {
545    /// Returns the `CategoryLimit` for the event.
546    ///
547    /// `None` if the event is not rate limited.
548    pub fn active_event(&self) -> Option<&CategoryLimit> {
549        if self.event.is_active() {
550            Some(&self.event)
551        } else if self.event_indexed.is_active() {
552            Some(&self.event_indexed)
553        } else {
554            None
555        }
556    }
557
558    /// Returns `true` if the event is rate limited.
559    pub fn is_event_active(&self) -> bool {
560        self.active_event().is_some()
561    }
562
563    /// Helper for `track_outcomes`.
564    fn get_outcomes(self) -> impl Iterator<Item = (Outcome, DataCategory, usize)> {
565        let Self {
566            event,
567            event_indexed,
568            attachments_limits:
569                AttachmentsLimits {
570                    event:
571                        AttachmentLimits {
572                            bytes: event_attachment_bytes,
573                            count: event_attachment_item,
574                        },
575                    trace:
576                        AttachmentLimits {
577                            bytes: trace_attachment_bytes,
578                            count: trace_attachment_item,
579                        },
580                    span:
581                        AttachmentLimits {
582                            bytes: span_attachment_bytes,
583                            count: span_attachment_item,
584                        },
585                },
586            sessions: _, // Do not report outcomes for sessions.
587            profiles,
588            profiles_backend,
589            profiles_ui,
590            profiles_indexed,
591            replays,
592            check_ins,
593            log_items,
594            log_bytes,
595            spans,
596            spans_indexed,
597            user_reports,
598            profile_chunks,
599            profile_chunks_ui,
600            trace_metrics,
601            trace_metrics_bytes,
602        } = self;
603
604        let limits = [
605            event,
606            event_indexed,
607            event_attachment_bytes,
608            event_attachment_item,
609            trace_attachment_bytes,
610            trace_attachment_item,
611            span_attachment_bytes,
612            span_attachment_item,
613            profiles,
614            profiles_backend,
615            profiles_ui,
616            profiles_indexed,
617            replays,
618            check_ins,
619            log_items,
620            log_bytes,
621            spans,
622            spans_indexed,
623            user_reports,
624            profile_chunks,
625            profile_chunks_ui,
626            trace_metrics,
627            trace_metrics_bytes,
628        ];
629
630        limits.into_iter().flat_map(|limit| limit.outcomes())
631    }
632
633    /// Applies the [`Enforcement`] on the [`Envelope`] by removing all items that were rate limited
634    /// and emits outcomes for each rate limited category.
635    ///
636    /// # Example
637    ///
638    /// ## Interaction between Events and Attachments
639    ///
640    /// An envelope with an `Error` event and an `Attachment`. Two quotas specify to drop all
641    /// attachments (reason `"a"`) and all errors (reason `"e"`). The result of enforcement will be:
642    ///
643    /// 1. All items are removed from the envelope.
644    /// 2. Enforcements report both the event and the attachment dropped with reason `"e"`, since
645    ///    dropping an event automatically drops all attachments with the same reason.
646    /// 3. Rate limits report the single event limit `"e"`, since attachment limits do not need to
647    ///    be checked in this case.
648    ///
649    /// ## Required Attachments
650    ///
651    /// An envelope with a single Minidump `Attachment`, and a single quota specifying to drop all
652    /// attachments with reason `"a"`:
653    ///
654    /// 1. Since the minidump creates an event and is required for processing, it remains in the
655    ///    envelope and is marked as `rate_limited`.
656    /// 2. Enforcements report the attachment dropped with reason `"a"`.
657    /// 3. Rate limits are empty since it is allowed to send required attachments even when rate
658    ///    limited.
659    ///
660    /// ## Previously Rate Limited Attachments
661    ///
662    /// An envelope with a single item marked as `rate_limited`, and a quota specifying to drop
663    /// everything with reason `"d"`:
664    ///
665    /// 1. The item remains in the envelope.
666    /// 2. Enforcements are empty. Rate limiting has occurred at an earlier stage in the pipeline.
667    /// 3. Rate limits are empty.
668    pub fn apply_to_managed(self, envelope: &mut Managed<Box<Envelope>>) {
669        envelope.modify(|envelope, records| {
670            envelope.retain_items(|item| self.retain_item(item));
671
672            // Sessions currently do not emit any outcomes, but may be dropped.
673            records.lenient(DataCategory::Session);
674            // This is an existing bug in how user reports handle rate limits and emit outcomes.
675            //
676            // User report v1 and v2 (feedback) are counting into the same category, but that is not
677            // completely consistent leading to some mismatches when emitting outcomes from rate
678            // limiting vs how outcomes are counted on the `Managed` instance.
679            //
680            // Issue: <https://github.com/getsentry/relay/issues/5524>.
681            records.lenient(DataCategory::UserReportV2);
682
683            for (outcome, category, quantity) in self.get_outcomes() {
684                records.reject_err(outcome, (category, quantity))
685            }
686        });
687    }
688
689    /// Returns `true` when an [`Item`] can be retained, `false` otherwise.
690    fn retain_item(&self, item: &mut Item) -> bool {
691        // Remove event items and all items that depend on this event
692        if self.event.is_active() && item.requires_event() {
693            return false;
694        }
695
696        // When checking limits for categories that have an indexed variant,
697        // we only have to check the more specific, the indexed, variant
698        // to determine whether an item is limited.
699        match item.ty() {
700            ItemType::Attachment => {
701                match item.attachment_parent_type() {
702                    AttachmentParentType::Span => !self.attachments_limits.span.is_active(),
703                    AttachmentParentType::Trace => !self.attachments_limits.trace.is_active(),
704                    AttachmentParentType::Event => {
705                        if !self.attachments_limits.event.is_active() {
706                            return true;
707                        }
708                        if item.creates_event() {
709                            item.set_rate_limited(true);
710                            true
711                        } else {
712                            false
713                        }
714                    }
715                }
716            }
717            ItemType::Session => !self.sessions.is_active(),
718            ItemType::Profile => {
719                if self.profiles_indexed.is_active() {
720                    false
721                } else if let Some(platform) = item.profile_type() {
722                    match platform {
723                        ProfileType::Backend => !self.profiles_backend.is_active(),
724                        ProfileType::Ui => !self.profiles_ui.is_active(),
725                    }
726                } else {
727                    true
728                }
729            }
730            ItemType::ReplayEvent => !self.replays.is_active(),
731            ItemType::ReplayVideo => !self.replays.is_active(),
732            ItemType::ReplayRecording => !self.replays.is_active(),
733            ItemType::UserReport => !self.user_reports.is_active(),
734            ItemType::CheckIn => !self.check_ins.is_active(),
735            ItemType::Log => {
736                !(self.log_items.is_active() || self.log_bytes.is_active())
737            }
738            ItemType::Span => !self.spans_indexed.is_active(),
739            ItemType::ProfileChunk => match item.profile_type() {
740                Some(ProfileType::Backend) => !self.profile_chunks.is_active(),
741                Some(ProfileType::Ui) => !self.profile_chunks_ui.is_active(),
742                None => true,
743            },
744            ItemType::TraceMetric => !(self.trace_metrics.is_active() || self.trace_metrics_bytes.is_active()),
745            ItemType::Integration => match item.integration() {
746                Some(Integration::Logs(_)) => !(self.log_items.is_active() || self.log_bytes.is_active()),
747                Some(Integration::Spans(_)) => !self.spans_indexed.is_active(),
748                None => true,
749            },
750            ItemType::Event
751            | ItemType::Transaction
752            | ItemType::Security
753            | ItemType::FormData
754            | ItemType::RawSecurity
755            | ItemType::UnrealReport
756            | ItemType::Sessions
757            | ItemType::MetricBuckets
758            | ItemType::ClientReport
759            | ItemType::UserReportV2  // This is an event type.
760            | ItemType::Unknown(_) => true,
761        }
762    }
763}
764
765/// Which limits to check with the [`EnvelopeLimiter`].
766#[derive(Debug, Copy, Clone)]
767pub enum CheckLimits {
768    /// Checks all limits except indexed categories.
769    ///
770    /// In the fast path it is necessary to apply cached rate limits but to not enforce indexed rate limits.
771    /// Because at the time of the check the decision whether an envelope is sampled or not is not yet known.
772    /// Additionally even if the item is later dropped by dynamic sampling, it must still be around to extract metrics
773    /// and cannot be dropped too early.
774    NonIndexed,
775}
776
777struct Check<F, E, R> {
778    limits: CheckLimits,
779    check: F,
780    _1: PhantomData<E>,
781    _2: PhantomData<R>,
782}
783
784impl<F, E, R> Check<F, E, R>
785where
786    F: FnMut(ItemScoping, usize) -> R,
787    R: Future<Output = Result<RateLimits, E>>,
788{
789    async fn apply(&mut self, scoping: ItemScoping, quantity: usize) -> Result<RateLimits, E> {
790        if matches!(self.limits, CheckLimits::NonIndexed) && scoping.category.is_indexed() {
791            return Ok(RateLimits::default());
792        }
793
794        (self.check)(scoping, quantity).await
795    }
796}
797
798/// Enforces rate limits with the given `check` function on items in the envelope.
799///
800/// The `check` function is called with the following rules:
801///  - Once for a single event, if present in the envelope.
802///  - Once for all comprised attachments, unless the event was rate limited.
803///  - Once for all comprised sessions.
804///
805/// Items violating the rate limit are removed from the envelope. This follows a set of rules:
806///  - If the event is removed, all items depending on the event are removed (e.g. attachments).
807///  - Attachments are not removed if they create events (e.g. minidumps).
808///  - Sessions are handled separately from all of the above.
809pub struct EnvelopeLimiter<F, E, R> {
810    check: Check<F, E, R>,
811    event_category: Option<DataCategory>,
812}
813
814impl<'a, F, E, R> EnvelopeLimiter<F, E, R>
815where
816    F: FnMut(ItemScoping, usize) -> R,
817    R: Future<Output = Result<RateLimits, E>>,
818{
819    /// Create a new `EnvelopeLimiter` with the given `check` function.
820    pub fn new(limits: CheckLimits, check: F) -> Self {
821        Self {
822            check: Check {
823                check,
824                limits,
825                _1: PhantomData,
826                _2: PhantomData,
827            },
828            event_category: None,
829        }
830    }
831
832    /// Process rate limits for the envelope, returning applied limits.
833    ///
834    /// Returns a tuple of `Enforcement` and `RateLimits`:
835    ///
836    /// - Enforcements declare the quantities of categories that have been rate limited with the
837    ///   individual reason codes that caused rate limiting. If multiple rate limits applied to a
838    ///   category, then the longest limit is reported.
839    /// - Rate limits declare all active rate limits, regardless of whether they have been applied
840    ///   to items in the envelope. This excludes rate limits applied to required attachments, since
841    ///   clients are allowed to continue sending them.
842    pub async fn compute(
843        mut self,
844        envelope: &Envelope,
845        scoping: &'a Scoping,
846    ) -> Result<(Enforcement, RateLimits), E> {
847        let mut summary = EnvelopeSummary::compute(envelope);
848        summary.event_category = self.event_category.or(summary.event_category);
849
850        let (enforcement, rate_limits) = self.execute(&summary, scoping).await?;
851        Ok((enforcement, rate_limits))
852    }
853
854    async fn execute(
855        &mut self,
856        summary: &EnvelopeSummary,
857        scoping: &'a Scoping,
858    ) -> Result<(Enforcement, RateLimits), E> {
859        let mut rate_limits = RateLimits::new();
860        let mut enforcement = Enforcement::default();
861
862        // Handle event.
863        if let Some(category) = summary.event_category {
864            // Check the broad category for limits.
865            let mut event_limits = self.check.apply(scoping.item(category), 1).await?;
866            enforcement.event = CategoryLimit::new(category, 1, event_limits.longest());
867
868            if let Some(index_category) = category.index_category() {
869                // Check the specific/indexed category for limits only if the specific one has not already
870                // an enforced limit.
871                if event_limits.is_empty() {
872                    event_limits.merge(self.check.apply(scoping.item(index_category), 1).await?);
873                }
874
875                enforcement.event_indexed =
876                    CategoryLimit::new(index_category, 1, event_limits.longest());
877            };
878
879            rate_limits.merge(event_limits);
880        }
881
882        // Handle spans.
883        if enforcement.is_event_active() {
884            enforcement.spans = enforcement
885                .event
886                .clone_for(DataCategory::Span, summary.span_quantity);
887
888            enforcement.spans_indexed = enforcement
889                .event_indexed
890                .clone_for(DataCategory::SpanIndexed, summary.span_quantity);
891        } else if summary.span_quantity > 0 || summary.has_span_dependent_items() {
892            let mut span_limits = self
893                .check
894                .apply(scoping.item(DataCategory::Span), summary.span_quantity)
895                .await?;
896            enforcement.spans = CategoryLimit::new(
897                DataCategory::Span,
898                summary.span_quantity,
899                span_limits.longest(),
900            );
901
902            if span_limits.is_empty() {
903                span_limits.merge(
904                    self.check
905                        .apply(
906                            scoping.item(DataCategory::SpanIndexed),
907                            summary.span_quantity,
908                        )
909                        .await?,
910                );
911            }
912
913            enforcement.spans_indexed = CategoryLimit::new(
914                DataCategory::SpanIndexed,
915                summary.span_quantity,
916                span_limits.longest(),
917            );
918
919            rate_limits.merge(span_limits);
920        }
921
922        // Handle span attachments
923        if enforcement.spans_indexed.is_active() {
924            enforcement.attachments_limits.span.bytes = enforcement.spans_indexed.clone_for(
925                DataCategory::Attachment,
926                summary.attachment_quantities.span.bytes,
927            );
928            enforcement.attachments_limits.span.count = enforcement.spans_indexed.clone_for(
929                DataCategory::AttachmentItem,
930                summary.attachment_quantities.span.count,
931            );
932        } else if !summary.attachment_quantities.span.is_empty() {
933            // While we could combine this check with the check that we do for event and trace
934            // attachments, this would complicate the logic so we opted against doing that.
935            // In practice the performance impact should be negligible since different types
936            // of attachments should rarely be send together.
937            enforcement.attachments_limits.span = self
938                .check_attachment_limits(scoping, &summary.attachment_quantities.span)
939                .await?;
940        }
941
942        // Handle attachments.
943        if let Some(limit) = enforcement.active_event() {
944            let limit1 = limit.clone_for(
945                DataCategory::Attachment,
946                summary.attachment_quantities.event.bytes,
947            );
948            let limit2 = limit.clone_for(
949                DataCategory::AttachmentItem,
950                summary.attachment_quantities.event.count,
951            );
952
953            enforcement.attachments_limits.event.bytes = limit1;
954            enforcement.attachments_limits.event.count = limit2;
955        } else {
956            let mut attachment_limits = RateLimits::new();
957            if summary.attachment_quantities.event.bytes > 0 {
958                let item_scoping = scoping.item(DataCategory::Attachment);
959
960                let attachment_byte_limits = self
961                    .check
962                    .apply(item_scoping, summary.attachment_quantities.event.bytes)
963                    .await?;
964
965                enforcement.attachments_limits.event.bytes = CategoryLimit::new(
966                    DataCategory::Attachment,
967                    summary.attachment_quantities.event.bytes,
968                    attachment_byte_limits.longest(),
969                );
970                enforcement.attachments_limits.event.count =
971                    enforcement.attachments_limits.event.bytes.clone_for(
972                        DataCategory::AttachmentItem,
973                        summary.attachment_quantities.event.count,
974                    );
975                attachment_limits.merge(attachment_byte_limits);
976            }
977            if !attachment_limits.is_limited() && summary.attachment_quantities.event.count > 0 {
978                let item_scoping = scoping.item(DataCategory::AttachmentItem);
979
980                let attachment_item_limits = self
981                    .check
982                    .apply(item_scoping, summary.attachment_quantities.event.count)
983                    .await?;
984
985                enforcement.attachments_limits.event.count = CategoryLimit::new(
986                    DataCategory::AttachmentItem,
987                    summary.attachment_quantities.event.count,
988                    attachment_item_limits.longest(),
989                );
990                enforcement.attachments_limits.event.bytes =
991                    enforcement.attachments_limits.event.count.clone_for(
992                        DataCategory::Attachment,
993                        summary.attachment_quantities.event.bytes,
994                    );
995                attachment_limits.merge(attachment_item_limits);
996            }
997
998            // Only record rate limits for plain attachments. For all other attachments, it's
999            // perfectly "legal" to send them. They will still be discarded in Sentry, but clients
1000            // can continue to send them.
1001            if summary.has_plain_attachments {
1002                rate_limits.merge(attachment_limits);
1003            }
1004        }
1005
1006        // Handle trace attachments.
1007        if !summary.attachment_quantities.trace.is_empty() {
1008            enforcement.attachments_limits.trace = self
1009                .check_attachment_limits(scoping, &summary.attachment_quantities.trace)
1010                .await?;
1011        }
1012
1013        // Handle sessions.
1014        if summary.session_quantity > 0 {
1015            let item_scoping = scoping.item(DataCategory::Session);
1016            let session_limits = self
1017                .check
1018                .apply(item_scoping, summary.session_quantity)
1019                .await?;
1020            enforcement.sessions = CategoryLimit::new(
1021                DataCategory::Session,
1022                summary.session_quantity,
1023                session_limits.longest(),
1024            );
1025            rate_limits.merge(session_limits);
1026        }
1027
1028        // Handle trace metrics.
1029        let mut trace_metric_limits = RateLimits::new();
1030        if summary.trace_metric_quantity > 0 {
1031            let item_scoping = scoping.item(DataCategory::TraceMetric);
1032            trace_metric_limits = self
1033                .check
1034                .apply(item_scoping, summary.trace_metric_quantity)
1035                .await?;
1036            enforcement.trace_metrics = CategoryLimit::new(
1037                DataCategory::TraceMetric,
1038                summary.trace_metric_quantity,
1039                trace_metric_limits.longest(),
1040            );
1041            enforcement.trace_metrics_bytes = CategoryLimit::new(
1042                DataCategory::TraceMetricByte,
1043                summary.trace_metric_byte_quantity,
1044                trace_metric_limits.longest(),
1045            );
1046        }
1047        if !trace_metric_limits.is_limited() && summary.trace_metric_byte_quantity > 0 {
1048            let item_scoping = scoping.item(DataCategory::TraceMetricByte);
1049            trace_metric_limits = self
1050                .check
1051                .apply(item_scoping, summary.trace_metric_byte_quantity)
1052                .await?;
1053            enforcement.trace_metrics = CategoryLimit::new(
1054                DataCategory::TraceMetric,
1055                summary.trace_metric_quantity,
1056                trace_metric_limits.longest(),
1057            );
1058            enforcement.trace_metrics_bytes = CategoryLimit::new(
1059                DataCategory::TraceMetricByte,
1060                summary.trace_metric_byte_quantity,
1061                trace_metric_limits.longest(),
1062            );
1063        }
1064        rate_limits.merge(trace_metric_limits);
1065
1066        // Handle logs.
1067        let mut log_limits = RateLimits::new();
1068        if summary.log_item_quantity > 0 {
1069            let item_scoping = scoping.item(DataCategory::LogItem);
1070            log_limits = self
1071                .check
1072                .apply(item_scoping, summary.log_item_quantity)
1073                .await?;
1074            enforcement.log_bytes = CategoryLimit::new(
1075                DataCategory::LogByte,
1076                summary.log_byte_quantity,
1077                log_limits.longest(),
1078            );
1079            enforcement.log_items = CategoryLimit::new(
1080                DataCategory::LogItem,
1081                summary.log_item_quantity,
1082                log_limits.longest(),
1083            );
1084        }
1085        if !log_limits.is_limited() && summary.log_byte_quantity > 0 {
1086            let item_scoping = scoping.item(DataCategory::LogByte);
1087            log_limits = self
1088                .check
1089                .apply(item_scoping, summary.log_byte_quantity)
1090                .await?;
1091            enforcement.log_bytes = CategoryLimit::new(
1092                DataCategory::LogByte,
1093                summary.log_byte_quantity,
1094                log_limits.longest(),
1095            );
1096            enforcement.log_items = CategoryLimit::new(
1097                DataCategory::LogItem,
1098                summary.log_item_quantity,
1099                log_limits.longest(),
1100            );
1101        }
1102        rate_limits.merge(log_limits);
1103
1104        // Handle profiles.
1105        if enforcement.is_event_active() {
1106            enforcement.profiles = enforcement
1107                .event
1108                .clone_for(DataCategory::Profile, summary.profile_quantity.total);
1109            enforcement.profiles_indexed = enforcement
1110                .event_indexed
1111                .clone_for(DataCategory::ProfileIndexed, summary.profile_quantity.total);
1112
1113            enforcement.profiles_backend = enforcement.event.clone_for(
1114                DataCategory::ProfileBackend,
1115                summary.profile_quantity.backend,
1116            );
1117            enforcement.profiles_ui = enforcement
1118                .event
1119                .clone_for(DataCategory::ProfileUi, summary.profile_quantity.ui);
1120        } else if summary.profile_quantity.total > 0 {
1121            let mut profile_limits = self
1122                .check
1123                .apply(
1124                    scoping.item(DataCategory::Profile),
1125                    summary.profile_quantity.total,
1126                )
1127                .await?;
1128
1129            // Profiles can persist in envelopes without transaction if the transaction item
1130            // was dropped by dynamic sampling.
1131            if profile_limits.is_empty() && summary.event_category.is_none() {
1132                profile_limits = self
1133                    .check
1134                    .apply(scoping.item(DataCategory::Transaction), 0)
1135                    .await?;
1136            }
1137
1138            enforcement.profiles = CategoryLimit::new(
1139                DataCategory::Profile,
1140                summary.profile_quantity.total,
1141                profile_limits.longest(),
1142            );
1143
1144            if enforcement.profiles.quantity == 0 {
1145                if summary.profile_quantity.backend > 0 {
1146                    let limit = self
1147                        .check
1148                        .apply(
1149                            scoping.item(DataCategory::ProfileBackend),
1150                            summary.profile_quantity.backend,
1151                        )
1152                        .await?;
1153
1154                    enforcement.profiles_backend = CategoryLimit::new(
1155                        DataCategory::ProfileBackend,
1156                        summary.profile_quantity.backend,
1157                        limit.longest(),
1158                    )
1159                    .add_outcome_category(DataCategory::Profile);
1160
1161                    profile_limits.merge(limit);
1162                }
1163                if summary.profile_quantity.ui > 0 {
1164                    let limit = self
1165                        .check
1166                        .apply(
1167                            scoping.item(DataCategory::ProfileUi),
1168                            summary.profile_quantity.ui,
1169                        )
1170                        .await?;
1171
1172                    enforcement.profiles_ui = CategoryLimit::new(
1173                        DataCategory::ProfileUi,
1174                        summary.profile_quantity.ui,
1175                        limit.longest(),
1176                    )
1177                    .add_outcome_category(DataCategory::Profile);
1178
1179                    profile_limits.merge(limit);
1180                }
1181            } else {
1182                enforcement.profiles_backend = CategoryLimit::new(
1183                    DataCategory::ProfileBackend,
1184                    summary.profile_quantity.backend,
1185                    profile_limits.longest(),
1186                );
1187                enforcement.profiles_ui = CategoryLimit::new(
1188                    DataCategory::ProfileUi,
1189                    summary.profile_quantity.ui,
1190                    profile_limits.longest(),
1191                );
1192            }
1193
1194            if enforcement.profiles.quantity > 0 {
1195                enforcement.profiles_indexed = enforcement
1196                    .profiles
1197                    .clone_for(DataCategory::ProfileIndexed, summary.profile_quantity.total);
1198            } else {
1199                let limit = self
1200                    .check
1201                    .apply(
1202                        scoping.item(DataCategory::ProfileIndexed),
1203                        summary.profile_quantity.total,
1204                    )
1205                    .await?;
1206
1207                if !limit.is_empty() {
1208                    enforcement.profiles_indexed = CategoryLimit::new(
1209                        DataCategory::ProfileIndexed,
1210                        summary.profile_quantity.total,
1211                        limit.longest(),
1212                    );
1213
1214                    profile_limits.merge(limit);
1215                } else {
1216                    enforcement.profiles_backend = enforcement
1217                        .profiles_backend
1218                        .add_outcome_category(DataCategory::ProfileIndexed);
1219                    enforcement.profiles_ui = enforcement
1220                        .profiles_ui
1221                        .add_outcome_category(DataCategory::ProfileIndexed);
1222                }
1223            }
1224
1225            rate_limits.merge(profile_limits);
1226        }
1227
1228        // Handle replays.
1229        if summary.replay_quantity > 0 {
1230            let item_scoping = scoping.item(DataCategory::Replay);
1231            let replay_limits = self
1232                .check
1233                .apply(item_scoping, summary.replay_quantity)
1234                .await?;
1235            enforcement.replays = CategoryLimit::new(
1236                DataCategory::Replay,
1237                summary.replay_quantity,
1238                replay_limits.longest(),
1239            );
1240            rate_limits.merge(replay_limits);
1241        }
1242
1243        // Handle user report v1s, which share limits with v2.
1244        if summary.user_report_quantity > 0 {
1245            let item_scoping = scoping.item(DataCategory::UserReportV2);
1246            let user_report_v2_limits = self
1247                .check
1248                .apply(item_scoping, summary.user_report_quantity)
1249                .await?;
1250            enforcement.user_reports = CategoryLimit::new(
1251                DataCategory::UserReportV2,
1252                summary.user_report_quantity,
1253                user_report_v2_limits.longest(),
1254            );
1255            rate_limits.merge(user_report_v2_limits);
1256        }
1257
1258        // Handle monitor checkins.
1259        if summary.monitor_quantity > 0 {
1260            let item_scoping = scoping.item(DataCategory::Monitor);
1261            let checkin_limits = self
1262                .check
1263                .apply(item_scoping, summary.monitor_quantity)
1264                .await?;
1265            enforcement.check_ins = CategoryLimit::new(
1266                DataCategory::Monitor,
1267                summary.monitor_quantity,
1268                checkin_limits.longest(),
1269            );
1270            rate_limits.merge(checkin_limits);
1271        }
1272
1273        // Handle profile chunks.
1274        if summary.profile_chunk_quantity > 0 {
1275            let item_scoping = scoping.item(DataCategory::ProfileChunk);
1276            let limits = self
1277                .check
1278                .apply(item_scoping, summary.profile_chunk_quantity)
1279                .await?;
1280            enforcement.profile_chunks = CategoryLimit::new(
1281                DataCategory::ProfileChunk,
1282                summary.profile_chunk_quantity,
1283                limits.longest(),
1284            );
1285            rate_limits.merge(limits);
1286        }
1287
1288        if summary.profile_chunk_ui_quantity > 0 {
1289            let item_scoping = scoping.item(DataCategory::ProfileChunkUi);
1290            let limits = self
1291                .check
1292                .apply(item_scoping, summary.profile_chunk_ui_quantity)
1293                .await?;
1294            enforcement.profile_chunks_ui = CategoryLimit::new(
1295                DataCategory::ProfileChunkUi,
1296                summary.profile_chunk_ui_quantity,
1297                limits.longest(),
1298            );
1299            rate_limits.merge(limits);
1300        }
1301
1302        Ok((enforcement, rate_limits))
1303    }
1304
1305    async fn check_attachment_limits(
1306        &mut self,
1307        scoping: &Scoping,
1308        quantities: &AttachmentQuantity,
1309    ) -> Result<AttachmentLimits, E> {
1310        let mut attachment_limits = self
1311            .check
1312            .apply(scoping.item(DataCategory::Attachment), quantities.bytes)
1313            .await?;
1314
1315        // Note: The check here is taken from the attachments logic for consistency I think just
1316        // checking `is_empty` should be fine?
1317        if !attachment_limits.is_limited() && quantities.count > 0 {
1318            attachment_limits.merge(
1319                self.check
1320                    .apply(scoping.item(DataCategory::AttachmentItem), quantities.count)
1321                    .await?,
1322            );
1323        }
1324
1325        Ok(AttachmentLimits {
1326            bytes: CategoryLimit::new(
1327                DataCategory::Attachment,
1328                quantities.bytes,
1329                attachment_limits.longest(),
1330            ),
1331            count: CategoryLimit::new(
1332                DataCategory::AttachmentItem,
1333                quantities.count,
1334                attachment_limits.longest(),
1335            ),
1336        })
1337    }
1338}
1339
1340impl<F, E, R> fmt::Debug for EnvelopeLimiter<F, E, R> {
1341    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1342        f.debug_struct("EnvelopeLimiter")
1343            .field("event_category", &self.event_category)
1344            .finish()
1345    }
1346}
1347
1348#[cfg(test)]
1349mod tests {
1350
1351    use std::collections::{BTreeMap, BTreeSet};
1352    use std::sync::Arc;
1353
1354    use relay_base_schema::organization::OrganizationId;
1355    use relay_base_schema::project::{ProjectId, ProjectKey};
1356    use relay_metrics::MetricNamespace;
1357    use relay_quotas::RetryAfter;
1358    use relay_system::Addr;
1359    use smallvec::smallvec;
1360    use tokio::sync::Mutex;
1361
1362    use super::*;
1363    use crate::envelope::ParentId;
1364    use crate::{
1365        envelope::{AttachmentType, ContentType, SourceQuantities},
1366        extractors::RequestMeta,
1367    };
1368
1369    struct RateLimitTestCase {
1370        name: &'static str,
1371        denied_categories: &'static [DataCategory],
1372        expect_attachment_limit_active: bool,
1373        expected_limiter_calls: &'static [(DataCategory, usize)],
1374        expected_outcomes: &'static [(DataCategory, usize)],
1375    }
1376
1377    #[tokio::test]
1378    async fn test_format_rate_limits() {
1379        let mut rate_limits = RateLimits::new();
1380
1381        // Add a generic rate limit for all categories.
1382        rate_limits.add(RateLimit {
1383            categories: Default::default(),
1384            scope: RateLimitScope::Organization(OrganizationId::new(42)),
1385            reason_code: Some(ReasonCode::new("my_limit")),
1386            retry_after: RetryAfter::from_secs(42),
1387            namespaces: smallvec![],
1388        });
1389
1390        // Add a more specific rate limit for just one category.
1391        rate_limits.add(RateLimit {
1392            categories: [DataCategory::Transaction, DataCategory::Security].into(),
1393            scope: RateLimitScope::Project(ProjectId::new(21)),
1394            reason_code: None,
1395            retry_after: RetryAfter::from_secs(4711),
1396            namespaces: smallvec![],
1397        });
1398
1399        let formatted = format_rate_limits(&rate_limits);
1400        let expected = "42::organization:my_limit, 4711:transaction;security:project";
1401        assert_eq!(formatted, expected);
1402    }
1403
1404    #[tokio::test]
1405    async fn test_format_rate_limits_namespace() {
1406        let mut rate_limits = RateLimits::new();
1407
1408        // Rate limit with reason code and namespace.
1409        rate_limits.add(RateLimit {
1410            categories: [DataCategory::MetricBucket].into(),
1411            scope: RateLimitScope::Organization(OrganizationId::new(42)),
1412            reason_code: Some(ReasonCode::new("my_limit")),
1413            retry_after: RetryAfter::from_secs(42),
1414            namespaces: smallvec![MetricNamespace::Transactions, MetricNamespace::Spans],
1415        });
1416
1417        // Rate limit without reason code.
1418        rate_limits.add(RateLimit {
1419            categories: [DataCategory::MetricBucket].into(),
1420            scope: RateLimitScope::Organization(OrganizationId::new(42)),
1421            reason_code: None,
1422            retry_after: RetryAfter::from_secs(42),
1423            namespaces: smallvec![MetricNamespace::Spans],
1424        });
1425
1426        let formatted = format_rate_limits(&rate_limits);
1427        let expected = "42:metric_bucket:organization:my_limit:transactions;spans, 42:metric_bucket:organization::spans";
1428        assert_eq!(formatted, expected);
1429    }
1430
1431    #[tokio::test]
1432    async fn test_parse_invalid_rate_limits() {
1433        let scoping = Scoping {
1434            organization_id: OrganizationId::new(42),
1435            project_id: ProjectId::new(21),
1436            project_key: ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
1437            key_id: Some(17),
1438        };
1439
1440        assert!(parse_rate_limits(&scoping, "").is_ok());
1441        assert!(parse_rate_limits(&scoping, "invalid").is_ok());
1442        assert!(parse_rate_limits(&scoping, ",,,").is_ok());
1443    }
1444
1445    #[tokio::test]
1446    async fn test_parse_rate_limits() {
1447        let scoping = Scoping {
1448            organization_id: OrganizationId::new(42),
1449            project_id: ProjectId::new(21),
1450            project_key: ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
1451            key_id: Some(17),
1452        };
1453
1454        // contains "foobar", an unknown scope that should be mapped to Unknown
1455        let formatted =
1456            "42::organization:my_limit, invalid, 4711:foobar;transaction;security:project";
1457        let rate_limits: Vec<RateLimit> =
1458            parse_rate_limits(&scoping, formatted).into_iter().collect();
1459
1460        assert_eq!(
1461            rate_limits,
1462            vec![
1463                RateLimit {
1464                    categories: Default::default(),
1465                    scope: RateLimitScope::Organization(OrganizationId::new(42)),
1466                    reason_code: Some(ReasonCode::new("my_limit")),
1467                    retry_after: rate_limits[0].retry_after,
1468                    namespaces: smallvec![],
1469                },
1470                RateLimit {
1471                    categories: [
1472                        DataCategory::Unknown,
1473                        DataCategory::Transaction,
1474                        DataCategory::Security,
1475                    ]
1476                    .into(),
1477                    scope: RateLimitScope::Project(ProjectId::new(21)),
1478                    reason_code: None,
1479                    retry_after: rate_limits[1].retry_after,
1480                    namespaces: smallvec![],
1481                }
1482            ]
1483        );
1484
1485        assert_eq!(42, rate_limits[0].retry_after.remaining_seconds());
1486        assert_eq!(4711, rate_limits[1].retry_after.remaining_seconds());
1487    }
1488
1489    #[tokio::test]
1490    async fn test_parse_rate_limits_namespace() {
1491        let scoping = Scoping {
1492            organization_id: OrganizationId::new(42),
1493            project_id: ProjectId::new(21),
1494            project_key: ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
1495            key_id: Some(17),
1496        };
1497
1498        let formatted = "42:metric_bucket:organization::transactions;spans";
1499        let rate_limits: Vec<RateLimit> =
1500            parse_rate_limits(&scoping, formatted).into_iter().collect();
1501
1502        assert_eq!(
1503            rate_limits,
1504            vec![RateLimit {
1505                categories: [DataCategory::MetricBucket].into(),
1506                scope: RateLimitScope::Organization(OrganizationId::new(42)),
1507                reason_code: None,
1508                retry_after: rate_limits[0].retry_after,
1509                namespaces: smallvec![MetricNamespace::Transactions, MetricNamespace::Spans],
1510            }]
1511        );
1512    }
1513
1514    #[tokio::test]
1515    async fn test_parse_rate_limits_empty_namespace() {
1516        let scoping = Scoping {
1517            organization_id: OrganizationId::new(42),
1518            project_id: ProjectId::new(21),
1519            project_key: ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
1520            key_id: Some(17),
1521        };
1522
1523        // notice the trailing colon
1524        let formatted = "42:metric_bucket:organization:some_reason:";
1525        let rate_limits: Vec<RateLimit> =
1526            parse_rate_limits(&scoping, formatted).into_iter().collect();
1527
1528        assert_eq!(
1529            rate_limits,
1530            vec![RateLimit {
1531                categories: [DataCategory::MetricBucket].into(),
1532                scope: RateLimitScope::Organization(OrganizationId::new(42)),
1533                reason_code: Some(ReasonCode::new("some_reason")),
1534                retry_after: rate_limits[0].retry_after,
1535                namespaces: smallvec![],
1536            }]
1537        );
1538    }
1539
1540    #[tokio::test]
1541    async fn test_parse_rate_limits_only_unknown() {
1542        let scoping = Scoping {
1543            organization_id: OrganizationId::new(42),
1544            project_id: ProjectId::new(21),
1545            project_key: ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
1546            key_id: Some(17),
1547        };
1548
1549        let formatted = "42:foo;bar:organization";
1550        let rate_limits: Vec<RateLimit> =
1551            parse_rate_limits(&scoping, formatted).into_iter().collect();
1552
1553        assert_eq!(
1554            rate_limits,
1555            vec![RateLimit {
1556                categories: [DataCategory::Unknown, DataCategory::Unknown].into(),
1557                scope: RateLimitScope::Organization(OrganizationId::new(42)),
1558                reason_code: None,
1559                retry_after: rate_limits[0].retry_after,
1560                namespaces: smallvec![],
1561            },]
1562        );
1563    }
1564
1565    macro_rules! envelope {
1566        ($( $item_type:ident $( :: $attachment_type:ident )? ),*) => {{
1567            let bytes = "{\"dsn\":\"https://e12d836b15bb49d7bbf99e64295d995b:@sentry.io/42\"}";
1568            #[allow(unused_mut)]
1569            let mut envelope = Envelope::parse_bytes(bytes.into()).unwrap();
1570            $(
1571                let mut item = Item::new(ItemType::$item_type);
1572                item.set_payload(ContentType::OctetStream, "0123456789");
1573                $( item.set_attachment_type(AttachmentType::$attachment_type); )?
1574                envelope.add_item(item);
1575            )*
1576
1577            envelope
1578        }}
1579    }
1580
1581    fn rate_limit(category: DataCategory) -> RateLimit {
1582        RateLimit {
1583            categories: [category].into(),
1584            scope: RateLimitScope::Organization(OrganizationId::new(42)),
1585            reason_code: None,
1586            retry_after: RetryAfter::from_secs(60),
1587            namespaces: smallvec![],
1588        }
1589    }
1590
1591    fn trace_attachment_item(bytes: usize, parent_id: Option<ParentId>) -> Item {
1592        let mut item = Item::new(ItemType::Attachment);
1593        item.set_payload(ContentType::TraceAttachment, "0".repeat(bytes));
1594        item.set_parent_id(parent_id);
1595        item
1596    }
1597
1598    #[derive(Debug, Default)]
1599    struct MockLimiter {
1600        denied: Vec<DataCategory>,
1601        called: BTreeMap<DataCategory, usize>,
1602        checked: BTreeSet<DataCategory>,
1603    }
1604
1605    impl MockLimiter {
1606        pub fn deny(mut self, category: DataCategory) -> Self {
1607            self.denied.push(category);
1608            self
1609        }
1610
1611        pub fn check(&mut self, scoping: ItemScoping, quantity: usize) -> Result<RateLimits, ()> {
1612            let cat = scoping.category;
1613            let previous = self.called.insert(cat, quantity);
1614            assert!(previous.is_none(), "rate limiter invoked twice for {cat}");
1615
1616            let mut limits = RateLimits::new();
1617            if self.denied.contains(&cat) {
1618                limits.add(rate_limit(cat));
1619            }
1620            Ok(limits)
1621        }
1622
1623        #[track_caller]
1624        pub fn assert_call(&mut self, category: DataCategory, expected: usize) {
1625            self.checked.insert(category);
1626
1627            let quantity = self.called.get(&category).copied();
1628            assert_eq!(
1629                quantity,
1630                Some(expected),
1631                "Expected quantity `{expected}` for data category `{category}`, got {quantity:?}."
1632            );
1633        }
1634    }
1635
1636    impl Drop for MockLimiter {
1637        fn drop(&mut self) {
1638            if std::thread::panicking() {
1639                return;
1640            }
1641
1642            for checked in &self.checked {
1643                self.called.remove(checked);
1644            }
1645
1646            if self.called.is_empty() {
1647                return;
1648            }
1649
1650            let not_asserted = self
1651                .called
1652                .iter()
1653                .map(|(k, v)| format!("- {k}: {v}"))
1654                .collect::<Vec<_>>()
1655                .join("\n");
1656
1657            panic!("Following calls to the limiter were not asserted:\n{not_asserted}");
1658        }
1659    }
1660
1661    async fn enforce_and_apply(
1662        mock: Arc<Mutex<MockLimiter>>,
1663        envelope: Box<Envelope>,
1664    ) -> (Box<Envelope>, Enforcement, RateLimits) {
1665        let mut envelope = Managed::from_envelope(envelope, Addr::custom().0);
1666        let scoping = envelope.scoping();
1667
1668        #[allow(unused_mut)]
1669        let mut limiter = EnvelopeLimiter::new(CheckLimits::NonIndexed, move |s, q| {
1670            let mock = mock.clone();
1671            async move {
1672                let mut mock = mock.lock().await;
1673                mock.check(s, q)
1674            }
1675        });
1676
1677        let (enforcement, limits) = limiter.compute(&envelope, &scoping).await.unwrap();
1678
1679        // We implemented `clone` only for tests because we don't want to make `apply_with_outcomes`
1680        // &self because we want move semantics to prevent double tracking.
1681        enforcement.clone().apply_to_managed(&mut envelope);
1682        let envelope = envelope.accept(|envelope| envelope);
1683
1684        (envelope, enforcement, limits)
1685    }
1686
1687    fn mock_limiter(categories: &[DataCategory]) -> Arc<Mutex<MockLimiter>> {
1688        let mut mock = MockLimiter::default();
1689        for &category in categories {
1690            mock = mock.deny(category);
1691        }
1692
1693        Arc::new(Mutex::new(mock))
1694    }
1695
1696    #[tokio::test]
1697    async fn test_enforce_pass_empty() {
1698        let envelope = envelope![];
1699
1700        let mock = mock_limiter(&[]);
1701        let (envelope, _, limits) = enforce_and_apply(mock, envelope).await;
1702
1703        assert!(!limits.is_limited());
1704        assert!(envelope.is_empty());
1705    }
1706
1707    #[tokio::test]
1708    async fn test_enforce_limit_error_event() {
1709        let envelope = envelope![Event];
1710
1711        let mock = mock_limiter(&[DataCategory::Error]);
1712        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1713
1714        assert!(limits.is_limited());
1715        assert!(envelope.is_empty());
1716        mock.lock().await.assert_call(DataCategory::Error, 1);
1717    }
1718
1719    #[tokio::test]
1720    async fn test_enforce_limit_error_with_attachments() {
1721        let envelope = envelope![Event, Attachment];
1722
1723        let mock = mock_limiter(&[DataCategory::Error]);
1724        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1725
1726        assert!(limits.is_limited());
1727        assert!(envelope.is_empty());
1728        mock.lock().await.assert_call(DataCategory::Error, 1);
1729    }
1730
1731    #[tokio::test]
1732    async fn test_enforce_limit_minidump() {
1733        let envelope = envelope![Attachment::Minidump];
1734
1735        let mock = mock_limiter(&[DataCategory::Error]);
1736        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1737
1738        assert!(limits.is_limited());
1739        assert!(envelope.is_empty());
1740        mock.lock().await.assert_call(DataCategory::Error, 1);
1741    }
1742
1743    #[tokio::test]
1744    async fn test_enforce_limit_attachments() {
1745        let envelope = envelope![Attachment::Minidump, Attachment];
1746
1747        let mock = mock_limiter(&[DataCategory::Attachment]);
1748        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1749
1750        // Attachments would be limited, but crash reports create events and are thus allowed.
1751        assert!(limits.is_limited());
1752        assert_eq!(envelope.len(), 1);
1753        mock.lock().await.assert_call(DataCategory::Error, 1);
1754        mock.lock().await.assert_call(DataCategory::Attachment, 20);
1755    }
1756
1757    /// Limit stand-alone profiles.
1758    #[tokio::test]
1759    async fn test_enforce_limit_profiles() {
1760        let envelope = envelope![Profile, Profile];
1761
1762        let mock = mock_limiter(&[DataCategory::Profile]);
1763        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1764
1765        assert!(limits.is_limited());
1766        assert_eq!(envelope.len(), 0);
1767        mock.lock().await.assert_call(DataCategory::Profile, 2);
1768
1769        assert_eq!(
1770            get_outcomes(enforcement),
1771            vec![
1772                (DataCategory::Profile, 2),
1773                (DataCategory::ProfileIndexed, 2)
1774            ]
1775        );
1776    }
1777
1778    /// Limit profile chunks.
1779    #[tokio::test]
1780    async fn test_enforce_limit_profile_chunks_no_profile_type() {
1781        // In this test we have profile chunks which have not yet been classified, which means they
1782        // should not be rate limited.
1783        let envelope = envelope![ProfileChunk, ProfileChunk];
1784
1785        let mock = mock_limiter(&[DataCategory::ProfileChunk]);
1786        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1787        assert!(!limits.is_limited());
1788        assert_eq!(get_outcomes(enforcement), vec![]);
1789
1790        let mock = mock_limiter(&[DataCategory::ProfileChunkUi]);
1791        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1792        assert!(!limits.is_limited());
1793        assert_eq!(get_outcomes(enforcement), vec![]);
1794
1795        assert_eq!(envelope.len(), 2);
1796    }
1797
1798    #[tokio::test]
1799    async fn test_enforce_limit_profile_chunks_ui() {
1800        let mut envelope = envelope![];
1801
1802        let mut item = Item::new(ItemType::ProfileChunk);
1803        item.set_platform("python".to_owned());
1804        envelope.add_item(item);
1805        let mut item = Item::new(ItemType::ProfileChunk);
1806        item.set_platform("javascript".to_owned());
1807        envelope.add_item(item);
1808
1809        let mock = mock_limiter(&[DataCategory::ProfileChunkUi]);
1810        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1811
1812        assert!(limits.is_limited());
1813        assert_eq!(envelope.len(), 1);
1814        mock.lock()
1815            .await
1816            .assert_call(DataCategory::ProfileChunkUi, 1);
1817        mock.lock().await.assert_call(DataCategory::ProfileChunk, 1);
1818
1819        assert_eq!(
1820            get_outcomes(enforcement),
1821            vec![(DataCategory::ProfileChunkUi, 1)]
1822        );
1823    }
1824
1825    #[tokio::test]
1826    async fn test_enforce_limit_profile_chunks_backend() {
1827        let mut envelope = envelope![];
1828
1829        let mut item = Item::new(ItemType::ProfileChunk);
1830        item.set_platform("python".to_owned());
1831        envelope.add_item(item);
1832        let mut item = Item::new(ItemType::ProfileChunk);
1833        item.set_platform("javascript".to_owned());
1834        envelope.add_item(item);
1835
1836        let mock = mock_limiter(&[DataCategory::ProfileChunk]);
1837        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1838
1839        assert!(limits.is_limited());
1840        assert_eq!(envelope.len(), 1);
1841        mock.lock()
1842            .await
1843            .assert_call(DataCategory::ProfileChunkUi, 1);
1844        mock.lock().await.assert_call(DataCategory::ProfileChunk, 1);
1845
1846        assert_eq!(
1847            get_outcomes(enforcement),
1848            vec![(DataCategory::ProfileChunk, 1)]
1849        );
1850    }
1851
1852    /// Limit replays.
1853    #[tokio::test]
1854    async fn test_enforce_limit_replays() {
1855        let envelope = envelope![ReplayEvent, ReplayRecording, ReplayVideo];
1856
1857        let mock = mock_limiter(&[DataCategory::Replay]);
1858        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1859
1860        assert!(limits.is_limited());
1861        assert_eq!(envelope.len(), 0);
1862        mock.lock().await.assert_call(DataCategory::Replay, 3);
1863
1864        assert_eq!(get_outcomes(enforcement), vec![(DataCategory::Replay, 3),]);
1865    }
1866
1867    /// Limit monitor checkins.
1868    #[tokio::test]
1869    async fn test_enforce_limit_monitor_checkins() {
1870        let envelope = envelope![CheckIn];
1871
1872        let mock = mock_limiter(&[DataCategory::Monitor]);
1873        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1874
1875        assert!(limits.is_limited());
1876        assert_eq!(envelope.len(), 0);
1877        mock.lock().await.assert_call(DataCategory::Monitor, 1);
1878
1879        assert_eq!(get_outcomes(enforcement), vec![(DataCategory::Monitor, 1)])
1880    }
1881
1882    #[tokio::test]
1883    async fn test_enforce_pass_minidump() {
1884        let envelope = envelope![Attachment::Minidump];
1885
1886        let mock = mock_limiter(&[DataCategory::Attachment]);
1887        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1888
1889        // If only crash report attachments are present, we don't emit a rate limit.
1890        assert!(!limits.is_limited());
1891        assert_eq!(envelope.len(), 1);
1892        mock.lock().await.assert_call(DataCategory::Error, 1);
1893        mock.lock().await.assert_call(DataCategory::Attachment, 10);
1894    }
1895
1896    #[tokio::test]
1897    async fn test_enforce_skip_rate_limited() {
1898        let mut envelope = envelope![];
1899
1900        let mut item = Item::new(ItemType::Attachment);
1901        item.set_payload(ContentType::OctetStream, "0123456789");
1902        item.set_rate_limited(true);
1903        envelope.add_item(item);
1904
1905        let mock = mock_limiter(&[DataCategory::Error]);
1906        let (envelope, _, limits) = enforce_and_apply(mock, envelope).await;
1907
1908        assert!(!limits.is_limited()); // No new rate limits applied.
1909        assert_eq!(envelope.len(), 1); // The item was retained
1910    }
1911
1912    #[tokio::test]
1913    async fn test_enforce_pass_sessions() {
1914        let envelope = envelope![Session, Session, Session];
1915
1916        let mock = mock_limiter(&[DataCategory::Error]);
1917        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1918
1919        // If only crash report attachments are present, we don't emit a rate limit.
1920        assert!(!limits.is_limited());
1921        assert_eq!(envelope.len(), 3);
1922        mock.lock().await.assert_call(DataCategory::Session, 3);
1923    }
1924
1925    #[tokio::test]
1926    async fn test_enforce_limit_sessions() {
1927        let envelope = envelope![Session, Session, Event];
1928
1929        let mock = mock_limiter(&[DataCategory::Session]);
1930        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1931
1932        // If only crash report attachments are present, we don't emit a rate limit.
1933        assert!(limits.is_limited());
1934        assert_eq!(envelope.len(), 1);
1935        mock.lock().await.assert_call(DataCategory::Error, 1);
1936        mock.lock().await.assert_call(DataCategory::Session, 2);
1937    }
1938
1939    #[tokio::test]
1940    #[cfg(feature = "processing")]
1941    async fn test_enforce_limit_assumed_event() {
1942        let envelope = envelope![Transaction];
1943
1944        let mock = mock_limiter(&[DataCategory::Transaction]);
1945        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1946
1947        assert!(limits.is_limited());
1948        assert!(envelope.is_empty()); // obviously
1949        mock.lock().await.assert_call(DataCategory::Transaction, 1);
1950    }
1951
1952    #[tokio::test]
1953    #[cfg(feature = "processing")]
1954    async fn test_enforce_limit_assumed_attachments() {
1955        let envelope = envelope![Event, Attachment, Attachment];
1956
1957        let mock = mock_limiter(&[DataCategory::Error]);
1958        let (envelope, _, limits) = enforce_and_apply(mock.clone(), envelope).await;
1959
1960        assert!(limits.is_limited());
1961        assert!(envelope.is_empty());
1962        mock.lock().await.assert_call(DataCategory::Error, 1);
1963    }
1964
1965    #[tokio::test]
1966    async fn test_enforce_transaction() {
1967        let envelope = envelope![Transaction];
1968
1969        let mock = mock_limiter(&[DataCategory::Transaction]);
1970        let (_, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
1971
1972        assert!(limits.is_limited());
1973        assert!(enforcement.event_indexed.is_active());
1974        assert!(enforcement.event.is_active());
1975        mock.lock().await.assert_call(DataCategory::Transaction, 1);
1976
1977        assert_eq!(
1978            get_outcomes(enforcement),
1979            vec![
1980                (DataCategory::Transaction, 1),
1981                (DataCategory::TransactionIndexed, 1),
1982                (DataCategory::Span, 1),
1983                (DataCategory::SpanIndexed, 1),
1984            ]
1985        );
1986    }
1987
1988    #[tokio::test]
1989    async fn test_enforce_transaction_non_indexed() {
1990        let envelope = envelope![Transaction, Profile];
1991        let scoping = envelope
1992            .headers()
1993            .meta()
1994            .get_partial_scoping()
1995            .into_scoping();
1996
1997        let mock = mock_limiter(&[DataCategory::TransactionIndexed]);
1998
1999        let mock_clone = mock.clone();
2000        let limiter = EnvelopeLimiter::new(CheckLimits::NonIndexed, move |s, q| {
2001            let mock_clone = mock_clone.clone();
2002            async move {
2003                let mut mock = mock_clone.lock().await;
2004                mock.check(s, q)
2005            }
2006        });
2007        let (enforcement, limits) = limiter.compute(&envelope, &scoping).await.unwrap();
2008
2009        assert!(!limits.is_limited());
2010        assert!(!enforcement.event_indexed.is_active());
2011        assert!(!enforcement.event.is_active());
2012        assert!(!enforcement.profiles_indexed.is_active());
2013        assert!(!enforcement.profiles.is_active());
2014        assert!(!enforcement.spans.is_active());
2015        assert!(!enforcement.spans_indexed.is_active());
2016        mock.lock().await.assert_call(DataCategory::Transaction, 1);
2017        mock.lock().await.assert_call(DataCategory::Profile, 1);
2018        mock.lock().await.assert_call(DataCategory::Span, 1);
2019    }
2020
2021    #[tokio::test]
2022    async fn test_enforce_transaction_no_indexing_quota() {
2023        let envelope = envelope![Transaction];
2024
2025        let mock = mock_limiter(&[DataCategory::TransactionIndexed]);
2026        let (_, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
2027
2028        assert!(!limits.is_limited());
2029        assert!(!enforcement.event_indexed.is_active());
2030        assert!(!enforcement.event.is_active());
2031        mock.lock().await.assert_call(DataCategory::Transaction, 1);
2032        mock.lock().await.assert_call(DataCategory::Span, 1);
2033    }
2034
2035    #[tokio::test]
2036    async fn test_enforce_transaction_attachment_enforced() {
2037        let envelope = envelope![Transaction, Attachment];
2038
2039        let mock = mock_limiter(&[DataCategory::Transaction]);
2040        let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2041
2042        assert!(enforcement.event.is_active());
2043        assert!(enforcement.attachments_limits.event.is_active());
2044        mock.lock().await.assert_call(DataCategory::Transaction, 1);
2045    }
2046
2047    fn get_outcomes(enforcement: Enforcement) -> Vec<(DataCategory, usize)> {
2048        enforcement
2049            .get_outcomes()
2050            .map(|(_, data_category, quantity)| (data_category, quantity))
2051            .collect::<Vec<_>>()
2052    }
2053
2054    #[tokio::test]
2055    async fn test_enforce_transaction_profile_enforced() {
2056        let envelope = envelope![Transaction, Profile];
2057
2058        let mock = mock_limiter(&[DataCategory::Transaction]);
2059        let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2060
2061        assert!(enforcement.event.is_active());
2062        assert!(enforcement.profiles.is_active());
2063        mock.lock().await.assert_call(DataCategory::Transaction, 1);
2064
2065        assert_eq!(
2066            get_outcomes(enforcement),
2067            vec![
2068                (DataCategory::Transaction, 1),
2069                (DataCategory::TransactionIndexed, 1),
2070                (DataCategory::Profile, 1),
2071                (DataCategory::ProfileIndexed, 1),
2072                (DataCategory::Span, 1),
2073                (DataCategory::SpanIndexed, 1),
2074            ]
2075        );
2076    }
2077
2078    #[tokio::test]
2079    async fn test_enforce_transaction_standalone_profile_enforced() {
2080        // When the transaction is sampled, the profile survives as standalone.
2081        let envelope = envelope![Profile];
2082
2083        let mock = mock_limiter(&[DataCategory::Transaction]);
2084        let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2085
2086        assert!(enforcement.profiles.is_active());
2087        mock.lock().await.assert_call(DataCategory::Profile, 1);
2088        mock.lock().await.assert_call(DataCategory::Transaction, 0);
2089
2090        assert_eq!(
2091            get_outcomes(enforcement),
2092            vec![
2093                (DataCategory::Profile, 1),
2094                (DataCategory::ProfileIndexed, 1),
2095            ]
2096        );
2097    }
2098
2099    #[tokio::test]
2100    async fn test_enforce_transaction_attachment_enforced_indexing_quota() {
2101        let envelope = envelope![Transaction, Attachment];
2102
2103        let mock = mock_limiter(&[DataCategory::TransactionIndexed]);
2104        let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2105
2106        assert!(!enforcement.event.is_active());
2107        assert!(!enforcement.event_indexed.is_active());
2108        assert!(!enforcement.attachments_limits.event.is_active());
2109        mock.lock().await.assert_call(DataCategory::Transaction, 1);
2110        mock.lock().await.assert_call(DataCategory::Span, 1);
2111        mock.lock().await.assert_call(DataCategory::Attachment, 10);
2112        mock.lock()
2113            .await
2114            .assert_call(DataCategory::AttachmentItem, 1);
2115
2116        assert_eq!(get_outcomes(enforcement), vec![]);
2117    }
2118
2119    #[tokio::test]
2120    async fn test_enforce_span() {
2121        let envelope = envelope![Span, Span];
2122
2123        let mock = mock_limiter(&[DataCategory::Span]);
2124        let (_, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
2125
2126        assert!(limits.is_limited());
2127        assert!(enforcement.spans_indexed.is_active());
2128        assert!(enforcement.spans.is_active());
2129        mock.lock().await.assert_call(DataCategory::Span, 2);
2130
2131        assert_eq!(
2132            get_outcomes(enforcement),
2133            vec![(DataCategory::Span, 2), (DataCategory::SpanIndexed, 2)]
2134        );
2135    }
2136
2137    #[tokio::test]
2138    async fn test_enforce_span_no_indexing_quota() {
2139        let envelope = envelope![Span, Span];
2140
2141        let mock = mock_limiter(&[DataCategory::SpanIndexed]);
2142        let (_, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
2143
2144        assert!(!limits.is_limited());
2145        assert!(!enforcement.spans_indexed.is_active());
2146        assert!(!enforcement.spans.is_active());
2147        mock.lock().await.assert_call(DataCategory::Span, 2);
2148
2149        assert_eq!(get_outcomes(enforcement), vec![]);
2150    }
2151
2152    #[test]
2153    fn test_source_quantity_for_total_quantity() {
2154        let dsn = "https://e12d836b15bb49d7bbf99e64295d995b:@sentry.io/42"
2155            .parse()
2156            .unwrap();
2157        let request_meta = RequestMeta::new(dsn);
2158
2159        let mut envelope = Envelope::from_request(None, request_meta);
2160
2161        let mut item = Item::new(ItemType::MetricBuckets);
2162        item.set_source_quantities(SourceQuantities {
2163            transactions: 5,
2164            spans: 0,
2165            buckets: 5,
2166        });
2167        envelope.add_item(item);
2168
2169        let mut item = Item::new(ItemType::MetricBuckets);
2170        item.set_source_quantities(SourceQuantities {
2171            transactions: 2,
2172            spans: 0,
2173            buckets: 3,
2174        });
2175        envelope.add_item(item);
2176
2177        let summary = EnvelopeSummary::compute(&envelope);
2178
2179        assert_eq!(summary.secondary_transaction_quantity, 7);
2180    }
2181
2182    #[tokio::test]
2183    async fn test_enforce_limit_logs_count() {
2184        let envelope = envelope![Log, Log];
2185
2186        let mock = mock_limiter(&[DataCategory::LogItem]);
2187        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
2188
2189        assert!(limits.is_limited());
2190        assert_eq!(envelope.len(), 0);
2191        mock.lock().await.assert_call(DataCategory::LogItem, 2);
2192
2193        assert_eq!(
2194            get_outcomes(enforcement),
2195            vec![(DataCategory::LogItem, 2), (DataCategory::LogByte, 20)]
2196        );
2197    }
2198
2199    #[tokio::test]
2200    async fn test_enforce_limit_logs_bytes() {
2201        let envelope = envelope![Log, Log];
2202
2203        let mock = mock_limiter(&[DataCategory::LogByte]);
2204        let (envelope, enforcement, limits) = enforce_and_apply(mock.clone(), envelope).await;
2205
2206        assert!(limits.is_limited());
2207        assert_eq!(envelope.len(), 0);
2208        mock.lock().await.assert_call(DataCategory::LogItem, 2);
2209        mock.lock().await.assert_call(DataCategory::LogByte, 20);
2210
2211        assert_eq!(
2212            get_outcomes(enforcement),
2213            vec![(DataCategory::LogItem, 2), (DataCategory::LogByte, 20)]
2214        );
2215    }
2216
2217    #[tokio::test]
2218    async fn test_enforce_standalone_span_attachment() {
2219        let test_cases = &[
2220            RateLimitTestCase {
2221                name: "span_limit",
2222                denied_categories: &[DataCategory::Span],
2223                expect_attachment_limit_active: true,
2224                expected_limiter_calls: &[(DataCategory::Span, 0)],
2225                expected_outcomes: &[
2226                    (DataCategory::Attachment, 7),
2227                    (DataCategory::AttachmentItem, 1),
2228                ],
2229            },
2230            RateLimitTestCase {
2231                name: "span_indexed_limit",
2232                denied_categories: &[DataCategory::SpanIndexed],
2233                expect_attachment_limit_active: false,
2234                expected_limiter_calls: &[
2235                    (DataCategory::Span, 0),
2236                    (DataCategory::Attachment, 7),
2237                    (DataCategory::AttachmentItem, 1),
2238                ],
2239                expected_outcomes: &[],
2240            },
2241            RateLimitTestCase {
2242                name: "attachment_limit",
2243                denied_categories: &[DataCategory::Attachment],
2244                expect_attachment_limit_active: true,
2245                expected_limiter_calls: &[(DataCategory::Span, 0), (DataCategory::Attachment, 7)],
2246                expected_outcomes: &[
2247                    (DataCategory::Attachment, 7),
2248                    (DataCategory::AttachmentItem, 1),
2249                ],
2250            },
2251            RateLimitTestCase {
2252                name: "attachment_indexed_limit",
2253                denied_categories: &[DataCategory::AttachmentItem],
2254                expect_attachment_limit_active: true,
2255                expected_limiter_calls: &[
2256                    (DataCategory::Span, 0),
2257                    (DataCategory::Attachment, 7),
2258                    (DataCategory::AttachmentItem, 1),
2259                ],
2260                expected_outcomes: &[
2261                    (DataCategory::Attachment, 7),
2262                    (DataCategory::AttachmentItem, 1),
2263                ],
2264            },
2265            RateLimitTestCase {
2266                name: "transaction_limit",
2267                denied_categories: &[DataCategory::Transaction],
2268                expect_attachment_limit_active: false,
2269                expected_limiter_calls: &[
2270                    (DataCategory::Span, 0),
2271                    (DataCategory::Attachment, 7),
2272                    (DataCategory::AttachmentItem, 1),
2273                ],
2274                expected_outcomes: &[],
2275            },
2276            RateLimitTestCase {
2277                name: "error_limit",
2278                denied_categories: &[DataCategory::Error],
2279                expect_attachment_limit_active: false,
2280                expected_limiter_calls: &[
2281                    (DataCategory::Span, 0),
2282                    (DataCategory::Attachment, 7),
2283                    (DataCategory::AttachmentItem, 1),
2284                ],
2285                expected_outcomes: &[],
2286            },
2287            RateLimitTestCase {
2288                name: "no_limits",
2289                denied_categories: &[],
2290                expect_attachment_limit_active: false,
2291                expected_limiter_calls: &[
2292                    (DataCategory::Span, 0),
2293                    (DataCategory::Attachment, 7),
2294                    (DataCategory::AttachmentItem, 1),
2295                ],
2296                expected_outcomes: &[],
2297            },
2298        ];
2299
2300        for RateLimitTestCase {
2301            name,
2302            denied_categories,
2303            expect_attachment_limit_active,
2304            expected_limiter_calls,
2305            expected_outcomes,
2306        } in test_cases
2307        {
2308            let mut envelope = envelope![];
2309            envelope.add_item(trace_attachment_item(7, Some(ParentId::SpanId(None))));
2310
2311            let mock = mock_limiter(denied_categories);
2312            let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2313
2314            for &(category, quantity) in *expected_limiter_calls {
2315                mock.lock().await.assert_call(category, quantity);
2316            }
2317
2318            assert_eq!(
2319                enforcement.attachments_limits.span.bytes.is_active(),
2320                *expect_attachment_limit_active,
2321                "{name}: span_attachment byte limit mismatch"
2322            );
2323            assert_eq!(
2324                enforcement.attachments_limits.span.count.is_active(),
2325                *expect_attachment_limit_active,
2326                "{name}: span_attachment count limit mismatch"
2327            );
2328
2329            assert_eq!(
2330                get_outcomes(enforcement),
2331                *expected_outcomes,
2332                "{name}: outcome mismatch"
2333            );
2334        }
2335    }
2336
2337    #[tokio::test]
2338    async fn test_enforce_span_with_span_attachment() {
2339        let test_cases = &[
2340            RateLimitTestCase {
2341                name: "span_limit",
2342                denied_categories: &[DataCategory::Span, DataCategory::Attachment], // Attachment here has no effect
2343                expect_attachment_limit_active: true,
2344                expected_limiter_calls: &[(DataCategory::Span, 1)],
2345                expected_outcomes: &[
2346                    (DataCategory::Attachment, 7),
2347                    (DataCategory::AttachmentItem, 1),
2348                    (DataCategory::Span, 1),
2349                    (DataCategory::SpanIndexed, 1),
2350                ],
2351            },
2352            RateLimitTestCase {
2353                name: "span_indexed_limit",
2354                denied_categories: &[DataCategory::SpanIndexed, DataCategory::Attachment],
2355                expect_attachment_limit_active: true,
2356                expected_limiter_calls: &[(DataCategory::Span, 1), (DataCategory::Attachment, 7)],
2357                expected_outcomes: &[
2358                    (DataCategory::Attachment, 7),
2359                    (DataCategory::AttachmentItem, 1),
2360                ],
2361            },
2362            RateLimitTestCase {
2363                name: "attachment_limit",
2364                denied_categories: &[DataCategory::Attachment],
2365                expect_attachment_limit_active: true,
2366                expected_limiter_calls: &[(DataCategory::Span, 1), (DataCategory::Attachment, 7)],
2367                expected_outcomes: &[
2368                    (DataCategory::Attachment, 7),
2369                    (DataCategory::AttachmentItem, 1),
2370                ],
2371            },
2372            RateLimitTestCase {
2373                name: "attachment_indexed_limit",
2374                denied_categories: &[DataCategory::AttachmentItem],
2375                expect_attachment_limit_active: true,
2376                expected_limiter_calls: &[
2377                    (DataCategory::Span, 1),
2378                    (DataCategory::Attachment, 7),
2379                    (DataCategory::AttachmentItem, 1),
2380                ],
2381                expected_outcomes: &[
2382                    (DataCategory::Attachment, 7),
2383                    (DataCategory::AttachmentItem, 1),
2384                ],
2385            },
2386            RateLimitTestCase {
2387                name: "transaction_limit",
2388                denied_categories: &[DataCategory::Transaction],
2389                expect_attachment_limit_active: false,
2390                expected_limiter_calls: &[
2391                    (DataCategory::Span, 1),
2392                    (DataCategory::Attachment, 7),
2393                    (DataCategory::AttachmentItem, 1),
2394                ],
2395                expected_outcomes: &[],
2396            },
2397            RateLimitTestCase {
2398                name: "error_limit",
2399                denied_categories: &[DataCategory::Error],
2400                expect_attachment_limit_active: false,
2401                expected_limiter_calls: &[
2402                    (DataCategory::Span, 1),
2403                    (DataCategory::Attachment, 7),
2404                    (DataCategory::AttachmentItem, 1),
2405                ],
2406                expected_outcomes: &[],
2407            },
2408            RateLimitTestCase {
2409                name: "no_limits",
2410                denied_categories: &[],
2411                expect_attachment_limit_active: false,
2412                expected_limiter_calls: &[
2413                    (DataCategory::Span, 1),
2414                    (DataCategory::Attachment, 7),
2415                    (DataCategory::AttachmentItem, 1),
2416                ],
2417                expected_outcomes: &[],
2418            },
2419        ];
2420
2421        for RateLimitTestCase {
2422            name,
2423            denied_categories,
2424            expect_attachment_limit_active,
2425            expected_limiter_calls,
2426            expected_outcomes,
2427        } in test_cases
2428        {
2429            let mut envelope = envelope![Span];
2430            envelope.add_item(trace_attachment_item(7, Some(ParentId::SpanId(None))));
2431
2432            let mock = mock_limiter(denied_categories);
2433            let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2434
2435            for &(category, quantity) in *expected_limiter_calls {
2436                mock.lock().await.assert_call(category, quantity);
2437            }
2438
2439            assert_eq!(
2440                enforcement.attachments_limits.span.bytes.is_active(),
2441                *expect_attachment_limit_active,
2442                "{name}: span_attachment byte limit mismatch"
2443            );
2444            assert_eq!(
2445                enforcement.attachments_limits.span.count.is_active(),
2446                *expect_attachment_limit_active,
2447                "{name}: span_attachment count limit mismatch"
2448            );
2449
2450            assert_eq!(
2451                get_outcomes(enforcement),
2452                *expected_outcomes,
2453                "{name}: outcome mismatch"
2454            );
2455        }
2456    }
2457
2458    #[tokio::test]
2459    async fn test_enforce_transaction_span_attachment() {
2460        let test_cases = &[
2461            RateLimitTestCase {
2462                name: "transaction_limit",
2463                denied_categories: &[DataCategory::Transaction],
2464                expect_attachment_limit_active: true,
2465                expected_limiter_calls: &[(DataCategory::Transaction, 1)],
2466                expected_outcomes: &[
2467                    (DataCategory::Transaction, 1),
2468                    (DataCategory::TransactionIndexed, 1),
2469                    (DataCategory::Attachment, 7),
2470                    (DataCategory::AttachmentItem, 1),
2471                    (DataCategory::Span, 1),
2472                    (DataCategory::SpanIndexed, 1),
2473                ],
2474            },
2475            RateLimitTestCase {
2476                name: "error_limit",
2477                denied_categories: &[DataCategory::Error],
2478                expect_attachment_limit_active: false,
2479                expected_limiter_calls: &[
2480                    (DataCategory::Transaction, 1),
2481                    (DataCategory::Span, 1),
2482                    (DataCategory::Attachment, 7),
2483                    (DataCategory::AttachmentItem, 1),
2484                ],
2485                expected_outcomes: &[],
2486            },
2487            RateLimitTestCase {
2488                name: "no_limits",
2489                denied_categories: &[],
2490                expect_attachment_limit_active: false,
2491                expected_limiter_calls: &[
2492                    (DataCategory::Transaction, 1),
2493                    (DataCategory::Span, 1),
2494                    (DataCategory::Attachment, 7),
2495                    (DataCategory::AttachmentItem, 1),
2496                ],
2497                expected_outcomes: &[],
2498            },
2499        ];
2500
2501        for RateLimitTestCase {
2502            name,
2503            denied_categories,
2504            expect_attachment_limit_active,
2505            expected_limiter_calls,
2506            expected_outcomes,
2507        } in test_cases
2508        {
2509            let mut envelope = envelope![Transaction];
2510            envelope.add_item(trace_attachment_item(7, Some(ParentId::SpanId(None))));
2511
2512            let mock = mock_limiter(denied_categories);
2513            let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2514
2515            for &(category, quantity) in *expected_limiter_calls {
2516                mock.lock().await.assert_call(category, quantity);
2517            }
2518
2519            assert_eq!(
2520                enforcement.attachments_limits.span.bytes.is_active(),
2521                *expect_attachment_limit_active,
2522                "{name}: span_attachment byte limit mismatch"
2523            );
2524            assert_eq!(
2525                enforcement.attachments_limits.span.count.is_active(),
2526                *expect_attachment_limit_active,
2527                "{name}: span_attachment count limit mismatch"
2528            );
2529
2530            assert_eq!(
2531                get_outcomes(enforcement),
2532                *expected_outcomes,
2533                "{name}: outcome mismatch"
2534            );
2535        }
2536    }
2537
2538    #[tokio::test]
2539    async fn test_enforce_standalone_trace_attachment() {
2540        let test_cases = &[
2541            RateLimitTestCase {
2542                name: "attachment_limit",
2543                denied_categories: &[DataCategory::Attachment],
2544                expect_attachment_limit_active: true,
2545                expected_limiter_calls: &[(DataCategory::Attachment, 7)],
2546                expected_outcomes: &[
2547                    (DataCategory::Attachment, 7),
2548                    (DataCategory::AttachmentItem, 1),
2549                ],
2550            },
2551            RateLimitTestCase {
2552                name: "attachment_limit_and_attachment_item_limit",
2553                denied_categories: &[DataCategory::Attachment, DataCategory::AttachmentItem],
2554                expect_attachment_limit_active: true,
2555                expected_limiter_calls: &[(DataCategory::Attachment, 7)],
2556                expected_outcomes: &[
2557                    (DataCategory::Attachment, 7),
2558                    (DataCategory::AttachmentItem, 1),
2559                ],
2560            },
2561            RateLimitTestCase {
2562                name: "attachment_item_limit",
2563                denied_categories: &[DataCategory::AttachmentItem],
2564                expect_attachment_limit_active: true,
2565                expected_limiter_calls: &[
2566                    (DataCategory::Attachment, 7),
2567                    (DataCategory::AttachmentItem, 1),
2568                ],
2569                expected_outcomes: &[
2570                    (DataCategory::Attachment, 7),
2571                    (DataCategory::AttachmentItem, 1),
2572                ],
2573            },
2574            RateLimitTestCase {
2575                name: "no_limits",
2576                denied_categories: &[],
2577                expect_attachment_limit_active: false,
2578                expected_limiter_calls: &[
2579                    (DataCategory::Attachment, 7),
2580                    (DataCategory::AttachmentItem, 1),
2581                ],
2582                expected_outcomes: &[],
2583            },
2584        ];
2585
2586        for RateLimitTestCase {
2587            name,
2588            denied_categories,
2589            expect_attachment_limit_active,
2590            expected_limiter_calls,
2591            expected_outcomes,
2592        } in test_cases
2593        {
2594            let mut envelope = envelope![];
2595            envelope.add_item(trace_attachment_item(7, None));
2596
2597            let mock = mock_limiter(denied_categories);
2598            let (_, enforcement, _) = enforce_and_apply(mock.clone(), envelope).await;
2599
2600            for &(category, quantity) in *expected_limiter_calls {
2601                mock.lock().await.assert_call(category, quantity);
2602            }
2603
2604            assert_eq!(
2605                enforcement.attachments_limits.trace.bytes.is_active(),
2606                *expect_attachment_limit_active,
2607                "{name}: trace_attachment byte limit mismatch"
2608            );
2609            assert_eq!(
2610                enforcement.attachments_limits.trace.count.is_active(),
2611                *expect_attachment_limit_active,
2612                "{name}: trace_attachment count limit mismatch"
2613            );
2614
2615            assert_eq!(
2616                get_outcomes(enforcement),
2617                *expected_outcomes,
2618                "{name}: outcome mismatch"
2619            );
2620        }
2621    }
2622}