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
70pub const MINIMUM_CLOCK_DRIFT: Duration = Duration::from_secs(55 * 60);
72
73#[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#[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 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 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#[derive(Debug)]
256pub struct ProcessEnvelope {
257 pub envelope: ManagedEnvelope,
259 pub project_info: Arc<ProjectInfo>,
261 pub rate_limits: Arc<RateLimits>,
263 pub sampling_project_info: Option<Arc<ProjectInfo>>,
265}
266
267#[derive(Debug)]
275pub struct ProcessMetrics {
276 pub data: MetricData,
278 pub project_key: ProjectKey,
280 pub source: BucketSource,
282 pub received_at: DateTime<Utc>,
284 pub sent_at: Option<DateTime<Utc>>,
287}
288
289#[derive(Debug)]
291pub enum MetricData {
292 Raw(Vec<Item>),
294 Parsed(Vec<Bucket>),
296}
297
298impl MetricData {
299 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 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 pub payload: Bytes,
344 pub source: BucketSource,
346 pub received_at: DateTime<Utc>,
348 pub sent_at: Option<DateTime<Utc>>,
350}
351
352#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
354pub enum BucketSource {
355 Internal,
361 External,
366}
367
368impl BucketSource {
369 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#[derive(Debug)]
380pub struct SubmitClientReports {
381 pub client_reports: Vec<ClientReport>,
383 pub scoping: Scoping,
385}
386
387#[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 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
452pub type EnvelopeProcessorServicePool = AsyncPool<BoxFuture<'static, ()>>;
454
455#[derive(Clone)]
459pub struct EnvelopeProcessorService {
460 inner: Arc<InnerProcessor>,
461}
462
463pub 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 #[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 "a_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 if let Some(sampling_state) = ctx.sampling_project_info {
579 envelope
582 .envelope_mut()
583 .parametrize_dsc_transaction(&sampling_state.config.tx_name_rules);
584 }
585
586 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 async fn process<'a>(
604 &self,
605 mut envelope: ManagedEnvelope,
606 ctx: processing::Context<'a>,
607 ) -> Vec<Output<Outputs>> {
608 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 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 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 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 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 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 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 upstream: Option<UpstreamDescriptor>,
857 ) {
858 if envelope.envelope_mut().is_empty() {
859 envelope.accept();
860 return;
861 }
862
863 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 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 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 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 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 #[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 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 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 #[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 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 store_forwarder.send(StoreMetrics {
1192 buckets,
1193 scoping,
1194 retention,
1195 });
1196 }
1197 }
1198
1199 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 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 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 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 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 if config.emit_outcomes() == EmitOutcomes::AsClientReports {
1398 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 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 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 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 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 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#[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 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#[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 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 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 fn is_empty(&self) -> bool {
1653 self.views.is_empty()
1654 }
1655
1656 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#[derive(Debug)]
1681struct SendMetricsRequest {
1682 upstream: Option<UpstreamDescriptor>,
1684 partition_key: String,
1686 unencoded: Bytes,
1688 encoded: Bytes,
1690 project_info: HashMap<ProjectKey, Scoping>,
1694 http_encoding: HttpEncoding,
1696 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 error.is_received() {
1783 return;
1784 }
1785
1786 self.create_error_outcomes()
1787 }
1788 }
1789 })
1790 }
1791}
1792
1793#[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 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 #[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 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 datascrubbing_settings.scrub_data = true;
2054 datascrubbing_settings.scrub_defaults = true;
2055 datascrubbing_settings.scrub_ip_addresses = true;
2056
2057 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 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 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}