1use std::collections::{BTreeMap, HashMap};
2use std::error::Error;
3use std::io::Write;
4use std::net::{IpAddr, SocketAddr};
5use std::num::{NonZeroU8, NonZeroU16};
6use std::path::{Path, PathBuf};
7use std::str::FromStr;
8use std::time::Duration;
9use std::{env, fmt, fs, io};
10
11use anyhow::Context;
12use relay_auth::{PublicKey, RelayId, SecretKey, generate_key_pair, generate_relay_id};
13use relay_common::Dsn;
14use relay_kafka::{
15 ConfigError as KafkaConfigError, KafkaConfigParam, KafkaTopic, KafkaTopicConfig,
16 TopicAssignments,
17};
18use relay_metrics::MetricNamespace;
19use serde::de::{DeserializeOwned, Unexpected, Visitor};
20use serde::{Deserialize, Deserializer, Serialize, Serializer};
21use uuid::Uuid;
22
23use crate::aggregator::{AggregatorServiceConfig, ScopedAggregatorConfig};
24use crate::byte_size::ByteSize;
25use crate::upstream::UpstreamDescriptor;
26use crate::{RedisConfig, RedisConfigs, RedisConfigsRef, build_redis_configs};
27
28const DEFAULT_NETWORK_OUTAGE_GRACE_PERIOD: u64 = 10;
29
30static CONFIG_YAML_HEADER: &str = r###"# Please see the relevant documentation.
31# Performance tuning: https://docs.sentry.io/product/relay/operating-guidelines/
32# All config options: https://docs.sentry.io/product/relay/options/
33"###;
34
35#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
37#[non_exhaustive]
38pub enum ConfigErrorKind {
39 CouldNotOpenFile,
41 CouldNotWriteFile,
43 BadYaml,
45 BadJson,
47 InvalidValue,
49 ProcessingNotAvailable,
52}
53
54impl fmt::Display for ConfigErrorKind {
55 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
56 match self {
57 Self::CouldNotOpenFile => write!(f, "could not open config file"),
58 Self::CouldNotWriteFile => write!(f, "could not write config file"),
59 Self::BadYaml => write!(f, "could not parse yaml config file"),
60 Self::BadJson => write!(f, "could not parse json config file"),
61 Self::InvalidValue => write!(f, "invalid config value"),
62 Self::ProcessingNotAvailable => write!(
63 f,
64 "was not compiled with processing, cannot enable processing"
65 ),
66 }
67 }
68}
69
70#[derive(Debug, Default)]
72enum ConfigErrorSource {
73 #[default]
75 None,
76 File(PathBuf),
78 FieldOverride(String),
80}
81
82impl fmt::Display for ConfigErrorSource {
83 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
84 match self {
85 ConfigErrorSource::None => Ok(()),
86 ConfigErrorSource::File(file_name) => {
87 write!(f, " (file {})", file_name.display())
88 }
89 ConfigErrorSource::FieldOverride(name) => write!(f, " (field {name})"),
90 }
91 }
92}
93
94#[derive(Debug)]
96pub struct ConfigError {
97 source: ConfigErrorSource,
98 kind: ConfigErrorKind,
99}
100
101impl ConfigError {
102 #[inline]
103 fn new(kind: ConfigErrorKind) -> Self {
104 Self {
105 source: ConfigErrorSource::None,
106 kind,
107 }
108 }
109
110 #[inline]
111 fn field(field: &'static str) -> Self {
112 Self {
113 source: ConfigErrorSource::FieldOverride(field.to_owned()),
114 kind: ConfigErrorKind::InvalidValue,
115 }
116 }
117
118 #[inline]
119 fn file(kind: ConfigErrorKind, p: impl AsRef<Path>) -> Self {
120 Self {
121 source: ConfigErrorSource::File(p.as_ref().to_path_buf()),
122 kind,
123 }
124 }
125
126 pub fn kind(&self) -> ConfigErrorKind {
128 self.kind
129 }
130}
131
132impl fmt::Display for ConfigError {
133 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
134 write!(f, "{}{}", self.kind(), self.source)
135 }
136}
137
138impl Error for ConfigError {}
139
140enum ConfigFormat {
141 Yaml,
142 Json,
143}
144
145impl ConfigFormat {
146 pub fn extension(&self) -> &'static str {
147 match self {
148 ConfigFormat::Yaml => "yml",
149 ConfigFormat::Json => "json",
150 }
151 }
152}
153
154trait ConfigObject: DeserializeOwned + Serialize {
155 fn format() -> ConfigFormat;
157
158 fn name() -> &'static str;
160
161 fn path(base: &Path) -> PathBuf {
163 base.join(format!("{}.{}", Self::name(), Self::format().extension()))
164 }
165
166 fn load(base: &Path) -> anyhow::Result<Self> {
168 let path = Self::path(base);
169
170 let f = fs::File::open(&path)
171 .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotOpenFile, &path))?;
172 let f = io::BufReader::new(f);
173
174 let mut source = {
175 let file = serde_vars::FileSource::default()
176 .with_variable_prefix("${file:")
177 .with_variable_suffix("}")
178 .with_base_path(base);
179 let env = serde_vars::EnvSource::default()
180 .with_variable_prefix("${")
181 .with_variable_suffix("}");
182 (file, env)
183 };
184 match Self::format() {
185 ConfigFormat::Yaml => {
186 serde_vars::deserialize(serde_yaml::Deserializer::from_reader(f), &mut source)
187 .with_context(|| ConfigError::file(ConfigErrorKind::BadYaml, &path))
188 }
189 ConfigFormat::Json => {
190 serde_vars::deserialize(&mut serde_json::Deserializer::from_reader(f), &mut source)
191 .with_context(|| ConfigError::file(ConfigErrorKind::BadJson, &path))
192 }
193 }
194 }
195
196 fn save(&self, base: &Path) -> anyhow::Result<()> {
198 let path = Self::path(base);
199 let mut options = fs::OpenOptions::new();
200 options.write(true).truncate(true).create(true);
201
202 #[cfg(unix)]
204 {
205 use std::os::unix::fs::OpenOptionsExt;
206 options.mode(0o600);
207 }
208
209 let mut f = options
210 .open(&path)
211 .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path))?;
212
213 match Self::format() {
214 ConfigFormat::Yaml => {
215 f.write_all(CONFIG_YAML_HEADER.as_bytes())?;
216 serde_yaml::to_writer(&mut f, self)
217 .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path))?
218 }
219 ConfigFormat::Json => serde_json::to_writer_pretty(&mut f, self)
220 .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path))?,
221 }
222
223 f.write_all(b"\n").ok();
224
225 Ok(())
226 }
227}
228
229#[derive(Debug, Default)]
232pub struct OverridableConfig {
233 pub mode: Option<String>,
235 pub instance: Option<String>,
237 pub log_level: Option<String>,
239 pub log_format: Option<String>,
241 pub upstream: Option<String>,
243 pub upstream_dsn: Option<String>,
245 pub host: Option<String>,
247 pub port: Option<String>,
249 pub processing: Option<String>,
251 pub kafka_url: Option<String>,
253 pub redis_url: Option<String>,
255 pub id: Option<String>,
257 pub secret_key: Option<String>,
259 pub public_key: Option<String>,
261 pub outcome_source: Option<String>,
263 pub shutdown_timeout: Option<String>,
265 pub server_name: Option<String>,
267}
268
269#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)]
271pub struct Credentials {
272 pub secret_key: SecretKey,
274 pub public_key: PublicKey,
276 pub id: RelayId,
278}
279
280impl Credentials {
281 pub fn generate() -> Self {
283 relay_log::info!("generating new relay credentials");
284 let (secret_key, public_key) = generate_key_pair();
285 Self {
286 secret_key,
287 public_key,
288 id: generate_relay_id(),
289 }
290 }
291
292 pub fn to_json_string(&self) -> anyhow::Result<String> {
294 serde_json::to_string(self)
295 .with_context(|| ConfigError::new(ConfigErrorKind::CouldNotWriteFile))
296 }
297}
298
299impl ConfigObject for Credentials {
300 fn format() -> ConfigFormat {
301 ConfigFormat::Json
302 }
303 fn name() -> &'static str {
304 "credentials"
305 }
306}
307
308#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
310#[serde(rename_all = "camelCase")]
311pub struct RelayInfo {
312 pub public_key: PublicKey,
314
315 #[serde(default)]
317 pub internal: bool,
318}
319
320impl RelayInfo {
321 pub fn new(public_key: PublicKey) -> Self {
323 Self {
324 public_key,
325 internal: false,
326 }
327 }
328}
329
330#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
332#[serde(rename_all = "camelCase")]
333pub enum RelayMode {
334 Proxy,
340
341 Managed,
347}
348
349impl<'de> Deserialize<'de> for RelayMode {
350 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
351 where
352 D: Deserializer<'de>,
353 {
354 let s = String::deserialize(deserializer)?;
355 match s.as_str() {
356 "proxy" => Ok(RelayMode::Proxy),
357 "managed" => Ok(RelayMode::Managed),
358 "static" => Err(serde::de::Error::custom(
359 "Relay mode 'static' has been removed. Please use 'managed' or 'proxy' instead.",
360 )),
361 other => Err(serde::de::Error::unknown_variant(
362 other,
363 &["proxy", "managed"],
364 )),
365 }
366 }
367}
368
369impl fmt::Display for RelayMode {
370 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
371 match self {
372 RelayMode::Proxy => write!(f, "proxy"),
373 RelayMode::Managed => write!(f, "managed"),
374 }
375 }
376}
377
378#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)]
380#[serde(rename_all = "camelCase")]
381pub enum RelayInstance {
382 Default,
384
385 Canary,
387}
388
389impl RelayInstance {
390 pub fn is_canary(&self) -> bool {
392 matches!(self, RelayInstance::Canary)
393 }
394}
395
396impl fmt::Display for RelayInstance {
397 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
398 match self {
399 RelayInstance::Default => write!(f, "default"),
400 RelayInstance::Canary => write!(f, "canary"),
401 }
402 }
403}
404
405impl FromStr for RelayInstance {
406 type Err = fmt::Error;
407
408 fn from_str(s: &str) -> Result<Self, Self::Err> {
409 match s {
410 "canary" => Ok(RelayInstance::Canary),
411 _ => Ok(RelayInstance::Default),
412 }
413 }
414}
415
416#[derive(Clone, Copy, Debug, Eq, PartialEq)]
418pub struct ParseRelayModeError;
419
420impl fmt::Display for ParseRelayModeError {
421 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
422 write!(f, "Relay mode must be one of: managed or proxy")
423 }
424}
425
426impl Error for ParseRelayModeError {}
427
428impl FromStr for RelayMode {
429 type Err = ParseRelayModeError;
430
431 fn from_str(s: &str) -> Result<Self, Self::Err> {
432 match s {
433 "proxy" => Ok(RelayMode::Proxy),
434 "managed" => Ok(RelayMode::Managed),
435 _ => Err(ParseRelayModeError),
436 }
437 }
438}
439
440fn is_default<T: Default + PartialEq>(t: &T) -> bool {
442 *t == T::default()
443}
444
445fn is_docker() -> bool {
447 if fs::metadata("/.dockerenv").is_ok() {
448 return true;
449 }
450
451 fs::read_to_string("/proc/self/cgroup").is_ok_and(|s| s.contains("/docker"))
452}
453
454fn default_host() -> IpAddr {
456 if is_docker() {
457 "0.0.0.0".parse().unwrap()
459 } else {
460 "127.0.0.1".parse().unwrap()
461 }
462}
463
464#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)]
468#[serde(rename_all = "lowercase")]
469#[derive(Default)]
470pub enum ReadinessCondition {
471 #[default]
480 Authenticated,
481 Always,
483}
484
485#[derive(Serialize, Deserialize, Debug)]
487#[serde(default)]
488pub struct Relay {
489 pub mode: RelayMode,
491 pub instance: RelayInstance,
493 pub upstream: UpstreamDescriptor,
495 pub advertised_upstream: Option<UpstreamDescriptor>,
504 pub host: IpAddr,
506 pub port: u16,
508 pub internal_host: Option<IpAddr>,
522 pub internal_port: Option<u16>,
526 #[serde(skip_serializing)]
528 pub tls_port: Option<u16>,
529 #[serde(skip_serializing)]
531 pub tls_identity_path: Option<PathBuf>,
532 #[serde(skip_serializing)]
534 pub tls_identity_password: Option<String>,
535 #[serde(skip_serializing_if = "is_default")]
540 pub override_project_ids: bool,
541}
542
543impl Default for Relay {
544 fn default() -> Self {
545 Relay {
546 mode: RelayMode::Managed,
547 instance: RelayInstance::Default,
548 upstream: "https://sentry.io/".parse().unwrap(),
549 advertised_upstream: None,
550 host: default_host(),
551 port: 3000,
552 internal_host: None,
553 internal_port: None,
554 tls_port: None,
555 tls_identity_path: None,
556 tls_identity_password: None,
557 override_project_ids: false,
558 }
559 }
560}
561
562#[derive(Serialize, Deserialize, Debug)]
564#[serde(default)]
565pub struct Metrics {
566 pub statsd: Option<String>,
570 pub statsd_buffer_size: Option<usize>,
574 pub prefix: String,
578 pub default_tags: BTreeMap<String, String>,
580 pub hostname_tag: Option<String>,
582 pub periodic_secs: u64,
587}
588
589impl Default for Metrics {
590 fn default() -> Self {
591 Metrics {
592 statsd: None,
593 statsd_buffer_size: None,
594 prefix: "sentry.relay".into(),
595 default_tags: BTreeMap::new(),
596 hostname_tag: None,
597 periodic_secs: 5,
598 }
599 }
600}
601
602#[derive(Serialize, Deserialize, Debug)]
604#[serde(default)]
605pub struct Limits {
606 pub max_concurrent_requests: usize,
609 pub max_concurrent_queries: usize,
614 pub max_event_size: ByteSize,
616 pub max_attachment_size: ByteSize,
618 pub max_attachment_count: usize,
620 pub max_attachments_size: ByteSize,
622 pub max_upload_size: ByteSize,
624 pub max_client_reports_size: ByteSize,
626 pub max_client_reports_count: usize,
628 pub max_check_in_size: ByteSize,
630 pub max_envelope_size: ByteSize,
632 pub max_sessions_size: ByteSize,
634 pub max_session_count: usize,
636 pub max_api_payload_size: ByteSize,
638 pub max_api_file_upload_size: ByteSize,
640 pub max_api_chunk_upload_size: ByteSize,
642 pub max_profile_size: ByteSize,
644 pub max_trace_metric_size: ByteSize,
646 pub max_log_size: ByteSize,
648 pub max_span_size: ByteSize,
650 pub max_standalone_span_count: usize,
652 pub max_container_size: ByteSize,
654 pub max_statsd_size: ByteSize,
656 pub max_metric_buckets_size: ByteSize,
658 pub max_replay_compressed_size: ByteSize,
660 #[serde(alias = "max_replay_size")]
662 max_replay_uncompressed_size: ByteSize,
663 pub max_replay_message_size: ByteSize,
665 pub max_removed_attribute_key_size: ByteSize,
675 pub max_thread_count: usize,
680 pub max_pool_concurrency: usize,
687 pub query_timeout: u64,
690 pub shutdown_timeout: u64,
693 pub keepalive_timeout: u64,
697 pub idle_timeout: Option<u64>,
704 pub max_connections: Option<usize>,
710 pub tcp_listen_backlog: u32,
718}
719
720impl Default for Limits {
721 fn default() -> Self {
722 Limits {
723 max_concurrent_requests: 100,
724 max_concurrent_queries: 5,
725 max_event_size: ByteSize::mebibytes(1),
726 max_attachment_size: ByteSize::mebibytes(200),
727 max_attachment_count: 30,
728 max_attachments_size: ByteSize::mebibytes(200),
729 max_upload_size: ByteSize::mebibytes(1024),
730 max_client_reports_size: ByteSize::kibibytes(100),
731 max_client_reports_count: 100,
732 max_check_in_size: ByteSize::kibibytes(100),
733 max_envelope_size: ByteSize::mebibytes(200),
734 max_sessions_size: ByteSize::mebibytes(10),
735 max_session_count: 100,
736 max_api_payload_size: ByteSize::mebibytes(20),
737 max_api_file_upload_size: ByteSize::mebibytes(40),
738 max_api_chunk_upload_size: ByteSize::mebibytes(100),
739 max_profile_size: ByteSize::mebibytes(50),
740 max_trace_metric_size: ByteSize::mebibytes(1),
741 max_log_size: ByteSize::mebibytes(1),
742 max_span_size: ByteSize::mebibytes(10),
743 max_standalone_span_count: 25,
744 max_container_size: ByteSize::mebibytes(12),
745 max_statsd_size: ByteSize::mebibytes(1),
746 max_metric_buckets_size: ByteSize::mebibytes(1),
747 max_replay_compressed_size: ByteSize::mebibytes(10),
748 max_replay_uncompressed_size: ByteSize::mebibytes(100),
749 max_replay_message_size: ByteSize::mebibytes(15),
750 max_thread_count: num_cpus::get(),
751 max_pool_concurrency: 1,
752 query_timeout: 30,
753 shutdown_timeout: 10,
754 keepalive_timeout: 5,
755 idle_timeout: None,
756 max_connections: None,
757 tcp_listen_backlog: 1024,
758 max_removed_attribute_key_size: ByteSize::kibibytes(10),
759 }
760 }
761}
762
763#[derive(Debug, Default, Deserialize, Serialize)]
765#[serde(default)]
766pub struct Routing {
767 pub accept_unknown_items: Option<bool>,
777}
778
779#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize)]
781#[serde(rename_all = "lowercase")]
782pub enum HttpEncoding {
783 #[default]
788 Identity,
789 Deflate,
795 Gzip,
802 Br,
804 Zstd,
806}
807
808impl HttpEncoding {
809 pub fn parse(str: &str) -> Self {
811 let str = str.trim();
812 if str.eq_ignore_ascii_case("zstd") {
813 Self::Zstd
814 } else if str.eq_ignore_ascii_case("br") {
815 Self::Br
816 } else if str.eq_ignore_ascii_case("gzip") || str.eq_ignore_ascii_case("x-gzip") {
817 Self::Gzip
818 } else if str.eq_ignore_ascii_case("deflate") {
819 Self::Deflate
820 } else {
821 Self::Identity
822 }
823 }
824
825 pub fn name(&self) -> Option<&'static str> {
829 match self {
830 Self::Identity => None,
831 Self::Deflate => Some("deflate"),
832 Self::Gzip => Some("gzip"),
833 Self::Br => Some("br"),
834 Self::Zstd => Some("zstd"),
835 }
836 }
837}
838
839#[derive(Serialize, Deserialize, Debug)]
841#[serde(default)]
842pub struct Http {
843 pub timeout: u32,
849 pub connection_timeout: u32,
854 pub max_retry_interval: u32,
856 pub host_header: Option<String>,
858 pub auth_interval: Option<u64>,
866 pub outage_grace_period: u64,
872 pub retry_delay: u64,
876 pub project_failure_interval: u64,
881 pub encoding: HttpEncoding,
897 pub global_metrics: bool,
904 pub forward: bool,
911 pub dns_cache: bool,
915}
916
917impl Default for Http {
918 fn default() -> Self {
919 Http {
920 timeout: 5,
921 connection_timeout: 3,
922 max_retry_interval: 60, host_header: None,
924 auth_interval: Some(600), outage_grace_period: DEFAULT_NETWORK_OUTAGE_GRACE_PERIOD,
926 retry_delay: 1,
927 project_failure_interval: 90,
928 encoding: HttpEncoding::Zstd,
929 global_metrics: false,
930 forward: true,
931 dns_cache: true,
932 }
933 }
934}
935
936#[derive(Clone, Copy, Debug, Eq, PartialEq, Default, Deserialize, Serialize)]
938#[serde(rename_all = "snake_case")]
939pub enum EnvelopeSpoolPartitioning {
940 ProjectKeyPair,
944 #[default]
951 RoundRobin,
952}
953
954#[derive(Debug, Serialize, Deserialize)]
956#[serde(default)]
957pub struct EnvelopeSpool {
958 pub path: Option<PathBuf>,
964 pub max_disk_size: ByteSize,
970 pub batch_size_bytes: ByteSize,
977 pub max_envelope_delay_secs: u64,
984 pub disk_usage_refresh_frequency_ms: u64,
989 pub max_backpressure_memory_percent: f32,
1019 pub partitions: NonZeroU8,
1026 pub partitioning: EnvelopeSpoolPartitioning,
1032 pub ephemeral: bool,
1039}
1040
1041impl Default for EnvelopeSpool {
1042 fn default() -> Self {
1043 Self {
1044 path: None,
1045 max_disk_size: ByteSize::mebibytes(500),
1046 batch_size_bytes: ByteSize::kibibytes(10),
1047 max_envelope_delay_secs: 24 * 60 * 60,
1048 disk_usage_refresh_frequency_ms: 100,
1049 max_backpressure_memory_percent: 0.8,
1050 partitions: NonZeroU8::new(1).unwrap(),
1051 partitioning: EnvelopeSpoolPartitioning::default(),
1052 ephemeral: false,
1053 }
1054 }
1055}
1056
1057#[derive(Debug, Serialize, Deserialize, Default)]
1059#[serde(default)]
1060pub struct Spool {
1061 pub envelopes: EnvelopeSpool,
1063}
1064
1065#[derive(Serialize, Deserialize, Debug)]
1067#[serde(default)]
1068pub struct Cache {
1069 pub project_request_full_config: bool,
1071 pub project_expiry: u32,
1073 pub project_grace_period: u32,
1078 pub project_refresh_interval: Option<u32>,
1084 pub relay_expiry: u32,
1086 #[serde(alias = "event_expiry")]
1092 envelope_expiry: u32,
1093 #[serde(alias = "event_buffer_size")]
1095 envelope_buffer_size: u32,
1096 pub miss_expiry: u32,
1098 pub batch_interval: u32,
1100 pub downstream_relays_batch_interval: u32,
1102 pub batch_size: usize,
1106 pub file_interval: u32,
1108 pub global_config_fetch_interval: u32,
1110}
1111
1112impl Default for Cache {
1113 fn default() -> Self {
1114 Cache {
1115 project_request_full_config: false,
1116 project_expiry: 300, project_grace_period: 120, project_refresh_interval: None,
1119 relay_expiry: 3600, envelope_expiry: 600, envelope_buffer_size: 1000,
1122 miss_expiry: 60, batch_interval: 100, downstream_relays_batch_interval: 100, batch_size: 500,
1126 file_interval: 10, global_config_fetch_interval: 10, }
1129 }
1130}
1131
1132#[derive(Serialize, Deserialize, Debug)]
1134#[serde(default)]
1135pub struct Processing {
1136 pub enabled: bool,
1138 pub geoip_path: Option<PathBuf>,
1140 pub max_secs_in_future: u32,
1142 pub max_session_secs_in_past: u32,
1144 pub kafka_config: Vec<KafkaConfigParam>,
1146 pub secondary_kafka_configs: BTreeMap<String, Vec<KafkaConfigParam>>,
1166 pub topics: TopicAssignments,
1168 pub kafka_validate_topics: bool,
1170 pub redis: Option<RedisConfigs>,
1172 pub attachment_chunk_size: ByteSize,
1174 pub projectconfig_cache_prefix: String,
1176 pub max_rate_limit: Option<u32>,
1178 pub quota_cache_ratio: Option<f32>,
1189 pub quota_cache_max: Option<f32>,
1196 #[serde(alias = "upload")]
1198 pub objectstore: ObjectstoreServiceConfig,
1199}
1200
1201impl Default for Processing {
1202 fn default() -> Self {
1204 Self {
1205 enabled: false,
1206 geoip_path: None,
1207 max_secs_in_future: 60, max_session_secs_in_past: 5 * 24 * 3600, kafka_config: Vec::new(),
1210 secondary_kafka_configs: BTreeMap::new(),
1211 topics: TopicAssignments::default(),
1212 kafka_validate_topics: false,
1213 redis: None,
1214 attachment_chunk_size: ByteSize::mebibytes(1),
1215 projectconfig_cache_prefix: "relayconfig".to_owned(),
1216 max_rate_limit: Some(300), quota_cache_ratio: None,
1218 quota_cache_max: None,
1219 objectstore: ObjectstoreServiceConfig::default(),
1220 }
1221 }
1222}
1223
1224#[derive(Debug, Default, Serialize, Deserialize)]
1226#[serde(default)]
1227pub struct Normalization {
1228 pub level: NormalizationLevel,
1230}
1231
1232#[derive(Copy, Clone, Debug, Default, Serialize, Deserialize, Eq, PartialEq)]
1234#[serde(rename_all = "lowercase")]
1235pub enum NormalizationLevel {
1236 #[default]
1240 Default,
1241 Full,
1246}
1247
1248#[derive(Serialize, Deserialize)]
1250pub struct ObjectstoreAuthConfig {
1251 pub key_id: String,
1254
1255 pub signing_key: String,
1257}
1258
1259impl fmt::Debug for ObjectstoreAuthConfig {
1260 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1261 f.debug_struct("ObjectstoreAuthConfig")
1262 .field("key_id", &self.key_id)
1263 .field("signing_key", &"[redacted]")
1264 .finish()
1265 }
1266}
1267
1268#[derive(Serialize, Deserialize, Debug)]
1270#[serde(default)]
1271pub struct ObjectstoreServiceConfig {
1272 pub objectstore_url: Option<String>,
1277
1278 pub max_concurrent_requests: usize,
1280
1281 pub max_backlog: usize,
1285
1286 pub timeout: u64,
1291
1292 pub stream_timeout: u64,
1297
1298 pub retry_delay: f64,
1300
1301 pub max_attempts: NonZeroU16,
1303
1304 pub fallback_to_kafka: bool,
1309
1310 pub auth: Option<ObjectstoreAuthConfig>,
1312}
1313
1314impl Default for ObjectstoreServiceConfig {
1315 fn default() -> Self {
1316 Self {
1317 objectstore_url: None,
1318 max_concurrent_requests: 10,
1319 max_backlog: 20,
1320 timeout: 60,
1321 stream_timeout: 5 * 60, retry_delay: 1.0,
1323 max_attempts: NonZeroU16::new(5).unwrap(),
1324 fallback_to_kafka: true,
1325 auth: None,
1326 }
1327 }
1328}
1329
1330#[derive(Copy, Clone, Debug, PartialEq, Eq)]
1333
1334pub enum EmitOutcomes {
1335 None,
1337 AsClientReports,
1339 AsOutcomes,
1341}
1342
1343impl EmitOutcomes {
1344 pub fn any(&self) -> bool {
1346 !matches!(self, EmitOutcomes::None)
1347 }
1348}
1349
1350impl Serialize for EmitOutcomes {
1351 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1352 where
1353 S: Serializer,
1354 {
1355 match self {
1357 Self::None => serializer.serialize_bool(false),
1358 Self::AsClientReports => serializer.serialize_str("as_client_reports"),
1359 Self::AsOutcomes => serializer.serialize_bool(true),
1360 }
1361 }
1362}
1363
1364struct EmitOutcomesVisitor;
1365
1366impl Visitor<'_> for EmitOutcomesVisitor {
1367 type Value = EmitOutcomes;
1368
1369 fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
1370 formatter.write_str("true, false, 'as_client_reports'")
1371 }
1372
1373 fn visit_bool<E>(self, v: bool) -> Result<Self::Value, E>
1374 where
1375 E: serde::de::Error,
1376 {
1377 Ok(if v {
1378 EmitOutcomes::AsOutcomes
1379 } else {
1380 EmitOutcomes::None
1381 })
1382 }
1383
1384 fn visit_str<E>(self, v: &str) -> Result<Self::Value, E>
1385 where
1386 E: serde::de::Error,
1387 {
1388 match v {
1389 "as_client_reports" => Ok(EmitOutcomes::AsClientReports),
1390 _ => Err(E::invalid_value(Unexpected::Str(v), &self)),
1391 }
1392 }
1393}
1394
1395impl<'de> Deserialize<'de> for EmitOutcomes {
1396 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1397 where
1398 D: Deserializer<'de>,
1399 {
1400 deserializer.deserialize_any(EmitOutcomesVisitor)
1401 }
1402}
1403
1404#[derive(Serialize, Deserialize, Debug)]
1406#[serde(default)]
1407pub struct Outcomes {
1408 pub emit_outcomes: EmitOutcomes,
1412 pub batch_size: usize,
1415 pub batch_interval: u64,
1418 pub source: Option<String>,
1421}
1422
1423impl Default for Outcomes {
1424 fn default() -> Self {
1425 Outcomes {
1426 emit_outcomes: EmitOutcomes::AsClientReports,
1427 batch_size: 1000,
1428 batch_interval: 500,
1429 source: None,
1430 }
1431 }
1432}
1433
1434#[derive(Serialize, Deserialize, Debug, Default)]
1436pub struct MinimalConfig {
1437 pub relay: Relay,
1439}
1440
1441impl MinimalConfig {
1442 pub fn save_in_folder<P: AsRef<Path>>(&self, p: P) -> anyhow::Result<()> {
1444 let path = p.as_ref();
1445 if fs::metadata(path).is_err() {
1446 fs::create_dir_all(path)
1447 .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotOpenFile, path))?;
1448 }
1449 self.save(path)
1450 }
1451}
1452
1453impl ConfigObject for MinimalConfig {
1454 fn format() -> ConfigFormat {
1455 ConfigFormat::Yaml
1456 }
1457
1458 fn name() -> &'static str {
1459 "config"
1460 }
1461}
1462
1463mod config_relay_info {
1465 use serde::ser::SerializeMap;
1466
1467 use super::*;
1468
1469 #[derive(Debug, Serialize, Deserialize, Clone)]
1471 struct RelayInfoConfig {
1472 public_key: PublicKey,
1473 #[serde(default)]
1474 internal: bool,
1475 }
1476
1477 impl From<RelayInfoConfig> for RelayInfo {
1478 fn from(v: RelayInfoConfig) -> Self {
1479 RelayInfo {
1480 public_key: v.public_key,
1481 internal: v.internal,
1482 }
1483 }
1484 }
1485
1486 impl From<RelayInfo> for RelayInfoConfig {
1487 fn from(v: RelayInfo) -> Self {
1488 RelayInfoConfig {
1489 public_key: v.public_key,
1490 internal: v.internal,
1491 }
1492 }
1493 }
1494
1495 pub(super) fn deserialize<'de, D>(des: D) -> Result<HashMap<RelayId, RelayInfo>, D::Error>
1496 where
1497 D: Deserializer<'de>,
1498 {
1499 let map = HashMap::<RelayId, RelayInfoConfig>::deserialize(des)?;
1500 Ok(map.into_iter().map(|(k, v)| (k, v.into())).collect())
1501 }
1502
1503 pub(super) fn serialize<S>(elm: &HashMap<RelayId, RelayInfo>, ser: S) -> Result<S::Ok, S::Error>
1504 where
1505 S: Serializer,
1506 {
1507 let mut map = ser.serialize_map(Some(elm.len()))?;
1508
1509 for (k, v) in elm {
1510 map.serialize_entry(k, &RelayInfoConfig::from(v.clone()))?;
1511 }
1512
1513 map.end()
1514 }
1515}
1516
1517#[derive(Serialize, Deserialize, Debug)]
1519#[serde(default)]
1520pub struct AuthConfig {
1521 #[serde(skip_serializing_if = "is_default")]
1523 pub ready: ReadinessCondition,
1524
1525 #[serde(with = "config_relay_info")]
1527 pub static_relays: HashMap<RelayId, RelayInfo>,
1528
1529 pub signature_max_age: u64,
1533}
1534
1535impl Default for AuthConfig {
1536 fn default() -> Self {
1537 Self {
1538 ready: ReadinessCondition::default(),
1539 static_relays: HashMap::new(),
1540 signature_max_age: 300, }
1542 }
1543}
1544
1545#[derive(Serialize, Deserialize, Debug, Default)]
1547pub struct GeoIpConfig {
1548 pub path: Option<PathBuf>,
1550}
1551
1552#[derive(Serialize, Deserialize, Debug)]
1554#[serde(default)]
1555pub struct CardinalityLimiter {
1556 pub cache_vacuum_interval: u64,
1562}
1563
1564impl Default for CardinalityLimiter {
1565 fn default() -> Self {
1566 Self {
1567 cache_vacuum_interval: 180,
1568 }
1569 }
1570}
1571
1572#[derive(Serialize, Deserialize, Debug)]
1577#[serde(default)]
1578pub struct Health {
1579 pub refresh_interval_ms: u64,
1586 pub max_memory_bytes: Option<ByteSize>,
1591 pub max_memory_percent: f32,
1595 pub probe_timeout_ms: u64,
1602 pub memory_stat_refresh_frequency_ms: u64,
1608}
1609
1610impl Default for Health {
1611 fn default() -> Self {
1612 Self {
1613 refresh_interval_ms: 3000,
1614 max_memory_bytes: None,
1615 max_memory_percent: 0.95,
1616 probe_timeout_ms: 900,
1617 memory_stat_refresh_frequency_ms: 100,
1618 }
1619 }
1620}
1621
1622#[derive(Serialize, Deserialize, Debug)]
1624#[serde(default)]
1625pub struct Cogs {
1626 pub max_queue_size: u64,
1632 pub relay_resource_id: String,
1638}
1639
1640impl Default for Cogs {
1641 fn default() -> Self {
1642 Self {
1643 max_queue_size: 10_000,
1644 relay_resource_id: "relay_service".to_owned(),
1645 }
1646 }
1647}
1648
1649#[derive(Debug, Clone, Serialize, Deserialize)]
1651#[serde(default)]
1652pub struct Upload {
1653 pub max_concurrent_requests: usize,
1657 pub timeout: u64,
1659 pub max_age: i64,
1663
1664 pub credentials: Option<UploadCredentials>,
1668}
1669
1670impl Default for Upload {
1671 fn default() -> Self {
1672 Self {
1673 max_concurrent_requests: 100,
1674 timeout: 5 * 60, max_age: 60 * 60, credentials: None,
1677 }
1678 }
1679}
1680
1681#[derive(Clone, Serialize, Deserialize)]
1683pub struct UploadCredentials {
1684 #[cfg(feature = "processing")]
1686 pub signing_key: SecretKey,
1687
1688 pub verification_key: PublicKey,
1690}
1691
1692impl fmt::Debug for UploadCredentials {
1693 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1694 let Self {
1695 #[cfg(feature = "processing")]
1696 signing_key: _,
1697 verification_key,
1698 } = self;
1699 let mut b = f.debug_struct("UploadCredentials");
1700 #[cfg(feature = "processing")]
1701 b.field("signing_key", &"[redacted]");
1702 b.field("verification_key", verification_key).finish()
1703 }
1704}
1705
1706#[derive(Serialize, Deserialize, Debug, Default)]
1708#[serde(default)]
1709#[allow(missing_docs)]
1710pub struct ConfigValues {
1711 pub relay: Relay,
1712 pub http: Http,
1713 pub cache: Cache,
1714 pub spool: Spool,
1715 pub limits: Limits,
1716 pub logging: relay_log::LogConfig,
1717 pub routing: Routing,
1718 pub metrics: Metrics,
1719 pub sentry: relay_log::SentryConfig,
1720 pub processing: Processing,
1721 pub outcomes: Outcomes,
1722 pub aggregator: AggregatorServiceConfig,
1723 pub secondary_aggregators: Vec<ScopedAggregatorConfig>,
1724 pub auth: AuthConfig,
1725 pub geoip: GeoIpConfig,
1726 pub normalization: Normalization,
1727 pub cardinality_limiter: CardinalityLimiter,
1728 pub health: Health,
1729 pub cogs: Cogs,
1730 pub upload: Upload,
1731}
1732
1733impl ConfigObject for ConfigValues {
1734 fn format() -> ConfigFormat {
1735 ConfigFormat::Yaml
1736 }
1737
1738 fn name() -> &'static str {
1739 "config"
1740 }
1741}
1742
1743pub struct Config {
1745 values: ConfigValues,
1746 credentials: Option<Credentials>,
1747 path: PathBuf,
1748}
1749
1750impl fmt::Debug for Config {
1751 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1752 f.debug_struct("Config")
1753 .field("path", &self.path)
1754 .field("values", &self.values)
1755 .finish()
1756 }
1757}
1758
1759impl Config {
1760 pub fn from_path<P: AsRef<Path>>(path: P) -> anyhow::Result<Config> {
1762 let path = env::current_dir()
1763 .map(|x| x.join(path.as_ref()))
1764 .unwrap_or_else(|_| path.as_ref().to_path_buf());
1765
1766 let config = Config {
1767 values: ConfigValues::load(&path)?,
1768 credentials: if Credentials::path(&path).exists() {
1769 Some(Credentials::load(&path)?)
1770 } else {
1771 None
1772 },
1773 path: path.clone(),
1774 };
1775
1776 if cfg!(not(feature = "processing")) && config.processing_enabled() {
1777 return Err(ConfigError::file(ConfigErrorKind::ProcessingNotAvailable, &path).into());
1778 }
1779
1780 Ok(config)
1781 }
1782
1783 pub fn from_json_value(value: serde_json::Value) -> anyhow::Result<Config> {
1787 Ok(Config {
1788 values: serde_json::from_value(value)
1789 .with_context(|| ConfigError::new(ConfigErrorKind::BadJson))?,
1790 credentials: None,
1791 path: PathBuf::new(),
1792 })
1793 }
1794
1795 pub fn apply_override(
1798 &mut self,
1799 mut overrides: OverridableConfig,
1800 ) -> anyhow::Result<&mut Self> {
1801 let relay = &mut self.values.relay;
1802
1803 if let Some(mode) = overrides.mode {
1804 relay.mode = mode
1805 .parse::<RelayMode>()
1806 .with_context(|| ConfigError::field("mode"))?;
1807 }
1808
1809 if let Some(deployment) = overrides.instance {
1810 relay.instance = deployment
1811 .parse::<RelayInstance>()
1812 .with_context(|| ConfigError::field("deployment"))?;
1813 }
1814
1815 if let Some(log_level) = overrides.log_level {
1816 self.values.logging.level = log_level.parse()?;
1817 }
1818
1819 if let Some(log_format) = overrides.log_format {
1820 self.values.logging.format = log_format.parse()?;
1821 }
1822
1823 if let Some(upstream) = overrides.upstream {
1824 relay.upstream = upstream
1825 .parse::<UpstreamDescriptor>()
1826 .with_context(|| ConfigError::field("upstream"))?;
1827 } else if let Some(upstream_dsn) = overrides.upstream_dsn {
1828 relay.upstream = upstream_dsn
1829 .parse::<Dsn>()
1830 .map(|dsn| UpstreamDescriptor::from_dsn(&dsn))
1831 .with_context(|| ConfigError::field("upstream_dsn"))?;
1832 }
1833
1834 if let Some(host) = overrides.host {
1835 relay.host = host
1836 .parse::<IpAddr>()
1837 .with_context(|| ConfigError::field("host"))?;
1838 }
1839
1840 if let Some(port) = overrides.port {
1841 relay.port = port
1842 .as_str()
1843 .parse()
1844 .with_context(|| ConfigError::field("port"))?;
1845 }
1846
1847 let processing = &mut self.values.processing;
1848 if let Some(enabled) = overrides.processing {
1849 match enabled.to_lowercase().as_str() {
1850 "true" | "1" => processing.enabled = true,
1851 "false" | "0" | "" => processing.enabled = false,
1852 _ => return Err(ConfigError::field("processing").into()),
1853 }
1854 }
1855
1856 if let Some(redis) = overrides.redis_url {
1857 processing.redis = Some(RedisConfigs::Unified(RedisConfig::single(redis)))
1858 }
1859
1860 if let Some(kafka_url) = overrides.kafka_url {
1861 let existing = processing
1862 .kafka_config
1863 .iter_mut()
1864 .find(|e| e.name == "bootstrap.servers");
1865
1866 if let Some(config_param) = existing {
1867 config_param.value = kafka_url;
1868 } else {
1869 processing.kafka_config.push(KafkaConfigParam {
1870 name: "bootstrap.servers".to_owned(),
1871 value: kafka_url,
1872 })
1873 }
1874 }
1875 let id = if let Some(id) = overrides.id {
1877 let id = Uuid::parse_str(&id).with_context(|| ConfigError::field("id"))?;
1878 Some(id)
1879 } else {
1880 None
1881 };
1882 let public_key = if let Some(public_key) = overrides.public_key {
1883 let public_key = public_key
1884 .parse::<PublicKey>()
1885 .with_context(|| ConfigError::field("public_key"))?;
1886 Some(public_key)
1887 } else {
1888 None
1889 };
1890
1891 let secret_key = if let Some(secret_key) = overrides.secret_key {
1892 let secret_key = secret_key
1893 .parse::<SecretKey>()
1894 .with_context(|| ConfigError::field("secret_key"))?;
1895 Some(secret_key)
1896 } else {
1897 None
1898 };
1899 let outcomes = &mut self.values.outcomes;
1900 if overrides.outcome_source.is_some() {
1901 outcomes.source = overrides.outcome_source.take();
1902 }
1903
1904 if let Some(credentials) = &mut self.credentials {
1905 if let Some(id) = id {
1907 credentials.id = id;
1908 }
1909 if let Some(public_key) = public_key {
1910 credentials.public_key = public_key;
1911 }
1912 if let Some(secret_key) = secret_key {
1913 credentials.secret_key = secret_key
1914 }
1915 } else {
1916 match (id, public_key, secret_key) {
1918 (Some(id), Some(public_key), Some(secret_key)) => {
1919 self.credentials = Some(Credentials {
1920 secret_key,
1921 public_key,
1922 id,
1923 })
1924 }
1925 (None, None, None) => {
1926 }
1929 _ => {
1930 return Err(ConfigError::field("incomplete credentials").into());
1931 }
1932 }
1933 }
1934
1935 let limits = &mut self.values.limits;
1936 if let Some(shutdown_timeout) = overrides.shutdown_timeout
1937 && let Ok(shutdown_timeout) = shutdown_timeout.parse::<u64>()
1938 {
1939 limits.shutdown_timeout = shutdown_timeout;
1940 }
1941
1942 if let Some(server_name) = overrides.server_name {
1943 self.values.sentry.server_name = Some(server_name.into());
1944 }
1945
1946 Ok(self)
1947 }
1948
1949 pub fn config_exists<P: AsRef<Path>>(path: P) -> bool {
1951 fs::metadata(ConfigValues::path(path.as_ref())).is_ok()
1952 }
1953
1954 pub fn path(&self) -> &Path {
1956 &self.path
1957 }
1958
1959 pub fn to_yaml_string(&self) -> anyhow::Result<String> {
1961 serde_yaml::to_string(&self.values)
1962 .with_context(|| ConfigError::new(ConfigErrorKind::CouldNotWriteFile))
1963 }
1964
1965 pub fn regenerate_credentials(&mut self, save: bool) -> anyhow::Result<()> {
1969 let creds = Credentials::generate();
1970 if save {
1971 creds.save(&self.path)?;
1972 }
1973 self.credentials = Some(creds);
1974 Ok(())
1975 }
1976
1977 pub fn credentials(&self) -> Option<&Credentials> {
1979 self.credentials.as_ref()
1980 }
1981
1982 pub fn replace_credentials(
1986 &mut self,
1987 credentials: Option<Credentials>,
1988 ) -> anyhow::Result<bool> {
1989 if self.credentials == credentials {
1990 return Ok(false);
1991 }
1992
1993 match credentials {
1994 Some(ref creds) => {
1995 creds.save(&self.path)?;
1996 }
1997 None => {
1998 let path = Credentials::path(&self.path);
1999 if fs::metadata(&path).is_ok() {
2000 fs::remove_file(&path).with_context(|| {
2001 ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path)
2002 })?;
2003 }
2004 }
2005 }
2006
2007 self.credentials = credentials;
2008 Ok(true)
2009 }
2010
2011 pub fn has_credentials(&self) -> bool {
2013 self.credentials.is_some()
2014 }
2015
2016 pub fn secret_key(&self) -> Option<&SecretKey> {
2018 self.credentials.as_ref().map(|x| &x.secret_key)
2019 }
2020
2021 pub fn public_key(&self) -> Option<&PublicKey> {
2023 self.credentials.as_ref().map(|x| &x.public_key)
2024 }
2025
2026 pub fn relay_id(&self) -> Option<&RelayId> {
2028 self.credentials.as_ref().map(|x| &x.id)
2029 }
2030
2031 pub fn relay_mode(&self) -> RelayMode {
2033 self.values.relay.mode
2034 }
2035
2036 pub fn relay_instance(&self) -> RelayInstance {
2038 self.values.relay.instance
2039 }
2040
2041 pub fn upstream(&self) -> &UpstreamDescriptor {
2043 &self.values.relay.upstream
2044 }
2045
2046 pub fn advertised_upstream(&self) -> Option<&UpstreamDescriptor> {
2048 self.values.relay.advertised_upstream.as_ref()
2049 }
2050
2051 pub fn http_host_header(&self) -> Option<&str> {
2053 self.values.http.host_header.as_deref()
2054 }
2055
2056 pub fn listen_addr(&self) -> SocketAddr {
2058 (self.values.relay.host, self.values.relay.port).into()
2059 }
2060
2061 pub fn listen_addr_internal(&self) -> Option<SocketAddr> {
2069 match (
2070 self.values.relay.internal_host,
2071 self.values.relay.internal_port,
2072 ) {
2073 (Some(host), None) => Some((host, self.values.relay.port).into()),
2074 (None, Some(port)) => Some((self.values.relay.host, port).into()),
2075 (Some(host), Some(port)) => Some((host, port).into()),
2076 (None, None) => None,
2077 }
2078 }
2079
2080 pub fn tls_listen_addr(&self) -> Option<SocketAddr> {
2082 if self.values.relay.tls_identity_path.is_some() {
2083 let port = self.values.relay.tls_port.unwrap_or(3443);
2084 Some((self.values.relay.host, port).into())
2085 } else {
2086 None
2087 }
2088 }
2089
2090 pub fn tls_identity_path(&self) -> Option<&Path> {
2092 self.values.relay.tls_identity_path.as_deref()
2093 }
2094
2095 pub fn tls_identity_password(&self) -> Option<&str> {
2097 self.values.relay.tls_identity_password.as_deref()
2098 }
2099
2100 pub fn override_project_ids(&self) -> bool {
2104 self.values.relay.override_project_ids
2105 }
2106
2107 pub fn requires_auth(&self) -> bool {
2111 match self.values.auth.ready {
2112 ReadinessCondition::Authenticated => self.relay_mode() == RelayMode::Managed,
2113 ReadinessCondition::Always => false,
2114 }
2115 }
2116
2117 pub fn http_auth_interval(&self) -> Option<Duration> {
2121 if self.processing_enabled() {
2122 return None;
2123 }
2124
2125 match self.values.http.auth_interval {
2126 None | Some(0) => None,
2127 Some(secs) => Some(Duration::from_secs(secs)),
2128 }
2129 }
2130
2131 pub fn http_outage_grace_period(&self) -> Duration {
2134 Duration::from_secs(self.values.http.outage_grace_period)
2135 }
2136
2137 pub fn http_retry_delay(&self) -> Duration {
2142 Duration::from_secs(self.values.http.retry_delay)
2143 }
2144
2145 pub fn http_project_failure_interval(&self) -> Duration {
2147 Duration::from_secs(self.values.http.project_failure_interval)
2148 }
2149
2150 pub fn http_encoding(&self) -> HttpEncoding {
2152 self.values.http.encoding
2153 }
2154
2155 pub fn http_global_metrics(&self) -> bool {
2157 self.values.http.global_metrics
2158 }
2159
2160 pub fn http_forward(&self) -> bool {
2165 self.values.http.forward && !self.processing_enabled()
2166 }
2167
2168 pub fn emit_outcomes(&self) -> EmitOutcomes {
2173 if self.processing_enabled() {
2174 return EmitOutcomes::AsOutcomes;
2175 }
2176 self.values.outcomes.emit_outcomes
2177 }
2178
2179 pub fn outcome_batch_size(&self) -> usize {
2181 self.values.outcomes.batch_size
2182 }
2183
2184 pub fn outcome_batch_interval(&self) -> Duration {
2186 Duration::from_millis(self.values.outcomes.batch_interval)
2187 }
2188
2189 pub fn outcome_source(&self) -> Option<&str> {
2191 self.values.outcomes.source.as_deref()
2192 }
2193
2194 pub fn logging(&self) -> &relay_log::LogConfig {
2196 &self.values.logging
2197 }
2198
2199 pub fn sentry(&self) -> &relay_log::SentryConfig {
2201 &self.values.sentry
2202 }
2203
2204 pub fn statsd_addr(&self) -> Option<&str> {
2206 self.values.metrics.statsd.as_deref()
2207 }
2208
2209 pub fn statsd_buffer_size(&self) -> Option<usize> {
2211 self.values.metrics.statsd_buffer_size
2212 }
2213
2214 pub fn metrics_prefix(&self) -> &str {
2216 &self.values.metrics.prefix
2217 }
2218
2219 pub fn metrics_default_tags(&self) -> &BTreeMap<String, String> {
2221 &self.values.metrics.default_tags
2222 }
2223
2224 pub fn metrics_hostname_tag(&self) -> Option<&str> {
2226 self.values.metrics.hostname_tag.as_deref()
2227 }
2228
2229 pub fn metrics_periodic_interval(&self) -> Option<Duration> {
2233 match self.values.metrics.periodic_secs {
2234 0 => None,
2235 secs => Some(Duration::from_secs(secs)),
2236 }
2237 }
2238
2239 pub fn http_timeout(&self) -> Duration {
2241 Duration::from_secs(self.values.http.timeout.into())
2242 }
2243
2244 pub fn http_connection_timeout(&self) -> Duration {
2246 Duration::from_secs(self.values.http.connection_timeout.into())
2247 }
2248
2249 pub fn http_max_retry_interval(&self) -> Duration {
2251 Duration::from_secs(self.values.http.max_retry_interval.into())
2252 }
2253
2254 pub fn http_dns_cache(&self) -> bool {
2256 self.values.http.dns_cache
2257 }
2258
2259 pub fn project_cache_expiry(&self) -> Duration {
2261 Duration::from_secs(self.values.cache.project_expiry.into())
2262 }
2263
2264 pub fn request_full_project_config(&self) -> bool {
2266 self.values.cache.project_request_full_config
2267 }
2268
2269 pub fn relay_cache_expiry(&self) -> Duration {
2271 Duration::from_secs(self.values.cache.relay_expiry.into())
2272 }
2273
2274 pub fn envelope_buffer_size(&self) -> usize {
2276 self.values
2277 .cache
2278 .envelope_buffer_size
2279 .try_into()
2280 .unwrap_or(usize::MAX)
2281 }
2282
2283 pub fn cache_miss_expiry(&self) -> Duration {
2285 Duration::from_secs(self.values.cache.miss_expiry.into())
2286 }
2287
2288 pub fn project_grace_period(&self) -> Duration {
2290 Duration::from_secs(self.values.cache.project_grace_period.into())
2291 }
2292
2293 pub fn project_refresh_interval(&self) -> Option<Duration> {
2297 self.values
2298 .cache
2299 .project_refresh_interval
2300 .map(Into::into)
2301 .map(Duration::from_secs)
2302 }
2303
2304 pub fn query_batch_interval(&self) -> Duration {
2307 Duration::from_millis(self.values.cache.batch_interval.into())
2308 }
2309
2310 pub fn downstream_relays_batch_interval(&self) -> Duration {
2312 Duration::from_millis(self.values.cache.downstream_relays_batch_interval.into())
2313 }
2314
2315 pub fn local_cache_interval(&self) -> Duration {
2317 Duration::from_secs(self.values.cache.file_interval.into())
2318 }
2319
2320 pub fn global_config_fetch_interval(&self) -> Duration {
2323 Duration::from_secs(self.values.cache.global_config_fetch_interval.into())
2324 }
2325
2326 pub fn spool_envelopes_path(&self, partition_id: u8) -> Option<PathBuf> {
2331 let mut path = self
2332 .values
2333 .spool
2334 .envelopes
2335 .path
2336 .as_ref()
2337 .map(|path| path.to_owned())?;
2338
2339 if partition_id == 0 {
2340 return Some(path);
2341 }
2342
2343 let file_name = path.file_name().and_then(|f| f.to_str())?;
2344 let new_file_name = format!("{file_name}.{partition_id}");
2345 path.set_file_name(new_file_name);
2346
2347 Some(path)
2348 }
2349
2350 pub fn spool_envelopes_max_disk_size(&self) -> usize {
2352 self.values.spool.envelopes.max_disk_size.as_bytes()
2353 }
2354
2355 pub fn spool_envelopes_batch_size_bytes(&self) -> usize {
2358 self.values.spool.envelopes.batch_size_bytes.as_bytes()
2359 }
2360
2361 pub fn spool_envelopes_max_age(&self) -> Duration {
2363 Duration::from_secs(self.values.spool.envelopes.max_envelope_delay_secs)
2364 }
2365
2366 pub fn spool_disk_usage_refresh_frequency_ms(&self) -> Duration {
2368 Duration::from_millis(self.values.spool.envelopes.disk_usage_refresh_frequency_ms)
2369 }
2370
2371 pub fn spool_max_backpressure_memory_percent(&self) -> f32 {
2373 self.values.spool.envelopes.max_backpressure_memory_percent
2374 }
2375
2376 pub fn spool_partitions(&self) -> NonZeroU8 {
2378 self.values.spool.envelopes.partitions
2379 }
2380
2381 pub fn spool_partitioning(&self) -> EnvelopeSpoolPartitioning {
2383 self.values.spool.envelopes.partitioning
2384 }
2385
2386 pub fn spool_ephemeral(&self) -> bool {
2388 self.values.spool.envelopes.ephemeral
2389 }
2390
2391 pub fn max_event_size(&self) -> usize {
2393 self.values.limits.max_event_size.as_bytes()
2394 }
2395
2396 pub fn max_attachment_size(&self) -> usize {
2398 self.values.limits.max_attachment_size.as_bytes()
2399 }
2400
2401 pub fn max_attachment_count(&self) -> usize {
2403 self.values.limits.max_attachment_count
2404 }
2405
2406 pub fn max_attachments_size(&self) -> usize {
2409 self.values.limits.max_attachments_size.as_bytes()
2410 }
2411
2412 pub fn max_upload_size(&self) -> usize {
2414 self.values.limits.max_upload_size.as_bytes()
2415 }
2416
2417 pub fn max_client_reports_count(&self) -> usize {
2419 self.values.limits.max_client_reports_count
2420 }
2421
2422 pub fn max_client_reports_size(&self) -> usize {
2424 self.values.limits.max_client_reports_size.as_bytes()
2425 }
2426
2427 pub fn max_check_in_size(&self) -> usize {
2429 self.values.limits.max_check_in_size.as_bytes()
2430 }
2431
2432 pub fn max_log_size(&self) -> usize {
2434 self.values.limits.max_log_size.as_bytes()
2435 }
2436
2437 pub fn max_span_size(&self) -> usize {
2439 self.values.limits.max_span_size.as_bytes()
2440 }
2441
2442 pub fn max_standalone_span_count(&self) -> usize {
2444 self.values.limits.max_standalone_span_count
2445 }
2446
2447 pub fn max_container_size(&self) -> usize {
2449 self.values.limits.max_container_size.as_bytes()
2450 }
2451
2452 pub fn max_envelope_size(&self) -> usize {
2456 self.values.limits.max_envelope_size.as_bytes()
2457 }
2458
2459 pub fn max_session_count(&self) -> usize {
2461 self.values.limits.max_session_count
2462 }
2463
2464 pub fn max_sessions_size(&self) -> usize {
2466 self.values.limits.max_sessions_size.as_bytes()
2467 }
2468
2469 pub fn max_statsd_size(&self) -> usize {
2471 self.values.limits.max_statsd_size.as_bytes()
2472 }
2473
2474 pub fn max_metric_buckets_size(&self) -> usize {
2476 self.values.limits.max_metric_buckets_size.as_bytes()
2477 }
2478
2479 pub fn max_api_payload_size(&self) -> usize {
2481 self.values.limits.max_api_payload_size.as_bytes()
2482 }
2483
2484 pub fn max_api_file_upload_size(&self) -> usize {
2486 self.values.limits.max_api_file_upload_size.as_bytes()
2487 }
2488
2489 pub fn max_api_chunk_upload_size(&self) -> usize {
2491 self.values.limits.max_api_chunk_upload_size.as_bytes()
2492 }
2493
2494 pub fn max_profile_size(&self) -> usize {
2496 self.values.limits.max_profile_size.as_bytes()
2497 }
2498
2499 pub fn max_trace_metric_size(&self) -> usize {
2501 self.values.limits.max_trace_metric_size.as_bytes()
2502 }
2503
2504 pub fn max_replay_compressed_size(&self) -> usize {
2506 self.values.limits.max_replay_compressed_size.as_bytes()
2507 }
2508
2509 pub fn max_replay_uncompressed_size(&self) -> usize {
2511 self.values.limits.max_replay_uncompressed_size.as_bytes()
2512 }
2513
2514 pub fn max_replay_message_size(&self) -> usize {
2520 self.values.limits.max_replay_message_size.as_bytes()
2521 }
2522
2523 pub fn max_concurrent_requests(&self) -> usize {
2525 self.values.limits.max_concurrent_requests
2526 }
2527
2528 pub fn max_concurrent_queries(&self) -> usize {
2530 self.values.limits.max_concurrent_queries
2531 }
2532
2533 pub fn max_removed_attribute_key_size(&self) -> usize {
2535 self.values.limits.max_removed_attribute_key_size.as_bytes()
2536 }
2537
2538 pub fn query_timeout(&self) -> Duration {
2540 Duration::from_secs(self.values.limits.query_timeout)
2541 }
2542
2543 pub fn shutdown_timeout(&self) -> Duration {
2546 Duration::from_secs(self.values.limits.shutdown_timeout)
2547 }
2548
2549 pub fn keepalive_timeout(&self) -> Duration {
2553 Duration::from_secs(self.values.limits.keepalive_timeout)
2554 }
2555
2556 pub fn idle_timeout(&self) -> Option<Duration> {
2558 self.values.limits.idle_timeout.map(Duration::from_secs)
2559 }
2560
2561 pub fn max_connections(&self) -> Option<usize> {
2563 self.values.limits.max_connections
2564 }
2565
2566 pub fn tcp_listen_backlog(&self) -> u32 {
2568 self.values.limits.tcp_listen_backlog
2569 }
2570
2571 pub fn cpu_concurrency(&self) -> usize {
2573 self.values.limits.max_thread_count
2574 }
2575
2576 pub fn pool_concurrency(&self) -> usize {
2578 self.values.limits.max_pool_concurrency
2579 }
2580
2581 pub fn query_batch_size(&self) -> usize {
2583 self.values.cache.batch_size
2584 }
2585
2586 pub fn project_configs_path(&self) -> PathBuf {
2588 self.path.join("projects")
2589 }
2590
2591 pub fn processing_enabled(&self) -> bool {
2593 self.values.processing.enabled
2594 }
2595
2596 pub fn normalization_level(&self) -> NormalizationLevel {
2598 self.values.normalization.level
2599 }
2600
2601 pub fn geoip_path(&self) -> Option<&Path> {
2603 self.values
2604 .geoip
2605 .path
2606 .as_deref()
2607 .or(self.values.processing.geoip_path.as_deref())
2608 }
2609
2610 pub fn max_secs_in_future(&self) -> i64 {
2614 self.values.processing.max_secs_in_future.into()
2615 }
2616
2617 pub fn max_session_secs_in_past(&self) -> i64 {
2619 self.values.processing.max_session_secs_in_past.into()
2620 }
2621
2622 pub fn kafka_configs(
2624 &self,
2625 topic: KafkaTopic,
2626 ) -> Result<KafkaTopicConfig<'_>, KafkaConfigError> {
2627 self.values.processing.topics.get(topic).kafka_configs(
2628 &self.values.processing.kafka_config,
2629 &self.values.processing.secondary_kafka_configs,
2630 )
2631 }
2632
2633 pub fn kafka_validate_topics(&self) -> bool {
2635 self.values.processing.kafka_validate_topics
2636 }
2637
2638 pub fn unused_topic_assignments(&self) -> &relay_kafka::Unused {
2640 &self.values.processing.topics.unused
2641 }
2642
2643 pub fn objectstore(&self) -> &ObjectstoreServiceConfig {
2645 &self.values.processing.objectstore
2646 }
2647
2648 pub fn upload(&self) -> &Upload {
2650 &self.values.upload
2651 }
2652
2653 #[cfg(feature = "processing")]
2655 pub fn upload_signing_key(&self) -> Option<&SecretKey> {
2656 self.upload()
2657 .credentials
2658 .as_ref()
2659 .map(|c| &c.signing_key)
2660 .or(self.credentials().map(|c| &c.secret_key))
2661 }
2662
2663 #[cfg(feature = "processing")]
2665 pub fn upload_verification_key(&self) -> Option<&PublicKey> {
2666 self.upload()
2667 .credentials
2668 .as_ref()
2669 .map(|c| &c.verification_key)
2670 .or(self.credentials().map(|c| &c.public_key))
2671 }
2672
2673 pub fn redis(&self) -> Option<RedisConfigsRef<'_>> {
2676 let redis_configs = self.values.processing.redis.as_ref()?;
2677
2678 Some(build_redis_configs(
2679 redis_configs,
2680 self.cpu_concurrency() as u32,
2681 self.pool_concurrency() as u32,
2682 ))
2683 }
2684
2685 pub fn attachment_chunk_size(&self) -> usize {
2687 self.values.processing.attachment_chunk_size.as_bytes()
2688 }
2689
2690 pub fn metrics_max_batch_size_bytes(&self) -> usize {
2692 self.values.aggregator.max_flush_bytes
2693 }
2694
2695 pub fn projectconfig_cache_prefix(&self) -> &str {
2698 &self.values.processing.projectconfig_cache_prefix
2699 }
2700
2701 pub fn max_rate_limit(&self) -> Option<u64> {
2703 self.values.processing.max_rate_limit.map(u32::into)
2704 }
2705
2706 pub fn quota_cache_ratio(&self) -> Option<f32> {
2708 self.values.processing.quota_cache_ratio
2709 }
2710
2711 pub fn quota_cache_max(&self) -> Option<f32> {
2713 self.values.processing.quota_cache_max
2714 }
2715
2716 pub fn cardinality_limiter_cache_vacuum_interval(&self) -> Duration {
2720 Duration::from_secs(self.values.cardinality_limiter.cache_vacuum_interval)
2721 }
2722
2723 pub fn health_refresh_interval(&self) -> Duration {
2725 Duration::from_millis(self.values.health.refresh_interval_ms)
2726 }
2727
2728 pub fn health_max_memory_watermark_bytes(&self) -> u64 {
2730 self.values
2731 .health
2732 .max_memory_bytes
2733 .as_ref()
2734 .map_or(u64::MAX, |b| b.as_bytes() as u64)
2735 }
2736
2737 pub fn health_max_memory_watermark_percent(&self) -> f32 {
2739 self.values.health.max_memory_percent
2740 }
2741
2742 pub fn health_probe_timeout(&self) -> Duration {
2744 Duration::from_millis(self.values.health.probe_timeout_ms)
2745 }
2746
2747 pub fn memory_stat_refresh_frequency_ms(&self) -> u64 {
2749 self.values.health.memory_stat_refresh_frequency_ms
2750 }
2751
2752 pub fn cogs_max_queue_size(&self) -> u64 {
2754 self.values.cogs.max_queue_size
2755 }
2756
2757 pub fn cogs_relay_resource_id(&self) -> &str {
2759 &self.values.cogs.relay_resource_id
2760 }
2761
2762 pub fn default_aggregator_config(&self) -> &AggregatorServiceConfig {
2764 &self.values.aggregator
2765 }
2766
2767 pub fn secondary_aggregator_configs(&self) -> &Vec<ScopedAggregatorConfig> {
2769 &self.values.secondary_aggregators
2770 }
2771
2772 pub fn aggregator_config_for(&self, namespace: MetricNamespace) -> &AggregatorServiceConfig {
2774 for entry in &self.values.secondary_aggregators {
2775 if entry.condition.matches(Some(namespace)) {
2776 return &entry.config;
2777 }
2778 }
2779 &self.values.aggregator
2780 }
2781
2782 pub fn static_relays(&self) -> &HashMap<RelayId, RelayInfo> {
2784 &self.values.auth.static_relays
2785 }
2786
2787 pub fn signature_max_age(&self) -> Duration {
2789 Duration::from_secs(self.values.auth.signature_max_age)
2790 }
2791
2792 pub fn accept_unknown_items(&self) -> bool {
2794 let forward = self.values.routing.accept_unknown_items;
2795 forward.unwrap_or_else(|| !self.processing_enabled())
2796 }
2797}
2798
2799impl Default for Config {
2800 fn default() -> Self {
2801 Self {
2802 values: ConfigValues::default(),
2803 credentials: None,
2804 path: PathBuf::new(),
2805 }
2806 }
2807}
2808
2809#[cfg(test)]
2810mod tests {
2811 use super::*;
2812
2813 #[test]
2815 fn test_event_buffer_size() {
2816 let yaml = r###"
2817cache:
2818 event_buffer_size: 1000000
2819 event_expiry: 1800
2820"###;
2821
2822 let values: ConfigValues = serde_yaml::from_str(yaml).unwrap();
2823 assert_eq!(values.cache.envelope_buffer_size, 1_000_000);
2824 assert_eq!(values.cache.envelope_expiry, 1800);
2825 }
2826
2827 #[cfg(feature = "processing")]
2828 #[test]
2829 fn test_upload_secret_key_from_file() {
2830 let path = env::temp_dir().join(Uuid::new_v4().to_string());
2831 fs::create_dir(&path).unwrap();
2832 fs::write(
2833 path.join("my_secret.txt"),
2834 "U3LSQM5NorvgnoYHW_aZpc_43nuuh3lhs3zjjcBwaks",
2835 )
2836 .unwrap();
2837 fs::write(
2838 ConfigValues::path(&path),
2839 r#"
2840upload:
2841 credentials:
2842 signing_key: ${file:my_secret.txt}
2843 verification_key: "VNS8haF0VTnuMMDR2t-f7AgnmUcXmcdzV3SVksSk34s""#,
2844 )
2845 .unwrap();
2846
2847 let config = Config::from_path(&path).unwrap();
2848
2849 fs::remove_dir_all(path).unwrap();
2850
2851 let signing_key = &config.upload().credentials.as_ref().unwrap().signing_key;
2852 assert_eq!(
2853 signing_key.to_string(),
2854 "U3LSQM5NorvgnoYHW_aZpc_43nuuh3lhs3zjjcBwaks"
2855 );
2856 }
2857
2858 #[test]
2859 fn test_emit_outcomes() {
2860 for (serialized, deserialized) in &[
2861 ("true", EmitOutcomes::AsOutcomes),
2862 ("false", EmitOutcomes::None),
2863 ("\"as_client_reports\"", EmitOutcomes::AsClientReports),
2864 ] {
2865 let value: EmitOutcomes = serde_json::from_str(serialized).unwrap();
2866 assert_eq!(value, *deserialized);
2867 assert_eq!(serde_json::to_string(&value).unwrap(), *serialized);
2868 }
2869 }
2870
2871 #[test]
2872 fn test_emit_outcomes_invalid() {
2873 assert!(serde_json::from_str::<EmitOutcomes>("asdf").is_err());
2874 }
2875}