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