1use 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
54const 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#[derive(Debug)]
114pub struct StoreEvent {
115 pub event_category: DataCategory,
117 pub event: Annotated<Event>,
119 pub attachments: Vec<Item>,
121 pub user_reports: Vec<Item>,
123 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#[derive(Clone, Debug)]
138pub struct StoreMetrics {
139 pub buckets: Vec<Bucket>,
140 pub scoping: Scoping,
141 pub retention: u16,
142}
143
144#[derive(Debug)]
146pub struct StoreTraceItem {
147 pub trace_item: TraceItem,
149}
150
151impl Counted for StoreTraceItem {
152 fn quantities(&self) -> Quantities {
153 self.trace_item.quantities()
154 }
155}
156
157#[derive(Debug)]
159pub struct StoreSpanV2 {
160 pub routing_key: Option<Uuid>,
162 pub retention_days: u16,
164 pub downsampled_retention_days: u16,
166 pub event_id: Option<EventId>,
168 pub item: SpanV2,
170 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#[derive(Debug)]
183pub struct StoreProfileChunk {
184 pub retention_days: u16,
186 pub payload: Bytes,
188 pub attachments: Vec<ProfileAttachment>,
190 pub quantities: Quantities,
194}
195
196#[derive(Debug)]
198pub struct ProfileAttachment {
199 pub name: String,
201 pub content_type: ContentType,
203 pub stored_id: ObjectstoreKey,
207}
208
209impl Counted for StoreProfileChunk {
210 fn quantities(&self) -> Quantities {
211 self.quantities.clone()
212 }
213}
214
215#[derive(Debug)]
217pub struct StoreReplay {
218 pub event_id: EventId,
220 pub retention_days: u16,
222 pub recording: Bytes,
224 pub event: Option<Bytes>,
226 pub video: Option<Bytes>,
228 pub quantities: Quantities,
232}
233
234impl Counted for StoreReplay {
235 fn quantities(&self) -> Quantities {
236 self.quantities.clone()
237 }
238}
239
240#[derive(Debug)]
242pub struct StoreAttachment {
243 pub event_id: EventId,
245 pub attachment: Item,
247 pub quantities: Quantities,
249 pub retention: u16,
251}
252
253impl Counted for StoreAttachment {
254 fn quantities(&self) -> Quantities {
255 self.quantities.clone()
256 }
257}
258
259#[derive(Debug)]
261pub struct StoreUserReport {
262 pub event_id: EventId,
264 pub report: Item,
266}
267
268impl Counted for StoreUserReport {
269 fn quantities(&self) -> Quantities {
270 smallvec::smallvec![(DataCategory::UserReportV2, 1)]
271 }
272}
273
274#[derive(Debug)]
276pub struct StoreProfile {
277 pub retention_days: u16,
279 pub profile: Item,
281 pub quantities: Quantities,
283}
284
285impl Counted for StoreProfile {
286 fn quantities(&self) -> Quantities {
287 self.quantities.clone()
288 }
289}
290
291#[derive(Debug)]
293pub struct StoreCheckIn {
294 pub check_in: Item,
296 pub sdk: Option<String>,
298 pub retention_days: u16,
300}
301
302impl Counted for StoreCheckIn {
303 fn quantities(&self) -> Quantities {
304 self.check_in.quantities()
305 }
306}
307
308pub type StoreServicePool = AsyncPool<StoreTask>;
310
311#[derive(Debug)]
313pub enum Store {
314 Event(Managed<Box<StoreEvent>>),
316 Metrics(StoreMetrics),
318 TraceItem(Managed<StoreTraceItem>),
320 Span(Managed<Box<StoreSpanV2>>),
322 ProfileChunk(Managed<StoreProfileChunk>),
324 Replay(Managed<StoreReplay>),
326 Attachment(Managed<StoreAttachment>),
328 UserReport(Managed<StoreUserReport>),
330 Profile(Managed<StoreProfile>),
332 CheckIn(Managed<StoreCheckIn>),
334}
335
336impl Store {
337 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
436pub 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>, }
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 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 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 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 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 true,
830 attachment.retention,
831 );
832 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 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 AttachmentPayload::Inline(payload)
1034 } else {
1035 let mut chunk_index = 0;
1036 let mut offset = 0;
1037 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 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 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 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
1267pub 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#[derive(Debug, Serialize)]
1288enum AttachmentPayload {
1289 #[serde(rename = "chunks")]
1294 Chunked(usize),
1295
1296 #[serde(rename = "data")]
1298 Inline(Bytes),
1299
1300 #[serde(rename = "stored_id")]
1302 Stored(String),
1303}
1304
1305#[derive(Debug, Serialize)]
1307struct ChunkedAttachment {
1308 id: String,
1312
1313 name: String,
1315
1316 rate_limited: bool,
1323
1324 #[serde(skip_serializing_if = "Option::is_none")]
1326 content_type: Option<String>,
1327
1328 #[serde(serialize_with = "serialize_attachment_type")]
1330 attachment_type: AttachmentType,
1331
1332 size: usize,
1334
1335 retention_days: u16,
1337
1338 #[serde(flatten)]
1340 payload: AttachmentPayload,
1341}
1342
1343fn 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#[derive(Debug, Serialize)]
1360struct EventKafkaMessage {
1361 payload: Bytes,
1363 start_time: u64,
1365 event_id: EventId,
1367 project_id: ProjectId,
1369 remote_addr: Option<String>,
1371 attachments: Vec<ChunkedAttachment>,
1373
1374 #[serde(skip)]
1376 org_id: OrganizationId,
1377}
1378
1379#[derive(Debug, Serialize)]
1381struct AttachmentChunkKafkaMessage {
1382 payload: Bytes,
1384 event_id: EventId,
1386 project_id: ProjectId,
1388 id: String,
1392 chunk_index: usize,
1394
1395 #[serde(skip)]
1397 org_id: OrganizationId,
1398}
1399
1400#[derive(Debug, Serialize)]
1405struct AttachmentKafkaMessage {
1406 event_id: EventId,
1408 project_id: ProjectId,
1410 attachment: ChunkedAttachment,
1412
1413 #[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#[derive(Debug, Serialize)]
1439struct UserReportKafkaMessage {
1440 project_id: ProjectId,
1442 start_time: u64,
1443 payload: Bytes,
1444
1445 #[serde(skip)]
1447 event_id: EventId,
1448 #[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#[derive(Debug, Serialize, Clone)]
1503pub struct OutcomeMessage<'a> {
1504 timestamp: String,
1506 #[serde(skip_serializing_if = "Option::is_none")]
1508 org_id: Option<OrganizationId>,
1509 project_id: ProjectId,
1511 #[serde(skip_serializing_if = "Option::is_none")]
1513 key_id: Option<u64>,
1514 outcome: OutcomeId,
1516 #[serde(skip_serializing_if = "Option::is_none")]
1518 reason: Option<&'a str>,
1519 #[serde(skip_serializing_if = "Option::is_none")]
1521 event_id: Option<&'a str>,
1522 #[serde(skip_serializing_if = "Option::is_none")]
1524 remote_addr: Option<&'a str>,
1525 #[serde(skip_serializing_if = "Option::is_none")]
1527 source: Option<&'a str>,
1528 #[serde(skip_serializing_if = "Option::is_none")]
1530 category: Option<u8>,
1531 #[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#[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 message_type: CheckInMessageType,
1565 payload: Bytes,
1567 start_time: u64,
1569 sdk: Option<String>,
1571 project_id: ProjectId,
1573 retention_days: u16,
1575
1576 #[serde(skip)]
1578 routing_key_hint: Option<Uuid>,
1579 #[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 #[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 received: f64,
1603 retention_days: u16,
1605 downsampled_retention_days: u16,
1607 #[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#[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 #[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 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 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 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 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 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
1823fn safe_timestamp(timestamp: DateTime<Utc>) -> u64 {
1827 let ts = timestamp.timestamp();
1828 if ts >= 0 {
1829 return ts as u64;
1830 }
1831
1832 Utc::now().timestamp() as u64
1834}