Skip to main content

relay_config/
config.rs

1use std::collections::{BTreeMap, HashMap};
2use std::error::Error;
3use std::io::Write;
4use std::net::{IpAddr, SocketAddr};
5use std::num::{NonZeroU8, NonZeroU16, NonZeroU32};
6use std::path::{Path, PathBuf};
7use std::str::FromStr;
8use std::sync::{Arc, Mutex, PoisonError};
9use std::time::Duration;
10use std::{env, fmt, fs, io};
11
12use anyhow::Context;
13use arc_swap::ArcSwap;
14use relay_auth::{PublicKey, RelayId, SecretKey, generate_key_pair, generate_relay_id};
15use relay_common::Dsn;
16use relay_kafka::{
17    ConfigError as KafkaConfigError, KafkaConfigParam, KafkaTopic, KafkaTopicConfig,
18    TopicAssignments,
19};
20use relay_metrics::MetricNamespace;
21use serde::de::{DeserializeOwned, Unexpected, Visitor};
22use serde::{Deserialize, Deserializer, Serialize, Serializer};
23use uuid::Uuid;
24
25use crate::aggregator::{AggregatorServiceConfig, ScopedAggregatorConfig};
26use crate::byte_size::ByteSize;
27use crate::upstream::UpstreamDescriptor;
28use crate::{RedisConfig, RedisConfigs, RedisConfigsRef, build_redis_configs};
29
30const DEFAULT_NETWORK_OUTAGE_GRACE_PERIOD: u64 = 10;
31
32static CONFIG_YAML_HEADER: &str = r###"# Please see the relevant documentation.
33# Performance tuning: https://docs.sentry.io/product/relay/operating-guidelines/
34# All config options: https://docs.sentry.io/product/relay/options/
35"###;
36
37/// Indicates config related errors.
38#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
39#[non_exhaustive]
40pub enum ConfigErrorKind {
41    /// Failed to open the file.
42    CouldNotOpenFile,
43    /// Failed to save a file.
44    CouldNotWriteFile,
45    /// Parsing YAML failed.
46    BadYaml,
47    /// Parsing JSON failed.
48    BadJson,
49    /// Invalid config value
50    InvalidValue,
51    /// The user attempted to run Relay with processing enabled, but uses a binary that was
52    /// compiled without the processing feature.
53    ProcessingNotAvailable,
54}
55
56impl fmt::Display for ConfigErrorKind {
57    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
58        match self {
59            Self::CouldNotOpenFile => write!(f, "could not open config file"),
60            Self::CouldNotWriteFile => write!(f, "could not write config file"),
61            Self::BadYaml => write!(f, "could not parse yaml config file"),
62            Self::BadJson => write!(f, "could not parse json config file"),
63            Self::InvalidValue => write!(f, "invalid config value"),
64            Self::ProcessingNotAvailable => write!(
65                f,
66                "was not compiled with processing, cannot enable processing"
67            ),
68        }
69    }
70}
71
72/// Defines the source of a config error
73#[derive(Debug, Default)]
74enum ConfigErrorSource {
75    /// An error occurring independently.
76    #[default]
77    None,
78    /// An error originating from a configuration file.
79    File(PathBuf),
80    /// An error originating in a field override (an env var, or a CLI parameter).
81    FieldOverride(String),
82}
83
84impl fmt::Display for ConfigErrorSource {
85    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
86        match self {
87            ConfigErrorSource::None => Ok(()),
88            ConfigErrorSource::File(file_name) => {
89                write!(f, " (file {})", file_name.display())
90            }
91            ConfigErrorSource::FieldOverride(name) => write!(f, " (field {name})"),
92        }
93    }
94}
95
96/// Indicates config related errors.
97#[derive(Debug)]
98pub struct ConfigError {
99    source: ConfigErrorSource,
100    kind: ConfigErrorKind,
101}
102
103impl ConfigError {
104    #[inline]
105    fn new(kind: ConfigErrorKind) -> Self {
106        Self {
107            source: ConfigErrorSource::None,
108            kind,
109        }
110    }
111
112    #[inline]
113    fn field(field: &'static str) -> Self {
114        Self {
115            source: ConfigErrorSource::FieldOverride(field.to_owned()),
116            kind: ConfigErrorKind::InvalidValue,
117        }
118    }
119
120    #[inline]
121    fn file(kind: ConfigErrorKind, p: impl AsRef<Path>) -> Self {
122        Self {
123            source: ConfigErrorSource::File(p.as_ref().to_path_buf()),
124            kind,
125        }
126    }
127
128    /// Returns the error kind of the error.
129    pub fn kind(&self) -> ConfigErrorKind {
130        self.kind
131    }
132}
133
134impl fmt::Display for ConfigError {
135    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
136        write!(f, "{}{}", self.kind(), self.source)
137    }
138}
139
140impl Error for ConfigError {}
141
142enum ConfigFormat {
143    Yaml,
144    Json,
145}
146
147impl ConfigFormat {
148    pub fn extension(&self) -> &'static str {
149        match self {
150            ConfigFormat::Yaml => "yml",
151            ConfigFormat::Json => "json",
152        }
153    }
154}
155
156trait ConfigObject: DeserializeOwned + Serialize {
157    /// The format in which to serialize this configuration.
158    fn format() -> ConfigFormat;
159
160    /// The basename of the config file.
161    fn name() -> &'static str;
162
163    /// The full filename of the config file, including the file extension.
164    fn path(base: &Path) -> PathBuf {
165        base.join(format!("{}.{}", Self::name(), Self::format().extension()))
166    }
167
168    /// Loads the config file from a file within the given directory location.
169    fn load(base: &Path) -> anyhow::Result<Self> {
170        let path = Self::path(base);
171
172        let f = fs::File::open(&path)
173            .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotOpenFile, &path))?;
174        let f = io::BufReader::new(f);
175
176        let mut source = {
177            let file = serde_vars::FileSource::default()
178                .with_variable_prefix("${file:")
179                .with_variable_suffix("}")
180                .with_base_path(base);
181            let env = serde_vars::EnvSource::default()
182                .with_variable_prefix("${")
183                .with_variable_suffix("}");
184            (file, env)
185        };
186        match Self::format() {
187            ConfigFormat::Yaml => {
188                serde_vars::deserialize(serde_yaml::Deserializer::from_reader(f), &mut source)
189                    .with_context(|| ConfigError::file(ConfigErrorKind::BadYaml, &path))
190            }
191            ConfigFormat::Json => {
192                serde_vars::deserialize(&mut serde_json::Deserializer::from_reader(f), &mut source)
193                    .with_context(|| ConfigError::file(ConfigErrorKind::BadJson, &path))
194            }
195        }
196    }
197
198    /// Writes the configuration to a file within the given directory location.
199    fn save(&self, base: &Path) -> anyhow::Result<()> {
200        let path = Self::path(base);
201        let mut options = fs::OpenOptions::new();
202        options.write(true).truncate(true).create(true);
203
204        // Remove all non-user permissions for the newly created file
205        #[cfg(unix)]
206        {
207            use std::os::unix::fs::OpenOptionsExt;
208            options.mode(0o600);
209        }
210
211        let mut f = options
212            .open(&path)
213            .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path))?;
214
215        match Self::format() {
216            ConfigFormat::Yaml => {
217                f.write_all(CONFIG_YAML_HEADER.as_bytes())?;
218                serde_yaml::to_writer(&mut f, self)
219                    .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path))?
220            }
221            ConfigFormat::Json => serde_json::to_writer_pretty(&mut f, self)
222                .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path))?,
223        }
224
225        f.write_all(b"\n").ok();
226
227        Ok(())
228    }
229}
230
231/// Structure used to hold information about configuration overrides via
232/// CLI parameters or environment variables
233#[derive(Debug, Default, Clone)]
234pub struct OverridableConfig {
235    /// The operation mode of this relay.
236    pub mode: Option<String>,
237    /// The instance type of this relay.
238    pub instance: Option<String>,
239    /// The log level of this relay.
240    pub log_level: Option<String>,
241    /// The log format of this relay.
242    pub log_format: Option<String>,
243    /// The upstream relay or sentry instance.
244    pub upstream: Option<String>,
245    /// Alternate upstream provided through a Sentry DSN. Key and project will be ignored.
246    pub upstream_dsn: Option<String>,
247    /// The host the relay should bind to (network interface).
248    pub host: Option<String>,
249    /// The port to bind for the unencrypted relay HTTP server.
250    pub port: Option<String>,
251    /// "true" if processing is enabled "false" otherwise
252    pub processing: Option<String>,
253    /// the kafka bootstrap.servers configuration string
254    pub kafka_url: Option<String>,
255    /// the redis server url
256    pub redis_url: Option<String>,
257    /// The globally unique ID of the relay.
258    pub id: Option<String>,
259    /// The secret key of the relay
260    pub secret_key: Option<String>,
261    /// The public key of the relay
262    pub public_key: Option<String>,
263    /// Outcome source
264    pub outcome_source: Option<String>,
265    /// shutdown timeout
266    pub shutdown_timeout: Option<String>,
267    /// Server name reported in the Sentry SDK.
268    pub server_name: Option<String>,
269}
270
271/// The relay credentials
272#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)]
273pub struct Credentials {
274    /// The secret key of the relay
275    pub secret_key: SecretKey,
276    /// The public key of the relay
277    pub public_key: PublicKey,
278    /// The globally unique ID of the relay.
279    pub id: RelayId,
280}
281
282impl Credentials {
283    /// Generates new random credentials.
284    pub fn generate() -> Self {
285        relay_log::info!("generating new relay credentials");
286        let (secret_key, public_key) = generate_key_pair();
287        Self {
288            secret_key,
289            public_key,
290            id: generate_relay_id(),
291        }
292    }
293
294    /// Serializes this configuration to JSON.
295    pub fn to_json_string(&self) -> anyhow::Result<String> {
296        serde_json::to_string(self)
297            .with_context(|| ConfigError::new(ConfigErrorKind::CouldNotWriteFile))
298    }
299}
300
301impl ConfigObject for Credentials {
302    fn format() -> ConfigFormat {
303        ConfigFormat::Json
304    }
305    fn name() -> &'static str {
306        "credentials"
307    }
308}
309
310/// Information on a downstream Relay.
311#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
312#[serde(rename_all = "camelCase")]
313pub struct RelayInfo {
314    /// The public key that this Relay uses to authenticate and sign requests.
315    pub public_key: PublicKey,
316
317    /// Marks an internal relay that has privileged access to more project configuration.
318    #[serde(default)]
319    pub internal: bool,
320}
321
322impl RelayInfo {
323    /// Creates a new RelayInfo
324    pub fn new(public_key: PublicKey) -> Self {
325        Self {
326            public_key,
327            internal: false,
328        }
329    }
330}
331
332/// The operation mode of a relay.
333#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
334#[serde(rename_all = "camelCase")]
335pub enum RelayMode {
336    /// This relay acts as a proxy for all requests and events.
337    ///
338    /// Events are normalized and rate limits from the upstream are enforced, but the relay will not
339    /// fetch project configurations from the upstream or perform PII stripping. All events are
340    /// accepted unless overridden on the file system.
341    Proxy,
342
343    /// Project configurations are managed by the upstream.
344    ///
345    /// Project configurations are always fetched from the upstream, unless they are statically
346    /// overridden in the file system. This relay must be allowed in the upstream Sentry. This is
347    /// only possible, if the upstream is Sentry directly, or another managed Relay.
348    Managed,
349}
350
351impl<'de> Deserialize<'de> for RelayMode {
352    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
353    where
354        D: Deserializer<'de>,
355    {
356        let s = String::deserialize(deserializer)?;
357        match s.as_str() {
358            "proxy" => Ok(RelayMode::Proxy),
359            "managed" => Ok(RelayMode::Managed),
360            "static" => Err(serde::de::Error::custom(
361                "Relay mode 'static' has been removed. Please use 'managed' or 'proxy' instead.",
362            )),
363            other => Err(serde::de::Error::unknown_variant(
364                other,
365                &["proxy", "managed"],
366            )),
367        }
368    }
369}
370
371impl fmt::Display for RelayMode {
372    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
373        match self {
374            RelayMode::Proxy => write!(f, "proxy"),
375            RelayMode::Managed => write!(f, "managed"),
376        }
377    }
378}
379
380/// The instance type of Relay.
381#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)]
382#[serde(rename_all = "camelCase")]
383pub enum RelayInstance {
384    /// This Relay is run as a default instance.
385    Default,
386
387    /// This Relay is run as a canary instance where experiments can be run.
388    Canary,
389}
390
391impl RelayInstance {
392    /// Returns `true` if the [`RelayInstance`] is of type [`RelayInstance::Canary`].
393    pub fn is_canary(&self) -> bool {
394        matches!(self, RelayInstance::Canary)
395    }
396}
397
398impl fmt::Display for RelayInstance {
399    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
400        match self {
401            RelayInstance::Default => write!(f, "default"),
402            RelayInstance::Canary => write!(f, "canary"),
403        }
404    }
405}
406
407impl FromStr for RelayInstance {
408    type Err = fmt::Error;
409
410    fn from_str(s: &str) -> Result<Self, Self::Err> {
411        match s {
412            "canary" => Ok(RelayInstance::Canary),
413            _ => Ok(RelayInstance::Default),
414        }
415    }
416}
417
418/// Error returned when parsing an invalid [`RelayMode`].
419#[derive(Clone, Copy, Debug, Eq, PartialEq)]
420pub struct ParseRelayModeError;
421
422impl fmt::Display for ParseRelayModeError {
423    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
424        write!(f, "Relay mode must be one of: managed or proxy")
425    }
426}
427
428impl Error for ParseRelayModeError {}
429
430impl FromStr for RelayMode {
431    type Err = ParseRelayModeError;
432
433    fn from_str(s: &str) -> Result<Self, Self::Err> {
434        match s {
435            "proxy" => Ok(RelayMode::Proxy),
436            "managed" => Ok(RelayMode::Managed),
437            _ => Err(ParseRelayModeError),
438        }
439    }
440}
441
442/// Returns `true` if this value is equal to `Default::default()`.
443fn is_default<T: Default + PartialEq>(t: &T) -> bool {
444    *t == T::default()
445}
446
447/// Checks if we are running in docker.
448fn is_docker() -> bool {
449    if fs::metadata("/.dockerenv").is_ok() {
450        return true;
451    }
452
453    fs::read_to_string("/proc/self/cgroup").is_ok_and(|s| s.contains("/docker"))
454}
455
456/// Default value for the "bind" configuration.
457fn default_host() -> IpAddr {
458    if is_docker() {
459        // Docker images rely on this service being exposed
460        "0.0.0.0".parse().unwrap()
461    } else {
462        "127.0.0.1".parse().unwrap()
463    }
464}
465
466/// Controls responses from the readiness health check endpoint based on authentication.
467///
468/// Independent of the the readiness condition, shutdown always switches Relay into unready state.
469#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)]
470#[serde(rename_all = "lowercase")]
471#[derive(Default)]
472pub enum ReadinessCondition {
473    /// (default) Relay is ready when authenticated and connected to the upstream.
474    ///
475    /// Before authentication has succeeded and during network outages, Relay responds as not ready.
476    /// Relay reauthenticates based on the `http.auth_interval` parameter. During reauthentication,
477    /// Relay remains ready until authentication fails.
478    ///
479    /// Authentication is only required for Relays in managed mode. Other Relays will only check for
480    /// network outages.
481    #[default]
482    Authenticated,
483    /// Relay reports readiness regardless of the authentication and networking state.
484    Always,
485}
486
487/// Relay specific configuration values.
488#[derive(Serialize, Deserialize, Debug, Clone)]
489#[serde(default)]
490pub struct Relay {
491    /// The operation mode of this Relay.
492    pub mode: RelayMode,
493    /// The instance type of this Relay.
494    pub instance: RelayInstance,
495    /// The upstream Relay or Sentry instance.
496    pub upstream: UpstreamDescriptor,
497    /// The upstream advertised to downstream Relay instances.
498    ///
499    /// This value will be advertised to downstream Relays as the upstream to use when forwarding
500    /// data. It can be used for traffic routing and balancing, it must not redirect to a different
501    /// Sentry instance.
502    ///
503    /// Downstream Relays will treat the advertised upstream as the same logical component as this instance
504    /// and re-use already established authentication keys.
505    pub advertised_upstream: Option<UpstreamDescriptor>,
506    /// The host the relay should bind to (network interface).
507    pub host: IpAddr,
508    /// The port to bind for the unencrypted relay HTTP server.
509    pub port: u16,
510    /// The host the relay should bind to (network interface) for internally exposed APIs, like
511    /// health checks.
512    ///
513    /// If not configured, internal routes are exposed on the main HTTP server.
514    ///
515    /// Note: configuring the internal http server on an address which overlaps with the main
516    /// server (e.g. main on `0.0.0.0:3000` and internal on `127.0.0.1:3000`) is a misconfiguration
517    /// resulting in approximately half of the requests sent to `127.0.0.1:3000` to fail, as the handling
518    /// http server is chosen by the operating system 'at random'.
519    ///
520    /// As a best practice you should always choose different ports to avoid this issue.
521    ///
522    /// Defaults to [`Self::host`].
523    pub internal_host: Option<IpAddr>,
524    /// The port to bind for internally exposed APIs.
525    ///
526    /// Defaults to [`Self::port`].
527    pub internal_port: Option<u16>,
528    /// Always override project IDs from the URL and DSN with the identifier used at the upstream.
529    ///
530    /// Enable this setting for Relays used to redirect traffic to a migrated Sentry instance.
531    /// Validation of project identifiers can be safely skipped in these cases.
532    #[serde(skip_serializing_if = "is_default")]
533    pub override_project_ids: bool,
534    /// Interval in seconds for Relay to check if its configuration changed.
535    ///
536    /// If configured Relay will periodically check its configuration for changes
537    /// and hot reload it.
538    ///
539    /// Hot reloading is only supported for a limited set of values.
540    ///
541    /// Defaults to `None` / off.
542    pub config_reload_interval: Option<u64>,
543}
544
545impl Default for Relay {
546    fn default() -> Self {
547        Relay {
548            mode: RelayMode::Managed,
549            instance: RelayInstance::Default,
550            upstream: "https://sentry.io/".parse().unwrap(),
551            advertised_upstream: None,
552            host: default_host(),
553            port: 3000,
554            internal_host: None,
555            internal_port: None,
556            override_project_ids: false,
557            config_reload_interval: None,
558        }
559    }
560}
561
562/// Control the metrics.
563#[derive(Serialize, Deserialize, Debug, Clone)]
564#[serde(default)]
565pub struct Metrics {
566    /// Hostname and port of the statsd server.
567    ///
568    /// Defaults to `None`.
569    pub statsd: Option<String>,
570    /// Buffer size used for metrics sent to the statsd socket.
571    ///
572    /// Defaults to `None`.
573    pub statsd_buffer_size: Option<usize>,
574    /// Common prefix that should be added to all metrics.
575    ///
576    /// Defaults to `"sentry.relay"`.
577    pub prefix: String,
578    /// Default tags to apply to all metrics.
579    pub default_tags: BTreeMap<String, String>,
580    /// Tag name to report the hostname to for each metric. Defaults to not sending such a tag.
581    pub hostname_tag: Option<String>,
582    /// Interval for periodic metrics emitted from Relay.
583    ///
584    /// Setting it to `0` seconds disables the periodic metrics.
585    /// Defaults to 5 seconds.
586    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/// Controls various limits
603#[derive(Serialize, Deserialize, Debug, Clone)]
604#[serde(default)]
605pub struct Limits {
606    /// How many requests can be sent concurrently from Relay to the upstream before Relay starts
607    /// buffering.
608    pub max_concurrent_requests: usize,
609    /// How many queries can be sent concurrently from Relay to the upstream before Relay starts
610    /// buffering.
611    ///
612    /// The concurrency of queries is additionally constrained by `max_concurrent_requests`.
613    pub max_concurrent_queries: usize,
614    /// The maximum payload size for events.
615    pub max_event_size: ByteSize,
616    /// The maximum size for each attachment.
617    pub max_attachment_size: ByteSize,
618    /// The maximum amount of attachments in a single envelope.
619    pub max_attachment_count: usize,
620    /// The maximum combined size for all attachments in an envelope or request.
621    pub max_attachments_size: ByteSize,
622    /// The maximum size for a TUS upload request body.
623    pub max_upload_size: ByteSize,
624    /// The maximum combined size for all client reports in an envelope or request.
625    pub max_client_reports_size: ByteSize,
626    /// The maximum number of client report items per envelope.
627    pub max_client_reports_count: usize,
628    /// The maximum payload size for a monitor check-in.
629    pub max_check_in_size: ByteSize,
630    /// The maximum payload size for an entire envelopes. Individual limits still apply.
631    pub max_envelope_size: ByteSize,
632    /// The maximum combined size for all sessions in an envelope in bytes.
633    pub max_sessions_size: ByteSize,
634    /// The maximum number of session items per envelope.
635    pub max_session_count: usize,
636    /// The maximum payload size for general API requests.
637    pub max_api_payload_size: ByteSize,
638    /// The maximum payload size for file uploads and chunks.
639    pub max_api_file_upload_size: ByteSize,
640    /// The maximum payload size for chunks
641    pub max_api_chunk_upload_size: ByteSize,
642    /// The maximum payload size for a profile
643    pub max_profile_size: ByteSize,
644    /// The maximum payload size for a trace metric.
645    pub max_trace_metric_size: ByteSize,
646    /// The maximum payload size for a log.
647    pub max_log_size: ByteSize,
648    /// The maximum number of operations that can occur in a span expansion.
649    pub max_expanded_log_operations: usize,
650    /// The maximum number of operations that can occur in a span expansion.
651    pub max_expanded_span_operations: usize,
652    /// The maximum number of operations that can occur in a trace metric expansion.
653    pub max_expanded_trace_metric_operations: usize,
654    /// The maximum payload size for a span.
655    pub max_span_size: ByteSize,
656    /// The maximum amount of standalone transaction spans per envelope.
657    pub max_standalone_span_count: usize,
658    /// The maximum payload size for an item container.
659    pub max_container_size: ByteSize,
660    /// The maximum payload size for metric buckets.
661    pub max_metric_buckets_size: ByteSize,
662    /// The maximum payload size for a compressed replay.
663    pub max_replay_compressed_size: ByteSize,
664    /// The maximum payload size for an uncompressed replay.
665    #[serde(alias = "max_replay_size")]
666    max_replay_uncompressed_size: ByteSize,
667    /// The maximum size for a replay recording Kafka message.
668    pub max_replay_message_size: ByteSize,
669    /// The byte size limit up to which Relay will retain
670    /// keys of invalid/removed attributes.
671    ///
672    /// This is only relevant for EAP items (spans, logs, …).
673    /// In principle, we want to record all deletions of attributes,
674    /// but we have to institute some limit to protect our infrastructure
675    /// against excessive metadata sizes.
676    ///
677    /// Defaults to 10KiB.
678    pub max_removed_attribute_key_size: ByteSize,
679    /// The maximum number of threads to spawn for CPU and web work, each.
680    ///
681    /// The total number of threads spawned will roughly be `2 * max_thread_count`. Defaults to
682    /// the number of logical CPU cores on the host.
683    pub max_thread_count: usize,
684    /// Controls the maximum concurrency of each worker thread.
685    ///
686    /// Increasing the concurrency, can lead to a better utilization of worker threads by
687    /// increasing the amount of I/O done concurrently.
688    //
689    /// Currently has no effect on defaults to `1`.
690    pub max_pool_concurrency: usize,
691    /// The maximum number of seconds a query is allowed to take across retries. Individual requests
692    /// have lower timeouts. Defaults to 30 seconds.
693    pub query_timeout: u64,
694    /// The maximum number of seconds to wait for pending envelopes after receiving a shutdown
695    /// signal.
696    pub shutdown_timeout: u64,
697    /// Server keep-alive timeout in seconds.
698    ///
699    /// By default, keep-alive is set to 5 seconds.
700    pub keepalive_timeout: u64,
701    /// Server idle timeout in seconds.
702    ///
703    /// The idle timeout limits the amount of time a connection is kept open without activity.
704    /// Setting this too short may abort connections before Relay is able to send a response.
705    ///
706    /// By default there is no idle timeout.
707    pub idle_timeout: Option<u64>,
708    /// Sets the maximum number of concurrent connections.
709    ///
710    /// Upon reaching the limit, the server will stop accepting connections.
711    ///
712    /// By default there is no limit.
713    pub max_connections: Option<usize>,
714    /// The TCP listen backlog.
715    ///
716    /// Configures the TCP listen backlog for the listening socket of Relay.
717    /// See [`man listen(2)`](https://man7.org/linux/man-pages/man2/listen.2.html)
718    /// for a more detailed description of the listen backlog.
719    ///
720    /// Defaults to `1024`, a value [google has been using for a long time](https://git.kernel.org/pub/scm/linux/kernel/git/torvalds/linux.git/commit/?id=19f92a030ca6d772ab44b22ee6a01378a8cb32d4).
721    pub tcp_listen_backlog: u32,
722}
723
724impl Default for Limits {
725    fn default() -> Self {
726        Limits {
727            max_concurrent_requests: 100,
728            max_concurrent_queries: 5,
729            max_event_size: ByteSize::mebibytes(1),
730            max_attachment_size: ByteSize::mebibytes(200),
731            max_attachment_count: 30,
732            max_attachments_size: ByteSize::mebibytes(200),
733            max_upload_size: ByteSize::mebibytes(1024),
734            max_client_reports_size: ByteSize::kibibytes(100),
735            max_client_reports_count: 100,
736            max_check_in_size: ByteSize::kibibytes(100),
737            max_envelope_size: ByteSize::mebibytes(200),
738            max_sessions_size: ByteSize::mebibytes(10),
739            max_session_count: 100,
740            max_api_payload_size: ByteSize::mebibytes(20),
741            max_api_file_upload_size: ByteSize::mebibytes(40),
742            max_api_chunk_upload_size: ByteSize::mebibytes(100),
743            max_profile_size: ByteSize::mebibytes(50),
744            max_trace_metric_size: ByteSize::mebibytes(1),
745            max_log_size: ByteSize::mebibytes(2),
746            max_expanded_log_operations: 3_000_000,
747            max_expanded_span_operations: 3_000_000,
748            max_expanded_trace_metric_operations: 3_000_000,
749            max_span_size: ByteSize::mebibytes(10),
750            max_standalone_span_count: 25,
751            max_container_size: ByteSize::mebibytes(12),
752            max_metric_buckets_size: ByteSize::mebibytes(1),
753            max_replay_compressed_size: ByteSize::mebibytes(10),
754            max_replay_uncompressed_size: ByteSize::mebibytes(100),
755            max_replay_message_size: ByteSize::mebibytes(15),
756            max_thread_count: num_cpus::get(),
757            max_pool_concurrency: 1,
758            query_timeout: 30,
759            shutdown_timeout: 10,
760            keepalive_timeout: 5,
761            idle_timeout: None,
762            max_connections: None,
763            tcp_listen_backlog: 1024,
764            max_removed_attribute_key_size: ByteSize::kibibytes(10),
765        }
766    }
767}
768
769/// Controls traffic steering.
770#[derive(Debug, Default, Deserialize, Serialize, Clone)]
771#[serde(default)]
772pub struct Routing {
773    /// Accept and forward unknown Envelope items to the upstream.
774    ///
775    /// Forwarding unknown items should be enabled in most cases to allow proxying traffic for newer
776    /// SDK versions. The upstream in Sentry makes the final decision on which items are valid. If
777    /// this is disabled, just the unknown items are removed from Envelopes, and the rest is
778    /// processed as usual.
779    ///
780    /// Defaults to `true` for all Relay modes other than processing mode. In processing mode, this
781    /// is disabled by default since the item cannot be handled.
782    pub accept_unknown_items: Option<bool>,
783}
784
785/// Http content encoding for both incoming and outgoing web requests.
786#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize)]
787#[serde(rename_all = "lowercase")]
788pub enum HttpEncoding {
789    /// Identity function without no compression.
790    ///
791    /// This is the default encoding and does not require the presence of the `content-encoding`
792    /// HTTP header.
793    #[default]
794    Identity,
795    /// Compression using a [zlib](https://en.wikipedia.org/wiki/Zlib) structure with
796    /// [deflate](https://en.wikipedia.org/wiki/DEFLATE) encoding.
797    ///
798    /// These structures are defined in [RFC 1950](https://datatracker.ietf.org/doc/html/rfc1950)
799    /// and [RFC 1951](https://datatracker.ietf.org/doc/html/rfc1951).
800    Deflate,
801    /// A format using the [Lempel-Ziv coding](https://en.wikipedia.org/wiki/LZ77_and_LZ78#LZ77)
802    /// (LZ77), with a 32-bit CRC.
803    ///
804    /// This is the original format of the UNIX gzip program. The HTTP/1.1 standard also recommends
805    /// that the servers supporting this content-encoding should recognize `x-gzip` as an alias, for
806    /// compatibility purposes.
807    Gzip,
808    /// A format using the [Brotli](https://en.wikipedia.org/wiki/Brotli) algorithm.
809    Br,
810    /// A format using the [Zstd](https://en.wikipedia.org/wiki/Zstd) compression algorithm.
811    Zstd,
812}
813
814impl HttpEncoding {
815    /// Parses a [`HttpEncoding`] from its `content-encoding` header value.
816    pub fn parse(str: &str) -> Self {
817        let str = str.trim();
818        if str.eq_ignore_ascii_case("zstd") {
819            Self::Zstd
820        } else if str.eq_ignore_ascii_case("br") {
821            Self::Br
822        } else if str.eq_ignore_ascii_case("gzip") || str.eq_ignore_ascii_case("x-gzip") {
823            Self::Gzip
824        } else if str.eq_ignore_ascii_case("deflate") {
825            Self::Deflate
826        } else {
827            Self::Identity
828        }
829    }
830
831    /// Returns the value for the `content-encoding` HTTP header.
832    ///
833    /// Returns `None` for [`Identity`](Self::Identity), and `Some` for other encodings.
834    pub fn name(&self) -> Option<&'static str> {
835        match self {
836            Self::Identity => None,
837            Self::Deflate => Some("deflate"),
838            Self::Gzip => Some("gzip"),
839            Self::Br => Some("br"),
840            Self::Zstd => Some("zstd"),
841        }
842    }
843}
844
845/// Controls authentication with upstream.
846#[derive(Serialize, Deserialize, Debug, Clone)]
847#[serde(default)]
848pub struct Http {
849    /// Timeout for upstream requests in seconds.
850    ///
851    /// This timeout covers the time from sending the request until receiving response headers.
852    /// Neither the connection process and handshakes, nor reading the response body is covered in
853    /// this timeout.
854    pub timeout: u32,
855    /// Timeout for establishing connections with the upstream in seconds.
856    ///
857    /// This includes SSL handshakes. Relay reuses connections when the upstream supports connection
858    /// keep-alive. Connections are retained for a maximum 75 seconds, or 15 seconds of inactivity.
859    pub connection_timeout: u32,
860    /// Maximum interval between failed request retries in seconds.
861    pub max_retry_interval: u32,
862    /// The custom HTTP Host header to send to the upstream.
863    pub host_header: Option<String>,
864    /// The interval in seconds at which Relay attempts to reauthenticate with the upstream server.
865    ///
866    /// Re-authentication happens even when Relay is idle. If authentication fails, Relay reverts
867    /// back into startup mode and tries to establish a connection. During this time, incoming
868    /// envelopes will be buffered.
869    ///
870    /// Defaults to `600` (10 minutes).
871    pub auth_interval: Option<u64>,
872    /// The maximum time of experiencing uninterrupted network failures until Relay considers that
873    /// it has encountered a network outage in seconds.
874    ///
875    /// During a network outage relay will try to reconnect and will buffer all upstream messages
876    /// until it manages to reconnect.
877    pub outage_grace_period: u64,
878    /// The time Relay waits before retrying an upstream request, in seconds.
879    ///
880    /// This time is only used before going into a network outage mode.
881    pub retry_delay: u64,
882    /// The interval in seconds for continued failed project fetches at which Relay will error.
883    ///
884    /// A successful fetch resets this interval. Relay does nothing during long
885    /// times without emitting requests.
886    pub project_failure_interval: u64,
887    /// Content encoding to apply to upstream store requests.
888    ///
889    /// By default, Relay applies `zstd` content encoding to compress upstream requests. Compression
890    /// can be disabled to reduce CPU consumption, but at the expense of increased network traffic.
891    ///
892    /// This setting applies to all store requests of SDK data, including events, transactions,
893    /// envelopes and sessions. At the moment, this does not apply to Relay's internal queries.
894    ///
895    /// Available options are:
896    ///
897    ///  - `identity`: Disables compression.
898    ///  - `deflate`: Compression using a zlib header with deflate encoding.
899    ///  - `gzip` (default): Compression using gzip.
900    ///  - `br`: Compression using the brotli algorithm.
901    ///  - `zstd`: Compression using the zstd algorithm.
902    pub encoding: HttpEncoding,
903    /// Submit metrics globally through a shared endpoint.
904    ///
905    /// As opposed to regular envelopes which are sent to an endpoint inferred from the project's
906    /// DSN, this submits metrics to the global endpoint with Relay authentication.
907    ///
908    /// This option does not have any effect on processing mode.
909    pub global_metrics: bool,
910    /// Controls whether the forward endpoint is enabled.
911    ///
912    /// The forward endpoint forwards unknown API requests to the upstream.
913    ///
914    /// Relay instances with processing enabled are expected to support the latest API and do never
915    /// support forwarding requests to Sentry.
916    pub forward: bool,
917    /// Enables an async DNS resolver through the `hickory-dns` crate, which uses an LRU cache for
918    /// the resolved entries. This helps to limit the amount of requests made to the upstream DNS
919    /// server (important for K8s infrastructure).
920    pub dns_cache: bool,
921}
922
923impl Default for Http {
924    fn default() -> Self {
925        Http {
926            timeout: 5,
927            connection_timeout: 3,
928            max_retry_interval: 60, // 1 minute
929            host_header: None,
930            auth_interval: Some(600), // 10 minutes
931            outage_grace_period: DEFAULT_NETWORK_OUTAGE_GRACE_PERIOD,
932            retry_delay: 1,
933            project_failure_interval: 90,
934            encoding: HttpEncoding::Zstd,
935            global_metrics: false,
936            forward: true,
937            dns_cache: true,
938        }
939    }
940}
941
942/// Strategy used to assign envelopes to buffer partitions.
943#[derive(Clone, Copy, Debug, Eq, PartialEq, Default, Deserialize, Serialize)]
944#[serde(rename_all = "snake_case")]
945pub enum EnvelopeSpoolPartitioning {
946    /// Envelopes with the same project key pair land on the same partition.
947    ///
948    /// Keeps per-project state, disk files, and event ordering co-located on one partition.
949    ProjectKeyPair,
950    /// Envelopes are distributed across partitions in a round-robin fashion (default).
951    ///
952    /// This prevents "hot" partitions when a single project pair dominates traffic, but has
953    /// trade-offs:
954    /// - Per-project LIFO ordering is no longer preserved across partitions.
955    /// - Per-partition memory footprint grows since every partition sees every project.
956    #[default]
957    RoundRobin,
958}
959
960/// Persistent buffering configuration for incoming envelopes.
961#[derive(Debug, Serialize, Deserialize, Clone)]
962#[serde(default)]
963pub struct EnvelopeSpool {
964    /// The path of the SQLite database file(s) which persist the data.
965    ///
966    /// Based on the number of partitions, more database files will be created within the same path.
967    ///
968    /// If not set, the envelopes will be buffered in memory.
969    pub path: Option<PathBuf>,
970    /// The maximum size of the buffer to keep, in bytes.
971    ///
972    /// When the on-disk buffer reaches this size, new envelopes will be dropped.
973    ///
974    /// Defaults to 500MB.
975    pub max_disk_size: ByteSize,
976    /// Size of the batch of compressed envelopes that are spooled to disk at once.
977    ///
978    /// Note that this is the size after which spooling will be triggered but it does not guarantee
979    /// that exactly this size will be spooled, it can be greater or equal.
980    ///
981    /// Defaults to 10 KiB.
982    pub batch_size_bytes: ByteSize,
983    /// Time after which a batch is flushed, regardless of batch size.
984    ///
985    /// The age of the batch is only checked when a new envelope comes in, but in practice this
986    /// has the desired effect: High-volume projects always form full batches, low-volume batches
987    /// flush individual envelopes to keep memory usage low.
988    pub flush_timeout_secs: Option<u64>,
989    /// Maximum time between receiving the envelope and processing it.
990    ///
991    /// When envelopes spend too much time in the buffer (e.g. because their project cannot be loaded),
992    /// they are dropped.
993    ///
994    /// Defaults to 24h.
995    pub max_envelope_delay_secs: u64,
996    /// The refresh frequency in ms of how frequently disk usage is updated by querying SQLite
997    /// internal page stats.
998    ///
999    /// Defaults to 100ms.
1000    pub disk_usage_refresh_frequency_ms: u64,
1001    /// The relative memory usage above which the buffer service will stop dequeueing envelopes.
1002    ///
1003    /// Only applies when [`Self::path`] is set.
1004    ///
1005    /// This value should be lower than [`Health::max_memory_percent`] to prevent flip-flopping.
1006    ///
1007    /// Warning: This threshold can cause the buffer service to deadlock when the buffer consumes
1008    /// excessive memory (as influenced by [`Self::batch_size_bytes`]).
1009    ///
1010    /// This scenario arises when the buffer stops spooling due to reaching the
1011    /// [`Self::max_backpressure_memory_percent`] limit, but the batch threshold for spooling
1012    /// ([`Self::batch_size_bytes`]) is never reached. As a result, no data is spooled, memory usage
1013    /// continues to grow, and the system becomes deadlocked.
1014    ///
1015    /// ### Example
1016    /// Suppose the system has 1GB of available memory and is configured to spool only after
1017    /// accumulating 10GB worth of envelopes. If Relay consumes 900MB of memory, it will stop
1018    /// unspooling due to reaching the [`Self::max_backpressure_memory_percent`] threshold.
1019    ///
1020    /// However, because the buffer hasn't accumulated the 10GB needed to trigger spooling,
1021    /// no data will be offloaded. Memory usage keeps increasing until it hits the
1022    /// [`Health::max_memory_percent`] threshold, e.g., at 950MB. At this point:
1023    ///
1024    /// - No more envelopes are accepted.
1025    /// - The buffer remains stuck, as unspooling won’t resume until memory drops below 900MB which
1026    ///   will not happen.
1027    /// - A deadlock occurs, with the system unable to recover without manual intervention.
1028    ///
1029    /// Defaults to 90% (5% less than max memory).
1030    pub max_backpressure_memory_percent: f32,
1031    /// Number of partitions of the buffer.
1032    ///
1033    /// A partition is a separate instance of the buffer which has its own isolated queue, stacks
1034    /// and other resources.
1035    ///
1036    /// Defaults to 1.
1037    pub partitions: NonZeroU8,
1038    /// Strategy used to assign envelopes to buffer partitions.
1039    ///
1040    /// Defaults to partitioning by `ProjectKeyPair`, which keeps all envelopes of a given project
1041    /// pair on the same partition. See [`EnvelopeSpoolPartitioning`] for alternatives and
1042    /// trade-offs.
1043    pub partitioning: EnvelopeSpoolPartitioning,
1044    /// Maximum number of envelopes unspooled from disk per second.
1045    ///
1046    /// The limit applies per Relay instance, not per spool partition.
1047    ///
1048    /// Defaults to `None`, which disables the limit.
1049    pub max_unspool_envelopes_per_second: Option<NonZeroU32>,
1050    /// Whether the database defined in `path` is on an ephemeral storage disk.
1051    ///
1052    /// With `ephemeral: true`, Relay does not spool in-flight data to disk
1053    /// during graceful shutdown. Instead, it attempts to process all data before it terminates.
1054    ///
1055    /// Defaults to `false`.
1056    pub ephemeral: bool,
1057}
1058
1059impl Default for EnvelopeSpool {
1060    fn default() -> Self {
1061        Self {
1062            path: None,
1063            max_disk_size: ByteSize::mebibytes(500),
1064            batch_size_bytes: ByteSize::kibibytes(10),
1065            max_envelope_delay_secs: 24 * 60 * 60,
1066            disk_usage_refresh_frequency_ms: 100,
1067            max_backpressure_memory_percent: 0.8,
1068            partitions: NonZeroU8::new(1).unwrap(),
1069            partitioning: EnvelopeSpoolPartitioning::default(),
1070            ephemeral: false,
1071            flush_timeout_secs: None,
1072            max_unspool_envelopes_per_second: None,
1073        }
1074    }
1075}
1076
1077/// Persistent buffering configuration.
1078#[derive(Debug, Serialize, Deserialize, Default, Clone)]
1079#[serde(default)]
1080pub struct Spool {
1081    /// Configuration for envelope spooling.
1082    pub envelopes: EnvelopeSpool,
1083}
1084
1085/// Controls internal caching behavior.
1086#[derive(Serialize, Deserialize, Debug, Clone)]
1087#[serde(default)]
1088pub struct Cache {
1089    /// The full project state will be requested by this Relay if set to `true`.
1090    ///
1091    /// Relay instances that receive the full project config have full access to quota config
1092    /// and perform dynamic sampling.
1093    pub project_request_full_config: bool,
1094    /// The cache timeout for project configurations in seconds.
1095    pub project_expiry: u32,
1096    /// Continue using project state this many seconds after cache expiry while a new state is
1097    /// being fetched. This is added on top of `project_expiry`.
1098    ///
1099    /// Default is 2 minutes.
1100    pub project_grace_period: u32,
1101    /// Refresh a project after the specified seconds.
1102    ///
1103    /// The time must be between expiry time and the grace period.
1104    ///
1105    /// By default there are no refreshes enabled.
1106    pub project_refresh_interval: Option<u32>,
1107    /// The cache timeout for downstream relay info (public keys) in seconds.
1108    pub relay_expiry: u32,
1109    /// The cache timeout for non-existing entries.
1110    pub miss_expiry: u32,
1111    /// The buffer timeout for batched project config queries before sending them upstream in ms.
1112    pub batch_interval: u32,
1113    /// The buffer timeout for batched queries of downstream relays in ms. Defaults to 100ms.
1114    pub downstream_relays_batch_interval: u32,
1115    /// The maximum number of project configs to fetch from Sentry at once. Defaults to 500.
1116    ///
1117    /// `cache.batch_interval` controls how quickly batches are sent, this controls the batch size.
1118    pub batch_size: usize,
1119    /// Interval for watching local cache override files in seconds.
1120    pub file_interval: u32,
1121    /// Interval for fetching new global configs from the upstream, in seconds.
1122    pub global_config_fetch_interval: u32,
1123}
1124
1125impl Default for Cache {
1126    fn default() -> Self {
1127        Cache {
1128            project_request_full_config: false,
1129            project_expiry: 300,       // 5 minutes
1130            project_grace_period: 120, // 2 minutes
1131            project_refresh_interval: None,
1132            relay_expiry: 3600,                    // 1 hour
1133            miss_expiry: 60,                       // 1 minute
1134            batch_interval: 100,                   // 100ms
1135            downstream_relays_batch_interval: 100, // 100ms
1136            batch_size: 500,
1137            file_interval: 10,                // 10 seconds
1138            global_config_fetch_interval: 10, // 10 seconds
1139        }
1140    }
1141}
1142
1143/// Controls Sentry-internal event processing.
1144#[derive(Serialize, Deserialize, Debug, Clone)]
1145#[serde(default)]
1146pub struct Processing {
1147    /// True if the Relay should do processing. Defaults to `false`.
1148    pub enabled: bool,
1149    /// GeoIp DB file source.
1150    pub geoip_path: Option<PathBuf>,
1151    /// Maximum future timestamp of ingested events.
1152    pub max_secs_in_future: u32,
1153    /// Maximum age of ingested sessions. Older sessions will be dropped.
1154    pub max_session_secs_in_past: u32,
1155    /// Kafka producer configurations.
1156    pub kafka_config: Vec<KafkaConfigParam>,
1157    /// Additional kafka producer configurations.
1158    ///
1159    /// The `kafka_config` is the default producer configuration used for all topics. A secondary
1160    /// kafka config can be referenced in `topics:` like this:
1161    ///
1162    /// ```yaml
1163    /// secondary_kafka_configs:
1164    ///   mycustomcluster:
1165    ///     - name: 'bootstrap.servers'
1166    ///       value: 'sentry_kafka_metrics:9093'
1167    ///
1168    /// topics:
1169    ///   transactions: ingest-transactions
1170    ///   metrics:
1171    ///     name: ingest-metrics
1172    ///     config: mycustomcluster
1173    /// ```
1174    ///
1175    /// Then metrics will be produced to an entirely different Kafka cluster.
1176    pub secondary_kafka_configs: BTreeMap<String, Vec<KafkaConfigParam>>,
1177    /// Kafka topic names.
1178    pub topics: TopicAssignments,
1179    /// Whether to validate the supplied topics by calling Kafka's metadata endpoints.
1180    pub kafka_validate_topics: bool,
1181    /// Redis hosts to connect to for storing state for rate limits.
1182    pub redis: Option<RedisConfigs>,
1183    /// Maximum chunk size of attachments for Kafka.
1184    pub attachment_chunk_size: ByteSize,
1185    /// Prefix to use when looking up project configs in Redis. Defaults to "relayconfig".
1186    pub projectconfig_cache_prefix: String,
1187    /// Maximum rate limit to report to clients.
1188    pub max_rate_limit: Option<u32>,
1189    /// Configures the quota cache ratio between `0.0` and `1.0`.
1190    ///
1191    /// The quota cache, caches the specified ratio of remaining quota in memory to reduce the
1192    /// amount of synchronizations required with Redis.
1193    ///
1194    /// The ratio is applied to the (per second) rate of the quota, not the total limit.
1195    /// For example a quota with limit 100 with a 10 second window is treated equally to a quota of
1196    /// 10 with a 1 second window.
1197    ///
1198    /// By default quota caching is disabled.
1199    pub quota_cache_ratio: Option<f32>,
1200    /// Relative amount of the total quota limit to which quota caching is applied.
1201    ///
1202    /// If exceeded, the rate limiter will no longer cache the quota and sync with Redis on every call instead.
1203    /// Lowering this value reduces the probability of incorrectly over-accepting.
1204    ///
1205    /// Must be between `0.0` and `1.0`, by default there is no limit configured.
1206    pub quota_cache_max: Option<f32>,
1207    /// Configuration for the objectstore service.
1208    #[serde(alias = "upload")]
1209    pub objectstore: ObjectstoreServiceConfig,
1210}
1211
1212impl Default for Processing {
1213    /// Constructs a disabled processing configuration.
1214    fn default() -> Self {
1215        Self {
1216            enabled: false,
1217            geoip_path: None,
1218            max_secs_in_future: 60,                  // 1 minute
1219            max_session_secs_in_past: 5 * 24 * 3600, // 5 days
1220            kafka_config: Vec::new(),
1221            secondary_kafka_configs: BTreeMap::new(),
1222            topics: TopicAssignments::default(),
1223            kafka_validate_topics: false,
1224            redis: None,
1225            attachment_chunk_size: ByteSize::mebibytes(1),
1226            projectconfig_cache_prefix: "relayconfig".to_owned(),
1227            max_rate_limit: Some(300), // 5 minutes
1228            quota_cache_ratio: None,
1229            quota_cache_max: None,
1230            objectstore: ObjectstoreServiceConfig::default(),
1231        }
1232    }
1233}
1234
1235/// Configuration for normalization in this Relay.
1236#[derive(Debug, Default, Serialize, Deserialize, Clone)]
1237#[serde(default)]
1238pub struct Normalization {
1239    /// Level of normalization for Relay to apply to incoming data.
1240    pub level: NormalizationLevel,
1241}
1242
1243/// Configuration for the level of normalization this Relay should do.
1244#[derive(Copy, Clone, Debug, Default, Serialize, Deserialize, Eq, PartialEq)]
1245#[serde(rename_all = "lowercase")]
1246pub enum NormalizationLevel {
1247    /// Runs normalization, excluding steps that break future compatibility.
1248    ///
1249    /// Processing Relays run [`NormalizationLevel::Full`] if this option is set.
1250    #[default]
1251    Default,
1252    /// Run full normalization.
1253    ///
1254    /// It includes steps that break future compatibility and should only run in
1255    /// the last layer of relays.
1256    Full,
1257}
1258
1259/// Configuration options for objectstore's auth scheme.
1260#[derive(Serialize, Deserialize, Clone)]
1261pub struct ObjectstoreAuthConfig {
1262    /// Identifier for the private key used to sign objectstore's tokens. Must correspond to a
1263    /// public key configured in objectstore.
1264    pub key_id: String,
1265
1266    /// EdDSA private key used to sign Objectstore's tokens, in PEM format.
1267    pub signing_key: String,
1268}
1269
1270impl fmt::Debug for ObjectstoreAuthConfig {
1271    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1272        f.debug_struct("ObjectstoreAuthConfig")
1273            .field("key_id", &self.key_id)
1274            .field("signing_key", &"[redacted]")
1275            .finish()
1276    }
1277}
1278
1279/// Configuration values for the objectstore service.
1280#[derive(Serialize, Deserialize, Debug, Clone)]
1281#[serde(default)]
1282pub struct ObjectstoreServiceConfig {
1283    /// The base URL for the objectstore service.
1284    ///
1285    /// This defaults to [`None`], which means that the service will be disabled,
1286    /// unless a proper configuration is provided.
1287    pub objectstore_url: Option<String>,
1288
1289    /// Maximum concurrency of uploads.
1290    pub max_concurrent_requests: usize,
1291
1292    /// Maximum size of the service input queue when `max_concurrent_requests` is saturated.
1293    ///
1294    /// The service will loadshed if this threshold is reached.
1295    pub max_backlog: usize,
1296
1297    /// Maximum duration of an attachment upload in seconds. Uploads that take longer are discarded.
1298    ///
1299    /// NOTE: This timeout applies to attachments that are already in-memory. Streaming uploads
1300    /// might take longer and are restricted independently by [`Self::stream_timeout`].
1301    pub timeout: u64,
1302
1303    /// Maximum duration of an upload stream.
1304    ///
1305    /// Streams get a larger default timeout because their duration depends on the client
1306    /// as well as the server.
1307    pub stream_timeout: u64,
1308
1309    /// Time between upload attempts.
1310    pub retry_delay: f64,
1311
1312    /// Maximum number of attempts made to upload.
1313    pub max_attempts: NonZeroU16,
1314
1315    /// Whether event attachment payloads may be sent through Kafka if objectstore upload fails.
1316    ///
1317    /// When disabled, failed event attachments are dropped with an `upload_failed`
1318    /// outcome instead of falling back to Store's Kafka attachment path.
1319    pub fallback_to_kafka: bool,
1320
1321    /// Configuration values for objectstore's auth scheme.
1322    pub auth: Option<ObjectstoreAuthConfig>,
1323}
1324
1325impl Default for ObjectstoreServiceConfig {
1326    fn default() -> Self {
1327        Self {
1328            objectstore_url: None,
1329            max_concurrent_requests: 10,
1330            max_backlog: 20,
1331            timeout: 60,
1332            stream_timeout: 5 * 60, // synced with `Upload::timeout`
1333            retry_delay: 1.0,
1334            max_attempts: NonZeroU16::new(5).unwrap(),
1335            fallback_to_kafka: true,
1336            auth: None,
1337        }
1338    }
1339}
1340
1341/// Determines how to emit outcomes.
1342/// For compatibility reasons, this can either be true, false or AsClientReports
1343#[derive(Copy, Clone, Debug, PartialEq, Eq)]
1344
1345pub enum EmitOutcomes {
1346    /// Do not emit any outcomes.
1347    None,
1348    /// Emit outcomes as client reports.
1349    AsClientReports,
1350    /// Emit outcomes as outcomes.
1351    AsOutcomes,
1352}
1353
1354impl EmitOutcomes {
1355    /// Returns true of outcomes are emitted via http, kafka, or client reports.
1356    pub fn any(&self) -> bool {
1357        !matches!(self, EmitOutcomes::None)
1358    }
1359}
1360
1361impl Serialize for EmitOutcomes {
1362    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1363    where
1364        S: Serializer,
1365    {
1366        // For compatibility, serialize None and AsOutcomes as booleans.
1367        match self {
1368            Self::None => serializer.serialize_bool(false),
1369            Self::AsClientReports => serializer.serialize_str("as_client_reports"),
1370            Self::AsOutcomes => serializer.serialize_bool(true),
1371        }
1372    }
1373}
1374
1375struct EmitOutcomesVisitor;
1376
1377impl Visitor<'_> for EmitOutcomesVisitor {
1378    type Value = EmitOutcomes;
1379
1380    fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
1381        formatter.write_str("true, false, 'as_client_reports'")
1382    }
1383
1384    fn visit_bool<E>(self, v: bool) -> Result<Self::Value, E>
1385    where
1386        E: serde::de::Error,
1387    {
1388        Ok(if v {
1389            EmitOutcomes::AsOutcomes
1390        } else {
1391            EmitOutcomes::None
1392        })
1393    }
1394
1395    fn visit_str<E>(self, v: &str) -> Result<Self::Value, E>
1396    where
1397        E: serde::de::Error,
1398    {
1399        match v {
1400            "as_client_reports" => Ok(EmitOutcomes::AsClientReports),
1401            _ => Err(E::invalid_value(Unexpected::Str(v), &self)),
1402        }
1403    }
1404}
1405
1406impl<'de> Deserialize<'de> for EmitOutcomes {
1407    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1408    where
1409        D: Deserializer<'de>,
1410    {
1411        deserializer.deserialize_any(EmitOutcomesVisitor)
1412    }
1413}
1414
1415/// Outcome generation specific configuration values.
1416#[derive(Serialize, Deserialize, Debug, Clone)]
1417#[serde(default)]
1418pub struct Outcomes {
1419    /// Controls whether outcomes will be emitted when processing is disabled.
1420    /// Processing relays always emit outcomes (for backwards compatibility).
1421    /// Can take the following values: false, "as_client_reports", true
1422    pub emit_outcomes: EmitOutcomes,
1423    /// Defines the source string registered in the outcomes originating from
1424    /// this Relay (typically something like the region or the layer).
1425    pub source: Option<String>,
1426}
1427
1428impl Default for Outcomes {
1429    fn default() -> Self {
1430        Outcomes {
1431            emit_outcomes: EmitOutcomes::AsClientReports,
1432            source: None,
1433        }
1434    }
1435}
1436
1437/// Minimal version of a config for dumping out.
1438#[derive(Serialize, Deserialize, Debug, Default)]
1439pub struct MinimalConfig {
1440    /// The relay part of the config.
1441    pub relay: Relay,
1442}
1443
1444impl MinimalConfig {
1445    /// Saves the config in the given config folder as config.yml
1446    pub fn save_in_folder<P: AsRef<Path>>(&self, p: P) -> anyhow::Result<()> {
1447        let path = p.as_ref();
1448        if fs::metadata(path).is_err() {
1449            fs::create_dir_all(path)
1450                .with_context(|| ConfigError::file(ConfigErrorKind::CouldNotOpenFile, path))?;
1451        }
1452        self.save(path)
1453    }
1454}
1455
1456impl ConfigObject for MinimalConfig {
1457    fn format() -> ConfigFormat {
1458        ConfigFormat::Yaml
1459    }
1460
1461    fn name() -> &'static str {
1462        "config"
1463    }
1464}
1465
1466/// Alternative serialization of RelayInfo for config file using snake case.
1467mod config_relay_info {
1468    use serde::ser::SerializeMap;
1469
1470    use super::*;
1471
1472    // Uses snake_case as opposed to camelCase.
1473    #[derive(Debug, Serialize, Deserialize, Clone)]
1474    struct RelayInfoConfig {
1475        public_key: PublicKey,
1476        #[serde(default)]
1477        internal: bool,
1478    }
1479
1480    impl From<RelayInfoConfig> for RelayInfo {
1481        fn from(v: RelayInfoConfig) -> Self {
1482            RelayInfo {
1483                public_key: v.public_key,
1484                internal: v.internal,
1485            }
1486        }
1487    }
1488
1489    impl From<RelayInfo> for RelayInfoConfig {
1490        fn from(v: RelayInfo) -> Self {
1491            RelayInfoConfig {
1492                public_key: v.public_key,
1493                internal: v.internal,
1494            }
1495        }
1496    }
1497
1498    pub(super) fn deserialize<'de, D>(des: D) -> Result<HashMap<RelayId, RelayInfo>, D::Error>
1499    where
1500        D: Deserializer<'de>,
1501    {
1502        let map = HashMap::<RelayId, RelayInfoConfig>::deserialize(des)?;
1503        Ok(map.into_iter().map(|(k, v)| (k, v.into())).collect())
1504    }
1505
1506    pub(super) fn serialize<S>(elm: &HashMap<RelayId, RelayInfo>, ser: S) -> Result<S::Ok, S::Error>
1507    where
1508        S: Serializer,
1509    {
1510        let mut map = ser.serialize_map(Some(elm.len()))?;
1511
1512        for (k, v) in elm {
1513            map.serialize_entry(k, &RelayInfoConfig::from(v.clone()))?;
1514        }
1515
1516        map.end()
1517    }
1518}
1519
1520/// Authentication options.
1521#[derive(Serialize, Deserialize, Debug, Clone)]
1522#[serde(default)]
1523pub struct AuthConfig {
1524    /// Controls responses from the readiness health check endpoint based on authentication.
1525    #[serde(skip_serializing_if = "is_default")]
1526    pub ready: ReadinessCondition,
1527
1528    /// Statically authenticated downstream relays.
1529    #[serde(with = "config_relay_info")]
1530    pub static_relays: HashMap<RelayId, RelayInfo>,
1531
1532    /// How old a signature can be before it is considered invalid, in seconds.
1533    ///
1534    /// Defaults to 5 minutes.
1535    pub signature_max_age: u64,
1536}
1537
1538impl Default for AuthConfig {
1539    fn default() -> Self {
1540        Self {
1541            ready: ReadinessCondition::default(),
1542            static_relays: HashMap::new(),
1543            signature_max_age: 300, // 5 minutes
1544        }
1545    }
1546}
1547
1548/// GeoIp database configuration options.
1549#[derive(Serialize, Deserialize, Debug, Default, Clone)]
1550pub struct GeoIpConfig {
1551    /// The path to GeoIP database.
1552    pub path: Option<PathBuf>,
1553}
1554
1555/// Settings to control Relay's health checks.
1556///
1557/// After breaching one of the configured thresholds, Relay will
1558/// return an `unhealthy` status from its health endpoint.
1559#[derive(Serialize, Deserialize, Debug, Clone, PartialEq)]
1560#[serde(default)]
1561pub struct Health {
1562    /// Interval to refresh internal health checks.
1563    ///
1564    /// Shorter intervals will decrease the time it takes the health check endpoint to report
1565    /// issues, but can also increase sporadic unhealthy responses.
1566    ///
1567    /// Defaults to `3000`` (3 seconds).
1568    pub refresh_interval_ms: u64,
1569    /// Maximum memory watermark in bytes.
1570    ///
1571    /// By default, there is no absolute limit set and the watermark
1572    /// is only controlled by setting [`Self::max_memory_percent`].
1573    pub max_memory_bytes: Option<ByteSize>,
1574    /// Maximum memory watermark as a percentage of maximum system memory.
1575    ///
1576    /// Defaults to `0.95` (95%).
1577    pub max_memory_percent: f32,
1578    /// Health check probe timeout in milliseconds.
1579    ///
1580    /// Any probe exceeding the timeout will be considered failed.
1581    /// This limits the max execution time of Relay health checks.
1582    ///
1583    /// Defaults to 900 milliseconds.
1584    pub probe_timeout_ms: u64,
1585    /// The refresh frequency of memory stats which are used to poll memory
1586    /// usage of Relay.
1587    ///
1588    /// The implementation of memory stats guarantees that the refresh will happen at
1589    /// least every `x` ms since memory readings are lazy and are updated only if needed.
1590    pub memory_stat_refresh_frequency_ms: u64,
1591}
1592
1593impl Default for Health {
1594    fn default() -> Self {
1595        Self {
1596            refresh_interval_ms: 3000,
1597            max_memory_bytes: None,
1598            max_memory_percent: 0.95,
1599            probe_timeout_ms: 900,
1600            memory_stat_refresh_frequency_ms: 100,
1601        }
1602    }
1603}
1604
1605/// COGS configuration.
1606#[derive(Serialize, Deserialize, Debug, Clone)]
1607#[serde(default)]
1608pub struct Cogs {
1609    /// Maximium amount of COGS measurements allowed to backlog.
1610    ///
1611    /// Any additional COGS measurements recorded will be dropped.
1612    ///
1613    /// Defaults to `10_000`.
1614    pub max_queue_size: u64,
1615    /// Relay COGS resource id.
1616    ///
1617    /// All Relay related COGS measurements are emitted with this resource id.
1618    ///
1619    /// Defaults to `relay_service`.
1620    pub relay_resource_id: String,
1621}
1622
1623impl Default for Cogs {
1624    fn default() -> Self {
1625        Self {
1626            max_queue_size: 10_000,
1627            relay_resource_id: "relay_service".to_owned(),
1628        }
1629    }
1630}
1631
1632/// Configuration for the upload service.
1633#[derive(Debug, Clone, Serialize, Deserialize)]
1634#[serde(default)]
1635pub struct Upload {
1636    /// Maximum number of uploads that the service accepts.
1637    ///
1638    /// Additional uploads will be rejected.
1639    pub max_concurrent_requests: usize,
1640    /// Maximum time spent trying to upload, in seconds.
1641    pub timeout: u64,
1642    /// The maximum time between creating the upload and uploading the data / the attachment placeholder.
1643    ///
1644    /// In seconds.
1645    pub max_age: i64,
1646
1647    /// Credentials used for signing & verifying upload locations.
1648    ///
1649    /// If omitted, relay's default [`Credentials`] are used.
1650    pub credentials: Option<UploadCredentials>,
1651}
1652
1653impl Default for Upload {
1654    fn default() -> Self {
1655        Self {
1656            max_concurrent_requests: 100,
1657            timeout: 5 * 60,  // five minutes
1658            max_age: 60 * 60, // 1h
1659            credentials: None,
1660        }
1661    }
1662}
1663
1664/// Credentials used for signing & verifying upload locations.
1665#[derive(Clone, Serialize, Deserialize)]
1666pub struct UploadCredentials {
1667    /// Key used to sign upload locations.
1668    #[cfg(feature = "processing")]
1669    pub signing_key: SecretKey,
1670
1671    /// Key used to verify upload locations.
1672    pub verification_key: PublicKey,
1673}
1674
1675impl fmt::Debug for UploadCredentials {
1676    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1677        let Self {
1678            #[cfg(feature = "processing")]
1679                signing_key: _,
1680            verification_key,
1681        } = self;
1682        let mut b = f.debug_struct("UploadCredentials");
1683        #[cfg(feature = "processing")]
1684        b.field("signing_key", &"[redacted]");
1685        b.field("verification_key", verification_key).finish()
1686    }
1687}
1688
1689/// All configuration values that can be deserialized from `config.yml`.
1690#[derive(Serialize, Deserialize, Debug, Default, Clone)]
1691#[serde(default)]
1692#[allow(missing_docs)]
1693pub struct ConfigValues {
1694    pub relay: Relay,
1695    pub http: Http,
1696    pub cache: Cache,
1697    pub spool: Spool,
1698    pub limits: Limits,
1699    pub logging: relay_log::LogConfig,
1700    pub routing: Routing,
1701    pub metrics: Metrics,
1702    pub sentry: relay_log::SentryConfig,
1703    pub processing: Processing,
1704    pub outcomes: Outcomes,
1705    pub aggregator: AggregatorServiceConfig,
1706    pub secondary_aggregators: Vec<ScopedAggregatorConfig>,
1707    pub auth: AuthConfig,
1708    pub geoip: GeoIpConfig,
1709    pub normalization: Normalization,
1710    pub health: Health,
1711    pub cogs: Cogs,
1712    pub upload: Upload,
1713}
1714
1715impl ConfigObject for ConfigValues {
1716    fn format() -> ConfigFormat {
1717        ConfigFormat::Yaml
1718    }
1719
1720    fn name() -> &'static str {
1721        "config"
1722    }
1723}
1724
1725#[derive(Default, Clone)]
1726struct ConfigInner {
1727    /// Relay's config values.
1728    values: ConfigValues,
1729    /// Configured Relay credentials.
1730    ///
1731    /// Credentials may be missing for proxy mode.
1732    credentials: Option<Credentials>,
1733}
1734
1735impl ConfigInner {
1736    fn apply_overrides(&mut self, overrides: &OverridableConfig) -> anyhow::Result<()> {
1737        if let Some(log_level) = &overrides.log_level {
1738            self.values.logging.level = log_level.parse()?;
1739        }
1740
1741        if let Some(log_format) = &overrides.log_format {
1742            self.values.logging.format = log_format.parse()?;
1743        }
1744
1745        let relay = &mut self.values.relay;
1746        if let Some(mode) = &overrides.mode {
1747            relay.mode = mode
1748                .parse::<RelayMode>()
1749                .with_context(|| ConfigError::field("mode"))?;
1750        }
1751        if let Some(deployment) = &overrides.instance {
1752            relay.instance = deployment
1753                .parse::<RelayInstance>()
1754                .with_context(|| ConfigError::field("deployment"))?;
1755        }
1756        if let Some(upstream) = &overrides.upstream {
1757            relay.upstream = upstream
1758                .parse::<UpstreamDescriptor>()
1759                .with_context(|| ConfigError::field("upstream"))?;
1760        } else if let Some(upstream_dsn) = &overrides.upstream_dsn {
1761            relay.upstream = upstream_dsn
1762                .parse::<Dsn>()
1763                .map(|dsn| UpstreamDescriptor::from_dsn(&dsn))
1764                .with_context(|| ConfigError::field("upstream_dsn"))?;
1765        }
1766        if let Some(host) = &overrides.host {
1767            relay.host = host
1768                .parse::<IpAddr>()
1769                .with_context(|| ConfigError::field("host"))?;
1770        }
1771        if let Some(port) = &overrides.port {
1772            relay.port = port
1773                .as_str()
1774                .parse()
1775                .with_context(|| ConfigError::field("port"))?;
1776        }
1777
1778        let processing = &mut self.values.processing;
1779        if let Some(enabled) = &overrides.processing {
1780            match enabled.to_lowercase().as_str() {
1781                "true" | "1" => processing.enabled = true,
1782                "false" | "0" | "" => processing.enabled = false,
1783                _ => return Err(ConfigError::field("processing").into()),
1784            }
1785        }
1786        if let Some(redis) = overrides.redis_url.clone() {
1787            processing.redis = Some(RedisConfigs::Unified(RedisConfig::single(redis)))
1788        }
1789        if let Some(kafka_url) = overrides.kafka_url.clone() {
1790            let existing = processing
1791                .kafka_config
1792                .iter_mut()
1793                .find(|e| e.name == "bootstrap.servers");
1794
1795            if let Some(config_param) = existing {
1796                config_param.value = kafka_url;
1797            } else {
1798                self.values.processing.kafka_config.push(KafkaConfigParam {
1799                    name: "bootstrap.servers".to_owned(),
1800                    value: kafka_url,
1801                })
1802            }
1803        }
1804
1805        if overrides.outcome_source.is_some() {
1806            self.values.outcomes.source = overrides.outcome_source.clone();
1807        }
1808
1809        if let Some(shutdown_timeout) = &overrides.shutdown_timeout
1810            && let Ok(shutdown_timeout) = shutdown_timeout.parse::<u64>()
1811        {
1812            self.values.limits.shutdown_timeout = shutdown_timeout;
1813        }
1814
1815        if let Some(server_name) = overrides.server_name.clone() {
1816            self.values.sentry.server_name = Some(server_name.into());
1817        }
1818
1819        let id = if let Some(id) = &overrides.id {
1820            let id = Uuid::parse_str(id).with_context(|| ConfigError::field("id"))?;
1821            Some(id)
1822        } else {
1823            None
1824        };
1825        let public_key = if let Some(public_key) = &overrides.public_key {
1826            let public_key = public_key
1827                .parse::<PublicKey>()
1828                .with_context(|| ConfigError::field("public_key"))?;
1829            Some(public_key)
1830        } else {
1831            None
1832        };
1833
1834        let secret_key = if let Some(secret_key) = &overrides.secret_key {
1835            let secret_key = secret_key
1836                .parse::<SecretKey>()
1837                .with_context(|| ConfigError::field("secret_key"))?;
1838            Some(secret_key)
1839        } else {
1840            None
1841        };
1842
1843        if let Some(credentials) = &mut self.credentials {
1844            //we have existing credentials we may override some entries
1845            if let Some(id) = id {
1846                credentials.id = id;
1847            }
1848            if let Some(public_key) = public_key {
1849                credentials.public_key = public_key;
1850            }
1851            if let Some(secret_key) = secret_key {
1852                credentials.secret_key = secret_key
1853            }
1854        } else {
1855            //no existing credentials we may only create the full credentials
1856            match (id, public_key, secret_key) {
1857                (Some(id), Some(public_key), Some(secret_key)) => {
1858                    self.credentials = Some(Credentials {
1859                        secret_key,
1860                        public_key,
1861                        id,
1862                    })
1863                }
1864                (None, None, None) => {
1865                    // nothing provided, we'll just leave the credentials None, maybe we
1866                    // don't need them in the current command or we'll override them later
1867                }
1868                _ => {
1869                    return Err(ConfigError::field("incomplete credentials").into());
1870                }
1871            }
1872        }
1873
1874        Ok(())
1875    }
1876
1877    /// Merges reloadable parts of the config from `other` into `self`.
1878    ///
1879    /// Returns `true` if any parts of the config were updated.
1880    fn reload_with(&mut self, other: &Self) -> bool {
1881        let mut changed = false;
1882
1883        if self.values.health != other.values.health {
1884            relay_log::debug!("updating health");
1885            self.values.health = other.values.health.clone();
1886            changed = true;
1887        }
1888
1889        changed
1890    }
1891}
1892
1893/// Relay's Configuration.
1894pub struct Config {
1895    /// A mutex to serialize all write accesses to `inner`.
1896    ///
1897    /// Accessing the arc swap is done with compare and swap, but we still want serialized
1898    /// access and operations to guarantee consistency when dealing with e.g. the filesystem.
1899    ///
1900    /// The mutex must be acquired before changing `inner`. Methods taking `&mut` can omit
1901    /// this, as the `&mut` requirement already satisfies that there is no concurrent access
1902    /// possible.
1903    inner_access: Mutex<()>,
1904    /// The actual config.
1905    inner: ArcSwap<ConfigInner>,
1906    /// A list of overrides applied to the config, in order.
1907    ///
1908    /// When re-loading the configuration these overrides need to be applied again
1909    /// in the same order as they were applied originally.
1910    overrides: Vec<OverridableConfig>,
1911    /// Path from which the config is loaded.
1912    ///
1913    /// This is Relay's configuration directory.
1914    path: PathBuf,
1915}
1916
1917impl fmt::Debug for Config {
1918    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1919        let inner = self.inner.load();
1920
1921        f.debug_struct("Config")
1922            .field("path", &self.path)
1923            // Only print specific parts of `inner` to not leak the credentials.
1924            .field("values", &inner.values)
1925            .finish()
1926    }
1927}
1928
1929impl Config {
1930    /// Loads a config from a given config folder.
1931    pub fn from_path<P: AsRef<Path>>(path: P) -> anyhow::Result<Config> {
1932        let path = env::current_dir()
1933            .map(|x| x.join(path.as_ref()))
1934            .unwrap_or_else(|_| path.as_ref().to_path_buf());
1935
1936        let inner = ConfigInner {
1937            values: ConfigValues::load(&path)?,
1938            credentials: match Credentials::path(&path).exists() {
1939                true => Some(Credentials::load(&path)?),
1940                false => None,
1941            },
1942        };
1943
1944        let config = Config {
1945            inner_access: Mutex::new(()),
1946            inner: ArcSwap::from_pointee(inner),
1947            overrides: Vec::new(),
1948            path: path.clone(),
1949        };
1950
1951        if cfg!(not(feature = "processing")) && config.current().processing_enabled() {
1952            return Err(ConfigError::file(ConfigErrorKind::ProcessingNotAvailable, &path).into());
1953        }
1954
1955        Ok(config)
1956    }
1957
1958    /// Creates a config from a JSON value.
1959    ///
1960    /// This is mostly useful for tests.
1961    pub fn from_json_value(value: serde_json::Value) -> anyhow::Result<Config> {
1962        Ok(Config {
1963            inner_access: Mutex::new(()),
1964            inner: ArcSwap::from_pointee(ConfigInner {
1965                values: serde_json::from_value(value)
1966                    .with_context(|| ConfigError::new(ConfigErrorKind::BadJson))?,
1967                credentials: None,
1968            }),
1969            overrides: Vec::new(),
1970            path: PathBuf::new(),
1971        })
1972    }
1973
1974    /// Override configuration with values coming from other sources (e.g. env variables or
1975    /// command line parameters).
1976    ///
1977    /// If applying the overrides fails, the config may be left in an inconsistent state.
1978    pub fn apply_override(&mut self, overrides: OverridableConfig) -> anyhow::Result<&mut Self> {
1979        // We could introduce a config builder which operates on mutable configs, which would eliminate
1980        // the need for this `try_rcu` dance here.
1981        crate::utils::try_rcu(&self.inner, |inner| {
1982            let mut new = ConfigInner::clone(inner);
1983            new.apply_overrides(&overrides)?;
1984            Ok::<_, anyhow::Error>(Arc::new(new))
1985        })?;
1986
1987        // Overrides successfully applied.
1988        self.overrides.push(overrides);
1989
1990        Ok(self)
1991    }
1992
1993    /// Checks if the config is already initialized.
1994    pub fn config_exists<P: AsRef<Path>>(path: P) -> bool {
1995        fs::metadata(ConfigValues::path(path.as_ref())).is_ok()
1996    }
1997
1998    /// Returns the filename of the config file.
1999    pub fn path(&self) -> &Path {
2000        &self.path
2001    }
2002
2003    /// Dumps out a YAML string of the values.
2004    pub fn to_yaml_string(&self) -> anyhow::Result<String> {
2005        serde_yaml::to_string(&self.inner.load().values)
2006            .with_context(|| ConfigError::new(ConfigErrorKind::CouldNotWriteFile))
2007    }
2008
2009    /// Set new credentials.
2010    ///
2011    /// This also writes the credentials back to the file, if this config was loaded from the file-system.
2012    pub fn replace_credentials(
2013        &mut self,
2014        credentials: Option<Credentials>,
2015    ) -> anyhow::Result<bool> {
2016        if self.inner.load().credentials == credentials {
2017            return Ok(false);
2018        }
2019
2020        if !self.path.is_empty() {
2021            match &credentials {
2022                Some(creds) => {
2023                    creds.save(&self.path)?;
2024                }
2025                None => {
2026                    let path = Credentials::path(&self.path);
2027                    if fs::metadata(&path).is_ok() {
2028                        fs::remove_file(&path).with_context(|| {
2029                            ConfigError::file(ConfigErrorKind::CouldNotWriteFile, &path)
2030                        })?;
2031                    }
2032                }
2033            }
2034        }
2035
2036        // Note: there is never anyone racing on the `ArcSwap` as long as `Self` borrowed mutably.
2037        //
2038        // We can improve this if we split out mutable operations into a separate struct and only
2039        // once `frozen()` we change to an `ArcSwap` internally.
2040        self.inner.rcu(|inner| {
2041            let mut inner = ConfigInner::clone(inner);
2042            inner.credentials = credentials.clone();
2043            Arc::new(inner)
2044        });
2045
2046        Ok(true)
2047    }
2048
2049    /// Reloads the configuration from disk.
2050    ///
2051    /// Returns `true` if the configuration changed.
2052    ///
2053    /// In order for a config to be reloadable it must've been loaded from a path. The original
2054    /// config will be re-read and reloadable parts of the config will be replaced with their updated
2055    /// values.
2056    ///
2057    /// If the reload fails for any reason, the current config is untouched.
2058    pub fn reload(&self) -> anyhow::Result<bool> {
2059        if self.path.is_empty() {
2060            return Ok(false);
2061        }
2062
2063        let _access = self
2064            .inner_access
2065            .lock()
2066            .unwrap_or_else(PoisonError::into_inner);
2067
2068        let mut new_config = Self::from_path(&self.path)?;
2069        for overrides in &self.overrides {
2070            new_config.apply_override(overrides.clone())?;
2071        }
2072        let new_config = new_config.current();
2073
2074        let mut changed = false;
2075
2076        // Since we do have the `_access` lock, this will always succeed on the first try.
2077        self.inner.rcu(|inner| {
2078            let mut new_inner = ConfigInner::clone(inner);
2079            changed = new_inner.reload_with(&new_config.inner);
2080            match changed {
2081                true => Arc::new(new_inner),
2082                false => Arc::clone(inner),
2083            }
2084        });
2085
2086        Ok(changed)
2087    }
2088
2089    /// Acquires a current [`snapshot`](ConfigSnapshot) of the config.
2090    ///
2091    /// A snapshot is the way to actually consume values from the config. A snapshot should ideally be acquired
2092    /// once per unit of work.
2093    pub fn current(&self) -> ConfigSnapshot {
2094        let inner = self.inner.load();
2095        ConfigSnapshot { inner }
2096    }
2097}
2098
2099impl Default for Config {
2100    fn default() -> Self {
2101        Self {
2102            inner_access: Mutex::new(()),
2103            inner: ArcSwap::from_pointee(Default::default()),
2104            overrides: Vec::new(),
2105            path: PathBuf::new(),
2106        }
2107    }
2108}
2109
2110/// A config snapshot is a point in time snapshot of the [`Config`].
2111///
2112/// The [`Config`] may change over time, to guarantee a consistent view of the config
2113/// a snapshot must be acquired first.
2114///
2115/// The snapshot should not be stored in a long lasting datastructure. As a rule of thumb it should
2116/// only exist on the stack.
2117pub struct ConfigSnapshot {
2118    inner: arc_swap::Guard<Arc<ConfigInner>>,
2119}
2120
2121impl ConfigSnapshot {
2122    /// Returns `true` if the config is ready to use.
2123    pub fn has_credentials(&self) -> bool {
2124        self.inner.credentials.is_some()
2125    }
2126
2127    /// Return the current credentials.
2128    pub fn credentials(&self) -> Option<&Credentials> {
2129        self.inner.credentials.as_ref()
2130    }
2131
2132    /// Returns the secret key if set.
2133    pub fn secret_key(&self) -> Option<&SecretKey> {
2134        self.inner.credentials.as_ref().map(|x| &x.secret_key)
2135    }
2136
2137    /// Returns the public key if set.
2138    pub fn public_key(&self) -> Option<&PublicKey> {
2139        self.inner.credentials.as_ref().map(|x| &x.public_key)
2140    }
2141
2142    /// Returns the relay ID.
2143    pub fn relay_id(&self) -> Option<&RelayId> {
2144        self.inner.credentials.as_ref().map(|x| &x.id)
2145    }
2146
2147    /// Returns the relay mode.
2148    pub fn relay_mode(&self) -> RelayMode {
2149        self.inner.values.relay.mode
2150    }
2151
2152    /// Returns the instance type of relay.
2153    pub fn relay_instance(&self) -> RelayInstance {
2154        self.inner.values.relay.instance
2155    }
2156
2157    /// Returns the upstream target as descriptor.
2158    pub fn upstream(&self) -> &UpstreamDescriptor {
2159        &self.inner.values.relay.upstream
2160    }
2161
2162    /// Returns the advertised upstream for downstream instances as descriptor.
2163    pub fn advertised_upstream(&self) -> Option<&UpstreamDescriptor> {
2164        self.inner.values.relay.advertised_upstream.as_ref()
2165    }
2166
2167    /// Returns the custom HTTP "Host" header.
2168    pub fn http_host_header(&self) -> Option<&str> {
2169        self.inner.values.http.host_header.as_deref()
2170    }
2171
2172    /// Returns the listen address.
2173    pub fn listen_addr(&self) -> SocketAddr {
2174        (self.inner.values.relay.host, self.inner.values.relay.port).into()
2175    }
2176
2177    /// Returns the listen address for internal APIs.
2178    ///
2179    /// Internal APIs are APIs which do not need to be publicly exposed,
2180    /// like health checks.
2181    ///
2182    /// Returns `None` when there is no explicit address configured for internal APIs,
2183    /// and they should instead be exposed on the main [`Self::listen_addr`].
2184    pub fn listen_addr_internal(&self) -> Option<SocketAddr> {
2185        match (
2186            self.inner.values.relay.internal_host,
2187            self.inner.values.relay.internal_port,
2188        ) {
2189            (Some(host), None) => Some((host, self.inner.values.relay.port).into()),
2190            (None, Some(port)) => Some((self.inner.values.relay.host, port).into()),
2191            (Some(host), Some(port)) => Some((host, port).into()),
2192            (None, None) => None,
2193        }
2194    }
2195
2196    /// Returns `true` when project IDs should be overriden rather than validated.
2197    ///
2198    /// Defaults to `false`, which requires project ID validation.
2199    pub fn override_project_ids(&self) -> bool {
2200        self.inner.values.relay.override_project_ids
2201    }
2202
2203    /// Returns the interval to check for configuration changes.
2204    ///
2205    /// `None` if the config should never be reloaded.
2206    pub fn config_reload_interval(&self) -> Option<Duration> {
2207        let interval = self.inner.values.relay.config_reload_interval?;
2208        Some(match interval {
2209            // Useful for tests to be able to configure this to a very small value,
2210            // but we also don't want it to be actually 0.
2211            0 => Duration::from_millis(50),
2212            secs => Duration::from_secs(secs),
2213        })
2214    }
2215
2216    /// Returns `true` if Relay requires authentication for readiness.
2217    ///
2218    /// See [`ReadinessCondition`] for more information.
2219    pub fn requires_auth(&self) -> bool {
2220        match self.inner.values.auth.ready {
2221            ReadinessCondition::Authenticated => self.relay_mode() == RelayMode::Managed,
2222            ReadinessCondition::Always => false,
2223        }
2224    }
2225
2226    /// Returns the interval at which Realy should try to re-authenticate with the upstream.
2227    ///
2228    /// Always disabled in processing mode.
2229    pub fn http_auth_interval(&self) -> Option<Duration> {
2230        if self.processing_enabled() {
2231            return None;
2232        }
2233
2234        match self.inner.values.http.auth_interval {
2235            None | Some(0) => None,
2236            Some(secs) => Some(Duration::from_secs(secs)),
2237        }
2238    }
2239
2240    /// The maximum time of experiencing uninterrupted network failures until Relay considers that
2241    /// it has encountered a network outage.
2242    pub fn http_outage_grace_period(&self) -> Duration {
2243        Duration::from_secs(self.inner.values.http.outage_grace_period)
2244    }
2245
2246    /// Time Relay waits before retrying an upstream request.
2247    ///
2248    /// Before going into a network outage, Relay may fail to make upstream
2249    /// requests. This is the time Relay waits before retrying the same request.
2250    pub fn http_retry_delay(&self) -> Duration {
2251        Duration::from_secs(self.inner.values.http.retry_delay)
2252    }
2253
2254    /// Time of continued project request failures before Relay emits an error.
2255    pub fn http_project_failure_interval(&self) -> Duration {
2256        Duration::from_secs(self.inner.values.http.project_failure_interval)
2257    }
2258
2259    /// Content encoding of upstream requests.
2260    pub fn http_encoding(&self) -> HttpEncoding {
2261        self.inner.values.http.encoding
2262    }
2263
2264    /// Returns whether metrics should be sent globally through a shared endpoint.
2265    pub fn http_global_metrics(&self) -> bool {
2266        self.inner.values.http.global_metrics
2267    }
2268
2269    /// Returns `true` if Relay supports forwarding unknown API requests.
2270    ///
2271    /// Relay instances with processing enabled are expected to support the latest API and do never
2272    /// support forwarding requests to Sentry.
2273    pub fn http_forward(&self) -> bool {
2274        self.inner.values.http.forward && !self.processing_enabled()
2275    }
2276
2277    /// Returns whether this Relay should emit outcomes.
2278    ///
2279    /// This is `true` either if `outcomes.emit_outcomes` is explicitly enabled, or if this Relay is
2280    /// in processing mode.
2281    pub fn emit_outcomes(&self) -> EmitOutcomes {
2282        if self.processing_enabled() {
2283            return EmitOutcomes::AsOutcomes;
2284        }
2285        self.inner.values.outcomes.emit_outcomes
2286    }
2287
2288    /// The originating source of the outcome
2289    pub fn outcome_source(&self) -> Option<&str> {
2290        self.inner.values.outcomes.source.as_deref()
2291    }
2292
2293    /// Returns logging configuration.
2294    pub fn logging(&self) -> &relay_log::LogConfig {
2295        &self.inner.values.logging
2296    }
2297
2298    /// Returns logging configuration.
2299    pub fn sentry(&self) -> &relay_log::SentryConfig {
2300        &self.inner.values.sentry
2301    }
2302
2303    /// Returns the addresses for statsd metrics.
2304    pub fn statsd_addr(&self) -> Option<&str> {
2305        self.inner.values.metrics.statsd.as_deref()
2306    }
2307
2308    /// Returns the addresses for statsd metrics.
2309    pub fn statsd_buffer_size(&self) -> Option<usize> {
2310        self.inner.values.metrics.statsd_buffer_size
2311    }
2312
2313    /// Return the prefix for statsd metrics.
2314    pub fn metrics_prefix(&self) -> &str {
2315        &self.inner.values.metrics.prefix
2316    }
2317
2318    /// Returns the default tags for statsd metrics.
2319    pub fn metrics_default_tags(&self) -> &BTreeMap<String, String> {
2320        &self.inner.values.metrics.default_tags
2321    }
2322
2323    /// Returns the name of the hostname tag that should be attached to each outgoing metric.
2324    pub fn metrics_hostname_tag(&self) -> Option<&str> {
2325        self.inner.values.metrics.hostname_tag.as_deref()
2326    }
2327
2328    /// Returns the interval for periodic metrics emitted from Relay.
2329    ///
2330    /// `None` if periodic metrics are disabled.
2331    pub fn metrics_periodic_interval(&self) -> Option<Duration> {
2332        match self.inner.values.metrics.periodic_secs {
2333            0 => None,
2334            secs => Some(Duration::from_secs(secs)),
2335        }
2336    }
2337
2338    /// Returns the default timeout for all upstream HTTP requests.
2339    pub fn http_timeout(&self) -> Duration {
2340        Duration::from_secs(self.inner.values.http.timeout.into())
2341    }
2342
2343    /// Returns the connection timeout for all upstream HTTP requests.
2344    pub fn http_connection_timeout(&self) -> Duration {
2345        Duration::from_secs(self.inner.values.http.connection_timeout.into())
2346    }
2347
2348    /// Returns the failed upstream request retry interval.
2349    pub fn http_max_retry_interval(&self) -> Duration {
2350        Duration::from_secs(self.inner.values.http.max_retry_interval.into())
2351    }
2352
2353    /// Returns `true` if relay should use an in-process cache for DNS lookups.
2354    pub fn http_dns_cache(&self) -> bool {
2355        self.inner.values.http.dns_cache
2356    }
2357
2358    /// Returns the expiry timeout for cached projects.
2359    pub fn project_cache_expiry(&self) -> Duration {
2360        Duration::from_secs(self.inner.values.cache.project_expiry.into())
2361    }
2362
2363    /// Returns `true` if the full project state should be requested from upstream.
2364    pub fn request_full_project_config(&self) -> bool {
2365        self.inner.values.cache.project_request_full_config
2366    }
2367
2368    /// Returns the expiry timeout for cached relay infos (public keys).
2369    pub fn relay_cache_expiry(&self) -> Duration {
2370        Duration::from_secs(self.inner.values.cache.relay_expiry.into())
2371    }
2372
2373    /// Returns the expiry timeout for cached misses before trying to refetch.
2374    pub fn cache_miss_expiry(&self) -> Duration {
2375        Duration::from_secs(self.inner.values.cache.miss_expiry.into())
2376    }
2377
2378    /// Returns the grace period for project caches.
2379    pub fn project_grace_period(&self) -> Duration {
2380        Duration::from_secs(self.inner.values.cache.project_grace_period.into())
2381    }
2382
2383    /// Returns the refresh interval for a project.
2384    ///
2385    /// Validates the refresh time to be between the grace period and expiry.
2386    pub fn project_refresh_interval(&self) -> Option<Duration> {
2387        self.inner
2388            .values
2389            .cache
2390            .project_refresh_interval
2391            .map(Into::into)
2392            .map(Duration::from_secs)
2393    }
2394
2395    /// Returns the duration in which batchable project config queries are
2396    /// collected before sending them in a single request.
2397    pub fn query_batch_interval(&self) -> Duration {
2398        Duration::from_millis(self.inner.values.cache.batch_interval.into())
2399    }
2400
2401    /// Returns the duration in which downstream relays are requested from upstream.
2402    pub fn downstream_relays_batch_interval(&self) -> Duration {
2403        Duration::from_millis(
2404            self.inner
2405                .values
2406                .cache
2407                .downstream_relays_batch_interval
2408                .into(),
2409        )
2410    }
2411
2412    /// Returns the interval in seconds in which local project configurations should be reloaded.
2413    pub fn local_cache_interval(&self) -> Duration {
2414        Duration::from_secs(self.inner.values.cache.file_interval.into())
2415    }
2416
2417    /// Returns the interval in seconds in which fresh global configs should be
2418    /// fetched from  upstream.
2419    pub fn global_config_fetch_interval(&self) -> Duration {
2420        Duration::from_secs(self.inner.values.cache.global_config_fetch_interval.into())
2421    }
2422
2423    /// Returns the path of the buffer file if the `cache.persistent_envelope_buffer.path` is configured.
2424    ///
2425    /// In case a partition with id > 0 is supplied, the filename of the envelopes path will be
2426    /// suffixed with `.{partition_id}`.
2427    pub fn spool_envelopes_path(&self, partition_id: u8) -> Option<PathBuf> {
2428        let mut path = self
2429            .inner
2430            .values
2431            .spool
2432            .envelopes
2433            .path
2434            .as_ref()
2435            .map(|path| path.to_owned())?;
2436
2437        if partition_id == 0 {
2438            return Some(path);
2439        }
2440
2441        let file_name = path.file_name().and_then(|f| f.to_str())?;
2442        let new_file_name = format!("{file_name}.{partition_id}");
2443        path.set_file_name(new_file_name);
2444
2445        Some(path)
2446    }
2447
2448    /// The maximum size of the buffer, in bytes.
2449    pub fn spool_envelopes_max_disk_size(&self) -> usize {
2450        self.inner.values.spool.envelopes.max_disk_size.as_bytes()
2451    }
2452
2453    /// Number of encoded envelope bytes that need to be accumulated before
2454    /// flushing one batch to disk.
2455    pub fn spool_envelopes_batch_size_bytes(&self) -> usize {
2456        self.inner
2457            .values
2458            .spool
2459            .envelopes
2460            .batch_size_bytes
2461            .as_bytes()
2462    }
2463
2464    /// Time after which a batch of envelopes is flushed to disk, regardless of its size.
2465    pub fn spool_envelopes_flush_timeout(&self) -> Option<Duration> {
2466        self.inner
2467            .values
2468            .spool
2469            .envelopes
2470            .flush_timeout_secs
2471            .map(Duration::from_secs)
2472    }
2473
2474    /// Maximum number of envelopes unspooled from disk per second.
2475    pub fn spool_max_unspool_envelopes_per_second(&self) -> Option<NonZeroU32> {
2476        self.inner
2477            .values
2478            .spool
2479            .envelopes
2480            .max_unspool_envelopes_per_second
2481    }
2482
2483    /// Returns the time after which we drop envelopes as a [`Duration`] object.
2484    pub fn spool_envelopes_max_age(&self) -> Duration {
2485        Duration::from_secs(self.inner.values.spool.envelopes.max_envelope_delay_secs)
2486    }
2487
2488    /// Returns the refresh frequency for disk usage monitoring as a [`Duration`] object.
2489    pub fn spool_disk_usage_refresh_frequency_ms(&self) -> Duration {
2490        Duration::from_millis(
2491            self.inner
2492                .values
2493                .spool
2494                .envelopes
2495                .disk_usage_refresh_frequency_ms,
2496        )
2497    }
2498
2499    /// Returns the relative memory usage up to which the disk buffer will unspool envelopes.
2500    pub fn spool_max_backpressure_memory_percent(&self) -> f32 {
2501        self.inner
2502            .values
2503            .spool
2504            .envelopes
2505            .max_backpressure_memory_percent
2506    }
2507
2508    /// Returns the number of partitions for the buffer.
2509    pub fn spool_partitions(&self) -> NonZeroU8 {
2510        self.inner.values.spool.envelopes.partitions
2511    }
2512
2513    /// Returns the strategy used to assign envelopes to buffer partitions.
2514    pub fn spool_partitioning(&self) -> EnvelopeSpoolPartitioning {
2515        self.inner.values.spool.envelopes.partitioning
2516    }
2517
2518    /// Returns `true` if the data is stored on ephemeral disks.
2519    pub fn spool_ephemeral(&self) -> bool {
2520        self.inner.values.spool.envelopes.ephemeral
2521    }
2522
2523    /// Returns the maximum size of an event payload in bytes.
2524    pub fn max_event_size(&self) -> usize {
2525        self.inner.values.limits.max_event_size.as_bytes()
2526    }
2527
2528    /// Returns the maximum size of each attachment.
2529    pub fn max_attachment_size(&self) -> usize {
2530        self.inner.values.limits.max_attachment_size.as_bytes()
2531    }
2532
2533    /// The maximum amount of attachments in a single envelope.
2534    pub fn max_attachment_count(&self) -> usize {
2535        self.inner.values.limits.max_attachment_count
2536    }
2537
2538    /// Returns the maximum combined size of attachments or payloads containing attachments
2539    /// (minidump, unreal, standalone attachments) in bytes.
2540    pub fn max_attachments_size(&self) -> usize {
2541        self.inner.values.limits.max_attachments_size.as_bytes()
2542    }
2543
2544    /// Returns the maximum size of a TUS upload request body.
2545    pub fn max_upload_size(&self) -> usize {
2546        self.inner.values.limits.max_upload_size.as_bytes()
2547    }
2548
2549    /// Returns the maximum number of client reports per envelope.
2550    pub fn max_client_reports_count(&self) -> usize {
2551        self.inner.values.limits.max_client_reports_count
2552    }
2553
2554    /// Returns the maximum combined size of client reports in bytes.
2555    pub fn max_client_reports_size(&self) -> usize {
2556        self.inner.values.limits.max_client_reports_size.as_bytes()
2557    }
2558
2559    /// Returns the maximum payload size of a monitor check-in in bytes.
2560    pub fn max_check_in_size(&self) -> usize {
2561        self.inner.values.limits.max_check_in_size.as_bytes()
2562    }
2563
2564    /// Returns the maximum payload size of a log in bytes.
2565    pub fn max_log_size(&self) -> usize {
2566        self.inner.values.limits.max_log_size.as_bytes()
2567    }
2568
2569    /// Returns the maximum number of operations to allow for a span expansion.
2570    pub fn max_expanded_log_operations(&self) -> usize {
2571        self.inner.values.limits.max_expanded_log_operations
2572    }
2573
2574    /// Returns the maximum number of operations to allow for a span expansion.
2575    pub fn max_expanded_span_operations(&self) -> usize {
2576        self.inner.values.limits.max_expanded_span_operations
2577    }
2578
2579    /// Returns the maximum number of operations to allow for a span expansion.
2580    pub fn max_expanded_trace_metric_operations(&self) -> usize {
2581        self.inner
2582            .values
2583            .limits
2584            .max_expanded_trace_metric_operations
2585    }
2586
2587    /// Returns the maximum payload size of a span in bytes.
2588    pub fn max_span_size(&self) -> usize {
2589        self.inner.values.limits.max_span_size.as_bytes()
2590    }
2591
2592    /// Returns the maximum amount of standalone transaction spans per envelope.
2593    pub fn max_standalone_span_count(&self) -> usize {
2594        self.inner.values.limits.max_standalone_span_count
2595    }
2596
2597    /// Returns the maximum payload size of an item container in bytes.
2598    pub fn max_container_size(&self) -> usize {
2599        self.inner.values.limits.max_container_size.as_bytes()
2600    }
2601
2602    /// Returns the maximum size of an envelope payload in bytes.
2603    ///
2604    /// Individual item size limits still apply.
2605    pub fn max_envelope_size(&self) -> usize {
2606        self.inner.values.limits.max_envelope_size.as_bytes()
2607    }
2608
2609    /// Returns the maximum number of sessions per envelope.
2610    pub fn max_session_count(&self) -> usize {
2611        self.inner.values.limits.max_session_count
2612    }
2613
2614    /// Returns the maximum combined size for all sessions in an envelope in bytes.
2615    pub fn max_sessions_size(&self) -> usize {
2616        self.inner.values.limits.max_sessions_size.as_bytes()
2617    }
2618
2619    /// Returns the maximum payload size of metric buckets in bytes.
2620    pub fn max_metric_buckets_size(&self) -> usize {
2621        self.inner.values.limits.max_metric_buckets_size.as_bytes()
2622    }
2623
2624    /// Returns the maximum payload size for general API requests.
2625    pub fn max_api_payload_size(&self) -> usize {
2626        self.inner.values.limits.max_api_payload_size.as_bytes()
2627    }
2628
2629    /// Returns the maximum payload size for file uploads and chunks.
2630    pub fn max_api_file_upload_size(&self) -> usize {
2631        self.inner.values.limits.max_api_file_upload_size.as_bytes()
2632    }
2633
2634    /// Returns the maximum payload size for chunks
2635    pub fn max_api_chunk_upload_size(&self) -> usize {
2636        self.inner
2637            .values
2638            .limits
2639            .max_api_chunk_upload_size
2640            .as_bytes()
2641    }
2642
2643    /// Returns the maximum payload size for a profile
2644    pub fn max_profile_size(&self) -> usize {
2645        self.inner.values.limits.max_profile_size.as_bytes()
2646    }
2647
2648    /// Returns the maximum payload size for a trace metric.
2649    pub fn max_trace_metric_size(&self) -> usize {
2650        self.inner.values.limits.max_trace_metric_size.as_bytes()
2651    }
2652
2653    /// Returns the maximum payload size for a compressed replay.
2654    pub fn max_replay_compressed_size(&self) -> usize {
2655        self.inner
2656            .values
2657            .limits
2658            .max_replay_compressed_size
2659            .as_bytes()
2660    }
2661
2662    /// Returns the maximum payload size for an uncompressed replay.
2663    pub fn max_replay_uncompressed_size(&self) -> usize {
2664        self.inner
2665            .values
2666            .limits
2667            .max_replay_uncompressed_size
2668            .as_bytes()
2669    }
2670
2671    /// Returns the maximum message size for an uncompressed replay.
2672    ///
2673    /// This is greater than max_replay_compressed_size because
2674    /// it can include additional metadata about the replay in
2675    /// addition to the recording.
2676    pub fn max_replay_message_size(&self) -> usize {
2677        self.inner.values.limits.max_replay_message_size.as_bytes()
2678    }
2679
2680    /// Returns the maximum number of active requests
2681    pub fn max_concurrent_requests(&self) -> usize {
2682        self.inner.values.limits.max_concurrent_requests
2683    }
2684
2685    /// Returns the maximum number of active queries
2686    pub fn max_concurrent_queries(&self) -> usize {
2687        self.inner.values.limits.max_concurrent_queries
2688    }
2689
2690    /// Returns the maximum combined size of keys of invalid attributes.
2691    pub fn max_removed_attribute_key_size(&self) -> usize {
2692        self.inner
2693            .values
2694            .limits
2695            .max_removed_attribute_key_size
2696            .as_bytes()
2697    }
2698
2699    /// The maximum number of seconds a query is allowed to take across retries.
2700    pub fn query_timeout(&self) -> Duration {
2701        Duration::from_secs(self.inner.values.limits.query_timeout)
2702    }
2703
2704    /// The maximum number of seconds to wait for pending envelopes after receiving a shutdown
2705    /// signal.
2706    pub fn shutdown_timeout(&self) -> Duration {
2707        Duration::from_secs(self.inner.values.limits.shutdown_timeout)
2708    }
2709
2710    /// Returns the server keep-alive timeout in seconds.
2711    ///
2712    /// By default keep alive is set to a 5 seconds.
2713    pub fn keepalive_timeout(&self) -> Duration {
2714        Duration::from_secs(self.inner.values.limits.keepalive_timeout)
2715    }
2716
2717    /// Returns the server idle timeout in seconds.
2718    pub fn idle_timeout(&self) -> Option<Duration> {
2719        self.inner
2720            .values
2721            .limits
2722            .idle_timeout
2723            .map(Duration::from_secs)
2724    }
2725
2726    /// Returns the maximum connections.
2727    pub fn max_connections(&self) -> Option<usize> {
2728        self.inner.values.limits.max_connections
2729    }
2730
2731    /// TCP listen backlog to configure on Relay's listening socket.
2732    pub fn tcp_listen_backlog(&self) -> u32 {
2733        self.inner.values.limits.tcp_listen_backlog
2734    }
2735
2736    /// Returns the number of cores to use for thread pools.
2737    pub fn cpu_concurrency(&self) -> usize {
2738        self.inner.values.limits.max_thread_count
2739    }
2740
2741    /// Returns the number of tasks that can run concurrently in the worker pool.
2742    pub fn pool_concurrency(&self) -> usize {
2743        self.inner.values.limits.max_pool_concurrency
2744    }
2745
2746    /// Returns the maximum size of a project config query.
2747    pub fn query_batch_size(&self) -> usize {
2748        self.inner.values.cache.batch_size
2749    }
2750
2751    /// True if the Relay should do processing.
2752    pub fn processing_enabled(&self) -> bool {
2753        self.inner.values.processing.enabled
2754    }
2755
2756    /// Level of normalization for Relay to apply to incoming data.
2757    pub fn normalization_level(&self) -> NormalizationLevel {
2758        self.inner.values.normalization.level
2759    }
2760
2761    /// The path to the GeoIp database required for event processing.
2762    pub fn geoip_path(&self) -> Option<&Path> {
2763        self.inner.values.geoip.path.as_deref().or(self
2764            .inner
2765            .values
2766            .processing
2767            .geoip_path
2768            .as_deref())
2769    }
2770
2771    /// Maximum future timestamp of ingested data.
2772    ///
2773    /// Events past this timestamp will be adjusted to `now()`. Sessions will be dropped.
2774    pub fn max_secs_in_future(&self) -> i64 {
2775        self.inner.values.processing.max_secs_in_future.into()
2776    }
2777
2778    /// Maximum age of ingested sessions. Older sessions will be dropped.
2779    pub fn max_session_secs_in_past(&self) -> i64 {
2780        self.inner.values.processing.max_session_secs_in_past.into()
2781    }
2782
2783    /// Configuration name and list of Kafka configuration parameters for a given topic.
2784    pub fn kafka_configs(
2785        &self,
2786        topic: KafkaTopic,
2787    ) -> Result<KafkaTopicConfig<'_>, KafkaConfigError> {
2788        self.inner
2789            .values
2790            .processing
2791            .topics
2792            .get(topic)
2793            .kafka_configs(
2794                &self.inner.values.processing.kafka_config,
2795                &self.inner.values.processing.secondary_kafka_configs,
2796            )
2797    }
2798
2799    /// Whether to validate the topics against Kafka.
2800    pub fn kafka_validate_topics(&self) -> bool {
2801        self.inner.values.processing.kafka_validate_topics
2802    }
2803
2804    /// All unused but configured topic assignments.
2805    pub fn unused_topic_assignments(&self) -> &relay_kafka::Unused {
2806        &self.inner.values.processing.topics.unused
2807    }
2808
2809    /// Configuration of the objectstore service.
2810    pub fn objectstore(&self) -> &ObjectstoreServiceConfig {
2811        &self.inner.values.processing.objectstore
2812    }
2813
2814    /// Configuration of the upload service.
2815    pub fn upload(&self) -> &Upload {
2816        &self.inner.values.upload
2817    }
2818
2819    /// Returns the key used to sign upload locations.
2820    #[cfg(feature = "processing")]
2821    pub fn upload_signing_key(&self) -> Option<&SecretKey> {
2822        self.upload()
2823            .credentials
2824            .as_ref()
2825            .map(|c| &c.signing_key)
2826            .or(self.credentials().map(|c| &c.secret_key))
2827    }
2828
2829    /// Returns the key used to verify upload locations.
2830    #[cfg(feature = "processing")]
2831    pub fn upload_verification_key(&self) -> Option<&PublicKey> {
2832        self.upload()
2833            .credentials
2834            .as_ref()
2835            .map(|c| &c.verification_key)
2836            .or(self.credentials().map(|c| &c.public_key))
2837    }
2838
2839    /// Redis servers to connect to for project configs, rate limiting, and metrics metadata.
2840    pub fn redis(&self) -> Option<RedisConfigsRef<'_>> {
2841        let redis_configs = self.inner.values.processing.redis.as_ref()?;
2842
2843        Some(build_redis_configs(
2844            redis_configs,
2845            self.cpu_concurrency() as u32,
2846            self.pool_concurrency() as u32,
2847        ))
2848    }
2849
2850    /// Chunk size of attachments in bytes.
2851    pub fn attachment_chunk_size(&self) -> usize {
2852        self.inner
2853            .values
2854            .processing
2855            .attachment_chunk_size
2856            .as_bytes()
2857    }
2858
2859    /// Maximum metrics batch size in bytes.
2860    pub fn metrics_max_batch_size_bytes(&self) -> usize {
2861        self.inner.values.aggregator.max_flush_bytes
2862    }
2863
2864    /// Default prefix to use when looking up project configs in Redis. This is only done when
2865    /// Relay is in processing mode.
2866    pub fn projectconfig_cache_prefix(&self) -> &str {
2867        &self.inner.values.processing.projectconfig_cache_prefix
2868    }
2869
2870    /// Maximum rate limit to report to clients in seconds.
2871    pub fn max_rate_limit(&self) -> Option<u64> {
2872        self.inner.values.processing.max_rate_limit.map(u32::into)
2873    }
2874
2875    /// Amount of remaining quota which is cached in memory.
2876    pub fn quota_cache_ratio(&self) -> Option<f32> {
2877        self.inner.values.processing.quota_cache_ratio
2878    }
2879
2880    /// Maximum limit (ratio) for the in memory quota cache.
2881    pub fn quota_cache_max(&self) -> Option<f32> {
2882        self.inner.values.processing.quota_cache_max
2883    }
2884
2885    /// Interval to refresh internal health checks.
2886    pub fn health_refresh_interval(&self) -> Duration {
2887        Duration::from_millis(self.inner.values.health.refresh_interval_ms)
2888    }
2889
2890    /// Maximum memory watermark in bytes.
2891    pub fn health_max_memory_watermark_bytes(&self) -> u64 {
2892        self.inner
2893            .values
2894            .health
2895            .max_memory_bytes
2896            .as_ref()
2897            .map_or(u64::MAX, |b| b.as_bytes() as u64)
2898    }
2899
2900    /// Maximum memory watermark as a percentage of maximum system memory.
2901    pub fn health_max_memory_watermark_percent(&self) -> f32 {
2902        self.inner.values.health.max_memory_percent
2903    }
2904
2905    /// Health check probe timeout.
2906    pub fn health_probe_timeout(&self) -> Duration {
2907        Duration::from_millis(self.inner.values.health.probe_timeout_ms)
2908    }
2909
2910    /// Refresh frequency for polling new memory stats.
2911    pub fn memory_stat_refresh_frequency_ms(&self) -> u64 {
2912        self.inner.values.health.memory_stat_refresh_frequency_ms
2913    }
2914
2915    /// Maximum amount of COGS measurements buffered in memory.
2916    pub fn cogs_max_queue_size(&self) -> u64 {
2917        self.inner.values.cogs.max_queue_size
2918    }
2919
2920    /// Resource ID to use for Relay COGS measurements.
2921    pub fn cogs_relay_resource_id(&self) -> &str {
2922        &self.inner.values.cogs.relay_resource_id
2923    }
2924
2925    /// Returns configuration for the default metrics aggregator.
2926    pub fn default_aggregator_config(&self) -> &AggregatorServiceConfig {
2927        &self.inner.values.aggregator
2928    }
2929
2930    /// Returns configuration for non-default metrics aggregator.
2931    pub fn secondary_aggregator_configs(&self) -> &Vec<ScopedAggregatorConfig> {
2932        &self.inner.values.secondary_aggregators
2933    }
2934
2935    /// Returns aggregator config for a given metrics namespace.
2936    pub fn aggregator_config_for(&self, namespace: MetricNamespace) -> &AggregatorServiceConfig {
2937        for entry in &self.inner.values.secondary_aggregators {
2938            if entry.condition.matches(Some(namespace)) {
2939                return &entry.config;
2940            }
2941        }
2942        &self.inner.values.aggregator
2943    }
2944
2945    /// Return the statically configured Relays.
2946    pub fn static_relays(&self) -> &HashMap<RelayId, RelayInfo> {
2947        &self.inner.values.auth.static_relays
2948    }
2949
2950    /// Returns the max age a signature is considered valid, in seconds.
2951    pub fn signature_max_age(&self) -> Duration {
2952        Duration::from_secs(self.inner.values.auth.signature_max_age)
2953    }
2954
2955    /// Returns `true` if unknown items should be accepted and forwarded.
2956    pub fn accept_unknown_items(&self) -> bool {
2957        let forward = self.inner.values.routing.accept_unknown_items;
2958        forward.unwrap_or_else(|| !self.processing_enabled())
2959    }
2960}
2961
2962impl fmt::Debug for ConfigSnapshot {
2963    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2964        f.debug_struct("ConfigSnapshot")
2965            .field("values", &self.inner.values)
2966            .finish()
2967    }
2968}
2969
2970#[cfg(test)]
2971mod tests {
2972    use super::*;
2973
2974    #[cfg(feature = "processing")]
2975    #[test]
2976    fn test_upload_secret_key_from_file() {
2977        let path = env::temp_dir().join(Uuid::new_v4().to_string());
2978        fs::create_dir(&path).unwrap();
2979        fs::write(
2980            path.join("my_secret.txt"),
2981            "U3LSQM5NorvgnoYHW_aZpc_43nuuh3lhs3zjjcBwaks",
2982        )
2983        .unwrap();
2984        fs::write(
2985            ConfigValues::path(&path),
2986            r#"
2987    upload:
2988        credentials:
2989            signing_key: ${file:my_secret.txt}
2990            verification_key: "VNS8haF0VTnuMMDR2t-f7AgnmUcXmcdzV3SVksSk34s""#,
2991        )
2992        .unwrap();
2993
2994        let config = Config::from_path(&path).unwrap().current();
2995
2996        fs::remove_dir_all(path).unwrap();
2997
2998        let signing_key = &config.upload().credentials.as_ref().unwrap().signing_key;
2999        assert_eq!(
3000            signing_key.to_string(),
3001            "U3LSQM5NorvgnoYHW_aZpc_43nuuh3lhs3zjjcBwaks"
3002        );
3003    }
3004
3005    #[test]
3006    fn test_emit_outcomes() {
3007        for (serialized, deserialized) in &[
3008            ("true", EmitOutcomes::AsOutcomes),
3009            ("false", EmitOutcomes::None),
3010            ("\"as_client_reports\"", EmitOutcomes::AsClientReports),
3011        ] {
3012            let value: EmitOutcomes = serde_json::from_str(serialized).unwrap();
3013            assert_eq!(value, *deserialized);
3014            assert_eq!(serde_json::to_string(&value).unwrap(), *serialized);
3015        }
3016    }
3017
3018    #[test]
3019    fn test_emit_outcomes_invalid() {
3020        assert!(serde_json::from_str::<EmitOutcomes>("asdf").is_err());
3021    }
3022}