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