Skip to main content

relay_server/services/
store.rs

1//! This module contains the service that forwards events and attachments to the Sentry store.
2//! The service uses Kafka topics to forward data to Sentry
3
4use std::borrow::Cow;
5use std::collections::BTreeMap;
6use std::error::Error;
7use std::pin::Pin;
8use std::sync::Arc;
9use std::task;
10use std::time::{Duration, Instant};
11
12use bytes::Bytes;
13use chrono::{DateTime, SecondsFormat, Utc};
14use prost::Message as _;
15use relay_base_schema::events::EventType;
16use relay_conventions::attributes::SENTRY__SEGMENT__ID;
17use sentry::protocol::{Attachment, SpanId};
18use sentry_protos::snuba::v1::{TraceItem, TraceItemType};
19use serde::Serialize;
20use tokio::sync::Mutex;
21use uuid::Uuid;
22
23use relay_base_schema::data_category::DataCategory;
24use relay_base_schema::organization::OrganizationId;
25use relay_base_schema::project::ProjectId;
26use relay_common::time::UnixTimestamp;
27use relay_config::{Config, ConfigSnapshot};
28use relay_event_schema::protocol::{Event, EventId, SpanV2, datetime_to_timestamp};
29use relay_kafka::{ClientError, KafkaClient, KafkaTopic, Message, SerializationOutput};
30use relay_metrics::{
31    Bucket, BucketView, BucketViewValue, BucketsView, ByNamespace, GaugeValue, MetricName,
32    MetricNamespace, SetView,
33};
34use relay_protocol::{Annotated, FiniteF64, SerializableAnnotated};
35use relay_quotas::Scoping;
36use relay_statsd::metric;
37use relay_system::{FromMessage, Interface, NoResponse, Service};
38use relay_threading::AsyncPool;
39
40use crate::envelope::{AttachmentPlaceholder, AttachmentType, ContentType, Item};
41use crate::managed::{Counted, Managed, OutcomeError, Quantities, Rejected};
42use crate::metrics::{ArrayEncoding, BucketEncoder, MetricOutcomes};
43
44use crate::service::ServiceError;
45use crate::services::global_config::GlobalConfigHandle;
46use crate::services::objectstore::ObjectstoreKey;
47use crate::services::outcome::{self, DiscardReason, Outcome, OutcomeId};
48use crate::services::upload::{Final, SignedLocation};
49use crate::statsd::{RelayCounters, RelayGauges, RelayTimers};
50use crate::utils;
51
52mod sessions;
53
54/// Fallback name used for attachment items without a `filename` header.
55const UNNAMED_ATTACHMENT: &str = "Unnamed Attachment";
56
57#[derive(Debug, thiserror::Error)]
58pub enum StoreError {
59    #[error("failed to send the message to kafka: {0}")]
60    SendFailed(#[from] ClientError),
61    #[error("failed to encode data: {0}")]
62    EncodingFailed(std::io::Error),
63    #[error("failed to serialize data: {0}")]
64    Serialize(#[from] serde_json::Error),
65    #[error("failed to store event because event id was missing")]
66    NoEventId,
67    #[error("invalid attachment reference")]
68    InvalidAttachmentRef,
69    #[error("invalid span id: {0}")]
70    InvalidSpanId(#[from] hex::FromHexError),
71}
72
73impl OutcomeError for StoreError {
74    type Error = Self;
75
76    fn consume(self) -> (Option<Outcome>, Self::Error) {
77        let outcome = match self {
78            StoreError::SendFailed(_)
79            | StoreError::EncodingFailed(_)
80            | StoreError::Serialize(_)
81            | StoreError::NoEventId => Some(Outcome::Invalid(DiscardReason::Internal)),
82            StoreError::InvalidAttachmentRef => {
83                Some(Outcome::Invalid(DiscardReason::InvalidAttachmentRef))
84            }
85            StoreError::InvalidSpanId(_) => Some(Outcome::Invalid(DiscardReason::Internal)),
86        };
87        (outcome, self)
88    }
89}
90
91struct Producer {
92    client: KafkaClient,
93}
94
95impl Producer {
96    pub fn create(config: &ConfigSnapshot) -> anyhow::Result<Self> {
97        let mut client_builder = KafkaClient::builder();
98
99        for topic in KafkaTopic::iter() {
100            let kafka_configs = config.kafka_configs(*topic)?;
101            client_builder = client_builder
102                .add_kafka_topic_config(*topic, &kafka_configs, config.kafka_validate_topics())
103                .map_err(|e| ServiceError::Kafka(e.to_string()))?;
104        }
105
106        Ok(Self {
107            client: client_builder.build(),
108        })
109    }
110}
111
112/// Publishes an [`Event`] and its attachments to Kafka.
113#[derive(Debug)]
114pub struct StoreEvent {
115    /// The data category the [`Self::event`] is counted in.
116    pub event_category: DataCategory,
117    /// The event to be stored.
118    pub event: Annotated<Event>,
119    /// A list of attachments associated with the event.
120    pub attachments: Vec<Item>,
121    /// A list of user reports associated with the event.
122    pub user_reports: Vec<Item>,
123    /// Event retention in days.
124    pub retention_days: u16,
125}
126
127impl Counted for StoreEvent {
128    fn quantities(&self) -> Quantities {
129        let mut quantities = smallvec::smallvec![(self.event_category, 1)];
130        quantities.extend(self.attachments.quantities());
131        quantities.extend(self.user_reports.quantities());
132        quantities
133    }
134}
135
136/// Publishes a list of [`Bucket`]s to the Sentry core application through Kafka topics.
137#[derive(Clone, Debug)]
138pub struct StoreMetrics {
139    pub buckets: Vec<Bucket>,
140    pub scoping: Scoping,
141    pub retention: u16,
142}
143
144/// Publishes a log item to the Sentry core application through Kafka.
145#[derive(Debug)]
146pub struct StoreTraceItem {
147    /// The final trace item which will be produced to Kafka.
148    pub trace_item: TraceItem,
149}
150
151impl Counted for StoreTraceItem {
152    fn quantities(&self) -> Quantities {
153        self.trace_item.quantities()
154    }
155}
156
157/// Publishes a span item to the Sentry core application through Kafka.
158#[derive(Debug)]
159pub struct StoreSpanV2 {
160    /// Routing key to assign a Kafka partition.
161    pub routing_key: Option<Uuid>,
162    /// Default retention of the span.
163    pub retention_days: u16,
164    /// Downsampled retention of the span.
165    pub downsampled_retention_days: u16,
166    /// Optional event id, if the span was extracted from a transaction.
167    pub event_id: Option<EventId>,
168    /// The final Sentry compatible span item.
169    pub item: SpanV2,
170    // Whether to run issue detection on the transaction or on the segment span.
171    // (only true for segment spans created from a transaction).
172    pub performance_issues_spans: bool,
173}
174
175impl Counted for StoreSpanV2 {
176    fn quantities(&self) -> Quantities {
177        smallvec::smallvec![(DataCategory::SpanIndexed, 1)]
178    }
179}
180
181/// Publishes a singular profile chunk to Kafka.
182#[derive(Debug)]
183pub struct StoreProfileChunk {
184    /// Default retention of the span.
185    pub retention_days: u16,
186    /// The serialized profile chunk payload.
187    pub payload: Bytes,
188    /// Additional attachments associated with this profile chunk.
189    pub attachments: Vec<ProfileAttachment>,
190    /// Outcome quantities associated with this profile.
191    ///
192    /// Quantities are different for backend and ui profile chunks.
193    pub quantities: Quantities,
194}
195
196/// Optional raw binary blob associated with [`profile chunk`](StoreProfileChunk).
197#[derive(Debug)]
198pub struct ProfileAttachment {
199    /// Name of the attachment,
200    pub name: String,
201    /// Content type of the attachment.
202    pub content_type: ContentType,
203    /// Objectstore id of the attachment.
204    ///
205    /// Using this id the attachment can be retrieved again.
206    pub stored_id: ObjectstoreKey,
207}
208
209impl Counted for StoreProfileChunk {
210    fn quantities(&self) -> Quantities {
211        self.quantities.clone()
212    }
213}
214
215/// A replay to be stored to Kafka.
216#[derive(Debug)]
217pub struct StoreReplay {
218    /// The event ID.
219    pub event_id: EventId,
220    /// Number of days to retain.
221    pub retention_days: u16,
222    /// The recording payload (rrweb data).
223    pub recording: Bytes,
224    /// Optional replay event payload (JSON).
225    pub event: Option<Bytes>,
226    /// Optional replay video.
227    pub video: Option<Bytes>,
228    /// Outcome quantities associated with this replay.
229    ///
230    /// Quantities are different for web and native replays.
231    pub quantities: Quantities,
232}
233
234impl Counted for StoreReplay {
235    fn quantities(&self) -> Quantities {
236        self.quantities.clone()
237    }
238}
239
240/// An attachment to be stored to Kafka.
241#[derive(Debug)]
242pub struct StoreAttachment {
243    /// The event ID.
244    pub event_id: EventId,
245    /// That attachment item.
246    pub attachment: Item,
247    /// Outcome quantities associated with this attachment.
248    pub quantities: Quantities,
249    /// Data retention in days for this attachment.
250    pub retention: u16,
251}
252
253impl Counted for StoreAttachment {
254    fn quantities(&self) -> Quantities {
255        self.quantities.clone()
256    }
257}
258
259/// A user report to be stored to Kafka.
260#[derive(Debug)]
261pub struct StoreUserReport {
262    /// The event ID.
263    pub event_id: EventId,
264    /// The user report.
265    pub report: Item,
266}
267
268impl Counted for StoreUserReport {
269    fn quantities(&self) -> Quantities {
270        smallvec::smallvec![(DataCategory::UserReportV2, 1)]
271    }
272}
273
274/// A profile to be stored to Kafka.
275#[derive(Debug)]
276pub struct StoreProfile {
277    /// Number of days to retain.
278    pub retention_days: u16,
279    /// The profile.
280    pub profile: Item,
281    /// Outcome quantities associated with this profile.
282    pub quantities: Quantities,
283}
284
285impl Counted for StoreProfile {
286    fn quantities(&self) -> Quantities {
287        self.quantities.clone()
288    }
289}
290
291/// A monitor check-in to be stored to Kafka.
292#[derive(Debug)]
293pub struct StoreCheckIn {
294    /// The serialized check-in.
295    pub check_in: Item,
296    /// The SDK client which produced the check-in.
297    pub sdk: Option<String>,
298    /// Check-in retention in days.
299    pub retention_days: u16,
300}
301
302impl Counted for StoreCheckIn {
303    fn quantities(&self) -> Quantities {
304        self.check_in.quantities()
305    }
306}
307
308/// The asynchronous thread pool used for scheduling storing tasks in the envelope store.
309pub type StoreServicePool = AsyncPool<StoreTask>;
310
311/// Service interface for storing items in Kafka.
312#[derive(Debug)]
313pub enum Store {
314    /// An [`Event`] and its associated items.
315    Event(Managed<Box<StoreEvent>>),
316    /// Aggregated generic metrics.
317    Metrics(StoreMetrics),
318    /// A singular [`TraceItem`].
319    TraceItem(Managed<StoreTraceItem>),
320    /// A singular Span.
321    Span(Managed<Box<StoreSpanV2>>),
322    /// A singular profile chunk.
323    ProfileChunk(Managed<StoreProfileChunk>),
324    /// A singular replay.
325    Replay(Managed<StoreReplay>),
326    /// A singular attachment.
327    Attachment(Managed<StoreAttachment>),
328    /// A singular user report.
329    UserReport(Managed<StoreUserReport>),
330    /// A single profile.
331    Profile(Managed<StoreProfile>),
332    /// A singular monitor check-in.
333    CheckIn(Managed<StoreCheckIn>),
334}
335
336impl Store {
337    /// Returns the name of the message variant.
338    fn variant(&self) -> &'static str {
339        match self {
340            Store::Event(_) => "event",
341            Store::Metrics(_) => "metrics",
342            Store::TraceItem(_) => "trace_item",
343            Store::Span(_) => "span",
344            Store::ProfileChunk(_) => "profile_chunk",
345            Store::Replay(_) => "replay",
346            Store::Attachment(_) => "attachment",
347            Store::UserReport(_) => "user_report",
348            Store::Profile(_) => "profile",
349            Store::CheckIn(_) => "check_in",
350        }
351    }
352}
353
354impl Interface for Store {}
355
356impl FromMessage<Managed<Box<StoreEvent>>> for Store {
357    type Response = NoResponse;
358
359    fn from_message(message: Managed<Box<StoreEvent>>, _: ()) -> Self {
360        Self::Event(message)
361    }
362}
363
364impl FromMessage<StoreMetrics> for Store {
365    type Response = NoResponse;
366
367    fn from_message(message: StoreMetrics, _: ()) -> Self {
368        Self::Metrics(message)
369    }
370}
371
372impl FromMessage<Managed<StoreTraceItem>> for Store {
373    type Response = NoResponse;
374
375    fn from_message(message: Managed<StoreTraceItem>, _: ()) -> Self {
376        Self::TraceItem(message)
377    }
378}
379
380impl FromMessage<Managed<Box<StoreSpanV2>>> for Store {
381    type Response = NoResponse;
382
383    fn from_message(message: Managed<Box<StoreSpanV2>>, _: ()) -> Self {
384        Self::Span(message)
385    }
386}
387
388impl FromMessage<Managed<StoreProfileChunk>> for Store {
389    type Response = NoResponse;
390
391    fn from_message(message: Managed<StoreProfileChunk>, _: ()) -> Self {
392        Self::ProfileChunk(message)
393    }
394}
395
396impl FromMessage<Managed<StoreReplay>> for Store {
397    type Response = NoResponse;
398
399    fn from_message(message: Managed<StoreReplay>, _: ()) -> Self {
400        Self::Replay(message)
401    }
402}
403
404impl FromMessage<Managed<StoreAttachment>> for Store {
405    type Response = NoResponse;
406
407    fn from_message(message: Managed<StoreAttachment>, _: ()) -> Self {
408        Self::Attachment(message)
409    }
410}
411
412impl FromMessage<Managed<StoreUserReport>> for Store {
413    type Response = NoResponse;
414
415    fn from_message(message: Managed<StoreUserReport>, _: ()) -> Self {
416        Self::UserReport(message)
417    }
418}
419
420impl FromMessage<Managed<StoreProfile>> for Store {
421    type Response = NoResponse;
422
423    fn from_message(message: Managed<StoreProfile>, _: ()) -> Self {
424        Self::Profile(message)
425    }
426}
427
428impl FromMessage<Managed<StoreCheckIn>> for Store {
429    type Response = NoResponse;
430
431    fn from_message(message: Managed<StoreCheckIn>, _: ()) -> Self {
432        Self::CheckIn(message)
433    }
434}
435
436/// Service implementing the [`Store`] interface.
437pub struct StoreService {
438    pool: StoreServicePool,
439    config: Arc<Config>,
440    global_config: GlobalConfigHandle,
441    metric_outcomes: MetricOutcomes,
442    producer: Producer,
443    last_span_report: Mutex<Instant>, // used to debounce expensive instrumentation.
444}
445
446impl StoreService {
447    pub fn create(
448        pool: StoreServicePool,
449        config: Arc<Config>,
450        global_config: GlobalConfigHandle,
451        metric_outcomes: MetricOutcomes,
452    ) -> anyhow::Result<Self> {
453        let producer = Producer::create(&config.current())?;
454        Ok(Self {
455            pool,
456            config,
457            global_config,
458            metric_outcomes,
459            producer,
460            last_span_report: Instant::now().into(),
461        })
462    }
463
464    fn handle_message(&self, message: Store) {
465        let ty = message.variant();
466        relay_statsd::metric!(timer(RelayTimers::StoreServiceDuration), message = ty, {
467            let result = match message {
468                Store::Event(message) => self.handle_store_event(message),
469                Store::Metrics(message) => {
470                    self.handle_store_metrics(message);
471                    Ok(())
472                }
473                Store::TraceItem(message) => self.handle_store_trace_item(message),
474                Store::Span(message) => self.handle_store_span(message),
475                Store::ProfileChunk(message) => self.handle_store_profile_chunk(message),
476                Store::Replay(message) => self.handle_store_replay(message),
477                Store::Attachment(message) => self.handle_store_attachment(message),
478                Store::UserReport(message) => self.handle_user_report(message),
479                Store::Profile(message) => self.handle_profile(message),
480                Store::CheckIn(message) => self.handle_check_in(message),
481            };
482            if let Err(error) = result {
483                relay_log::error!(
484                    error = &error as &dyn Error,
485                    tags.message = ty,
486                    "failed to store message"
487                );
488            }
489        })
490    }
491
492    fn handle_store_event(
493        &self,
494        message: Managed<Box<StoreEvent>>,
495    ) -> Result<(), Rejected<StoreError>> {
496        let received_at = message.received_at();
497        let scoping = message.scoping();
498        let remote_addr = message.remote_addr().map(|ip| ip.to_string());
499
500        message.try_accept(|m| self.do_store_event(*m, scoping, received_at, remote_addr))
501    }
502
503    fn do_store_event(
504        &self,
505        store: StoreEvent,
506        scoping: Scoping,
507        received_at: DateTime<Utc>,
508        remote_addr: Option<String>,
509    ) -> Result<(), StoreError> {
510        let event_id = store.event.value().and_then(|e| e.id.value()).copied();
511        let event_id = event_id.ok_or(StoreError::NoEventId)?;
512
513        let event_type = store.event.value().and_then(|e| e.ty.value());
514        let send_individual_attachments = matches!(
515            event_type,
516            Some(&EventType::Transaction) | Some(&EventType::UserReportV2)
517        );
518
519        let mut attachments = Vec::new();
520        for attachment in store.attachments {
521            // Note: Technically when an error occurs we may have already produced some items to Kafka,
522            // but since we don't keep track of that properly and just error out here, outcomes will
523            // also report items which were already produced to Kafka as dropped.
524            //
525            // We specifically accept this here, since this is an extreme edge case with minimal practical
526            // consequences.
527            if let Some(attachment) = self.produce_attachment(
528                event_id,
529                scoping.project_id,
530                scoping.organization_id,
531                &attachment,
532                send_individual_attachments,
533                store.retention_days,
534            )? {
535                attachments.push(attachment);
536            }
537        }
538
539        for user_report in &store.user_reports {
540            self.produce_user_report(
541                event_id,
542                scoping.project_id,
543                scoping.organization_id,
544                received_at,
545                user_report,
546            )?;
547        }
548
549        let event_topic = if event_type == Some(&EventType::Transaction) {
550            KafkaTopic::Transactions
551        } else if event_type == Some(&EventType::UserReportV2) {
552            KafkaTopic::Feedback
553        } else if !attachments.is_empty() || !store.user_reports.is_empty() {
554            KafkaTopic::Attachments
555        } else {
556            KafkaTopic::Events
557        };
558
559        let payload = store.event.to_json()?.into_bytes().into();
560        self.produce(
561            event_topic,
562            KafkaMessage::Event(EventKafkaMessage {
563                payload,
564                start_time: safe_timestamp(received_at),
565                event_id,
566                project_id: scoping.project_id,
567                org_id: scoping.organization_id,
568                remote_addr,
569                attachments,
570            }),
571        )
572    }
573
574    fn handle_store_metrics(&self, message: StoreMetrics) {
575        let StoreMetrics {
576            buckets,
577            scoping,
578            retention,
579        } = message;
580
581        let batch_size = self.config.current().metrics_max_batch_size_bytes();
582        let mut error = None;
583
584        let global_config = self.global_config.current().unwrap_or_default();
585        let mut encoder = BucketEncoder::new(&global_config);
586
587        let emit_sessions_to_eap = utils::is_rolled_out(
588            scoping.organization_id.value(),
589            global_config.options.sessions_eap_rollout_rate,
590        )
591        .is_keep();
592
593        let now = UnixTimestamp::now();
594        let mut delay_stats = ByNamespace::<(u64, u64, u64)>::default();
595
596        for mut bucket in buckets {
597            let namespace = encoder.prepare(&mut bucket);
598
599            if let Some(received_at) = bucket.metadata.received_at {
600                let delay = now.as_secs().saturating_sub(received_at.as_secs());
601                let (total, count, max) = delay_stats.get_mut(namespace);
602                *total += delay;
603                *count += 1;
604                *max = (*max).max(delay);
605            }
606
607            // Create a local bucket view to avoid splitting buckets unnecessarily. Since we produce
608            // each bucket separately, we only need to split buckets that exceed the size, but not
609            // batches.
610            for view in BucketsView::new(std::slice::from_ref(&bucket))
611                .by_size(batch_size)
612                .flatten()
613            {
614                let message =
615                    self.create_metric_message(&scoping, &mut encoder, namespace, &view, retention);
616
617                let result =
618                    message.and_then(|message| self.send_metric_message(namespace, message));
619
620                let outcome = match result {
621                    Ok(()) => Outcome::Accepted,
622                    Err(e) => {
623                        error.get_or_insert(e);
624                        Outcome::Invalid(DiscardReason::Internal)
625                    }
626                };
627
628                self.metric_outcomes.track(scoping, &[view], outcome);
629            }
630
631            if emit_sessions_to_eap
632                && let Some(trace_item) = sessions::to_trace_item(scoping, bucket, retention)
633            {
634                let message = KafkaMessage::for_item(scoping, trace_item);
635                let res = self.produce(KafkaTopic::Items, message);
636                if let Err(error) = res {
637                    relay_log::error!(
638                        error = &error as &dyn std::error::Error,
639                        "failed to produce session metrics to EAP"
640                    )
641                }
642            }
643        }
644
645        if let Some(error) = error {
646            relay_log::error!(
647                error = &error as &dyn std::error::Error,
648                "failed to produce metric buckets: {error}"
649            );
650        }
651
652        for (namespace, (total, count, max)) in delay_stats {
653            if count == 0 {
654                continue;
655            }
656            metric!(
657                counter(RelayCounters::MetricDelaySum) += total,
658                namespace = namespace.as_str()
659            );
660            metric!(
661                counter(RelayCounters::MetricDelayCount) += count,
662                namespace = namespace.as_str()
663            );
664            metric!(
665                gauge(RelayGauges::MetricDelayMax) = max,
666                namespace = namespace.as_str()
667            );
668        }
669    }
670
671    fn handle_store_trace_item(
672        &self,
673        message: Managed<StoreTraceItem>,
674    ) -> Result<(), Rejected<StoreError>> {
675        let scoping = message.scoping();
676
677        message.try_accept(|item| {
678            let message = KafkaMessage::for_item(scoping, item.trace_item);
679            self.produce(KafkaTopic::Items, message)
680        })?;
681
682        Ok(())
683    }
684
685    fn handle_store_span(
686        &self,
687        message: Managed<Box<StoreSpanV2>>,
688    ) -> Result<(), Rejected<StoreError>> {
689        let scoping = message.scoping();
690        let received_at = message.received_at();
691
692        let meta = SpanMeta {
693            organization_id: scoping.organization_id,
694            project_id: scoping.project_id,
695            key_id: scoping.key_id,
696            event_id: message.event_id,
697            retention_days: message.retention_days,
698            downsampled_retention_days: message.downsampled_retention_days,
699            received: datetime_to_timestamp(received_at),
700            performance_issues_spans: message.performance_issues_spans,
701        };
702
703        message.try_accept(|span| {
704            // Temporary validation of the segment ID.
705            if let Some(segment_id) = span
706                .item
707                .attributes
708                .value()
709                .and_then(|a| a.get_value(SENTRY__SEGMENT__ID))
710                .and_then(|v| v.as_str())
711                && let Err(e) = segment_id.parse::<SpanId>()
712            {
713                relay_log::configure_scope(|scope| {
714                    scope.set_tag("sentry_project_id", scoping.project_id);
715                    if let Ok(mut last_report) = self.last_span_report.try_lock() {
716                        let now = Instant::now();
717                        if (now.saturating_duration_since(*last_report)) > Duration::from_secs(1) {
718                            *last_report = now;
719                            if let Ok(json) = Annotated::new(span.item).to_json() {
720                                scope.add_attachment(Attachment {
721                                    buffer: json.into_bytes(),
722                                    filename: "span.json".to_owned(),
723                                    content_type: Some("application/json".to_owned()),
724                                    ty: None,
725                                });
726                            }
727                        }
728                    }
729                });
730                return Err(e.into());
731            }
732            let item = Annotated::new(span.item);
733            let message = KafkaMessage::SpanV2 {
734                routing_key: span.routing_key,
735                headers: BTreeMap::from([(
736                    "project_id".to_owned(),
737                    scoping.project_id.to_string(),
738                )]),
739                message: SpanKafkaMessage {
740                    meta,
741                    span: SerializableAnnotated(&item),
742                },
743                org_id: scoping.organization_id,
744            };
745
746            self.produce(KafkaTopic::Spans, message)
747        })?;
748
749        relay_statsd::metric!(
750            counter(RelayCounters::SpanV2Produced) += 1,
751            via = "processing"
752        );
753
754        Ok(())
755    }
756
757    fn handle_store_profile_chunk(
758        &self,
759        message: Managed<StoreProfileChunk>,
760    ) -> Result<(), Rejected<StoreError>> {
761        let scoping = message.scoping();
762        let received_at = message.received_at();
763
764        message.try_accept(|message| {
765            let message = ProfileChunkKafkaMessage {
766                organization_id: scoping.organization_id,
767                project_id: scoping.project_id,
768                received: safe_timestamp(received_at),
769                retention_days: message.retention_days,
770                headers: BTreeMap::from([(
771                    "project_id".to_owned(),
772                    scoping.project_id.to_string(),
773                )]),
774                payload: message.payload,
775                attachments: message
776                    .attachments
777                    .into_iter()
778                    .map(|attachment| ProfileChunkKafkaAttachment {
779                        name: attachment.name,
780                        content_type: attachment.content_type.as_str(),
781                        stored_id: attachment.stored_id.into_inner(),
782                    })
783                    .collect(),
784            };
785
786            self.produce(KafkaTopic::Profiles, KafkaMessage::ProfileChunk(message))
787        })
788    }
789
790    fn handle_store_replay(
791        &self,
792        message: Managed<StoreReplay>,
793    ) -> Result<(), Rejected<StoreError>> {
794        let scoping = message.scoping();
795        let received_at = message.received_at();
796
797        message.try_accept(|replay| {
798            let kafka_msg =
799                KafkaMessage::ReplayRecordingNotChunked(ReplayRecordingNotChunkedKafkaMessage {
800                    replay_id: replay.event_id,
801                    key_id: scoping.key_id,
802                    org_id: scoping.organization_id,
803                    project_id: scoping.project_id,
804                    received: safe_timestamp(received_at),
805                    retention_days: replay.retention_days,
806                    payload: &replay.recording,
807                    replay_event: replay.event.as_deref(),
808                    replay_video: replay.video.as_deref(),
809                    // Hardcoded to `true` to indicate to the consumer that it should always publish the
810                    // replay_event as relay no longer does it.
811                    relay_snuba_publish_disabled: true,
812                });
813            self.produce(KafkaTopic::ReplayRecordings, kafka_msg)
814        })
815    }
816
817    fn handle_store_attachment(
818        &self,
819        message: Managed<StoreAttachment>,
820    ) -> Result<(), Rejected<StoreError>> {
821        let scoping = message.scoping();
822        message.try_accept(|attachment| {
823            let result = self.produce_attachment(
824                attachment.event_id,
825                scoping.project_id,
826                scoping.organization_id,
827                &attachment.attachment,
828                // Hardcoded to `true` since standalone attachments are 'individual attachments'.
829                true,
830                attachment.retention,
831            );
832            // Since we are sending an 'individual attachment' this function should never return a
833            // `ChunkedAttachment`.
834            debug_assert!(!matches!(result, Ok(Some(_))));
835            result.map(|_| ())
836        })
837    }
838
839    fn handle_user_report(
840        &self,
841        message: Managed<StoreUserReport>,
842    ) -> Result<(), Rejected<StoreError>> {
843        let scoping = message.scoping();
844        let received_at = message.received_at();
845
846        message.try_accept(|report| {
847            let kafka_msg = KafkaMessage::UserReport(UserReportKafkaMessage {
848                project_id: scoping.project_id,
849                event_id: report.event_id,
850                start_time: safe_timestamp(received_at),
851                payload: report.report.payload(),
852                org_id: scoping.organization_id,
853            });
854            self.produce(KafkaTopic::Attachments, kafka_msg)
855        })
856    }
857
858    fn handle_profile(&self, message: Managed<StoreProfile>) -> Result<(), Rejected<StoreError>> {
859        let scoping = message.scoping();
860        let received_at = message.received_at();
861
862        message.try_accept(|profile| {
863            self.produce_profile(
864                scoping.organization_id,
865                scoping.project_id,
866                scoping.key_id,
867                received_at,
868                profile.retention_days,
869                &profile.profile,
870            )
871        })
872    }
873
874    fn handle_check_in(&self, message: Managed<StoreCheckIn>) -> Result<(), Rejected<StoreError>> {
875        let scoping = message.scoping();
876        let received_at = message.received_at();
877
878        message.try_accept(|check_in| {
879            let message = KafkaMessage::CheckIn(CheckInKafkaMessage {
880                message_type: CheckInMessageType::CheckIn,
881                project_id: scoping.project_id,
882                org_id: scoping.organization_id,
883                retention_days: check_in.retention_days,
884                start_time: safe_timestamp(received_at),
885                sdk: check_in.sdk,
886                payload: check_in.check_in.payload(),
887                routing_key_hint: check_in.check_in.routing_hint(),
888            });
889
890            self.produce(KafkaTopic::Monitors, message)
891        })
892    }
893
894    fn create_metric_message<'a>(
895        &self,
896        scoping: &Scoping,
897        encoder: &'a mut BucketEncoder,
898        namespace: MetricNamespace,
899        view: &BucketView<'a>,
900        retention_days: u16,
901    ) -> Result<MetricKafkaMessage<'a>, StoreError> {
902        let value = match view.value() {
903            BucketViewValue::Counter(c) => MetricValue::Counter(c),
904            BucketViewValue::Distribution(data) => MetricValue::Distribution(
905                encoder
906                    .encode_distribution(namespace, data)
907                    .map_err(StoreError::EncodingFailed)?,
908            ),
909            BucketViewValue::Set(data) => MetricValue::Set(
910                encoder
911                    .encode_set(namespace, data)
912                    .map_err(StoreError::EncodingFailed)?,
913            ),
914            BucketViewValue::Gauge(g) => MetricValue::Gauge(g),
915        };
916
917        Ok(MetricKafkaMessage {
918            org_id: scoping.organization_id,
919            project_id: scoping.project_id,
920            key_id: scoping.key_id,
921            name: view.name(),
922            value,
923            timestamp: view.timestamp(),
924            tags: view.tags(),
925            retention_days,
926            received_at: view.metadata().received_at,
927        })
928    }
929
930    fn produce(
931        &self,
932        topic: KafkaTopic,
933        // Takes message by value to ensure it is not being produced twice.
934        message: KafkaMessage,
935    ) -> Result<(), StoreError> {
936        relay_log::trace!(
937            "Sending kafka message of type {} to {topic:?}",
938            message.variant()
939        );
940
941        let topic_name = self.producer.client.send_message(topic, &message)?;
942
943        match &message {
944            KafkaMessage::Metric {
945                message: metric, ..
946            } => {
947                metric!(
948                    counter(RelayCounters::ProcessingMessageEnqueued) += 1,
949                    event_type = message.variant(),
950                    topic = topic_name,
951                    metric_type = metric.value.variant(),
952                    metric_encoding = metric.value.encoding().unwrap_or(""),
953                );
954            }
955            KafkaMessage::ReplayRecordingNotChunked(replay) => {
956                let has_video = replay.replay_video.is_some();
957
958                metric!(
959                    counter(RelayCounters::ProcessingMessageEnqueued) += 1,
960                    event_type = message.variant(),
961                    topic = topic_name,
962                    has_video = bool_to_str(has_video),
963                );
964            }
965            message => {
966                metric!(
967                    counter(RelayCounters::ProcessingMessageEnqueued) += 1,
968                    event_type = message.variant(),
969                    topic = topic_name,
970                );
971            }
972        }
973
974        Ok(())
975    }
976
977    fn chunked_attachment_from_placeholder(
978        &self,
979        item: &Item,
980        retention_days: u16,
981    ) -> Result<ChunkedAttachment, StoreError> {
982        debug_assert!(
983            item.stored_key().is_none(),
984            "AttachmentRef should not have been uploaded to objectstore"
985        );
986
987        let payload = item.payload();
988        let placeholder: AttachmentPlaceholder<'_> =
989            serde_json::from_slice(&payload).map_err(|_| StoreError::InvalidAttachmentRef)?;
990        let config = self.config.current();
991        let location = SignedLocation::<Final>::try_from_str(placeholder.location)
992            .ok_or(StoreError::InvalidAttachmentRef)?
993            .verify(Utc::now(), &config)
994            .map_err(|_| StoreError::InvalidAttachmentRef)?;
995
996        let store_key = location.key;
997
998        Ok(ChunkedAttachment {
999            id: Uuid::new_v4().to_string(),
1000            name: item.filename().unwrap_or(UNNAMED_ATTACHMENT).to_owned(),
1001            rate_limited: item.rate_limited(),
1002            content_type: placeholder.content_type,
1003            attachment_type: item.attachment_type().unwrap_or_default(),
1004            size: item.attachment_body_size(),
1005            retention_days,
1006            payload: AttachmentPayload::Stored(store_key),
1007        })
1008    }
1009
1010    fn chunked_attachment_from_attachment(
1011        &self,
1012        event_id: EventId,
1013        project_id: ProjectId,
1014        org_id: OrganizationId,
1015        item: &Item,
1016        send_individual_attachments: bool,
1017        retention_days: u16,
1018    ) -> Result<ChunkedAttachment, StoreError> {
1019        let id = Uuid::new_v4().to_string();
1020
1021        let payload = item.payload();
1022        let size = item.len();
1023        let max_chunk_size = self.config.current().attachment_chunk_size();
1024
1025        let payload = if size == 0 {
1026            AttachmentPayload::Chunked(0)
1027        } else if let Some(stored_key) = item.stored_key() {
1028            AttachmentPayload::Stored(stored_key.into())
1029        } else if send_individual_attachments && size < max_chunk_size {
1030            // When sending individual attachments, and we have a single chunk, we want to send the
1031            // `data` inline in the `attachment` message.
1032            // This avoids a needless roundtrip through the attachments cache on the Sentry side.
1033            AttachmentPayload::Inline(payload)
1034        } else {
1035            let mut chunk_index = 0;
1036            let mut offset = 0;
1037            // This skips chunks for empty attachments. The consumer does not require chunks for
1038            // empty attachments. `chunks` will be `0` in this case.
1039            while offset < size {
1040                let chunk_size = std::cmp::min(max_chunk_size, size - offset);
1041                let chunk_message = AttachmentChunkKafkaMessage {
1042                    payload: payload.slice(offset..offset + chunk_size),
1043                    event_id,
1044                    project_id,
1045                    id: id.clone(),
1046                    chunk_index,
1047                    org_id,
1048                };
1049
1050                self.produce(
1051                    KafkaTopic::Attachments,
1052                    KafkaMessage::AttachmentChunk(chunk_message),
1053                )?;
1054                offset += chunk_size;
1055                chunk_index += 1;
1056            }
1057
1058            // The chunk_index is incremented after every loop iteration. After we exit the loop, it
1059            // is one larger than the last chunk, so it is equal to the number of chunks.
1060            AttachmentPayload::Chunked(chunk_index)
1061        };
1062
1063        Ok(ChunkedAttachment {
1064            id,
1065            name: match item.filename() {
1066                Some(name) => name.to_owned(),
1067                None => UNNAMED_ATTACHMENT.to_owned(),
1068            },
1069            rate_limited: item.rate_limited(),
1070            content_type: item.raw_content_type().map(|s| s.to_ascii_lowercase()),
1071            attachment_type: item.attachment_type().unwrap_or_default(),
1072            size,
1073            retention_days,
1074            payload,
1075        })
1076    }
1077
1078    /// Produces Kafka messages for the content and metadata of an attachment item.
1079    ///
1080    /// The `send_individual_attachments` controls whether the metadata of an attachment
1081    /// is produced directly as an individual `attachment` message, or returned from this function
1082    /// to be later sent as part of an `event` message.
1083    ///
1084    /// Attachment contents are chunked and sent as multiple `attachment_chunk` messages,
1085    /// unless the `send_individual_attachments` flag is set, and the content is small enough
1086    /// to fit inside a message.
1087    /// In that case, no `attachment_chunk` is produced, but the content is sent as part
1088    /// of the `attachment` message instead.
1089    fn produce_attachment(
1090        &self,
1091        event_id: EventId,
1092        project_id: ProjectId,
1093        org_id: OrganizationId,
1094        item: &Item,
1095        send_individual_attachments: bool,
1096        retention_days: u16,
1097    ) -> Result<Option<ChunkedAttachment>, StoreError> {
1098        let attachment = if item.is_attachment_ref() {
1099            self.chunked_attachment_from_placeholder(item, retention_days)
1100        } else {
1101            self.chunked_attachment_from_attachment(
1102                event_id,
1103                project_id,
1104                org_id,
1105                item,
1106                send_individual_attachments,
1107                retention_days,
1108            )
1109        }?;
1110
1111        if send_individual_attachments {
1112            let message = KafkaMessage::Attachment(AttachmentKafkaMessage {
1113                event_id,
1114                project_id,
1115                attachment,
1116                org_id,
1117            });
1118            self.produce(KafkaTopic::Attachments, message)?;
1119            Ok(None)
1120        } else {
1121            Ok(Some(attachment))
1122        }
1123    }
1124
1125    fn produce_user_report(
1126        &self,
1127        event_id: EventId,
1128        project_id: ProjectId,
1129        org_id: OrganizationId,
1130        received_at: DateTime<Utc>,
1131        item: &Item,
1132    ) -> Result<(), StoreError> {
1133        let message = KafkaMessage::UserReport(UserReportKafkaMessage {
1134            project_id,
1135            event_id,
1136            start_time: safe_timestamp(received_at),
1137            payload: item.payload(),
1138            org_id,
1139        });
1140
1141        self.produce(KafkaTopic::Attachments, message)
1142    }
1143
1144    fn send_metric_message(
1145        &self,
1146        namespace: MetricNamespace,
1147        message: MetricKafkaMessage,
1148    ) -> Result<(), StoreError> {
1149        let topic = match namespace {
1150            MetricNamespace::Sessions => KafkaTopic::MetricsSessions,
1151            MetricNamespace::Outcomes => {
1152                return self.send_metric_based_outcome(message);
1153            }
1154            MetricNamespace::Spans | MetricNamespace::Transactions => {
1155                // Generic metrics (spans/transactions) are no longer ingested.
1156                return Ok(());
1157            }
1158            MetricNamespace::Unsupported => {
1159                relay_log::error!(
1160                    metric_message.name = message.name.as_ref(),
1161                    "store service dropping unknown metric usecase"
1162                );
1163                return Ok(());
1164            }
1165        };
1166
1167        let headers = BTreeMap::from([("namespace".to_owned(), namespace.to_string())]);
1168        self.produce(topic, KafkaMessage::Metric { headers, message })?;
1169        Ok(())
1170    }
1171
1172    fn send_metric_based_outcome(&self, message: MetricKafkaMessage) -> Result<(), StoreError> {
1173        let Some(outcome) = outcome::metric::to_outcome_id(message.name) else {
1174            relay_log::error!(
1175                mri = message.name.as_ref(),
1176                "invalid outcome metric, cannot infer outcome id from metric name"
1177            );
1178            return Ok(());
1179        };
1180        let quantity = match message.value {
1181            MetricValue::Counter(c) => c.to_f64() as _,
1182            v => {
1183                relay_log::error!(
1184                    mri = message.name.as_ref(),
1185                    "invalid outcome metric, expected a counter got '{}'",
1186                    v.variant()
1187                );
1188                return Ok(());
1189            }
1190        };
1191
1192        let outcome = OutcomeMessage {
1193            timestamp: message
1194                .timestamp
1195                .as_datetime()
1196                .unwrap_or_else(Utc::now)
1197                .to_rfc3339_opts(SecondsFormat::Micros, true),
1198            org_id: Some(message.org_id).filter(|id| id.value() != 0),
1199            project_id: message.project_id,
1200            key_id: message.key_id,
1201            outcome,
1202            reason: message.tags.get("reason").map(|s| s.as_str()),
1203            event_id: message.tags.get("event_id").map(|s| s.as_str()),
1204            remote_addr: message.tags.get("remote_addr").map(|s| s.as_str()),
1205            source: message.tags.get("source").map(|s| s.as_str()),
1206            category: message.tags.get("category").and_then(|s| s.parse().ok()),
1207            quantity: Some(quantity),
1208        };
1209
1210        let topic = match outcome.outcome.is_billing() {
1211            true => KafkaTopic::OutcomesBilling,
1212            false => KafkaTopic::Outcomes,
1213        };
1214
1215        self.produce(topic, KafkaMessage::Outcome(outcome))
1216    }
1217
1218    fn produce_profile(
1219        &self,
1220        organization_id: OrganizationId,
1221        project_id: ProjectId,
1222        key_id: Option<u64>,
1223        received_at: DateTime<Utc>,
1224        retention_days: u16,
1225        item: &Item,
1226    ) -> Result<(), StoreError> {
1227        let message = ProfileKafkaMessage {
1228            organization_id,
1229            project_id,
1230            key_id,
1231            received: safe_timestamp(received_at),
1232            retention_days,
1233            headers: BTreeMap::from([
1234                (
1235                    "sampled".to_owned(),
1236                    if item.sampled() { "true" } else { "false" }.to_owned(),
1237                ),
1238                ("project_id".to_owned(), project_id.to_string()),
1239            ]),
1240            payload: item.payload(),
1241        };
1242        self.produce(KafkaTopic::Profiles, KafkaMessage::Profile(message))?;
1243        Ok(())
1244    }
1245}
1246
1247impl Service for StoreService {
1248    type Interface = Store;
1249
1250    async fn run(self, mut rx: relay_system::Receiver<Self::Interface>) {
1251        let this = Arc::new(self);
1252
1253        relay_log::info!("store forwarder started");
1254
1255        while let Some(message) = rx.recv().await {
1256            let task = StoreTask {
1257                service: Arc::clone(&this),
1258                message: Some(message),
1259            };
1260            this.pool.spawn_async(task).await;
1261        }
1262
1263        relay_log::info!("store forwarder stopped");
1264    }
1265}
1266
1267/// Task executed by [`StoreService`].
1268pub struct StoreTask {
1269    service: Arc<StoreService>,
1270    message: Option<Store>,
1271}
1272
1273impl Future for StoreTask {
1274    type Output = ();
1275
1276    fn poll(mut self: Pin<&mut Self>, _: &mut task::Context<'_>) -> task::Poll<Self::Output> {
1277        let message = self
1278            .message
1279            .take()
1280            .expect("StoreTask polled after completion");
1281        let () = relay_log::with_scope(|_| {}, || self.service.handle_message(message));
1282        task::Poll::Ready(())
1283    }
1284}
1285
1286/// This signifies how the attachment payload is being transfered.
1287#[derive(Debug, Serialize)]
1288enum AttachmentPayload {
1289    /// The payload has been split into multiple chunks.
1290    ///
1291    /// The individual chunks are being sent as separate [`AttachmentChunkKafkaMessage`] messages.
1292    /// If the payload `size == 0`, the number of chunks will also be `0`.
1293    #[serde(rename = "chunks")]
1294    Chunked(usize),
1295
1296    /// The payload is inlined here directly, and thus into the [`ChunkedAttachment`].
1297    #[serde(rename = "data")]
1298    Inline(Bytes),
1299
1300    /// The attachment has already been stored into the objectstore, with the given Id.
1301    #[serde(rename = "stored_id")]
1302    Stored(String),
1303}
1304
1305/// Common attributes for both standalone attachments and processing-relevant attachments.
1306#[derive(Debug, Serialize)]
1307struct ChunkedAttachment {
1308    /// The attachment ID within the event.
1309    ///
1310    /// The triple `(project_id, event_id, id)` identifies an attachment uniquely.
1311    id: String,
1312
1313    /// File name of the attachment file.
1314    name: String,
1315
1316    /// Whether this attachment was rate limited and should be removed after processing.
1317    ///
1318    /// By default, rate limited attachments are immediately removed from Envelopes. For processing,
1319    /// native crash reports still need to be retained. These attachments are marked with the
1320    /// `rate_limited` header, which signals to the processing pipeline that the attachment should
1321    /// not be persisted after processing.
1322    rate_limited: bool,
1323
1324    /// Content type of the attachment payload.
1325    #[serde(skip_serializing_if = "Option::is_none")]
1326    content_type: Option<String>,
1327
1328    /// The Sentry-internal attachment type used in the processing pipeline.
1329    #[serde(serialize_with = "serialize_attachment_type")]
1330    attachment_type: AttachmentType,
1331
1332    /// The size of the attachment in bytes.
1333    size: usize,
1334
1335    /// The retention in days for this attachment.
1336    retention_days: u16,
1337
1338    /// The attachment payload, chunked, inlined, or already stored.
1339    #[serde(flatten)]
1340    payload: AttachmentPayload,
1341}
1342
1343/// A hack to make rmp-serde behave more like serde-json when serializing enums.
1344///
1345/// Cannot serialize bytes.
1346///
1347/// See <https://github.com/3Hren/msgpack-rust/pull/214>
1348fn serialize_attachment_type<S, T>(t: &T, serializer: S) -> Result<S::Ok, S::Error>
1349where
1350    S: serde::Serializer,
1351    T: serde::Serialize,
1352{
1353    serde_json::to_value(t)
1354        .map_err(|e| serde::ser::Error::custom(e.to_string()))?
1355        .serialize(serializer)
1356}
1357
1358/// Container payload for event messages.
1359#[derive(Debug, Serialize)]
1360struct EventKafkaMessage {
1361    /// Raw event payload.
1362    payload: Bytes,
1363    /// Time at which the event was received by Relay.
1364    start_time: u64,
1365    /// The event id.
1366    event_id: EventId,
1367    /// The project id for the current event.
1368    project_id: ProjectId,
1369    /// The client ip address.
1370    remote_addr: Option<String>,
1371    /// Attachments that are potentially relevant for processing.
1372    attachments: Vec<ChunkedAttachment>,
1373
1374    /// Used for [`KafkaMessage::key`]
1375    #[serde(skip)]
1376    org_id: OrganizationId,
1377}
1378
1379/// Container payload for chunks of attachments.
1380#[derive(Debug, Serialize)]
1381struct AttachmentChunkKafkaMessage {
1382    /// Chunk payload of the attachment.
1383    payload: Bytes,
1384    /// The event id.
1385    event_id: EventId,
1386    /// The project id for the current event.
1387    project_id: ProjectId,
1388    /// The attachment ID within the event.
1389    ///
1390    /// The triple `(project_id, event_id, id)` identifies an attachment uniquely.
1391    id: String,
1392    /// Sequence number of chunk. Starts at 0 and ends at `AttachmentKafkaMessage.num_chunks - 1`.
1393    chunk_index: usize,
1394
1395    /// Used for [`KafkaMessage::key`]
1396    #[serde(skip)]
1397    org_id: OrganizationId,
1398}
1399
1400/// A "standalone" attachment.
1401///
1402/// Still belongs to an event but can be sent independently (like UserReport) and is not
1403/// considered in processing.
1404#[derive(Debug, Serialize)]
1405struct AttachmentKafkaMessage {
1406    /// The event id.
1407    event_id: EventId,
1408    /// The project id for the current event.
1409    project_id: ProjectId,
1410    /// The attachment.
1411    attachment: ChunkedAttachment,
1412
1413    /// Used for [`KafkaMessage::key`]
1414    #[serde(skip)]
1415    org_id: OrganizationId,
1416}
1417
1418#[derive(Debug, Serialize)]
1419struct ReplayRecordingNotChunkedKafkaMessage<'a> {
1420    replay_id: EventId,
1421    key_id: Option<u64>,
1422    org_id: OrganizationId,
1423    project_id: ProjectId,
1424    received: u64,
1425    retention_days: u16,
1426    #[serde(with = "serde_bytes")]
1427    payload: &'a [u8],
1428    #[serde(with = "serde_bytes")]
1429    replay_event: Option<&'a [u8]>,
1430    #[serde(with = "serde_bytes")]
1431    replay_video: Option<&'a [u8]>,
1432    relay_snuba_publish_disabled: bool,
1433}
1434
1435/// User report for an event wrapped up in a message ready for consumption in Kafka.
1436///
1437/// Is always independent of an event and can be sent as part of any envelope.
1438#[derive(Debug, Serialize)]
1439struct UserReportKafkaMessage {
1440    /// The project id for the current event.
1441    project_id: ProjectId,
1442    start_time: u64,
1443    payload: Bytes,
1444
1445    /// Used for [`KafkaMessage::key`]
1446    #[serde(skip)]
1447    event_id: EventId,
1448    /// Used for [`KafkaMessage::key`]
1449    #[serde(skip)]
1450    org_id: OrganizationId,
1451}
1452
1453#[derive(Clone, Debug, Serialize)]
1454struct MetricKafkaMessage<'a> {
1455    org_id: OrganizationId,
1456    project_id: ProjectId,
1457    #[serde(skip)]
1458    key_id: Option<u64>,
1459    name: &'a MetricName,
1460    #[serde(flatten)]
1461    value: MetricValue<'a>,
1462    timestamp: UnixTimestamp,
1463    tags: &'a BTreeMap<String, String>,
1464    retention_days: u16,
1465    #[serde(skip_serializing_if = "Option::is_none")]
1466    received_at: Option<UnixTimestamp>,
1467}
1468
1469#[derive(Clone, Debug, Serialize)]
1470#[serde(tag = "type", content = "value")]
1471enum MetricValue<'a> {
1472    #[serde(rename = "c")]
1473    Counter(FiniteF64),
1474    #[serde(rename = "d")]
1475    Distribution(ArrayEncoding<'a, &'a [FiniteF64]>),
1476    #[serde(rename = "s")]
1477    Set(ArrayEncoding<'a, SetView<'a>>),
1478    #[serde(rename = "g")]
1479    Gauge(GaugeValue),
1480}
1481
1482impl MetricValue<'_> {
1483    fn variant(&self) -> &'static str {
1484        match self {
1485            Self::Counter(_) => "counter",
1486            Self::Distribution(_) => "distribution",
1487            Self::Set(_) => "set",
1488            Self::Gauge(_) => "gauge",
1489        }
1490    }
1491
1492    fn encoding(&self) -> Option<&'static str> {
1493        match self {
1494            Self::Distribution(ae) => Some(ae.name()),
1495            Self::Set(ae) => Some(ae.name()),
1496            _ => None,
1497        }
1498    }
1499}
1500
1501/// Raw representation of an outcome for Kafka.
1502#[derive(Debug, Serialize, Clone)]
1503pub struct OutcomeMessage<'a> {
1504    /// The timespan of the event outcome.
1505    timestamp: String,
1506    /// Organization id.
1507    #[serde(skip_serializing_if = "Option::is_none")]
1508    org_id: Option<OrganizationId>,
1509    /// Project id.
1510    project_id: ProjectId,
1511    /// The DSN project key id.
1512    #[serde(skip_serializing_if = "Option::is_none")]
1513    key_id: Option<u64>,
1514    /// The outcome.
1515    outcome: OutcomeId,
1516    /// Reason for the outcome.
1517    #[serde(skip_serializing_if = "Option::is_none")]
1518    reason: Option<&'a str>,
1519    /// The event id.
1520    #[serde(skip_serializing_if = "Option::is_none")]
1521    event_id: Option<&'a str>,
1522    /// The client ip address.
1523    #[serde(skip_serializing_if = "Option::is_none")]
1524    remote_addr: Option<&'a str>,
1525    /// The source of the outcome (which Relay sent it)
1526    #[serde(skip_serializing_if = "Option::is_none")]
1527    source: Option<&'a str>,
1528    /// The event's data category.
1529    #[serde(skip_serializing_if = "Option::is_none")]
1530    category: Option<u8>,
1531    /// The number of events or total attachment size in bytes.
1532    #[serde(skip_serializing_if = "Option::is_none")]
1533    quantity: Option<u64>,
1534}
1535
1536#[derive(Clone, Debug, Serialize)]
1537struct ProfileKafkaMessage {
1538    organization_id: OrganizationId,
1539    project_id: ProjectId,
1540    key_id: Option<u64>,
1541    received: u64,
1542    retention_days: u16,
1543    #[serde(skip)]
1544    headers: BTreeMap<String, String>,
1545    payload: Bytes,
1546}
1547
1548/// Used to discriminate cron monitor ingestion messages.
1549///
1550/// There are two types of messages that end up in the ingest-monitors kafka topic, "check_in" (the
1551/// ones produced here in relay) and "clock_pulse" messages, which are produced externally and are
1552/// intended to ensure the clock continues to run even when ingestion volume drops.
1553#[allow(dead_code)]
1554#[derive(Debug, Serialize)]
1555#[serde(rename_all = "snake_case")]
1556enum CheckInMessageType {
1557    ClockPulse,
1558    CheckIn,
1559}
1560
1561#[derive(Debug, Serialize)]
1562struct CheckInKafkaMessage {
1563    /// Used by the consumer to discrinminate the message.
1564    message_type: CheckInMessageType,
1565    /// Raw event payload.
1566    payload: Bytes,
1567    /// Time at which the event was received by Relay.
1568    start_time: u64,
1569    /// The SDK client which produced the event.
1570    sdk: Option<String>,
1571    /// The project id for the current event.
1572    project_id: ProjectId,
1573    /// Number of days to retain.
1574    retention_days: u16,
1575
1576    /// Used for [`KafkaMessage::key`]
1577    #[serde(skip)]
1578    routing_key_hint: Option<Uuid>,
1579    /// Used for [`KafkaMessage::key`]
1580    #[serde(skip)]
1581    org_id: OrganizationId,
1582}
1583
1584#[derive(Debug, Serialize)]
1585struct SpanKafkaMessage<'a> {
1586    #[serde(flatten)]
1587    meta: SpanMeta,
1588    #[serde(flatten)]
1589    span: SerializableAnnotated<'a, SpanV2>,
1590}
1591
1592#[derive(Debug, Serialize)]
1593struct SpanMeta {
1594    organization_id: OrganizationId,
1595    project_id: ProjectId,
1596    // Required for the buffer to emit outcomes scoped to the DSN.
1597    #[serde(skip_serializing_if = "Option::is_none")]
1598    key_id: Option<u64>,
1599    #[serde(skip_serializing_if = "Option::is_none")]
1600    event_id: Option<EventId>,
1601    /// Time at which the event was received by Relay. Not to be confused with `start_timestamp_ms`.
1602    received: f64,
1603    /// Number of days until these data should be deleted.
1604    retention_days: u16,
1605    /// Number of days until the downsampled version of this data should be deleted.
1606    downsampled_retention_days: u16,
1607    /// Whether the segment span should be used for issue detection instead of the transaction.
1608    #[serde(rename = "_performance_issues_spans", skip_serializing_if = "is_false")]
1609    performance_issues_spans: bool,
1610}
1611
1612fn is_false(val: &bool) -> bool {
1613    !val
1614}
1615
1616#[derive(Clone, Debug, Serialize)]
1617struct ProfileChunkKafkaMessage {
1618    organization_id: OrganizationId,
1619    project_id: ProjectId,
1620    received: u64,
1621    retention_days: u16,
1622    #[serde(skip)]
1623    headers: BTreeMap<String, String>,
1624    payload: Bytes,
1625    #[serde(skip_serializing_if = "Vec::is_empty")]
1626    attachments: Vec<ProfileChunkKafkaAttachment>,
1627}
1628
1629#[derive(Clone, Debug, Serialize)]
1630struct ProfileChunkKafkaAttachment {
1631    name: String,
1632    content_type: &'static str,
1633    stored_id: String,
1634}
1635
1636/// An enum over all possible ingest messages.
1637#[derive(Debug, Serialize)]
1638#[serde(tag = "type", rename_all = "snake_case")]
1639#[allow(clippy::large_enum_variant)]
1640enum KafkaMessage<'a> {
1641    Event(EventKafkaMessage),
1642    UserReport(UserReportKafkaMessage),
1643    Metric {
1644        #[serde(skip)]
1645        headers: BTreeMap<String, String>,
1646        #[serde(flatten)]
1647        message: MetricKafkaMessage<'a>,
1648    },
1649    CheckIn(CheckInKafkaMessage),
1650    Item {
1651        #[serde(skip)]
1652        headers: BTreeMap<String, String>,
1653        #[serde(skip)]
1654        item_type: TraceItemType,
1655        #[serde(skip)]
1656        message: TraceItem,
1657    },
1658    SpanV2 {
1659        #[serde(skip)]
1660        routing_key: Option<Uuid>,
1661        #[serde(skip)]
1662        headers: BTreeMap<String, String>,
1663        #[serde(flatten)]
1664        message: SpanKafkaMessage<'a>,
1665
1666        /// Used for [`KafkaMessage::key`]
1667        #[serde(skip)]
1668        org_id: OrganizationId,
1669    },
1670
1671    Attachment(AttachmentKafkaMessage),
1672    AttachmentChunk(AttachmentChunkKafkaMessage),
1673
1674    Profile(ProfileKafkaMessage),
1675    ProfileChunk(ProfileChunkKafkaMessage),
1676
1677    ReplayRecordingNotChunked(ReplayRecordingNotChunkedKafkaMessage<'a>),
1678
1679    Outcome(OutcomeMessage<'a>),
1680}
1681
1682impl KafkaMessage<'_> {
1683    /// Creates a [`KafkaMessage`] for a [`TraceItem`].
1684    fn for_item(scoping: Scoping, item: TraceItem) -> KafkaMessage<'static> {
1685        let item_type = item.item_type();
1686        KafkaMessage::Item {
1687            headers: BTreeMap::from([
1688                ("project_id".to_owned(), scoping.project_id.to_string()),
1689                ("item_type".to_owned(), (item_type as i32).to_string()),
1690            ]),
1691            message: item,
1692            item_type,
1693        }
1694    }
1695}
1696
1697impl Message for KafkaMessage<'_> {
1698    fn variant(&self) -> &'static str {
1699        match self {
1700            KafkaMessage::Event(_) => "event",
1701            KafkaMessage::UserReport(_) => "user_report",
1702            KafkaMessage::Metric { message, .. } => match message.name.namespace() {
1703                MetricNamespace::Sessions => "metric_sessions",
1704                MetricNamespace::Spans => "metric_spans",
1705                MetricNamespace::Transactions => "metric_transactions",
1706                MetricNamespace::Outcomes => "metric_outcomes",
1707                MetricNamespace::Unsupported => "metric_unsupported",
1708            },
1709            KafkaMessage::CheckIn(_) => "check_in",
1710            KafkaMessage::SpanV2 { .. } => "span",
1711            KafkaMessage::Item { item_type, .. } => item_type.as_str_name(),
1712
1713            KafkaMessage::Attachment(_) => "attachment",
1714            KafkaMessage::AttachmentChunk(_) => "attachment_chunk",
1715
1716            KafkaMessage::Profile(_) => "profile",
1717            KafkaMessage::ProfileChunk(_) => "profile_chunk",
1718
1719            KafkaMessage::ReplayRecordingNotChunked(_) => "replay_recording_not_chunked",
1720
1721            KafkaMessage::Outcome(_) => "outcome",
1722        }
1723    }
1724
1725    /// Returns the partitioning key for this Kafka message determining.
1726    fn key(&self) -> Option<relay_kafka::Key> {
1727        match self {
1728            Self::Event(message) => Some((message.event_id.0, message.org_id)),
1729            Self::UserReport(message) => Some((message.event_id.0, message.org_id)),
1730            Self::SpanV2 {
1731                routing_key,
1732                org_id,
1733                ..
1734            } => routing_key.map(|r| (r, *org_id)),
1735
1736            // Monitor check-ins use the hinted UUID passed through from the Envelope.
1737            //
1738            // XXX(epurkhiser): In the future it would be better if all KafkaMessage's would
1739            // recieve the routing_key_hint form their envelopes.
1740            Self::CheckIn(message) => message.routing_key_hint.map(|r| (r, message.org_id)),
1741
1742            Self::Attachment(message) => Some((message.event_id.0, message.org_id)),
1743            Self::AttachmentChunk(message) => Some((message.event_id.0, message.org_id)),
1744
1745            // Random partitioning
1746            Self::Metric { .. }
1747            | Self::Item { .. }
1748            | Self::Profile(_)
1749            | Self::ProfileChunk(_)
1750            | Self::ReplayRecordingNotChunked(_)
1751            | Self::Outcome(_) => None,
1752        }
1753        .filter(|(uuid, _)| !uuid.is_nil())
1754        .map(|(uuid, org_id)| {
1755            // mix key with org id for better paritioning in case of incidents where messages
1756            // across orgs get the same id
1757            let mut res = uuid.into_bytes();
1758            for (i, &b) in org_id.value().to_be_bytes().iter().enumerate() {
1759                res[i] ^= b;
1760            }
1761            u128::from_be_bytes(res)
1762        })
1763    }
1764
1765    fn headers(&self) -> Option<&BTreeMap<String, String>> {
1766        match &self {
1767            KafkaMessage::Metric { headers, .. }
1768            | KafkaMessage::SpanV2 { headers, .. }
1769            | KafkaMessage::Item { headers, .. }
1770            | KafkaMessage::Profile(ProfileKafkaMessage { headers, .. })
1771            | KafkaMessage::ProfileChunk(ProfileChunkKafkaMessage { headers, .. }) => Some(headers),
1772
1773            KafkaMessage::Event(_)
1774            | KafkaMessage::UserReport(_)
1775            | KafkaMessage::CheckIn(_)
1776            | KafkaMessage::Attachment(_)
1777            | KafkaMessage::AttachmentChunk(_)
1778            | KafkaMessage::ReplayRecordingNotChunked(_)
1779            | KafkaMessage::Outcome(_) => None,
1780        }
1781    }
1782
1783    fn serialize(&self) -> Result<SerializationOutput<'_>, ClientError> {
1784        match self {
1785            KafkaMessage::Metric { message, .. } => serialize_as_json(message),
1786            KafkaMessage::SpanV2 { message, .. } => serialize_as_json(message),
1787            KafkaMessage::Item { message, .. } => {
1788                let mut payload = Vec::new();
1789                match message.encode(&mut payload) {
1790                    Ok(_) => Ok(SerializationOutput::Protobuf(Cow::Owned(payload))),
1791                    Err(_) => Err(ClientError::ProtobufEncodingFailed),
1792                }
1793            }
1794            KafkaMessage::Outcome(outcome) => serialize_as_json(outcome),
1795            KafkaMessage::Event(_)
1796            | KafkaMessage::UserReport(_)
1797            | KafkaMessage::CheckIn(_)
1798            | KafkaMessage::Attachment(_)
1799            | KafkaMessage::AttachmentChunk(_)
1800            | KafkaMessage::Profile(_)
1801            | KafkaMessage::ProfileChunk(_)
1802            | KafkaMessage::ReplayRecordingNotChunked(_) => match rmp_serde::to_vec_named(&self) {
1803                Ok(x) => Ok(SerializationOutput::MsgPack(Cow::Owned(x))),
1804                Err(err) => Err(ClientError::InvalidMsgPack(err)),
1805            },
1806        }
1807    }
1808}
1809
1810fn serialize_as_json<T: serde::Serialize>(
1811    value: &T,
1812) -> Result<SerializationOutput<'_>, ClientError> {
1813    match serde_json::to_vec(value) {
1814        Ok(vec) => Ok(SerializationOutput::Json(Cow::Owned(vec))),
1815        Err(err) => Err(ClientError::InvalidJson(err)),
1816    }
1817}
1818
1819fn bool_to_str(value: bool) -> &'static str {
1820    if value { "true" } else { "false" }
1821}
1822
1823/// Returns a safe timestamp for Kafka.
1824///
1825/// Kafka expects timestamps to be in UTC and in seconds since epoch.
1826fn safe_timestamp(timestamp: DateTime<Utc>) -> u64 {
1827    let ts = timestamp.timestamp();
1828    if ts >= 0 {
1829        return ts as u64;
1830    }
1831
1832    // We assume this call can't return < 0.
1833    Utc::now().timestamp() as u64
1834}