Skip to main content

relay_server/services/
processor.rs

1use std::borrow::Cow;
2use std::collections::{BTreeMap, BTreeSet, HashMap};
3use std::error::Error;
4use std::fmt::Debug;
5use std::future::Future;
6use std::io::Write;
7use std::pin::Pin;
8use std::sync::Arc;
9use std::time::Duration;
10
11use anyhow::Context;
12use brotli::CompressorWriter as BrotliEncoder;
13use bytes::Bytes;
14use chrono::{DateTime, Utc};
15use flate2::Compression;
16use flate2::write::{GzEncoder, ZlibEncoder};
17use futures::future::BoxFuture;
18use relay_base_schema::project::{ProjectId, ProjectKey};
19use relay_cogs::{AppFeature, Cogs, FeatureWeights, ResourceId, Token};
20use relay_common::time::UnixTimestamp;
21use relay_config::{Config, ConfigSnapshot, EmitOutcomes, HttpEncoding, UpstreamDescriptor};
22use relay_event_normalization::{ClockDriftProcessor, GeoIpLookup};
23use relay_event_schema::processor::ProcessingAction;
24use relay_event_schema::protocol::ClientReport;
25use relay_filter::FilterStatKey;
26use relay_log::sentry::SentryFutureExt;
27use relay_metrics::{Bucket, BucketMetadata, BucketView, BucketsView, MetricNamespace};
28use relay_quotas::{RateLimits, Scoping};
29use relay_sampling::evaluation::SamplingDecision;
30use relay_statsd::metric;
31use relay_system::{Addr, FromMessage, NoResponse, Service};
32use reqwest::header;
33use zstd::stream::Encoder as ZstdEncoder;
34
35use crate::envelope::{self, ContentType, Envelope, EnvelopeError, Item, ItemType};
36use crate::extractors::{PartialDsn, RequestMeta, RequestTrust};
37use crate::managed::ManagedEnvelope;
38use crate::metrics::{MetricOutcomes, MetricsLimiter, MinimalTrackableBucket};
39use crate::metrics_extraction::ExtractedMetrics;
40use crate::processing::errors::SwitchProcessingError;
41use crate::processing::relay::RelayProcessor;
42use crate::processing::{Forward as _, Output, Outputs, QuotaRateLimiter};
43use crate::service::ServiceError;
44use crate::services::global_config::GlobalConfigHandle;
45use crate::services::metrics::{Aggregator, FlushBuckets, MergeBuckets, ProjectBuckets};
46use crate::services::outcome::{self, DiscardItemType, DiscardReason, Outcome, TrackOutcome};
47use crate::services::projects::cache::ProjectCacheHandle;
48use crate::services::projects::project::{ProjectInfo, ProjectState};
49use crate::services::upstream::{
50    SendRequest, Sign, SignatureType, UpstreamRelay, UpstreamRequest, UpstreamRequestError,
51};
52use crate::statsd::{RelayCounters, RelayDistributions, RelayTimers};
53use crate::utils;
54use crate::{http, processing};
55use relay_threading::AsyncPool;
56#[cfg(feature = "processing")]
57use {
58    crate::services::objectstore::Objectstore,
59    crate::services::store::Store,
60    itertools::Itertools,
61    relay_dynamic_config::GlobalConfig,
62    relay_quotas::{Quota, RateLimitingError, RedisRateLimiter},
63    relay_redis::RedisClients,
64    std::time::Instant,
65    symbolic_unreal::{Unreal4Error, Unreal4ErrorKind},
66};
67
68mod metrics;
69
70/// The minimum clock drift for correction to apply.
71pub const MINIMUM_CLOCK_DRIFT: Duration = Duration::from_secs(55 * 60);
72
73/// An error returned when handling [`ProcessEnvelope`].
74#[derive(Debug, thiserror::Error)]
75pub enum ProcessingError {
76    #[error("invalid json in event")]
77    InvalidJson(#[source] serde_json::Error),
78
79    #[error("invalid message pack event payload")]
80    InvalidMsgpack(#[from] rmp_serde::decode::Error),
81
82    #[error("event data too deeply nested")]
83    NestingTooDeep,
84
85    #[cfg(feature = "processing")]
86    #[error("invalid unreal crash report")]
87    InvalidUnrealReport(#[source] Unreal4Error),
88
89    #[error("event payload too large")]
90    PayloadTooLarge(DiscardItemType),
91
92    #[error("invalid transaction event")]
93    InvalidTransaction,
94
95    #[error("the item is not allowed/supported in this envelope")]
96    UnsupportedItem,
97
98    #[error("envelope processor failed")]
99    ProcessingFailed(#[from] ProcessingAction),
100
101    #[error("duplicate {0} in event")]
102    DuplicateItem(ItemType),
103
104    #[error("failed to extract event payload")]
105    NoEventPayload,
106
107    #[error("invalid security report type: {0:?}")]
108    InvalidSecurityType(Bytes),
109
110    #[error("unsupported security report type")]
111    UnsupportedSecurityType,
112
113    #[error("invalid security report")]
114    InvalidSecurityReport(#[source] serde_json::Error),
115
116    #[error("event filtered with reason: {0:?}")]
117    EventFiltered(FilterStatKey),
118
119    #[error("could not serialize event payload")]
120    SerializeFailed(#[source] serde_json::Error),
121
122    #[cfg(feature = "processing")]
123    #[error("failed to apply quotas")]
124    QuotasFailed(#[from] RateLimitingError),
125
126    #[error("nintendo switch dying message processing failed {0:?}")]
127    InvalidNintendoDyingMessage(#[source] SwitchProcessingError),
128
129    #[cfg(all(sentry, feature = "processing"))]
130    #[error("playstation dump processing failed: {0}")]
131    InvalidPlaystationDump(String),
132
133    #[cfg(feature = "processing")]
134    #[error("invalid attachment reference")]
135    InvalidAttachmentRef,
136}
137
138impl ProcessingError {
139    pub fn to_outcome(&self) -> Option<Outcome> {
140        match self {
141            Self::PayloadTooLarge(payload_type) => {
142                Some(Outcome::Invalid(DiscardReason::ItemTooLarge(*payload_type)))
143            }
144            Self::InvalidJson(_) => Some(Outcome::Invalid(DiscardReason::InvalidJson)),
145            Self::InvalidMsgpack(_) => Some(Outcome::Invalid(DiscardReason::InvalidMsgpack)),
146            Self::NestingTooDeep => Some(Outcome::Invalid(DiscardReason::NestingTooDeep)),
147            Self::InvalidSecurityType(_) => {
148                Some(Outcome::Invalid(DiscardReason::SecurityReportType))
149            }
150            Self::UnsupportedItem => Some(Outcome::Invalid(DiscardReason::InvalidEnvelope)),
151            Self::InvalidSecurityReport(_) => Some(Outcome::Invalid(DiscardReason::SecurityReport)),
152            Self::UnsupportedSecurityType => Some(Outcome::Filtered(FilterStatKey::InvalidCsp)),
153            Self::InvalidTransaction => Some(Outcome::Invalid(DiscardReason::InvalidTransaction)),
154            Self::DuplicateItem(_) => Some(Outcome::Invalid(DiscardReason::DuplicateItem)),
155            Self::NoEventPayload => Some(Outcome::Invalid(DiscardReason::NoEventPayload)),
156            Self::InvalidNintendoDyingMessage(_) => Some(Outcome::Invalid(DiscardReason::Payload)),
157            #[cfg(all(sentry, feature = "processing"))]
158            Self::InvalidPlaystationDump(_) => Some(Outcome::Invalid(DiscardReason::Payload)),
159            #[cfg(feature = "processing")]
160            Self::InvalidUnrealReport(err) if err.kind() == Unreal4ErrorKind::BadCompression => {
161                Some(Outcome::Invalid(DiscardReason::InvalidCompression))
162            }
163            #[cfg(feature = "processing")]
164            Self::InvalidUnrealReport(_) => Some(Outcome::Invalid(DiscardReason::ProcessUnreal)),
165            Self::SerializeFailed(_) | Self::ProcessingFailed(_) => {
166                Some(Outcome::Invalid(DiscardReason::Internal))
167            }
168            #[cfg(feature = "processing")]
169            Self::QuotasFailed(_) => Some(Outcome::Invalid(DiscardReason::Internal)),
170            Self::EventFiltered(key) => Some(Outcome::Filtered(key.clone())),
171
172            #[cfg(feature = "processing")]
173            Self::InvalidAttachmentRef => {
174                Some(Outcome::Invalid(DiscardReason::InvalidAttachmentRef))
175            }
176        }
177    }
178}
179
180#[cfg(feature = "processing")]
181impl From<Unreal4Error> for ProcessingError {
182    fn from(err: Unreal4Error) -> Self {
183        match err.kind() {
184            Unreal4ErrorKind::TooLarge => Self::PayloadTooLarge(ItemType::UnrealReport.into()),
185            _ => ProcessingError::InvalidUnrealReport(err),
186        }
187    }
188}
189
190/// A container for extracted metrics during processing.
191///
192/// The container enforces that the extracted metrics are correctly tagged
193/// with the dynamic sampling decision.
194#[derive(Debug)]
195pub struct ProcessingExtractedMetrics {
196    metrics: ExtractedMetrics,
197}
198
199impl ProcessingExtractedMetrics {
200    pub fn new() -> Self {
201        Self {
202            metrics: ExtractedMetrics::default(),
203        }
204    }
205
206    pub fn into_inner(self) -> ExtractedMetrics {
207        self.metrics
208    }
209
210    /// Extends the contained metrics with [`ExtractedMetrics`].
211    pub fn extend(
212        &mut self,
213        extracted: ExtractedMetrics,
214        sampling_decision: Option<SamplingDecision>,
215    ) {
216        self.extend_project_metrics(extracted.0, sampling_decision);
217    }
218
219    /// Extends the contained project metrics.
220    pub fn extend_project_metrics<I>(
221        &mut self,
222        buckets: I,
223        sampling_decision: Option<SamplingDecision>,
224    ) where
225        I: IntoIterator<Item = Bucket>,
226    {
227        self.metrics.0.extend(buckets.into_iter().map(|mut bucket| {
228            bucket.metadata.extracted_from_indexed =
229                sampling_decision == Some(SamplingDecision::Keep);
230            bucket
231        }));
232    }
233}
234
235fn send_metrics(metrics: ExtractedMetrics, project_key: ProjectKey, aggregator: &Addr<Aggregator>) {
236    let ExtractedMetrics(project_metrics) = metrics;
237
238    if !project_metrics.is_empty() {
239        aggregator.send(MergeBuckets {
240            project_key,
241            buckets: project_metrics,
242        });
243    }
244}
245
246/// Applies processing to all contents of the given envelope.
247///
248/// Depending on the contents of the envelope and Relay's mode, this includes:
249///
250///  - Basic normalization and validation for all item types.
251///  - Clock drift correction if the required `sent_at` header is present.
252///  - Expansion of certain item types (e.g. unreal).
253///  - Store normalization for event payloads in processing mode.
254///  - Rate limiters and inbound filters on events in processing mode.
255#[derive(Debug)]
256pub struct ProcessEnvelope {
257    /// Envelope to process.
258    pub envelope: ManagedEnvelope,
259    /// The project info.
260    pub project_info: Arc<ProjectInfo>,
261    /// Currently active cached rate limits for this project.
262    pub rate_limits: Arc<RateLimits>,
263    /// Root sampling project info.
264    pub sampling_project_info: Option<Arc<ProjectInfo>>,
265}
266
267/// Parses a list of metrics or metric buckets and pushes them to the project's aggregator.
268///
269/// This parses and validates the metrics:
270///  - For [`Metrics`](ItemType::Statsd), each metric is parsed separately, and invalid metrics are
271///    ignored independently.
272///  - For [`MetricBuckets`](ItemType::MetricBuckets), the entire list of buckets is parsed and
273///    dropped together on parsing failure.
274///  - Other envelope items will be ignored with an error message.
275///
276/// Additionally, processing applies clock drift correction using the system clock of this Relay, if
277/// the Envelope specifies the [`sent_at`](Envelope::sent_at) header.
278#[derive(Debug)]
279pub struct ProcessMetrics {
280    /// A list of metric items.
281    pub data: MetricData,
282    /// The target project.
283    pub project_key: ProjectKey,
284    /// Whether to keep or reset the metric metadata.
285    pub source: BucketSource,
286    /// The wall clock time at which the request was received.
287    pub received_at: DateTime<Utc>,
288    /// The value of the Envelope's [`sent_at`](Envelope::sent_at) header for clock drift
289    /// correction.
290    pub sent_at: Option<DateTime<Utc>>,
291}
292
293/// Raw unparsed metric data.
294#[derive(Debug)]
295pub enum MetricData {
296    /// Raw data, unparsed envelope items.
297    Raw(Vec<Item>),
298    /// Already parsed buckets but unprocessed.
299    Parsed(Vec<Bucket>),
300}
301
302impl MetricData {
303    /// Consumes the metric data and parses the contained buckets.
304    ///
305    /// If the contained data is already parsed the buckets are returned unchanged.
306    /// Raw buckets are parsed and created with the passed `timestamp`.
307    fn into_buckets(self, timestamp: UnixTimestamp) -> Vec<Bucket> {
308        let items = match self {
309            Self::Parsed(buckets) => return buckets,
310            Self::Raw(items) => items,
311        };
312
313        let mut buckets = Vec::new();
314        for item in items {
315            let payload = item.payload();
316            if item.ty() == &ItemType::Statsd {
317                for bucket_result in Bucket::parse_all(&payload, timestamp) {
318                    match bucket_result {
319                        Ok(bucket) => buckets.push(bucket),
320                        Err(error) => relay_log::debug!(
321                            error = &error as &dyn Error,
322                            "failed to parse metric bucket from statsd format",
323                        ),
324                    }
325                }
326            } else if item.ty() == &ItemType::MetricBuckets {
327                match serde_json::from_slice::<Vec<Bucket>>(&payload) {
328                    Ok(parsed_buckets) => {
329                        // Re-use the allocation of `b` if possible.
330                        if buckets.is_empty() {
331                            buckets = parsed_buckets;
332                        } else {
333                            buckets.extend(parsed_buckets);
334                        }
335                    }
336                    Err(error) => {
337                        relay_log::debug!(
338                            error = &error as &dyn Error,
339                            "failed to parse metric bucket",
340                        );
341                        metric!(counter(RelayCounters::MetricBucketsParsingFailed) += 1);
342                    }
343                }
344            } else {
345                relay_log::error!(
346                    "invalid item of type {} passed to ProcessMetrics",
347                    item.ty()
348                );
349            }
350        }
351        buckets
352    }
353}
354
355#[derive(Debug)]
356pub struct ProcessBatchedMetrics {
357    /// Metrics payload in JSON format.
358    pub payload: Bytes,
359    /// Whether to keep or reset the metric metadata.
360    pub source: BucketSource,
361    /// The wall clock time at which the request was received.
362    pub received_at: DateTime<Utc>,
363    /// The wall clock time at which the request was received.
364    pub sent_at: Option<DateTime<Utc>>,
365}
366
367/// Source information where a metric bucket originates from.
368#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
369pub enum BucketSource {
370    /// The metric bucket originated from an internal Relay use case.
371    ///
372    /// The metric bucket originates either from within the same Relay
373    /// or was accepted coming from another Relay which is registered as
374    /// an internal Relay via Relay's configuration.
375    Internal,
376    /// The bucket source originated from an untrusted source.
377    ///
378    /// Managed Relays sending extracted metrics are considered external,
379    /// it's a project use case but it comes from an untrusted source.
380    External,
381}
382
383impl BucketSource {
384    /// Infers the bucket source from [`RequestMeta::request_trust`].
385    pub fn from_meta(meta: &RequestMeta) -> Self {
386        match meta.request_trust() {
387            RequestTrust::Trusted => Self::Internal,
388            RequestTrust::Untrusted => Self::External,
389        }
390    }
391}
392
393/// Sends a client report to the upstream.
394#[derive(Debug)]
395pub struct SubmitClientReports {
396    /// The client report to be sent.
397    pub client_reports: Vec<ClientReport>,
398    /// Scoping information for the client report.
399    pub scoping: Scoping,
400}
401
402/// CPU-intensive processing tasks for envelopes.
403#[derive(Debug)]
404pub enum EnvelopeProcessor {
405    ProcessEnvelope(Box<ProcessEnvelope>),
406    ProcessProjectMetrics(Box<ProcessMetrics>),
407    ProcessBatchedMetrics(Box<ProcessBatchedMetrics>),
408    FlushBuckets(Box<FlushBuckets>),
409    SubmitClientReports(Box<SubmitClientReports>),
410}
411
412impl EnvelopeProcessor {
413    /// Returns the name of the message variant.
414    pub fn variant(&self) -> &'static str {
415        match self {
416            EnvelopeProcessor::ProcessEnvelope(_) => "ProcessEnvelope",
417            EnvelopeProcessor::ProcessProjectMetrics(_) => "ProcessProjectMetrics",
418            EnvelopeProcessor::ProcessBatchedMetrics(_) => "ProcessBatchedMetrics",
419            EnvelopeProcessor::FlushBuckets(_) => "FlushBuckets",
420            EnvelopeProcessor::SubmitClientReports(_) => "SubmitClientReports",
421        }
422    }
423}
424
425impl relay_system::Interface for EnvelopeProcessor {}
426
427impl FromMessage<ProcessEnvelope> for EnvelopeProcessor {
428    type Response = relay_system::NoResponse;
429
430    fn from_message(message: ProcessEnvelope, _sender: ()) -> Self {
431        Self::ProcessEnvelope(Box::new(message))
432    }
433}
434
435impl FromMessage<ProcessMetrics> for EnvelopeProcessor {
436    type Response = NoResponse;
437
438    fn from_message(message: ProcessMetrics, _: ()) -> Self {
439        Self::ProcessProjectMetrics(Box::new(message))
440    }
441}
442
443impl FromMessage<ProcessBatchedMetrics> for EnvelopeProcessor {
444    type Response = NoResponse;
445
446    fn from_message(message: ProcessBatchedMetrics, _: ()) -> Self {
447        Self::ProcessBatchedMetrics(Box::new(message))
448    }
449}
450
451impl FromMessage<FlushBuckets> for EnvelopeProcessor {
452    type Response = NoResponse;
453
454    fn from_message(message: FlushBuckets, _: ()) -> Self {
455        Self::FlushBuckets(Box::new(message))
456    }
457}
458
459impl FromMessage<SubmitClientReports> for EnvelopeProcessor {
460    type Response = NoResponse;
461
462    fn from_message(message: SubmitClientReports, _: ()) -> Self {
463        Self::SubmitClientReports(Box::new(message))
464    }
465}
466
467/// The asynchronous thread pool used for scheduling processing tasks in the processor.
468pub type EnvelopeProcessorServicePool = AsyncPool<BoxFuture<'static, ()>>;
469
470/// Service implementing the [`EnvelopeProcessor`] interface.
471///
472/// This service handles messages in a worker pool with configurable concurrency.
473#[derive(Clone)]
474pub struct EnvelopeProcessorService {
475    inner: Arc<InnerProcessor>,
476}
477
478/// Contains the addresses of services that the processor publishes to.
479pub struct Addrs {
480    pub outcome_aggregator: Addr<TrackOutcome>,
481    pub upstream_relay: Addr<UpstreamRelay>,
482    #[cfg(feature = "processing")]
483    pub objectstore: Option<Addr<Objectstore>>,
484    #[cfg(feature = "processing")]
485    pub store_forwarder: Option<Addr<Store>>,
486    pub aggregator: Addr<Aggregator>,
487}
488
489impl Default for Addrs {
490    fn default() -> Self {
491        Addrs {
492            outcome_aggregator: Addr::dummy(),
493            upstream_relay: Addr::dummy(),
494            #[cfg(feature = "processing")]
495            objectstore: None,
496            #[cfg(feature = "processing")]
497            store_forwarder: None,
498            aggregator: Addr::dummy(),
499        }
500    }
501}
502
503struct InnerProcessor {
504    pool: EnvelopeProcessorServicePool,
505    config: Arc<Config>,
506    global_config: GlobalConfigHandle,
507    project_cache: ProjectCacheHandle,
508    cogs: Cogs,
509    addrs: Addrs,
510    #[cfg(feature = "processing")]
511    rate_limiter: Option<Arc<RedisRateLimiter>>,
512    metric_outcomes: MetricOutcomes,
513    processor: RelayProcessor,
514}
515
516impl EnvelopeProcessorService {
517    /// Creates a multi-threaded envelope processor.
518    #[cfg_attr(feature = "processing", expect(clippy::too_many_arguments))]
519    pub fn new(
520        pool: EnvelopeProcessorServicePool,
521        config: Arc<Config>,
522        global_config: GlobalConfigHandle,
523        project_cache: ProjectCacheHandle,
524        cogs: Cogs,
525        #[cfg(feature = "processing")] redis: Option<RedisClients>,
526        addrs: Addrs,
527        metric_outcomes: MetricOutcomes,
528    ) -> Self {
529        let c = config.current();
530
531        let geoip_lookup = c
532            .geoip_path()
533            .and_then(
534                |p| match GeoIpLookup::open(p).context(ServiceError::GeoIp) {
535                    Ok(geoip) => Some(geoip),
536                    Err(err) => {
537                        relay_log::error!("failed to open GeoIP db {p:?}: {err:?}");
538                        None
539                    }
540                },
541            )
542            .unwrap_or_else(GeoIpLookup::empty);
543
544        if let Some(build_epoch) = geoip_lookup.build_epoch() {
545            relay_log::info!("Loaded GeoIP database (build: {build_epoch})");
546        }
547
548        #[cfg(feature = "processing")]
549        let rate_limiter = redis.map(|redis| {
550            RedisRateLimiter::new(redis.quotas)
551                .max_limit(c.max_rate_limit())
552                .cache(c.quota_cache_ratio(), c.quota_cache_max())
553        });
554
555        let quota_limiter = Arc::new(QuotaRateLimiter::new(
556            #[cfg(feature = "processing")]
557            project_cache.clone(),
558            #[cfg(feature = "processing")]
559            rate_limiter.clone(),
560        ));
561        #[cfg(feature = "processing")]
562        let rate_limiter = rate_limiter.map(Arc::new);
563        let inner = InnerProcessor {
564            pool,
565            global_config,
566            project_cache,
567            #[cfg(feature = "processing")]
568            rate_limiter,
569            processor: RelayProcessor::new(
570                cogs.clone(),
571                &quota_limiter,
572                &geoip_lookup,
573                addrs.outcome_aggregator.clone(),
574            ),
575            cogs,
576            addrs,
577            metric_outcomes,
578            config,
579        };
580
581        Self {
582            inner: Arc::new(inner),
583        }
584    }
585
586    async fn process_envelope(
587        &self,
588        project_id: ProjectId,
589        mut envelope: ManagedEnvelope,
590        ctx: processing::Context<'_>,
591    ) -> Vec<Output<Outputs>> {
592        // Pre-process the envelope headers.
593        if let Some(sampling_state) = ctx.sampling_project_info {
594            // Both transactions and standalone span envelopes need a normalized DSC header
595            // to make sampling rules based on the segment/transaction name work correctly.
596            envelope
597                .envelope_mut()
598                .parametrize_dsc_transaction(&sampling_state.config.tx_name_rules);
599        }
600
601        // Ensure the project ID is updated to the stored instance for this project cache. This can
602        // differ in two cases:
603        //  1. The envelope was sent to the legacy `/store/` endpoint without a project ID.
604        //  2. The DSN was moved and the envelope sent to the old project ID.
605        envelope
606            .envelope_mut()
607            .meta_mut()
608            .set_project_id(project_id);
609
610        self.inner.processor.run(envelope, ctx).await
611    }
612
613    /// Processes the envelope and returns the processed envelope back.
614    ///
615    /// Returns `Some` if the envelope passed inbound filtering and rate limiting. Invalid items are
616    /// removed from the envelope. Otherwise, if the envelope is empty or the entire envelope needs
617    /// to be dropped, this is `None`.
618    async fn process<'a>(
619        &self,
620        mut envelope: ManagedEnvelope,
621        ctx: processing::Context<'a>,
622    ) -> Vec<Output<Outputs>> {
623        // Prefer the project's project ID, and fall back to the stated project id from the
624        // envelope. The project ID is available in all modes, other than in proxy mode, where
625        // envelopes for unknown projects are forwarded blindly.
626        //
627        // Neither ID can be available in proxy mode on the /store/ endpoint. This is not supported,
628        // since we cannot process an envelope without project ID, so drop it.
629        let Some(project_id) = ctx
630            .project_info
631            .project_id
632            .or_else(|| envelope.envelope().meta().project_id())
633        else {
634            relay_log::error!(
635                tags.project_key = %envelope.envelope().meta().public_key(),
636                "project info does not contain project id"
637            );
638            envelope.reject(Outcome::Invalid(DiscardReason::Internal));
639            return Vec::new();
640        };
641
642        relay_log::configure_scope(|scope| {
643            scope.set_tag("project_id", project_id);
644        });
645
646        self.process_envelope(project_id, envelope, ctx).await
647    }
648
649    async fn handle_process_envelope(&self, cogs: &mut Token, message: ProcessEnvelope) {
650        let wait_time = message.envelope.age();
651        metric!(timer(RelayTimers::EnvelopeWaitTime) = wait_time);
652
653        // This COGS handling may need an overhaul in the future:
654        // Cancel the passed in token, to start individual measurements per processor instead.
655        cogs.cancel();
656
657        let global_config = self.inner.global_config.current().unwrap_or_default();
658        let config = self.inner.config.current();
659
660        let ctx = processing::Context {
661            config: &config,
662            global_config: &global_config,
663            project_info: &message.project_info,
664            sampling_project_info: message.sampling_project_info.as_deref(),
665            rate_limits: &message.rate_limits,
666        };
667
668        let project_key = message.envelope.meta().public_key();
669        // Only allow sending to the sampling key, if we successfully loaded a sampling project
670        // info relating to it. This filters out unknown/invalid project keys as well as project
671        // keys from different organizations.
672        let sampling_key = ctx
673            .sampling_project_info
674            .and_then(|p| p.get_public_key_config())
675            .map(|pkc| pkc.public_key);
676
677        relay_log::configure_scope(|scope| {
678            scope.set_tag("project_key", project_key);
679            if let Some(sampling_key) = sampling_key {
680                scope.set_tag("sampling_key", sampling_key);
681            }
682            let meta = message.envelope.envelope().meta();
683            scope.set_tag("sdk_name", meta.client_name());
684            if let Some(client) = meta.client() {
685                scope.set_tag("sdk", client);
686            }
687            if let Some(user_agent) = meta.user_agent() {
688                scope.set_extra("user_agent", user_agent.into());
689            }
690        });
691
692        let mut envelopes: smallvec::SmallVec<[ManagedEnvelope; 1]> =
693            smallvec::smallvec![message.envelope];
694
695        // The first envelope we process is not an intermediate.
696        let mut is_intermediate = false;
697
698        while let Some(envelope) = envelopes.pop() {
699            let outputs = metric!(
700                timer(RelayTimers::EnvelopeProcessingTime),
701                is_intermediate = if is_intermediate { "true" } else { "false" },
702                { self.process(envelope, ctx).await }
703            );
704
705            let ctx = ctx.to_forward();
706            for Output {
707                main,
708                metrics,
709                intermediates,
710            } in outputs
711            {
712                if let Some(metrics) = metrics {
713                    let agg = &self.inner.addrs.aggregator;
714                    metrics.accept(|metrics| {
715                        send_metrics(metrics, project_key, agg);
716                    });
717                }
718
719                if let Some(output) = main {
720                    // Only counting processing time for COGS at the moment.
721                    self.submit_upstream(&mut Token::noop(), output, ctx);
722                }
723
724                if let Some(intermediates) = intermediates {
725                    envelopes.push(intermediates)
726                }
727
728                // Every envelope past the first is an intermediate.
729                is_intermediate = true;
730            }
731        }
732    }
733
734    fn handle_process_metrics(&self, cogs: &mut Token, message: ProcessMetrics) {
735        let ProcessMetrics {
736            data,
737            project_key,
738            received_at,
739            sent_at,
740            source,
741        } = message;
742
743        let received_timestamp =
744            UnixTimestamp::from_datetime(received_at).unwrap_or(UnixTimestamp::now());
745
746        let mut buckets = data.into_buckets(received_timestamp);
747        if buckets.is_empty() {
748            return;
749        };
750        cogs.update(relay_metrics::cogs::BySize(&buckets));
751
752        let clock_drift_processor =
753            ClockDriftProcessor::new(sent_at, received_at).at_least(MINIMUM_CLOCK_DRIFT);
754
755        buckets.retain_mut(|bucket| {
756            if let Err(error) = relay_metrics::normalize_bucket(bucket) {
757                relay_log::debug!(error = &error as &dyn Error, "dropping bucket {bucket:?}");
758                return false;
759            }
760
761            if !self::metrics::is_valid_namespace(bucket, source) {
762                relay_log::debug!("dropping bucket in invalid namespace {bucket:?}");
763                return false;
764            }
765
766            clock_drift_processor.process_timestamp(&mut bucket.timestamp);
767
768            if !matches!(source, BucketSource::Internal) {
769                bucket.metadata = BucketMetadata::new(received_timestamp);
770            }
771
772            true
773        });
774
775        let project = self.inner.project_cache.get(project_key);
776
777        // Best effort check to filter and rate limit buckets, if there is no project state
778        // available at the current time, we will check again after flushing.
779        let buckets = match project.state() {
780            ProjectState::Enabled(project_info) => {
781                let rate_limits = project.rate_limits().current_limits();
782                self.check_buckets(project_key, project_info, &rate_limits, buckets)
783            }
784            _ => buckets,
785        };
786
787        relay_log::trace!("merging metric buckets into the aggregator");
788        self.inner
789            .addrs
790            .aggregator
791            .send(MergeBuckets::new(project_key, buckets));
792    }
793
794    fn handle_process_batched_metrics(&self, cogs: &mut Token, message: ProcessBatchedMetrics) {
795        let ProcessBatchedMetrics {
796            payload,
797            source,
798            received_at,
799            sent_at,
800        } = message;
801
802        #[derive(serde::Deserialize)]
803        struct Wrapper {
804            buckets: HashMap<ProjectKey, Vec<Bucket>>,
805        }
806
807        let buckets = match serde_json::from_slice(&payload) {
808            Ok(Wrapper { buckets }) => buckets,
809            Err(error) => {
810                relay_log::debug!(
811                    error = &error as &dyn Error,
812                    "failed to parse batched metrics",
813                );
814                metric!(counter(RelayCounters::MetricBucketsParsingFailed) += 1);
815                return;
816            }
817        };
818
819        for (project_key, buckets) in buckets {
820            self.handle_process_metrics(
821                cogs,
822                ProcessMetrics {
823                    data: MetricData::Parsed(buckets),
824                    project_key,
825                    source,
826                    received_at,
827                    sent_at,
828                },
829            )
830        }
831    }
832
833    /// Submits a processor [`Output`] to the appropriate upstream.
834    ///
835    /// If processing is enabled, the upstream is Kafka.
836    fn submit_upstream(
837        &self,
838        cogs: &mut Token,
839        output: Outputs,
840        ctx: processing::ForwardContext<'_>,
841    ) {
842        let _submit = cogs.start_category("submit");
843
844        #[cfg(feature = "processing")]
845        if ctx.config.processing_enabled()
846            && let Some(store_forwarder) = &self.inner.addrs.store_forwarder
847        {
848            use crate::processing::StoreHandle;
849
850            let objectstore = self.inner.addrs.objectstore.as_ref();
851            let handle = StoreHandle::new(store_forwarder, objectstore, ctx.global_config);
852
853            output
854                .forward_store(handle, ctx)
855                .unwrap_or_else(|err| err.into_inner());
856
857            return;
858        }
859
860        match output.serialize_envelope(ctx) {
861            Ok(envelope) => {
862                let envelope = ManagedEnvelope::from(envelope);
863                self.submit_envelope_upstream(
864                    envelope,
865                    ctx.config,
866                    ctx.project_info.upstream.clone(),
867                );
868            }
869            Err(_) => relay_log::error!("failed to serialize output to an envelope"),
870        };
871    }
872
873    fn submit_envelope_upstream(
874        &self,
875        mut envelope: ManagedEnvelope,
876        config: &ConfigSnapshot,
877        // Currently allowed to be optional as code is migrated to respect the upstream override
878        // provided from the project config. Eventually must be available and is required.
879        upstream: Option<UpstreamDescriptor>,
880    ) {
881        if envelope.envelope_mut().is_empty() {
882            envelope.accept();
883            return;
884        }
885
886        // No code path should hit this.
887        //
888        // Any item which is produced by processing is handled in `submit_upstream`,
889        // metrics are sent to the store directly and outcomes must be produced to Kafka
890        // instead of being sent onward as client report.
891        if config.processing_enabled() {
892            relay_log::error!(
893                "attempt to forward envelope to http upstream when processing is enabled"
894            );
895            return;
896        }
897
898        // Override the `sent_at` timestamp. Since the envelope went through basic
899        // normalization, all timestamps have been corrected. We propagate the new
900        // `sent_at` to allow the next Relay to double-check this timestamp and
901        // potentially apply correction again. This is done as close to sending as
902        // possible so that we avoid internal delays.
903        envelope.envelope_mut().set_sent_at(Utc::now());
904
905        relay_log::trace!("sending envelope to sentry endpoint");
906        let http_encoding = config.http_encoding();
907        let result = envelope.envelope().to_vec().and_then(|v| {
908            encode_payload(&v.into(), http_encoding).map_err(EnvelopeError::PayloadIoFailed)
909        });
910
911        match result {
912            Ok(body) => {
913                self.inner
914                    .addrs
915                    .upstream_relay
916                    .send(SendRequest(SendEnvelope {
917                        upstream,
918                        envelope,
919                        body,
920                        http_encoding,
921                        project_cache: self.inner.project_cache.clone(),
922                    }));
923            }
924            Err(error) => {
925                // Errors are only logged for what we consider an internal discard reason. These
926                // indicate errors in the infrastructure or implementation bugs.
927                relay_log::error!(
928                    error = &error as &dyn Error,
929                    tags.project_key = %envelope.scoping().project_key,
930                    "failed to serialize envelope payload"
931                );
932
933                envelope.reject(Outcome::Invalid(DiscardReason::Internal));
934            }
935        }
936    }
937
938    fn handle_submit_client_reports(&self, message: SubmitClientReports) {
939        let SubmitClientReports {
940            client_reports,
941            scoping,
942        } = message;
943
944        relay_log::trace!(
945            "sending {} client report(s) to project id {}",
946            client_reports.len(),
947            scoping.project_id
948        );
949
950        if client_reports.is_empty() {
951            return;
952        }
953
954        let config = self.inner.config.current();
955        let upstream = config.upstream();
956        let dsn = PartialDsn::outbound(&scoping, upstream);
957
958        let mut envelope = Envelope::from_request(None, RequestMeta::outbound(dsn));
959        for client_report in client_reports {
960            match client_report.serialize() {
961                Ok(payload) => {
962                    let mut item = Item::new(ItemType::ClientReport);
963                    item.set_payload(ContentType::Json, payload);
964                    envelope.add_item(item);
965                }
966                Err(error) => {
967                    relay_log::error!(
968                        error = &error as &dyn std::error::Error,
969                        "failed to serialize client report"
970                    );
971                }
972            }
973        }
974
975        let envelope = ManagedEnvelope::new(envelope, self.inner.addrs.outcome_aggregator.clone());
976        self.submit_envelope_upstream(envelope, &self.inner.config.current(), None);
977    }
978
979    fn check_buckets(
980        &self,
981        project_key: ProjectKey,
982        project_info: &ProjectInfo,
983        rate_limits: &RateLimits,
984        buckets: Vec<Bucket>,
985    ) -> Vec<Bucket> {
986        let Some(scoping) = project_info.scoping(project_key) else {
987            relay_log::error!(
988                tags.project_key = project_key.as_str(),
989                "there is no scoping: dropping {} buckets",
990                buckets.len(),
991            );
992            return Vec::new();
993        };
994
995        let mut buckets =
996            self::metrics::remove_invalid_namespaces(buckets, &self.inner.metric_outcomes, scoping);
997
998        let mut namespaces: BTreeSet<MetricNamespace> = buckets
999            .iter()
1000            .filter_map(|bucket| bucket.name.try_namespace())
1001            .collect();
1002
1003        // Never rate limit outcomes.
1004        namespaces.remove(&MetricNamespace::Outcomes);
1005
1006        for namespace in namespaces {
1007            let limits = rate_limits
1008                .check_with_quotas(project_info.get_quotas(), &scoping.metric_bucket(namespace));
1009
1010            if limits.is_limited() {
1011                let rejected;
1012                (buckets, rejected) = utils::split_off(buckets, |bucket| {
1013                    bucket.name.try_namespace() == Some(namespace)
1014                });
1015
1016                let reason_code = limits.longest().and_then(|limit| limit.reason_code.clone());
1017                self.inner.metric_outcomes.track(
1018                    scoping,
1019                    &rejected,
1020                    Outcome::RateLimited(reason_code),
1021                );
1022            }
1023        }
1024
1025        let quotas = project_info.config.quotas.clone();
1026        match MetricsLimiter::create(buckets, quotas, scoping) {
1027            Ok(mut bucket_limiter) => {
1028                bucket_limiter.enforce_limits(rate_limits, &self.inner.metric_outcomes);
1029                bucket_limiter.into_buckets()
1030            }
1031            Err(buckets) => buckets,
1032        }
1033    }
1034
1035    #[cfg(feature = "processing")]
1036    async fn rate_limit_buckets(
1037        &self,
1038        scoping: Scoping,
1039        project_info: &ProjectInfo,
1040        mut buckets: Vec<Bucket>,
1041    ) -> Vec<Bucket> {
1042        let Some(rate_limiter) = &self.inner.rate_limiter else {
1043            return buckets;
1044        };
1045
1046        let global_config = self.inner.global_config.current().unwrap_or_default();
1047        let mut namespaces = buckets
1048            .iter()
1049            .filter_map(|bucket| bucket.name.try_namespace())
1050            .counts();
1051
1052        // Never rate limit outcomes.
1053        namespaces.remove(&MetricNamespace::Outcomes);
1054
1055        let quotas = CombinedQuotas::new(&global_config, project_info.get_quotas());
1056
1057        for (namespace, quantity) in namespaces {
1058            let item_scoping = scoping.metric_bucket(namespace);
1059
1060            let limits = match rate_limiter
1061                .is_rate_limited(quotas, &item_scoping, quantity, false)
1062                .await
1063            {
1064                Ok(limits) => limits,
1065                Err(err) => {
1066                    relay_log::error!(
1067                        error = &err as &dyn std::error::Error,
1068                        "failed to check redis rate limits"
1069                    );
1070                    break;
1071                }
1072            };
1073
1074            if limits.is_limited() {
1075                let rejected;
1076                (buckets, rejected) = utils::split_off(buckets, |bucket| {
1077                    bucket.name.try_namespace() == Some(namespace)
1078                });
1079
1080                let reason_code = limits.longest().and_then(|limit| limit.reason_code.clone());
1081                self.inner.metric_outcomes.track(
1082                    scoping,
1083                    &rejected,
1084                    Outcome::RateLimited(reason_code),
1085                );
1086
1087                self.inner
1088                    .project_cache
1089                    .get(item_scoping.scoping.project_key)
1090                    .rate_limits()
1091                    .merge(limits);
1092            }
1093        }
1094
1095        match MetricsLimiter::create(buckets, project_info.config.quotas.clone(), scoping) {
1096            Err(buckets) => buckets,
1097            Ok(bucket_limiter) => self.apply_other_rate_limits(bucket_limiter).await,
1098        }
1099    }
1100
1101    /// Check and apply rate limits to metrics buckets for transactions and spans.
1102    #[cfg(feature = "processing")]
1103    async fn apply_other_rate_limits(&self, mut bucket_limiter: MetricsLimiter) -> Vec<Bucket> {
1104        relay_log::trace!("handle_rate_limit_buckets");
1105
1106        let scoping = *bucket_limiter.scoping();
1107
1108        if let Some(rate_limiter) = self.inner.rate_limiter.as_ref() {
1109            let global_config = self.inner.global_config.current().unwrap_or_default();
1110            let quotas = CombinedQuotas::new(&global_config, bucket_limiter.quotas());
1111
1112            // We set over_accept_once such that the limit is actually reached, which allows subsequent
1113            // calls with quantity=0 to be rate limited.
1114            let over_accept_once = true;
1115            let mut rate_limits = RateLimits::new();
1116
1117            let (category, count) = bucket_limiter.count();
1118
1119            let timer = Instant::now();
1120            let mut is_limited = false;
1121
1122            if let Some(count) = count {
1123                match rate_limiter
1124                    .is_rate_limited(quotas, &scoping.item(category), count, over_accept_once)
1125                    .await
1126                {
1127                    Ok(limits) => {
1128                        is_limited = limits.is_limited();
1129                        rate_limits.merge(limits)
1130                    }
1131                    Err(e) => {
1132                        relay_log::error!(error = &e as &dyn Error, "rate limiting error")
1133                    }
1134                }
1135            }
1136
1137            relay_statsd::metric!(
1138                timer(RelayTimers::RateLimitBucketsDuration) = timer.elapsed(),
1139                category = category.name(),
1140                limited = if is_limited { "true" } else { "false" },
1141                count = match count {
1142                    None => "none",
1143                    Some(0) => "0",
1144                    Some(1) => "1",
1145                    Some(1..=10) => "10",
1146                    Some(1..=25) => "25",
1147                    Some(1..=50) => "50",
1148                    Some(51..=100) => "100",
1149                    Some(101..=500) => "500",
1150                    _ => "> 500",
1151                },
1152            );
1153
1154            if rate_limits.is_limited() {
1155                let was_enforced =
1156                    bucket_limiter.enforce_limits(&rate_limits, &self.inner.metric_outcomes);
1157
1158                if was_enforced {
1159                    // Update the rate limits in the project cache.
1160                    self.inner
1161                        .project_cache
1162                        .get(scoping.project_key)
1163                        .rate_limits()
1164                        .merge(rate_limits);
1165                }
1166            }
1167        }
1168
1169        bucket_limiter.into_buckets()
1170    }
1171
1172    /// Processes metric buckets and sends them to Kafka.
1173    ///
1174    /// This function runs the following steps:
1175    ///  - rate limiting
1176    ///  - emit billing outcomes
1177    ///  - submit to `StoreForwarder`
1178    #[cfg(feature = "processing")]
1179    async fn encode_metrics_processing(
1180        &self,
1181        message: FlushBuckets,
1182        store_forwarder: &Addr<Store>,
1183    ) {
1184        use crate::constants::DEFAULT_EVENT_RETENTION;
1185        use crate::services::store::StoreMetrics;
1186
1187        for ProjectBuckets {
1188            buckets,
1189            scoping,
1190            project_info,
1191            ..
1192        } in message.buckets.into_values()
1193        {
1194            let mut buckets = self
1195                .rate_limit_buckets(scoping, &project_info, buckets)
1196                .await;
1197
1198            if buckets.is_empty() {
1199                continue;
1200            }
1201
1202            // Emit metric billing outcomes.
1203            self.inner
1204                .metric_outcomes
1205                .track_accepted_outcome(scoping, &mut buckets);
1206
1207            let retention = project_info
1208                .config
1209                .event_retention
1210                .unwrap_or(DEFAULT_EVENT_RETENTION);
1211
1212            // The store forwarder takes care of bucket splitting internally, so we can submit the
1213            // entire list of buckets. There is no batching needed here.
1214            store_forwarder.send(StoreMetrics {
1215                buckets,
1216                scoping,
1217                retention,
1218            });
1219        }
1220    }
1221
1222    /// Serializes metric buckets to JSON and sends them to the upstream.
1223    ///
1224    /// This function runs the following steps:
1225    ///  - partitioning
1226    ///  - batching by configured size limit
1227    ///  - serialize to JSON and pack in an envelope
1228    ///
1229    /// Rate limiting runs only in processing Relays as it requires access to the central Redis instance.
1230    /// Cached rate limits are applied in the project cache already.
1231    fn encode_metrics_envelope(&self, message: FlushBuckets) {
1232        let FlushBuckets {
1233            partition_key,
1234            buckets,
1235        } = message;
1236
1237        let config = self.inner.config.current();
1238        let batch_size = config.metrics_max_batch_size_bytes();
1239        let upstream = config.upstream();
1240
1241        for ProjectBuckets {
1242            buckets,
1243            scoping,
1244            project_info,
1245            ..
1246        } in buckets.values()
1247        {
1248            let dsn = PartialDsn::outbound(scoping, upstream);
1249
1250            relay_statsd::metric!(
1251                distribution(RelayDistributions::PartitionKeys) = u64::from(partition_key)
1252            );
1253
1254            let mut num_batches = 0;
1255            for batch in BucketsView::from(buckets).by_size(batch_size) {
1256                let mut envelope = Envelope::from_request(None, RequestMeta::outbound(dsn.clone()));
1257
1258                let mut item = Item::new(ItemType::MetricBuckets);
1259                item.set_source_quantities(crate::metrics::extract_quantities(batch));
1260                item.set_payload(ContentType::Json, serde_json::to_vec(&buckets).unwrap());
1261                envelope.add_item(item);
1262
1263                let mut envelope =
1264                    ManagedEnvelope::new(envelope, self.inner.addrs.outcome_aggregator.clone());
1265                envelope
1266                    .set_partition_key(Some(partition_key))
1267                    .scope(*scoping);
1268
1269                relay_statsd::metric!(
1270                    distribution(RelayDistributions::BucketsPerBatch) = batch.len() as u64
1271                );
1272
1273                self.submit_envelope_upstream(envelope, &config, project_info.upstream.clone());
1274                num_batches += 1;
1275            }
1276
1277            relay_statsd::metric!(
1278                distribution(RelayDistributions::BatchesPerPartition) = num_batches
1279            );
1280        }
1281    }
1282
1283    /// Creates a [`SendMetricsRequest`] and sends it to the upstream relay.
1284    fn send_global_partition(
1285        &self,
1286        upstream: Option<UpstreamDescriptor>,
1287        partition_key: u32,
1288        partition: &mut Partition<'_>,
1289    ) {
1290        if partition.is_empty() {
1291            return;
1292        }
1293
1294        let (unencoded, project_info) = partition.take();
1295        let http_encoding = self.inner.config.current().http_encoding();
1296        let encoded = match encode_payload(&unencoded, http_encoding) {
1297            Ok(payload) => payload,
1298            Err(error) => {
1299                let error = &error as &dyn std::error::Error;
1300                relay_log::error!(error, "failed to encode metrics payload");
1301                return;
1302            }
1303        };
1304
1305        let request = SendMetricsRequest {
1306            upstream,
1307            partition_key: partition_key.to_string(),
1308            unencoded,
1309            encoded,
1310            project_info,
1311            http_encoding,
1312            metric_outcomes: self.inner.metric_outcomes.clone(),
1313        };
1314
1315        self.inner.addrs.upstream_relay.send(SendRequest(request));
1316    }
1317
1318    /// Serializes metric buckets to JSON and sends them to the upstream via the global endpoint.
1319    ///
1320    /// This function is similar to [`Self::encode_metrics_envelope`], but sends a global batched
1321    /// payload directly instead of per-project Envelopes.
1322    ///
1323    /// This function runs the following steps:
1324    ///  - partitioning
1325    ///  - batching by configured size limit
1326    ///  - serialize to JSON
1327    ///  - submit directly to the upstream
1328    fn encode_metrics_global(&self, message: FlushBuckets) {
1329        let FlushBuckets {
1330            partition_key,
1331            buckets,
1332        } = message;
1333
1334        let batch_size = self.inner.config.current().metrics_max_batch_size_bytes();
1335        let mut partitions = BTreeMap::new();
1336        let mut partition_splits = 0;
1337
1338        for ProjectBuckets {
1339            buckets,
1340            scoping,
1341            project_info,
1342            ..
1343        } in buckets.values()
1344        {
1345            let partition = match partitions.get_mut(&project_info.upstream) {
1346                Some(partition) => partition,
1347                None => partitions
1348                    .entry(project_info.upstream.clone())
1349                    .or_insert_with(|| Partition::new(batch_size)),
1350            };
1351
1352            for bucket in buckets {
1353                let mut remaining = Some(BucketView::new(bucket));
1354
1355                while let Some(bucket) = remaining.take() {
1356                    if let Some(next) = partition.insert(bucket, *scoping) {
1357                        // A part of the bucket could not be inserted. Take the partition and submit
1358                        // it immediately. Repeat until the final part was inserted. This should
1359                        // always result in a request, otherwise we would enter an endless loop.
1360                        self.send_global_partition(
1361                            project_info.upstream.clone(),
1362                            partition_key,
1363                            partition,
1364                        );
1365                        remaining = Some(next);
1366                        partition_splits += 1;
1367                    }
1368                }
1369            }
1370        }
1371
1372        if partition_splits > 0 {
1373            metric!(distribution(RelayDistributions::PartitionSplits) = partition_splits);
1374        }
1375
1376        for (upstream, mut partition) in partitions {
1377            self.send_global_partition(upstream, partition_key, &mut partition);
1378        }
1379    }
1380
1381    /// Removes all outcome metrics from `message` and sends them as client reports.
1382    ///
1383    /// Returns a new [`FlushBuckets`] message, without any outcome metrics remaining.
1384    fn encode_metrics_client_reports(&self, mut message: FlushBuckets) -> FlushBuckets {
1385        for ProjectBuckets {
1386            buckets, scoping, ..
1387        } in message.buckets.values_mut()
1388        {
1389            let client_reports = outcome::metric::extract_client_reports(buckets).collect();
1390
1391            self.handle_submit_client_reports(SubmitClientReports {
1392                client_reports,
1393                scoping: *scoping,
1394            });
1395        }
1396
1397        message
1398    }
1399
1400    async fn handle_flush_buckets(&self, mut message: FlushBuckets) {
1401        for (project_key, pb) in message.buckets.iter_mut() {
1402            let buckets = std::mem::take(&mut pb.buckets);
1403            pb.buckets =
1404                self.check_buckets(*project_key, &pb.project_info, &pb.rate_limits, buckets);
1405        }
1406
1407        let config = self.inner.config.current();
1408
1409        #[cfg(feature = "processing")]
1410        if config.processing_enabled()
1411            && let Some(ref store_forwarder) = self.inner.addrs.store_forwarder
1412        {
1413            return self
1414                .encode_metrics_processing(message, store_forwarder)
1415                .await;
1416        }
1417
1418        // Processing Relays never send outcomes as client reports, which is why this check is after
1419        // the processing check.
1420        if config.emit_outcomes() == EmitOutcomes::AsClientReports {
1421            // Remove client reports from metrics to be sent, if configured as client reports
1422            // and send them separately.
1423            message = self.encode_metrics_client_reports(message);
1424        }
1425
1426        if config.http_global_metrics() {
1427            self.encode_metrics_global(message)
1428        } else {
1429            self.encode_metrics_envelope(message)
1430        }
1431    }
1432
1433    #[cfg(all(test, feature = "processing"))]
1434    fn redis_rate_limiter_enabled(&self) -> bool {
1435        self.inner.rate_limiter.is_some()
1436    }
1437
1438    async fn handle_message(self, message: EnvelopeProcessor) {
1439        let ty = message.variant();
1440        let feature_weights = self.feature_weights(&message);
1441
1442        metric!(timer(RelayTimers::ProcessMessageDuration), message = ty, {
1443            let mut cogs = self.inner.cogs.timed(ResourceId::Relay, feature_weights);
1444
1445            match message {
1446                EnvelopeProcessor::ProcessEnvelope(m) => {
1447                    self.handle_process_envelope(&mut cogs, *m).await
1448                }
1449                EnvelopeProcessor::ProcessProjectMetrics(m) => {
1450                    self.handle_process_metrics(&mut cogs, *m)
1451                }
1452                EnvelopeProcessor::ProcessBatchedMetrics(m) => {
1453                    self.handle_process_batched_metrics(&mut cogs, *m)
1454                }
1455                EnvelopeProcessor::FlushBuckets(m) => self.handle_flush_buckets(*m).await,
1456                EnvelopeProcessor::SubmitClientReports(m) => self.handle_submit_client_reports(*m),
1457            }
1458        });
1459    }
1460
1461    fn feature_weights(&self, message: &EnvelopeProcessor) -> FeatureWeights {
1462        match message {
1463            // Envelope is split later and tokens are attributed then.
1464            EnvelopeProcessor::ProcessEnvelope(_) => AppFeature::Unattributed.into(),
1465            EnvelopeProcessor::ProcessProjectMetrics(_) => AppFeature::Unattributed.into(),
1466            EnvelopeProcessor::ProcessBatchedMetrics(_) => AppFeature::Unattributed.into(),
1467            EnvelopeProcessor::FlushBuckets(v) => v
1468                .buckets
1469                .values()
1470                .map(|s| {
1471                    if self.inner.config.current().processing_enabled() {
1472                        // Processing does not encode the metrics but instead rate limit the metrics,
1473                        // which scales by count and not size.
1474                        relay_metrics::cogs::ByCount(&s.buckets).into()
1475                    } else {
1476                        relay_metrics::cogs::BySize(&s.buckets).into()
1477                    }
1478                })
1479                .fold(FeatureWeights::none(), FeatureWeights::merge),
1480            EnvelopeProcessor::SubmitClientReports(_) => AppFeature::ClientReports.into(),
1481        }
1482    }
1483}
1484
1485impl Service for EnvelopeProcessorService {
1486    type Interface = EnvelopeProcessor;
1487
1488    async fn run(self, mut rx: relay_system::Receiver<Self::Interface>) {
1489        while let Some(message) = rx.recv().await {
1490            let service = self.clone();
1491            // Create a new hub to prevent sentry scopes from bleeding to other tasks.
1492            let hub = relay_log::Hub::new_from_top(relay_log::Hub::current());
1493
1494            self.inner
1495                .pool
1496                .spawn_async(Box::pin(service.handle_message(message).bind_hub(hub)))
1497                .await;
1498        }
1499    }
1500}
1501
1502pub fn encode_payload(body: &Bytes, http_encoding: HttpEncoding) -> Result<Bytes, std::io::Error> {
1503    let envelope_body: Vec<u8> = match http_encoding {
1504        HttpEncoding::Identity => return Ok(body.clone()),
1505        HttpEncoding::Deflate => {
1506            let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
1507            encoder.write_all(body.as_ref())?;
1508            encoder.finish()?
1509        }
1510        HttpEncoding::Gzip => {
1511            let mut encoder = GzEncoder::new(Vec::new(), Compression::default());
1512            encoder.write_all(body.as_ref())?;
1513            encoder.finish()?
1514        }
1515        HttpEncoding::Br => {
1516            // Use default buffer size (via 0), medium quality (5), and the default lgwin (22).
1517            let mut encoder = BrotliEncoder::new(Vec::new(), 0, 5, 22);
1518            encoder.write_all(body.as_ref())?;
1519            encoder.into_inner()
1520        }
1521        HttpEncoding::Zstd => {
1522            // Use the fastest compression level, our main objective here is to get the best
1523            // compression ratio for least amount of time spent.
1524            let mut encoder = ZstdEncoder::new(Vec::new(), 1)?;
1525            encoder.write_all(body.as_ref())?;
1526            encoder.finish()?
1527        }
1528    };
1529
1530    Ok(envelope_body.into())
1531}
1532
1533/// An upstream request that submits an envelope via HTTP.
1534#[derive(Debug)]
1535pub struct SendEnvelope {
1536    pub upstream: Option<UpstreamDescriptor>,
1537    pub envelope: ManagedEnvelope,
1538    pub body: Bytes,
1539    pub http_encoding: HttpEncoding,
1540    pub project_cache: ProjectCacheHandle,
1541}
1542
1543impl UpstreamRequest for SendEnvelope {
1544    fn upstream(&self) -> Option<&UpstreamDescriptor> {
1545        self.upstream.as_ref()
1546    }
1547
1548    fn method(&self) -> reqwest::Method {
1549        reqwest::Method::POST
1550    }
1551
1552    fn path(&self) -> Cow<'_, str> {
1553        format!("/api/{}/envelope/", self.envelope.scoping().project_id).into()
1554    }
1555
1556    fn route(&self) -> &'static str {
1557        "envelope"
1558    }
1559
1560    fn build(&mut self, builder: &mut http::RequestBuilder) -> Result<(), http::HttpError> {
1561        let envelope_body = self.body.clone();
1562
1563        let meta = &self.envelope.meta();
1564        let shard = self.envelope.partition_key().map(|p| p.to_string());
1565        builder
1566            .content_encoding(self.http_encoding)
1567            .header_opt("Origin", meta.origin().map(|url| url.as_str()))
1568            .header_opt("User-Agent", meta.user_agent())
1569            .header("X-Sentry-Auth", meta.auth_header())
1570            .header("X-Forwarded-For", meta.forwarded_for())
1571            .header("Content-Type", envelope::CONTENT_TYPE)
1572            .header_opt("X-Sentry-Relay-Shard", shard)
1573            .body(envelope_body);
1574
1575        Ok(())
1576    }
1577
1578    fn sign(&mut self) -> Option<Sign> {
1579        Some(Sign::Optional(SignatureType::RequestSign))
1580    }
1581
1582    fn respond(
1583        self: Box<Self>,
1584        result: Result<http::Response, UpstreamRequestError>,
1585    ) -> Pin<Box<dyn Future<Output = ()> + Send + Sync>> {
1586        Box::pin(async move {
1587            let result = match result {
1588                Ok(mut response) => response.consume().await.map_err(UpstreamRequestError::Http),
1589                Err(error) => Err(error),
1590            };
1591
1592            match result {
1593                Ok(()) => self.envelope.accept(),
1594                Err(error) if error.is_received() => {
1595                    let scoping = self.envelope.scoping();
1596                    self.envelope.accept();
1597
1598                    if let UpstreamRequestError::RateLimited(limits) = error {
1599                        self.project_cache
1600                            .get(scoping.project_key)
1601                            .rate_limits()
1602                            .merge(limits.scope(&scoping));
1603                    }
1604                }
1605                Err(error) => {
1606                    // Errors are only logged for what we consider an internal discard reason. These
1607                    // indicate errors in the infrastructure or implementation bugs.
1608                    let mut envelope = self.envelope;
1609                    envelope.reject(Outcome::Invalid(DiscardReason::Internal));
1610                    relay_log::error!(
1611                        error = &error as &dyn Error,
1612                        tags.project_key = %envelope.scoping().project_key,
1613                        "error sending envelope"
1614                    );
1615                }
1616            }
1617        })
1618    }
1619}
1620
1621/// A container for metric buckets from multiple projects.
1622///
1623/// This container is used to send metrics to the upstream in global batches as part of the
1624/// [`FlushBuckets`] message if the `http.global_metrics` option is enabled. The container monitors
1625/// the size of all metrics and allows to split them into multiple batches. See
1626/// [`insert`](Self::insert) for more information.
1627#[derive(Debug)]
1628struct Partition<'a> {
1629    max_size: usize,
1630    remaining: usize,
1631    views: HashMap<ProjectKey, Vec<BucketView<'a>>>,
1632    project_info: HashMap<ProjectKey, Scoping>,
1633}
1634
1635impl<'a> Partition<'a> {
1636    /// Creates a new partition with the given maximum size in bytes.
1637    pub fn new(size: usize) -> Self {
1638        Self {
1639            max_size: size,
1640            remaining: size,
1641            views: HashMap::new(),
1642            project_info: HashMap::new(),
1643        }
1644    }
1645
1646    /// Inserts a bucket into the partition, splitting it if necessary.
1647    ///
1648    /// This function attempts to add the bucket to this partition. If the bucket does not fit
1649    /// entirely into the partition given its maximum size, the remaining part of the bucket is
1650    /// returned from this function call.
1651    ///
1652    /// If this function returns `Some(_)`, the partition is full and should be submitted to the
1653    /// upstream immediately. Use [`Self::take`] to retrieve the contents of the
1654    /// partition. Afterwards, the caller is responsible to call this function again with the
1655    /// remaining bucket until it is fully inserted.
1656    pub fn insert(&mut self, bucket: BucketView<'a>, scoping: Scoping) -> Option<BucketView<'a>> {
1657        let (current, next) = bucket.split(self.remaining, Some(self.max_size));
1658
1659        if let Some(current) = current {
1660            self.remaining = self.remaining.saturating_sub(current.estimated_size());
1661            self.views
1662                .entry(scoping.project_key)
1663                .or_default()
1664                .push(current);
1665
1666            self.project_info
1667                .entry(scoping.project_key)
1668                .or_insert(scoping);
1669        }
1670
1671        next
1672    }
1673
1674    /// Returns `true` if the partition does not hold any data.
1675    fn is_empty(&self) -> bool {
1676        self.views.is_empty()
1677    }
1678
1679    /// Returns the serialized buckets for this partition.
1680    ///
1681    /// This empties the partition, so that it can be reused.
1682    fn take(&mut self) -> (Bytes, HashMap<ProjectKey, Scoping>) {
1683        #[derive(serde::Serialize)]
1684        struct Wrapper<'a> {
1685            buckets: &'a HashMap<ProjectKey, Vec<BucketView<'a>>>,
1686        }
1687
1688        let buckets = &self.views;
1689        let payload = serde_json::to_vec(&Wrapper { buckets }).unwrap().into();
1690
1691        let scopings = std::mem::take(&mut self.project_info);
1692
1693        self.views.clear();
1694        self.remaining = self.max_size;
1695
1696        (payload, scopings)
1697    }
1698}
1699
1700/// An upstream request that submits metric buckets via HTTP.
1701///
1702/// This request is not awaited. It automatically tracks outcomes if the request is not received.
1703#[derive(Debug)]
1704struct SendMetricsRequest {
1705    /// Optional upstream override where the request will be sent to.
1706    upstream: Option<UpstreamDescriptor>,
1707    /// If the partition key is set, the request is marked with `X-Sentry-Relay-Shard`.
1708    partition_key: String,
1709    /// Serialized metric buckets without encoding applied, used for signing.
1710    unencoded: Bytes,
1711    /// Serialized metric buckets with the stated HTTP encoding applied.
1712    encoded: Bytes,
1713    /// Mapping of all contained project keys to their scoping and extraction mode.
1714    ///
1715    /// Used to track outcomes for transmission failures.
1716    project_info: HashMap<ProjectKey, Scoping>,
1717    /// Encoding (compression) of the payload.
1718    http_encoding: HttpEncoding,
1719    /// Metric outcomes instance to send outcomes on error.
1720    metric_outcomes: MetricOutcomes,
1721}
1722
1723impl SendMetricsRequest {
1724    fn create_error_outcomes(self) {
1725        #[derive(serde::Deserialize)]
1726        struct Wrapper {
1727            buckets: HashMap<ProjectKey, Vec<MinimalTrackableBucket>>,
1728        }
1729
1730        let buckets = match serde_json::from_slice(&self.unencoded) {
1731            Ok(Wrapper { buckets }) => buckets,
1732            Err(err) => {
1733                relay_log::error!(
1734                    error = &err as &dyn std::error::Error,
1735                    "failed to parse buckets from failed transmission"
1736                );
1737                return;
1738            }
1739        };
1740
1741        for (key, buckets) in buckets {
1742            let Some(&scoping) = self.project_info.get(&key) else {
1743                relay_log::error!("missing scoping for project key");
1744                continue;
1745            };
1746
1747            self.metric_outcomes.track(
1748                scoping,
1749                &buckets,
1750                Outcome::Invalid(DiscardReason::Internal),
1751            );
1752        }
1753    }
1754}
1755
1756impl UpstreamRequest for SendMetricsRequest {
1757    fn upstream(&self) -> Option<&UpstreamDescriptor> {
1758        self.upstream.as_ref()
1759    }
1760
1761    fn set_relay_id(&self) -> bool {
1762        true
1763    }
1764
1765    fn sign(&mut self) -> Option<Sign> {
1766        Some(Sign::Required(SignatureType::Body(self.unencoded.clone())))
1767    }
1768
1769    fn method(&self) -> reqwest::Method {
1770        reqwest::Method::POST
1771    }
1772
1773    fn path(&self) -> Cow<'_, str> {
1774        "/api/0/relays/metrics/".into()
1775    }
1776
1777    fn route(&self) -> &'static str {
1778        "global_metrics"
1779    }
1780
1781    fn build(&mut self, builder: &mut http::RequestBuilder) -> Result<(), http::HttpError> {
1782        builder
1783            .content_encoding(self.http_encoding)
1784            .header("X-Sentry-Relay-Shard", self.partition_key.as_bytes())
1785            .header(header::CONTENT_TYPE, b"application/json")
1786            .body(self.encoded.clone());
1787
1788        Ok(())
1789    }
1790
1791    fn respond(
1792        self: Box<Self>,
1793        result: Result<http::Response, UpstreamRequestError>,
1794    ) -> Pin<Box<dyn Future<Output = ()> + Send + Sync>> {
1795        Box::pin(async {
1796            match result {
1797                Ok(mut response) => {
1798                    response.consume().await.ok();
1799                }
1800                Err(error) => {
1801                    relay_log::error!(error = &error as &dyn Error, "Failed to send metrics batch");
1802
1803                    // If the request did not arrive at the upstream, we are responsible for outcomes.
1804                    // Otherwise, the upstream is responsible to log outcomes.
1805                    if error.is_received() {
1806                        return;
1807                    }
1808
1809                    self.create_error_outcomes()
1810                }
1811            }
1812        })
1813    }
1814}
1815
1816/// Container for global and project level [`Quota`].
1817#[derive(Copy, Clone, Debug)]
1818#[cfg(feature = "processing")]
1819struct CombinedQuotas<'a> {
1820    global_quotas: &'a [Quota],
1821    project_quotas: &'a [Quota],
1822}
1823
1824#[cfg(feature = "processing")]
1825impl<'a> CombinedQuotas<'a> {
1826    /// Returns a new [`CombinedQuotas`].
1827    pub fn new(global_config: &'a GlobalConfig, project_quotas: &'a [Quota]) -> Self {
1828        Self {
1829            global_quotas: &global_config.quotas,
1830            project_quotas,
1831        }
1832    }
1833}
1834
1835#[cfg(feature = "processing")]
1836impl<'a> IntoIterator for CombinedQuotas<'a> {
1837    type Item = &'a Quota;
1838    type IntoIter = std::iter::Chain<std::slice::Iter<'a, Quota>, std::slice::Iter<'a, Quota>>;
1839
1840    fn into_iter(self) -> Self::IntoIter {
1841        self.global_quotas.iter().chain(self.project_quotas.iter())
1842    }
1843}
1844
1845#[cfg(test)]
1846mod tests {
1847    use insta::assert_debug_snapshot;
1848    use relay_common::glob2::LazyGlob;
1849    use relay_dynamic_config::ProjectConfig;
1850    use relay_event_normalization::{
1851        NormalizationConfig, RedactionRule, TransactionNameConfig, TransactionNameRule,
1852    };
1853    use relay_event_schema::protocol::{Event, EventId, TransactionSource};
1854    use relay_pii::DataScrubbingConfig;
1855    use relay_protocol::Annotated;
1856    #[cfg(feature = "processing")]
1857    use relay_quotas::DataCategory;
1858    use similar_asserts::assert_eq;
1859
1860    use crate::testutils::{create_test_processor, create_test_processor_with_addrs};
1861
1862    #[cfg(feature = "processing")]
1863    use {
1864        relay_metrics::BucketValue,
1865        relay_quotas::{QuotaScope, ReasonCode},
1866        relay_test::mock_service,
1867    };
1868
1869    use super::*;
1870
1871    async fn process_to_single_envelope<'a>(
1872        processor: &EnvelopeProcessorService,
1873        envelope: ManagedEnvelope,
1874        ctx: processing::Context<'a>,
1875    ) -> Box<Envelope> {
1876        let mut outputs = processor.process(envelope, ctx).await;
1877        assert_eq!(outputs.len(), 1);
1878
1879        let Output {
1880            main,
1881            metrics,
1882            intermediates: _,
1883        } = outputs.pop().unwrap();
1884
1885        if let Some(metrics) = metrics {
1886            metrics.accept(drop);
1887        }
1888
1889        main.unwrap()
1890            .serialize_envelope(ctx.to_forward())
1891            .unwrap()
1892            .accept(|envelope| envelope)
1893    }
1894
1895    #[cfg(feature = "processing")]
1896    fn mock_quota(id: &str) -> Quota {
1897        Quota {
1898            id: Some(id.into()),
1899            categories: [DataCategory::MetricBucket].into(),
1900            scope: QuotaScope::Organization,
1901            scope_id: None,
1902            limit: Some(0),
1903            window: None,
1904            reason_code: None,
1905            namespace: None,
1906        }
1907    }
1908
1909    #[cfg(feature = "processing")]
1910    #[test]
1911    fn test_dynamic_quotas() {
1912        let global_config = relay_dynamic_config::GlobalConfig {
1913            quotas: vec![mock_quota("foo"), mock_quota("bar")],
1914            ..Default::default()
1915        };
1916
1917        let project_quotas = vec![mock_quota("baz"), mock_quota("qux")];
1918
1919        let dynamic_quotas = CombinedQuotas::new(&global_config, &project_quotas);
1920
1921        let quota_ids = dynamic_quotas.into_iter().filter_map(|q| q.id.as_deref());
1922        assert!(quota_ids.eq(["foo", "bar", "baz", "qux"]));
1923    }
1924
1925    /// Ensures that if we ratelimit one batch of buckets in [`FlushBuckets`] message, it won't
1926    /// also ratelimit the next batches in the same message automatically.
1927    #[cfg(feature = "processing")]
1928    #[tokio::test]
1929    async fn test_ratelimit_per_batch() {
1930        use relay_base_schema::organization::OrganizationId;
1931        use relay_protocol::FiniteF64;
1932
1933        let rate_limited_org = Scoping {
1934            organization_id: OrganizationId::new(1),
1935            project_id: ProjectId::new(21),
1936            project_key: ProjectKey::parse("00000000000000000000000000000000").unwrap(),
1937            key_id: Some(17),
1938        };
1939
1940        let not_rate_limited_org = Scoping {
1941            organization_id: OrganizationId::new(2),
1942            project_id: ProjectId::new(21),
1943            project_key: ProjectKey::parse("11111111111111111111111111111111").unwrap(),
1944            key_id: Some(17),
1945        };
1946
1947        let message = {
1948            let project_info = {
1949                let quota = Quota {
1950                    id: Some("testing".into()),
1951                    categories: [DataCategory::MetricBucket].into(),
1952                    scope: relay_quotas::QuotaScope::Organization,
1953                    scope_id: Some(rate_limited_org.organization_id.to_string().into()),
1954                    limit: Some(0),
1955                    window: None,
1956                    reason_code: Some(ReasonCode::new("test")),
1957                    namespace: None,
1958                };
1959
1960                let mut config = ProjectConfig::default();
1961                config.quotas.push(quota);
1962
1963                Arc::new(ProjectInfo {
1964                    config,
1965                    ..Default::default()
1966                })
1967            };
1968
1969            let project_metrics = |scoping| ProjectBuckets {
1970                buckets: vec![Bucket {
1971                    name: "d:spans/bar".into(),
1972                    value: BucketValue::Counter(FiniteF64::new(1.0).unwrap()),
1973                    timestamp: UnixTimestamp::now(),
1974                    tags: Default::default(),
1975                    width: 10,
1976                    metadata: BucketMetadata::default(),
1977                }],
1978                rate_limits: Default::default(),
1979                project_info: project_info.clone(),
1980                scoping,
1981            };
1982
1983            let buckets = hashbrown::HashMap::from([
1984                (
1985                    rate_limited_org.project_key,
1986                    project_metrics(rate_limited_org),
1987                ),
1988                (
1989                    not_rate_limited_org.project_key,
1990                    project_metrics(not_rate_limited_org),
1991                ),
1992            ]);
1993
1994            FlushBuckets {
1995                partition_key: 0,
1996                buckets,
1997            }
1998        };
1999
2000        // ensure the order of the map while iterating is as expected.
2001        assert_eq!(message.buckets.keys().count(), 2);
2002
2003        let config = {
2004            let config_json = serde_json::json!({
2005                "processing": {
2006                    "enabled": true,
2007                    "kafka_config": [],
2008                    "redis": {
2009                        "server": std::env::var("RELAY_REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_owned()),
2010                    }
2011                }
2012            });
2013            Config::from_json_value(config_json).unwrap()
2014        };
2015
2016        let (store, handle) = {
2017            let f = |org_ids: &mut Vec<OrganizationId>, msg: Store| {
2018                let org_id = match msg {
2019                    Store::Metrics(x) => x.scoping.organization_id,
2020                    _ => panic!("received envelope when expecting only metrics"),
2021                };
2022                org_ids.push(org_id);
2023            };
2024
2025            mock_service("store_forwarder", vec![], f)
2026        };
2027
2028        let processor = create_test_processor(config).await;
2029        assert!(processor.redis_rate_limiter_enabled());
2030
2031        processor.encode_metrics_processing(message, &store).await;
2032
2033        drop(store);
2034        let orgs_not_ratelimited = handle.await.unwrap();
2035
2036        assert_eq!(
2037            orgs_not_ratelimited,
2038            vec![not_rate_limited_org.organization_id]
2039        );
2040    }
2041
2042    #[tokio::test]
2043    async fn test_browser_version_extraction_with_pii_like_data() {
2044        let processor = create_test_processor(Default::default()).await;
2045        let outcome_aggregator = Addr::dummy();
2046        let event_id = EventId::new();
2047
2048        let dsn = "https://e12d836b15bb49d7bbf99e64295d995b:@sentry.io/42"
2049            .parse()
2050            .unwrap();
2051
2052        let request_meta = RequestMeta::new(dsn);
2053        let mut envelope = Envelope::from_request(Some(event_id), request_meta);
2054
2055        envelope.add_item({
2056                let mut item = Item::new(ItemType::Event);
2057                item.set_payload(
2058                    ContentType::Json,
2059                    r#"
2060                    {
2061                        "request": {
2062                            "headers": [
2063                                ["User-Agent", "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/103.0.0.0 Safari/537.36"]
2064                            ]
2065                        }
2066                    }
2067                "#,
2068                );
2069                item
2070            });
2071
2072        let mut datascrubbing_settings = DataScrubbingConfig::default();
2073        // enable all the default scrubbing
2074        datascrubbing_settings.scrub_data = true;
2075        datascrubbing_settings.scrub_defaults = true;
2076        datascrubbing_settings.scrub_ip_addresses = true;
2077
2078        // Make sure to mask any IP-like looking data
2079        let pii_config = serde_json::from_str(r#"{"applications": {"**": ["@ip:mask"]}}"#).unwrap();
2080
2081        let config = ProjectConfig {
2082            datascrubbing_settings,
2083            pii_config: Some(pii_config),
2084            ..Default::default()
2085        };
2086
2087        let project_info = ProjectInfo {
2088            config,
2089            ..Default::default()
2090        };
2091
2092        let envelope = ManagedEnvelope::new(envelope, outcome_aggregator);
2093
2094        let ctx = processing::Context {
2095            project_info: &project_info,
2096            ..processing::Context::for_test()
2097        };
2098
2099        let new_envelope = process_to_single_envelope(&processor, envelope, ctx).await;
2100
2101        let event_item = new_envelope.items().last().unwrap();
2102        let annotated_event: Annotated<Event> =
2103            Annotated::from_json_bytes(&event_item.payload()).unwrap();
2104        let event = annotated_event.into_value().unwrap();
2105        let headers = event
2106            .request
2107            .into_value()
2108            .unwrap()
2109            .headers
2110            .into_value()
2111            .unwrap();
2112
2113        // IP-like data must be masked
2114        assert_eq!(
2115            Some(
2116                "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/********* Safari/537.36"
2117            ),
2118            headers.get_header("User-Agent")
2119        );
2120        // But we still get correct browser and version number
2121        let contexts = event.contexts.into_value().unwrap();
2122        let browser = contexts.0.get("browser").unwrap();
2123        assert_eq!(
2124            r#"{"browser":"Chrome 103.0.0","name":"Chrome","version":"103.0.0","type":"browser"}"#,
2125            browser.to_json().unwrap()
2126        );
2127    }
2128
2129    #[tokio::test]
2130    #[cfg(feature = "processing")]
2131    async fn test_materialize_dsc() {
2132        use crate::services::projects::project::PublicKeyConfig;
2133
2134        let dsn = "https://e12d836b15bb49d7bbf99e64295d995b:@sentry.io/42"
2135            .parse()
2136            .unwrap();
2137        let request_meta = RequestMeta::new(dsn);
2138        let mut envelope = Envelope::from_request(None, request_meta);
2139
2140        let dsc = r#"{
2141            "trace_id": "00000000-0000-0000-0000-000000000001",
2142            "public_key": "e12d836b15bb49d7bbf99e64295d995b",
2143            "sample_rate": "0.2"
2144        }"#;
2145        envelope.set_dsc(serde_json::from_str(dsc).unwrap());
2146
2147        let mut item = Item::new(ItemType::Event);
2148        item.set_payload(ContentType::Json, r#"{}"#);
2149        envelope.add_item(item);
2150
2151        let outcome_aggregator = Addr::dummy();
2152        let managed_envelope = ManagedEnvelope::new(envelope, outcome_aggregator);
2153
2154        let mut project_info = ProjectInfo::default();
2155        project_info.public_keys.push(PublicKeyConfig {
2156            public_key: ProjectKey::parse("e12d836b15bb49d7bbf99e64295d995b").unwrap(),
2157            numeric_id: Some(1),
2158        });
2159
2160        let config = serde_json::json!({
2161            "processing": {
2162                "enabled": true,
2163                "kafka_config": [],
2164            }
2165        });
2166
2167        let processor =
2168            create_test_processor(Config::from_json_value(config.clone()).unwrap()).await;
2169        let config = Config::from_json_value(config).unwrap().current();
2170        let ctx = processing::Context {
2171            config: &config,
2172            project_info: &project_info,
2173            sampling_project_info: Some(&project_info),
2174            ..processing::Context::for_test()
2175        };
2176
2177        let envelope = process_to_single_envelope(&processor, managed_envelope, ctx).await;
2178        let event = envelope
2179            .get_item_by(|item| item.ty() == &ItemType::Event)
2180            .unwrap();
2181
2182        let event = Annotated::<Event>::from_json_bytes(&event.payload()).unwrap();
2183        insta::assert_debug_snapshot!(event.value().unwrap()._dsc, @r###"
2184        Object(
2185            {
2186                "environment": ~,
2187                "public_key": String(
2188                    "e12d836b15bb49d7bbf99e64295d995b",
2189                ),
2190                "release": ~,
2191                "replay_id": ~,
2192                "sample_rate": String(
2193                    "0.2",
2194                ),
2195                "trace_id": String(
2196                    "00000000000000000000000000000001",
2197                ),
2198                "transaction": ~,
2199            },
2200        )
2201        "###);
2202    }
2203
2204    fn capture_test_event(transaction_name: &str, source: TransactionSource) -> Vec<String> {
2205        let mut event = Annotated::<Event>::from_json(
2206            r#"
2207            {
2208                "type": "transaction",
2209                "transaction": "/foo/",
2210                "timestamp": 946684810.0,
2211                "start_timestamp": 946684800.0,
2212                "contexts": {
2213                    "trace": {
2214                        "trace_id": "4c79f60c11214eb38604f4ae0781bfb2",
2215                        "span_id": "fa90fdead5f74053",
2216                        "op": "http.server",
2217                        "type": "trace"
2218                    }
2219                },
2220                "transaction_info": {
2221                    "source": "url"
2222                }
2223            }
2224            "#,
2225        )
2226        .unwrap();
2227        let e = event.value_mut().as_mut().unwrap();
2228        e.transaction.set_value(Some(transaction_name.into()));
2229
2230        e.transaction_info
2231            .value_mut()
2232            .as_mut()
2233            .unwrap()
2234            .source
2235            .set_value(Some(source));
2236
2237        relay_statsd::with_capturing_test_client(|| {
2238            utils::log_transaction_name_metrics(&mut event, |event| {
2239                let config = NormalizationConfig {
2240                    transaction_name_config: TransactionNameConfig {
2241                        rules: &[TransactionNameRule {
2242                            pattern: LazyGlob::new("/foo/*/**".to_owned()),
2243                            expiry: DateTime::<Utc>::MAX_UTC,
2244                            redaction: RedactionRule::Replace {
2245                                substitution: "*".to_owned(),
2246                            },
2247                        }],
2248                    },
2249                    ..Default::default()
2250                };
2251                relay_event_normalization::normalize_event(event, &config)
2252            });
2253        })
2254    }
2255
2256    #[test]
2257    fn test_log_transaction_metrics_none() {
2258        let captures = capture_test_event("/nothing", TransactionSource::Url);
2259        insta::assert_debug_snapshot!(captures, @r###"
2260        [
2261            "event.transaction_name_changes:1|c|#source_in:url,changes:none,source_out:sanitized,is_404:false",
2262        ]
2263        "###);
2264    }
2265
2266    #[test]
2267    fn test_log_transaction_metrics_rule() {
2268        let captures = capture_test_event("/foo/john/denver", TransactionSource::Url);
2269        insta::assert_debug_snapshot!(captures, @r###"
2270        [
2271            "event.transaction_name_changes:1|c|#source_in:url,changes:rule,source_out:sanitized,is_404:false",
2272        ]
2273        "###);
2274    }
2275
2276    #[test]
2277    fn test_log_transaction_metrics_pattern() {
2278        let captures = capture_test_event("/something/12345", TransactionSource::Url);
2279        insta::assert_debug_snapshot!(captures, @r###"
2280        [
2281            "event.transaction_name_changes:1|c|#source_in:url,changes:pattern,source_out:sanitized,is_404:false",
2282        ]
2283        "###);
2284    }
2285
2286    #[test]
2287    fn test_log_transaction_metrics_both() {
2288        let captures = capture_test_event("/foo/john/12345", TransactionSource::Url);
2289        insta::assert_debug_snapshot!(captures, @r###"
2290        [
2291            "event.transaction_name_changes:1|c|#source_in:url,changes:both,source_out:sanitized,is_404:false",
2292        ]
2293        "###);
2294    }
2295
2296    #[test]
2297    fn test_log_transaction_metrics_no_match() {
2298        let captures = capture_test_event("/foo/john/12345", TransactionSource::Route);
2299        insta::assert_debug_snapshot!(captures, @r###"
2300        [
2301            "event.transaction_name_changes:1|c|#source_in:route,changes:none,source_out:route,is_404:false",
2302        ]
2303        "###);
2304    }
2305
2306    #[tokio::test]
2307    async fn test_process_metrics_bucket_metadata() {
2308        let mut token = Cogs::noop().timed(ResourceId::Relay, AppFeature::Unattributed);
2309        let project_key = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap();
2310        let received_at = Utc::now();
2311        let config = Config::default();
2312
2313        let (aggregator, mut aggregator_rx) = Addr::custom();
2314        let processor = create_test_processor_with_addrs(
2315            config,
2316            Addrs {
2317                aggregator,
2318                ..Default::default()
2319            },
2320        )
2321        .await;
2322
2323        let mut item = Item::new(ItemType::Statsd);
2324        item.set_payload(ContentType::Text, "sessions/foo:3182887624:4267882815|s");
2325        for (source, expected_received_at) in [
2326            (
2327                BucketSource::External,
2328                Some(UnixTimestamp::from_datetime(received_at).unwrap()),
2329            ),
2330            (BucketSource::Internal, None),
2331        ] {
2332            let message = ProcessMetrics {
2333                data: MetricData::Raw(vec![item.clone()]),
2334                project_key,
2335                source,
2336                received_at,
2337                sent_at: Some(Utc::now()),
2338            };
2339            processor.handle_process_metrics(&mut token, message);
2340
2341            let Aggregator::MergeBuckets(merge_buckets) = aggregator_rx.recv().await.unwrap();
2342            let buckets = merge_buckets.buckets;
2343            assert_eq!(buckets.len(), 1);
2344            assert_eq!(buckets[0].metadata.received_at, expected_received_at);
2345        }
2346    }
2347
2348    #[tokio::test]
2349    async fn test_process_batched_metrics() {
2350        let mut token = Cogs::noop().timed(ResourceId::Relay, AppFeature::Unattributed);
2351        let received_at = Utc::now();
2352        let config = Config::default();
2353
2354        let (aggregator, mut aggregator_rx) = Addr::custom();
2355        let processor = create_test_processor_with_addrs(
2356            config,
2357            Addrs {
2358                aggregator,
2359                ..Default::default()
2360            },
2361        )
2362        .await;
2363
2364        let payload = r#"{
2365    "buckets": {
2366        "11111111111111111111111111111111": [
2367            {
2368                "timestamp": 1615889440,
2369                "width": 0,
2370                "name": "d:transactions/endpoint.response_time@millisecond",
2371                "type": "d",
2372                "value": [
2373                  68.0
2374                ],
2375                "tags": {
2376                  "route": "user_index"
2377                }
2378            }
2379        ],
2380        "22222222222222222222222222222222": [
2381            {
2382                "timestamp": 1615889440,
2383                "width": 0,
2384                "name": "d:transactions/endpoint.cache_rate@none",
2385                "type": "d",
2386                "value": [
2387                  36.0
2388                ]
2389            }
2390        ]
2391    }
2392}
2393"#;
2394        let message = ProcessBatchedMetrics {
2395            payload: Bytes::from(payload),
2396            source: BucketSource::Internal,
2397            received_at,
2398            sent_at: Some(Utc::now()),
2399        };
2400        processor.handle_process_batched_metrics(&mut token, message);
2401
2402        let Aggregator::MergeBuckets(mb1) = aggregator_rx.recv().await.unwrap();
2403        let Aggregator::MergeBuckets(mb2) = aggregator_rx.recv().await.unwrap();
2404
2405        let mut messages = vec![mb1, mb2];
2406        messages.sort_by_key(|mb| mb.project_key);
2407
2408        let actual = messages
2409            .into_iter()
2410            .map(|mb| (mb.project_key, mb.buckets))
2411            .collect::<Vec<_>>();
2412
2413        assert_debug_snapshot!(actual, @r###"
2414        [
2415            (
2416                ProjectKey("11111111111111111111111111111111"),
2417                [
2418                    Bucket {
2419                        timestamp: UnixTimestamp(1615889440),
2420                        width: 0,
2421                        name: MetricName(
2422                            "d:transactions/endpoint.response_time@millisecond",
2423                        ),
2424                        value: Distribution(
2425                            [
2426                                68.0,
2427                            ],
2428                        ),
2429                        tags: {
2430                            "route": "user_index",
2431                        },
2432                        metadata: BucketMetadata {
2433                            merges: 1,
2434                            received_at: None,
2435                            extracted_from_indexed: false,
2436                        },
2437                    },
2438                ],
2439            ),
2440            (
2441                ProjectKey("22222222222222222222222222222222"),
2442                [
2443                    Bucket {
2444                        timestamp: UnixTimestamp(1615889440),
2445                        width: 0,
2446                        name: MetricName(
2447                            "d:transactions/endpoint.cache_rate@none",
2448                        ),
2449                        value: Distribution(
2450                            [
2451                                36.0,
2452                            ],
2453                        ),
2454                        tags: {},
2455                        metadata: BucketMetadata {
2456                            merges: 1,
2457                            received_at: None,
2458                            extracted_from_indexed: false,
2459                        },
2460                    },
2461                ],
2462            ),
2463        ]
2464        "###);
2465    }
2466}