Skip to main content

relay_server/
service.rs

1use std::convert::Infallible;
2use std::sync::Arc;
3use std::time::Duration;
4
5use crate::metrics::MetricOutcomes;
6use crate::services::autoscaling::{AutoscalingMetricService, AutoscalingMetrics};
7use crate::services::buffer::{
8    ObservableEnvelopeBuffer, PartitionedEnvelopeBuffer, ProjectKeyPair,
9};
10use crate::services::cogs::{CogsService, CogsServiceRecorder};
11use crate::services::config_reload::ConfigReloadService;
12use crate::services::global_config::{
13    GlobalConfigHandle, GlobalConfigManager, GlobalConfigService,
14};
15use crate::services::health_check::{HealthCheck, HealthCheckService};
16use crate::services::metrics::RouterService;
17#[cfg(feature = "processing")]
18use crate::services::objectstore::Objectstore;
19#[cfg(feature = "processing")]
20use crate::services::objectstore::ObjectstoreService;
21use crate::services::outcome::{
22    ClientReportOutcomeProducerService, NullOutcomeProducerService, OutcomeProducerService,
23    TrackOutcome,
24};
25use crate::services::processor::{
26    self, EnvelopeProcessor, EnvelopeProcessorService, EnvelopeProcessorServicePool,
27};
28use crate::services::projects::cache::{ProjectCacheHandle, ProjectCacheService};
29use crate::services::projects::source::ProjectSource;
30use crate::services::proxy_processor::{ProxyAddrs, ProxyProcessorService};
31use crate::services::relays::{RelayCache, RelayCacheService};
32use crate::services::stats::RelayStats;
33#[cfg(feature = "processing")]
34use crate::services::store::{StoreService, StoreServicePool};
35use crate::services::upload::{self, Upload};
36use crate::services::upstream::{UpstreamRelay, UpstreamRelayService};
37use crate::utils::{MemoryChecker, MemoryStat, ThreadKind};
38#[cfg(feature = "processing")]
39use anyhow::Context;
40use anyhow::Result;
41use axum::extract::FromRequestParts;
42use axum::http::request::Parts;
43use relay_cogs::Cogs;
44use relay_config::{Config, ConfigSnapshot, EmitOutcomes, RelayMode};
45#[cfg(feature = "processing")]
46use relay_config::{RedisConfigRef, RedisConfigsRef};
47#[cfg(feature = "processing")]
48use relay_redis::AsyncRedisClient;
49#[cfg(feature = "processing")]
50use relay_redis::redis::Script;
51#[cfg(feature = "processing")]
52use relay_redis::{RedisClients, RedisError, RedisScripts};
53#[cfg(feature = "processing")]
54use relay_system::ConcurrentService;
55use relay_system::{Addr, Service, ServiceSpawn, ServiceSpawnExt as _, channel};
56
57/// Indicates the type of failure of the server.
58#[derive(Debug, thiserror::Error)]
59pub enum ServiceError {
60    /// GeoIp construction failed.
61    #[error("could not load the Geoip Db")]
62    GeoIp,
63
64    /// Initializing the Kafka producer failed.
65    #[cfg(feature = "processing")]
66    #[error("could not initialize kafka producer: {0}")]
67    Kafka(String),
68
69    /// Initializing the Redis client failed.
70    #[cfg(feature = "processing")]
71    #[error("could not initialize redis client during startup")]
72    Redis,
73}
74
75#[derive(Clone, Debug)]
76pub struct Registry {
77    pub cogs: Cogs,
78    pub health_check: Addr<HealthCheck>,
79    pub outcome_aggregator: Addr<TrackOutcome>,
80    pub processor: Addr<EnvelopeProcessor>,
81    pub relay_cache: Addr<RelayCache>,
82    pub global_config: Addr<GlobalConfigManager>,
83    pub upstream_relay: Addr<UpstreamRelay>,
84    pub envelope_buffer: Arc<PartitionedEnvelopeBuffer>,
85    pub project_cache_handle: ProjectCacheHandle,
86    pub autoscaling: Option<Addr<AutoscalingMetrics>>,
87    #[cfg(feature = "processing")]
88    pub objectstore: Option<Addr<Objectstore>>,
89    pub upload: Addr<Upload>,
90    pub global_config_handle: GlobalConfigHandle,
91}
92
93/// Constructs a Tokio [`relay_system::Runtime`] configured for running [services](relay_system::Service).
94pub fn create_runtime(name: &'static str, threads: usize) -> relay_system::Runtime {
95    relay_system::Runtime::builder(name)
96        .worker_threads(threads)
97        // Relay uses `spawn_blocking` only for Redis connections within the project
98        // cache, those should never exceed 100 concurrent connections
99        // (limited by connection pool).
100        //
101        // Relay also does not use other blocking operations from Tokio which require
102        // this pool, no usage of `tokio::fs` and `tokio::io::{Stdin, Stdout, Stderr}`.
103        //
104        // We limit the maximum amount of threads here, we've seen that Tokio
105        // expands this pool very very aggressively and basically never shrinks it
106        // which leads to a massive resource waste.
107        .max_blocking_threads(150)
108        // We also lower down the default (10s) keep alive timeout for blocking
109        // threads to encourage the runtime to not keep too many idle blocking threads
110        // around.
111        .thread_keep_alive(Duration::from_secs(1))
112        .build()
113}
114
115fn create_processor_pool(config: &ConfigSnapshot) -> Result<EnvelopeProcessorServicePool> {
116    // Adjust thread count for small cpu counts to not have too many idle cores
117    // and distribute workload better.
118    let thread_count = match config.cpu_concurrency() {
119        conc @ 0..=2 => conc.max(1),
120        conc @ 3..=4 => conc - 1,
121        conc => conc - 2,
122    };
123    relay_log::info!("starting {thread_count} envelope processing workers");
124
125    let pool = crate::utils::ThreadPoolBuilder::new("processor", tokio::runtime::Handle::current())
126        .num_threads(thread_count)
127        .max_concurrency(config.pool_concurrency())
128        .thread_kind(ThreadKind::Worker)
129        .build()?;
130
131    Ok(pool)
132}
133
134#[cfg(feature = "processing")]
135fn create_store_pool(config: &ConfigSnapshot) -> Result<StoreServicePool> {
136    // Spawn a store worker for every 12 threads in the processor pool.
137    // This ratio was found empirically and may need adjustments in the future.
138    //
139    // Ideally in the future the store will be single threaded again, after we move
140    // all the heavy processing (de- and re-serialization) into the processor.
141    let thread_count = config.cpu_concurrency().div_ceil(8);
142    relay_log::info!("starting {thread_count} store workers");
143
144    let pool = crate::utils::ThreadPoolBuilder::new("store", tokio::runtime::Handle::current())
145        .num_threads(thread_count)
146        .max_concurrency(config.pool_concurrency())
147        .build()?;
148
149    Ok(pool)
150}
151
152#[derive(Debug)]
153struct StateInner {
154    config: Arc<Config>,
155    memory_checker: MemoryChecker,
156    registry: Registry,
157}
158
159/// Server state.
160#[derive(Clone, Debug)]
161pub struct ServiceState {
162    inner: Arc<StateInner>,
163}
164
165impl ServiceState {
166    /// Starts all services and returns addresses to all of them.
167    pub async fn start(
168        handle: &relay_system::Handle,
169        services: &dyn ServiceSpawn,
170        config: Arc<Config>,
171    ) -> Result<Self> {
172        let upstream_relay = services.start(UpstreamRelayService::new(config.clone()));
173        let current_config = config.current();
174
175        #[cfg(feature = "processing")]
176        let redis_clients = current_config
177            .redis()
178            .filter(|_| current_config.processing_enabled())
179            .map(create_redis_clients)
180            .transpose()
181            .context(ServiceError::Redis)?;
182
183        // If we have Redis configured, we want to initialize all the scripts by loading them in
184        // the scripts cache if not present. Our custom ConnectionLike implementation relies on this
185        // initialization to work properly since it assumes that scripts are loaded across all Redis
186        // instances.
187        #[cfg(feature = "processing")]
188        if let Some(redis_clients) = &redis_clients {
189            initialize_redis_scripts_for_client(redis_clients)
190                .await
191                .context(ServiceError::Redis)?;
192        }
193
194        // We create an instance of `MemoryStat` which can be supplied composed with any arbitrary
195        // configuration object down the line.
196        let memory_stat = MemoryStat::new(current_config.memory_stat_refresh_frequency_ms());
197
198        // Create an address for the `EnvelopeProcessor`, which can be injected into the
199        // other services.
200        let (processor, processor_rx) = match current_config.relay_mode() {
201            RelayMode::Proxy => channel(ProxyProcessorService::name()),
202            RelayMode::Managed => channel(EnvelopeProcessorService::name()),
203        };
204
205        let (aggregator, aggregator_rx) = channel(RouterService::name());
206
207        let outcome_aggregator = match current_config.emit_outcomes() {
208            EmitOutcomes::None => services.start(NullOutcomeProducerService::new()),
209            _ => match current_config.relay_mode() {
210                RelayMode::Proxy => services.start(ClientReportOutcomeProducerService::new(
211                    &current_config,
212                    processor.clone(),
213                )),
214                RelayMode::Managed => services.start(OutcomeProducerService::new(
215                    Arc::clone(&config),
216                    aggregator.clone(),
217                )),
218            },
219        };
220
221        let (global_config, global_config_rx) =
222            GlobalConfigService::new(config.clone(), upstream_relay.clone());
223        let global_config_handle = global_config.handle();
224        // The global config service must start before dependant services are
225        // started. Messages like subscription requests to the global config
226        // service fail if the service is not running.
227        let global_config = services.start(global_config);
228
229        let project_source = ProjectSource::start_in(
230            services,
231            Arc::clone(&config),
232            upstream_relay.clone(),
233            #[cfg(feature = "processing")]
234            redis_clients.clone(),
235        )
236        .await;
237        let project_cache_handle =
238            ProjectCacheService::new(Arc::clone(&config), project_source).start_in(services);
239
240        let cogs = CogsService::new(&current_config);
241        let cogs = Cogs::new(CogsServiceRecorder::new(
242            &current_config,
243            services.start(cogs),
244        ));
245
246        let metric_outcomes = MetricOutcomes::new(outcome_aggregator.clone());
247
248        #[cfg(feature = "processing")]
249        let store_pool = create_store_pool(&current_config)?;
250        #[cfg(feature = "processing")]
251        let store = current_config
252            .processing_enabled()
253            .then(|| {
254                StoreService::create(
255                    store_pool.clone(),
256                    config.clone(),
257                    global_config_handle.clone(),
258                    metric_outcomes.clone(),
259                )
260                .map(|s| services.start(s))
261            })
262            .transpose()?;
263
264        #[cfg(feature = "processing")]
265        let objectstore = ObjectstoreService::new(current_config.objectstore(), store.clone())?
266            .map(|s| {
267                let concurrent = ConcurrentService::new(s)
268                    .with_backlog_limit(current_config.objectstore().max_backlog)
269                    .with_concurrency_limit(current_config.objectstore().max_concurrent_requests);
270                services.start(concurrent)
271            });
272
273        let envelope_buffer = PartitionedEnvelopeBuffer::create(
274            current_config.spool_partitions(),
275            config.clone(),
276            memory_stat.clone(),
277            global_config_rx.clone(),
278            project_cache_handle.clone(),
279            processor.clone(),
280            outcome_aggregator.clone(),
281            services,
282        );
283
284        let (processor_pool, aggregator_handle, autoscaling) = match current_config.relay_mode() {
285            RelayMode::Proxy => {
286                services.start_with(
287                    ProxyProcessorService::new(
288                        config.clone(),
289                        project_cache_handle.clone(),
290                        ProxyAddrs {
291                            outcome_aggregator: outcome_aggregator.clone(),
292                            upstream_relay: upstream_relay.clone(),
293                        },
294                    ),
295                    processor_rx,
296                );
297                (None, None, None)
298            }
299            RelayMode::Managed => {
300                let processor_pool = create_processor_pool(&current_config)?;
301
302                let router = RouterService::new(
303                    handle.clone(),
304                    current_config.default_aggregator_config().clone(),
305                    current_config.secondary_aggregator_configs().clone(),
306                    Some(processor.clone().recipient()),
307                    project_cache_handle.clone(),
308                );
309                let router_handle = router.handle();
310                services.start_with(router, aggregator_rx);
311
312                services.start_with(
313                    EnvelopeProcessorService::new(
314                        processor_pool.clone(),
315                        config.clone(),
316                        global_config_handle.clone(),
317                        project_cache_handle.clone(),
318                        cogs.clone(),
319                        #[cfg(feature = "processing")]
320                        redis_clients.clone(),
321                        processor::Addrs {
322                            outcome_aggregator: outcome_aggregator.clone(),
323                            upstream_relay: upstream_relay.clone(),
324                            #[cfg(feature = "processing")]
325                            objectstore: objectstore.clone(),
326                            #[cfg(feature = "processing")]
327                            store_forwarder: store,
328                            aggregator: aggregator.clone(),
329                        },
330                        metric_outcomes.clone(),
331                    ),
332                    processor_rx,
333                );
334
335                let autoscaling = services.start(AutoscalingMetricService::new(
336                    memory_stat.clone(),
337                    envelope_buffer.clone(),
338                    handle.clone(),
339                    processor_pool.clone(),
340                ));
341
342                (Some(processor_pool), Some(router_handle), Some(autoscaling))
343            }
344        };
345
346        let health_check = services.start(HealthCheckService::new(
347            config.clone(),
348            MemoryChecker::new(memory_stat.clone(), config.clone()),
349            aggregator_handle,
350            upstream_relay.clone(),
351            envelope_buffer.clone(),
352        ));
353
354        services.start(RelayStats::new(
355            config.clone(),
356            handle.clone(),
357            upstream_relay.clone(),
358            #[cfg(feature = "processing")]
359            redis_clients.clone(),
360            processor_pool,
361            #[cfg(feature = "processing")]
362            store_pool,
363        ));
364
365        let relay_cache = services.start(RelayCacheService::new(
366            config.clone(),
367            upstream_relay.clone(),
368        ));
369
370        let upload = services.start(upload::create_service(
371            &config,
372            &upstream_relay,
373            #[cfg(feature = "processing")]
374            &objectstore,
375        ));
376
377        let _ = services.start(ConfigReloadService::new(config.clone()));
378
379        let registry = Registry {
380            cogs,
381            processor,
382            health_check,
383            outcome_aggregator,
384            relay_cache,
385            global_config,
386            project_cache_handle,
387            upstream_relay,
388            envelope_buffer,
389            autoscaling,
390            #[cfg(feature = "processing")]
391            objectstore,
392            upload,
393            global_config_handle,
394        };
395
396        let state = StateInner {
397            config: config.clone(),
398            memory_checker: MemoryChecker::new(memory_stat, config.clone()),
399            registry,
400        };
401
402        Ok(ServiceState {
403            inner: Arc::new(state),
404        })
405    }
406
407    /// Returns a snapshot of the Relay configuration.
408    pub fn config(&self) -> ConfigSnapshot {
409        self.inner.config.current()
410    }
411
412    /// Returns a reference to the [`Cogs`] tracking.
413    pub fn cogs(&self) -> &Cogs {
414        &self.inner.registry.cogs
415    }
416
417    /// Returns a reference to the [`MemoryChecker`] which is a [`Config`] aware wrapper on the
418    /// [`MemoryStat`] which gives utility methods to determine whether memory usage is above
419    /// thresholds set in the [`Config`].
420    pub fn memory_checker(&self) -> &MemoryChecker {
421        &self.inner.memory_checker
422    }
423
424    pub fn autoscaling(&self) -> Option<&Addr<AutoscalingMetrics>> {
425        self.inner.registry.autoscaling.as_ref()
426    }
427
428    /// Returns the V2 envelope buffer, if present.
429    pub fn envelope_buffer(&self, project_key_pair: ProjectKeyPair) -> &ObservableEnvelopeBuffer {
430        self.inner.registry.envelope_buffer.buffer(project_key_pair)
431    }
432
433    /// Returns a [`ProjectCacheHandle`].
434    pub fn project_cache_handle(&self) -> &ProjectCacheHandle {
435        &self.inner.registry.project_cache_handle
436    }
437
438    /// Returns the address of the [`RelayCache`] service.
439    pub fn relay_cache(&self) -> &Addr<RelayCache> {
440        &self.inner.registry.relay_cache
441    }
442
443    /// Returns the address of the [`HealthCheck`] service.
444    pub fn health_check(&self) -> &Addr<HealthCheck> {
445        &self.inner.registry.health_check
446    }
447
448    /// Returns the address of the [`UpstreamRelay`] service.
449    pub fn upstream_relay(&self) -> &Addr<UpstreamRelay> {
450        &self.inner.registry.upstream_relay
451    }
452
453    /// Returns the address of the [`EnvelopeProcessor`] service.
454    pub fn processor(&self) -> &Addr<EnvelopeProcessor> {
455        &self.inner.registry.processor
456    }
457
458    /// Returns the address of the [`GlobalConfigService`] service.
459    pub fn global_config(&self) -> &Addr<GlobalConfigManager> {
460        &self.inner.registry.global_config
461    }
462
463    /// Returns the address of the [`TrackOutcome`] service.
464    pub fn outcome_aggregator(&self) -> &Addr<TrackOutcome> {
465        &self.inner.registry.outcome_aggregator
466    }
467
468    #[cfg(feature = "processing")]
469    /// Returns the address of the [`Objectstore`] service.
470    pub fn objectstore(&self) -> Option<&Addr<Objectstore>> {
471        self.inner.registry.objectstore.as_ref()
472    }
473
474    /// Returns the address of the [`Upload`] service.
475    pub fn upload(&self) -> &Addr<Upload> {
476        &self.inner.registry.upload
477    }
478
479    pub fn global_config_handle(&self) -> &GlobalConfigHandle {
480        &self.inner.registry.global_config_handle
481    }
482}
483
484/// Creates Redis clients from the given `configs`.
485///
486/// If `configs` is [`Unified`](RedisConfigsRef::Unified), one client is created and then cloned
487/// for project configs, cardinality, and quotas, meaning that they really use the same client.
488///
489/// If it is [`Individual`](RedisConfigsRef::Individual), an actual separate client
490/// is created for each use case.
491#[cfg(feature = "processing")]
492pub fn create_redis_clients(configs: RedisConfigsRef<'_>) -> Result<RedisClients, RedisError> {
493    const PROJECT_CONFIG_REDIS_CLIENT: &str = "projectconfig";
494    const QUOTA_REDIS_CLIENT: &str = "quotas";
495    const UNIFIED_REDIS_CLIENT: &str = "unified";
496
497    match configs {
498        RedisConfigsRef::Unified(unified) => {
499            let client = create_async_redis_client(UNIFIED_REDIS_CLIENT, &unified)?;
500
501            Ok(RedisClients {
502                project_configs: client.clone(),
503                quotas: client,
504            })
505        }
506        RedisConfigsRef::Individual {
507            project_configs,
508            quotas,
509        } => {
510            let project_configs =
511                create_async_redis_client(PROJECT_CONFIG_REDIS_CLIENT, &project_configs)?;
512            let quotas = create_async_redis_client(QUOTA_REDIS_CLIENT, &quotas)?;
513
514            Ok(RedisClients {
515                project_configs,
516                quotas,
517            })
518        }
519    }
520}
521
522#[cfg(feature = "processing")]
523fn create_async_redis_client(
524    name: &'static str,
525    config: &RedisConfigRef<'_>,
526) -> Result<AsyncRedisClient, RedisError> {
527    match config {
528        RedisConfigRef::Cluster {
529            cluster_nodes,
530            options,
531        } => AsyncRedisClient::cluster(name, cluster_nodes.iter().map(|s| s.as_str()), options),
532        RedisConfigRef::Single { server, options } => {
533            AsyncRedisClient::single(name, server, options)
534        }
535    }
536}
537
538#[cfg(feature = "processing")]
539async fn initialize_redis_scripts_for_client(
540    redis_clients: &RedisClients,
541) -> Result<(), RedisError> {
542    let scripts = RedisScripts::all();
543
544    let RedisClients {
545        project_configs,
546        quotas,
547    } = redis_clients;
548
549    initialize_redis_scripts(project_configs, &scripts).await?;
550    initialize_redis_scripts(quotas, &scripts).await?;
551
552    Ok(())
553}
554
555#[cfg(feature = "processing")]
556async fn initialize_redis_scripts(
557    client: &AsyncRedisClient,
558    scripts: &[&Script],
559) -> Result<(), RedisError> {
560    let mut connection = client.get_connection().await?;
561
562    for script in scripts {
563        // We load on all instances without checking if the script is already in cache because of a
564        // limitation in the connection implementation.
565        script
566            .prepare_invoke()
567            .load_async(&mut connection)
568            .await
569            .map_err(RedisError::Redis)?;
570    }
571
572    Ok(())
573}
574
575impl FromRequestParts<Self> for ServiceState {
576    type Rejection = Infallible;
577
578    async fn from_request_parts(_: &mut Parts, state: &Self) -> Result<Self, Self::Rejection> {
579        Ok(state.clone())
580    }
581}