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