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#[derive(Debug, thiserror::Error)]
59pub enum ServiceError {
60 #[error("could not load the Geoip Db")]
62 GeoIp,
63
64 #[cfg(feature = "processing")]
66 #[error("could not initialize kafka producer: {0}")]
67 Kafka(String),
68
69 #[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
93pub fn create_runtime(name: &'static str, threads: usize) -> relay_system::Runtime {
95 relay_system::Runtime::builder(name)
96 .worker_threads(threads)
97 .max_blocking_threads(150)
108 .thread_keep_alive(Duration::from_secs(1))
112 .build()
113}
114
115fn create_processor_pool(config: &ConfigSnapshot) -> Result<EnvelopeProcessorServicePool> {
116 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 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#[derive(Clone, Debug)]
161pub struct ServiceState {
162 inner: Arc<StateInner>,
163}
164
165impl ServiceState {
166 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 #[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 let memory_stat = MemoryStat::new(current_config.memory_stat_refresh_frequency_ms());
197
198 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 ¤t_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 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(¤t_config);
241 let cogs = Cogs::new(CogsServiceRecorder::new(
242 ¤t_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(¤t_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(¤t_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 pub fn config(&self) -> ConfigSnapshot {
409 self.inner.config.current()
410 }
411
412 pub fn cogs(&self) -> &Cogs {
414 &self.inner.registry.cogs
415 }
416
417 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 pub fn envelope_buffer_is_empty(&self) -> bool {
430 self.inner.registry.envelope_buffer.is_empty()
431 }
432
433 pub fn envelope_buffer(&self, project_key_pair: ProjectKeyPair) -> &ObservableEnvelopeBuffer {
435 self.inner.registry.envelope_buffer.buffer(project_key_pair)
436 }
437
438 pub fn project_cache_handle(&self) -> &ProjectCacheHandle {
440 &self.inner.registry.project_cache_handle
441 }
442
443 pub fn relay_cache(&self) -> &Addr<RelayCache> {
445 &self.inner.registry.relay_cache
446 }
447
448 pub fn health_check(&self) -> &Addr<HealthCheck> {
450 &self.inner.registry.health_check
451 }
452
453 pub fn upstream_relay(&self) -> &Addr<UpstreamRelay> {
455 &self.inner.registry.upstream_relay
456 }
457
458 pub fn processor(&self) -> &Addr<EnvelopeProcessor> {
460 &self.inner.registry.processor
461 }
462
463 pub fn global_config(&self) -> &Addr<GlobalConfigManager> {
465 &self.inner.registry.global_config
466 }
467
468 pub fn outcome_aggregator(&self) -> &Addr<TrackOutcome> {
470 &self.inner.registry.outcome_aggregator
471 }
472
473 #[cfg(feature = "processing")]
474 pub fn objectstore(&self) -> Option<&Addr<Objectstore>> {
476 self.inner.registry.objectstore.as_ref()
477 }
478
479 pub fn upload(&self) -> &Addr<Upload> {
481 &self.inner.registry.upload
482 }
483
484 pub fn global_config_handle(&self) -> &GlobalConfigHandle {
485 &self.inner.registry.global_config_handle
486 }
487}
488
489#[cfg(feature = "processing")]
497pub fn create_redis_clients(configs: RedisConfigsRef<'_>) -> Result<RedisClients, RedisError> {
498 const PROJECT_CONFIG_REDIS_CLIENT: &str = "projectconfig";
499 const QUOTA_REDIS_CLIENT: &str = "quotas";
500 const UNIFIED_REDIS_CLIENT: &str = "unified";
501
502 match configs {
503 RedisConfigsRef::Unified(unified) => {
504 let client = create_async_redis_client(UNIFIED_REDIS_CLIENT, &unified)?;
505
506 Ok(RedisClients {
507 project_configs: client.clone(),
508 quotas: client,
509 })
510 }
511 RedisConfigsRef::Individual {
512 project_configs,
513 quotas,
514 } => {
515 let project_configs =
516 create_async_redis_client(PROJECT_CONFIG_REDIS_CLIENT, &project_configs)?;
517 let quotas = create_async_redis_client(QUOTA_REDIS_CLIENT, "as)?;
518
519 Ok(RedisClients {
520 project_configs,
521 quotas,
522 })
523 }
524 }
525}
526
527#[cfg(feature = "processing")]
528fn create_async_redis_client(
529 name: &'static str,
530 config: &RedisConfigRef<'_>,
531) -> Result<AsyncRedisClient, RedisError> {
532 match config {
533 RedisConfigRef::Cluster {
534 cluster_nodes,
535 options,
536 } => AsyncRedisClient::cluster(name, cluster_nodes.iter().map(|s| s.as_str()), options),
537 RedisConfigRef::Single { server, options } => {
538 AsyncRedisClient::single(name, server, options)
539 }
540 }
541}
542
543#[cfg(feature = "processing")]
544async fn initialize_redis_scripts_for_client(
545 redis_clients: &RedisClients,
546) -> Result<(), RedisError> {
547 let scripts = RedisScripts::all();
548
549 let RedisClients {
550 project_configs,
551 quotas,
552 } = redis_clients;
553
554 initialize_redis_scripts(project_configs, &scripts).await?;
555 initialize_redis_scripts(quotas, &scripts).await?;
556
557 Ok(())
558}
559
560#[cfg(feature = "processing")]
561async fn initialize_redis_scripts(
562 client: &AsyncRedisClient,
563 scripts: &[&Script],
564) -> Result<(), RedisError> {
565 let mut connection = client.get_connection().await?;
566
567 for script in scripts {
568 script
571 .prepare_invoke()
572 .load_async(&mut connection)
573 .await
574 .map_err(RedisError::Redis)?;
575 }
576
577 Ok(())
578}
579
580impl FromRequestParts<Self> for ServiceState {
581 type Rejection = Infallible;
582
583 async fn from_request_parts(_: &mut Parts, state: &Self) -> Result<Self, Self::Rejection> {
584 Ok(state.clone())
585 }
586}