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(&self, project_key_pair: ProjectKeyPair) -> &ObservableEnvelopeBuffer {
430 self.inner.registry.envelope_buffer.buffer(project_key_pair)
431 }
432
433 pub fn project_cache_handle(&self) -> &ProjectCacheHandle {
435 &self.inner.registry.project_cache_handle
436 }
437
438 pub fn relay_cache(&self) -> &Addr<RelayCache> {
440 &self.inner.registry.relay_cache
441 }
442
443 pub fn health_check(&self) -> &Addr<HealthCheck> {
445 &self.inner.registry.health_check
446 }
447
448 pub fn upstream_relay(&self) -> &Addr<UpstreamRelay> {
450 &self.inner.registry.upstream_relay
451 }
452
453 pub fn processor(&self) -> &Addr<EnvelopeProcessor> {
455 &self.inner.registry.processor
456 }
457
458 pub fn global_config(&self) -> &Addr<GlobalConfigManager> {
460 &self.inner.registry.global_config
461 }
462
463 pub fn outcome_aggregator(&self) -> &Addr<TrackOutcome> {
465 &self.inner.registry.outcome_aggregator
466 }
467
468 #[cfg(feature = "processing")]
469 pub fn objectstore(&self) -> Option<&Addr<Objectstore>> {
471 self.inner.registry.objectstore.as_ref()
472 }
473
474 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#[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, "as)?;
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 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}