1#![recursion_limit = "256"]
11
12use ::http::HeaderValue;
19use std::collections::{BTreeMap, BTreeSet};
20use std::panic::AssertUnwindSafe;
21use std::path::PathBuf;
22use std::pin::Pin;
23use std::sync::{Arc, LazyLock, Mutex};
24use std::time::{Duration, Instant};
25use std::{env, io};
26
27use anyhow::{Context, anyhow};
28use derivative::Derivative;
29use futures::FutureExt;
30use ipnet::IpNet;
31use mz_adapter::config::{
32 SystemParameterSyncClientConfig, SystemParameterSyncConfig, system_parameter_sync,
33};
34use mz_adapter::webhook::WebhookConcurrencyLimiter;
35use mz_adapter::{AdapterError, Client as AdapterClient, load_remote_system_parameters};
36use mz_adapter_types::bootstrap_builtin_cluster_config::BootstrapBuiltinClusterConfig;
37use mz_adapter_types::dyncfgs::{
38 ENABLE_0DT_DEPLOYMENT_PANIC_AFTER_TIMEOUT, WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL,
39 WITH_0DT_DEPLOYMENT_MAX_WAIT,
40};
41use mz_auth::password::Password;
42use mz_authenticator::GenericOidcAuthenticator;
43use mz_build_info::{BuildInfo, build_info};
44use mz_catalog::config::ClusterReplicaSizeMap;
45use mz_catalog::durable::BootstrapArgs;
46use mz_cloud_resources::CloudResourceController;
47use mz_controller::ControllerConfig;
48use mz_dyncfg::ConfigSet;
49use mz_frontegg_auth::Authenticator as FronteggAuthenticator;
50use mz_license_keys::ValidatedLicenseKey;
51use mz_ore::future::OreFutureExt;
52use mz_ore::metrics::MetricsRegistry;
53use mz_ore::now::NowFn;
54use mz_ore::tracing::TracingHandle;
55use mz_ore::url::SensitiveUrl;
56use mz_ore::{instrument, task};
57use mz_persist_client::cache::PersistClientCache;
58use mz_persist_client::usage::StorageUsageClient;
59use mz_pgwire::MetricsConfig;
60use mz_pgwire_common::ConnectionCounter;
61use mz_repr::strconv;
62use mz_secrets::SecretsController;
63use mz_server_core::listeners::v26_32_0::ListenersConfig;
64use mz_server_core::listeners::{HttpListenerConfig, ListenerConfig, SqlListenerConfig};
65use mz_server_core::{
66 ConnectionStream, ListenerHandle, ReloadTrigger, ReloadingSslContext, ServeConfig,
67 TlsCertConfig, TlsMode,
68};
69use mz_sql::catalog::EnvironmentId;
70use mz_sql::session::vars::{Value, VarInput};
71use tokio::sync::oneshot;
72use tower_http::cors::AllowOrigin;
73use tracing::{Instrument, info, info_span};
74
75use crate::deployment::preflight::{self, CatchupConfig};
76use crate::deployment::state::DeploymentState;
77use crate::http::{HttpConfig, HttpServer, InternalRouteConfig};
78
79pub use crate::http::{SqlResponse, WebSocketAuth, WebSocketResponse};
80
81mod deployment;
82pub mod environmentd;
83pub mod http;
84mod telemetry;
85#[cfg(feature = "test")]
86pub mod test_util;
87
88pub const BUILD_INFO: BuildInfo = build_info!();
89
90#[derive(Derivative)]
92#[derivative(Debug)]
93pub struct Config {
94 pub unsafe_mode: bool,
98 pub all_features: bool,
101
102 pub tls: Option<TlsCertConfig>,
105 #[derivative(Debug = "ignore")]
107 pub tls_reload_certs: ReloadTrigger,
108 pub external_login_password_mz_system: Option<Password>,
110 pub frontegg: Option<FronteggAuthenticator>,
112 pub frontegg_oauth_issuer_url: Option<String>,
114 pub cors_allowed_origin: AllowOrigin,
117 pub cors_allowed_origin_list: Vec<HeaderValue>,
122 pub egress_addresses: Vec<IpNet>,
125 pub http_host_name: Option<String>,
131 pub internal_console_redirect_url: Option<String>,
134
135 pub controller: ControllerConfig,
138 pub secrets_controller: Arc<dyn SecretsController>,
140 pub cloud_resource_controller: Option<Arc<dyn CloudResourceController>>,
142 pub system_dyncfgs: Arc<ConfigSet>,
144
145 pub storage_usage_collection_interval: Duration,
148 pub storage_usage_retention_period: Option<Duration>,
150
151 pub catalog_config: CatalogConfig,
154 pub availability_zones: Vec<String>,
157 pub cluster_replica_sizes: ClusterReplicaSizeMap,
159 pub timestamp_oracle_url: Option<SensitiveUrl>,
161 pub segment_api_key: Option<String>,
163 pub segment_client_side: bool,
166 pub test_only_dummy_segment_client: bool,
168 pub launchdarkly_sdk_key: Option<String>,
171 pub launchdarkly_base_uri: Option<String>,
174 pub launchdarkly_key_map: BTreeMap<String, String>,
177 pub config_sync_timeout: Duration,
179 pub config_sync_loop_interval: Option<Duration>,
181 pub config_sync_file_path: Option<PathBuf>,
183
184 pub environment_id: EnvironmentId,
187 pub bootstrap_role: Option<String>,
189 pub bootstrap_default_cluster_replica_size: String,
191 pub bootstrap_default_cluster_replication_factor: u32,
193 pub bootstrap_builtin_system_cluster_config: BootstrapBuiltinClusterConfig,
195 pub bootstrap_builtin_catalog_server_cluster_config: BootstrapBuiltinClusterConfig,
197 pub bootstrap_builtin_probe_cluster_config: BootstrapBuiltinClusterConfig,
199 pub bootstrap_builtin_support_cluster_config: BootstrapBuiltinClusterConfig,
201 pub bootstrap_builtin_analytics_cluster_config: BootstrapBuiltinClusterConfig,
203 pub system_parameter_defaults: BTreeMap<String, String>,
206 pub helm_chart_version: Option<String>,
208 pub license_key: ValidatedLicenseKey,
210
211 pub aws_account_id: Option<String>,
215 pub aws_privatelink_availability_zones: Option<Vec<String>>,
217
218 pub metrics_registry: MetricsRegistry,
221 pub tracing_handle: TracingHandle,
223
224 pub now: NowFn,
227 pub force_builtin_schema_migration: Option<String>,
230}
231
232#[derive(Debug, Clone)]
234pub struct CatalogConfig {
235 pub persist_clients: Arc<PersistClientCache>,
237 pub metrics: Arc<mz_catalog::durable::Metrics>,
239}
240
241pub struct Listener<C> {
242 pub handle: ListenerHandle,
243 connection_stream: Pin<Box<dyn ConnectionStream>>,
244 config: C,
245}
246impl<C> Listener<C>
247where
248 C: ListenerConfig,
249{
250 async fn bind(config: C) -> Result<Self, io::Error> {
262 let (handle, connection_stream) = mz_server_core::listen(&config.addr()).await?;
263 Ok(Self {
264 handle,
265 connection_stream,
266 config,
267 })
268 }
269}
270
271impl Listener<SqlListenerConfig> {
272 #[instrument(name = "environmentd::serve_sql")]
273 pub async fn serve_sql(
274 self,
275 name: String,
276 active_connection_counter: ConnectionCounter,
277 tls_reloading_context: Option<ReloadingSslContext>,
278 frontegg: Option<FronteggAuthenticator>,
279 adapter_client: AdapterClient,
280 oidc: GenericOidcAuthenticator,
281 metrics: MetricsConfig,
282 helm_chart_version: Option<String>,
283 ) -> ListenerHandle {
284 let label = leak_listener_name(&name);
285 let tls = tls_reloading_context.map(|context| mz_server_core::ReloadingTlsConfig {
286 context,
287 mode: if self.config.enable_tls {
288 TlsMode::Require
289 } else {
290 TlsMode::Allow
291 },
292 });
293
294 task::spawn(|| format!("{}_sql_server", label), {
295 let sql_server = mz_pgwire::Server::new(mz_pgwire::Config {
296 label,
297 tls,
298 adapter_client,
299 authenticator_kind: self.config.authenticator_kind,
300 frontegg,
301 oidc,
302 metrics,
303 active_connection_counter,
304 helm_chart_version,
305 allowed_roles: self.config.allowed_roles,
306 });
307 mz_server_core::serve(ServeConfig {
308 conns: self.connection_stream,
309 server: sql_server,
310 dyncfg: None,
313 })
314 });
315 self.handle
316 }
317}
318
319impl Listener<HttpListenerConfig> {
320 #[instrument(name = "environmentd::serve_http")]
321 pub async fn serve_http(self, config: HttpConfig) -> ListenerHandle {
322 let task_name = format!("{}_http_server", config.source);
323 task::spawn(|| task_name, {
324 let http_server = HttpServer::new(config);
325 mz_server_core::serve(ServeConfig {
326 conns: self.connection_stream,
327 server: http_server,
328 dyncfg: None,
331 })
332 });
333 self.handle
334 }
335}
336
337static LISTENER_NAMES: LazyLock<Mutex<BTreeSet<&'static str>>> =
340 LazyLock::new(|| Mutex::new(BTreeSet::new()));
341
342fn leak_listener_name(name: &str) -> &'static str {
351 let mut names = LISTENER_NAMES.lock().expect("lock poisoned");
352 if let Some(name) = names.get(name) {
353 return name;
354 }
355 let name: &'static str = Box::leak(name.to_owned().into_boxed_str());
356 names.insert(name);
357 name
358}
359
360pub struct Listeners {
361 pub http: BTreeMap<String, Listener<HttpListenerConfig>>,
362 pub sql: BTreeMap<String, Listener<SqlListenerConfig>>,
363}
364
365impl Listeners {
366 pub async fn bind(config: ListenersConfig) -> Result<Self, io::Error> {
367 let mut sql = BTreeMap::new();
368 for (name, config) in config.sql {
369 sql.insert(name, Listener::bind(config).await?);
370 }
371
372 let mut http = BTreeMap::new();
373 for (name, config) in config.http {
374 http.insert(name, Listener::bind(config).await?);
375 }
376
377 Ok(Listeners { http, sql })
378 }
379
380 #[instrument(name = "environmentd::serve")]
384 pub async fn serve(self, config: Config) -> Result<Server, AdapterError> {
385 let serve_start = Instant::now();
386 info!("startup: envd serve: beginning");
387 info!("startup: envd serve: preamble beginning");
388
389 let tls_reloading_context = match config.tls {
391 Some(tls_config) => Some(tls_config.reloading_context(config.tls_reload_certs)?),
392 None => None,
393 };
394
395 let active_connection_counter = ConnectionCounter::default();
396 let (deployment_state, deployment_state_handle) = DeploymentState::new();
397
398 let webhook_concurrency_limit = WebhookConcurrencyLimiter::default();
407 let internal_route_config = Arc::new(InternalRouteConfig {
408 deployment_state_handle,
409 internal_console_redirect_url: config.internal_console_redirect_url,
410 });
411
412 let (authenticator_oidc_tx, authenticator_oidc_rx) = oneshot::channel();
413 let authenticator_oidc_rx = authenticator_oidc_rx.shared();
414 let (adapter_client_tx, adapter_client_rx) = oneshot::channel();
415 let adapter_client_rx = adapter_client_rx.shared();
416
417 let metrics_registry = config.metrics_registry.clone();
418 let metrics = http::Metrics::register_into(&metrics_registry, "mz_http");
419 let mcp_metrics = http::mcp_metrics::McpMetrics::register_into(&metrics_registry);
420 let oauth_metadata_metrics =
421 http::oauth_metadata::OauthMetadataMetrics::register_into(&metrics_registry);
422 let mut http_listener_handles = BTreeMap::new();
423 for (name, listener) in self.http {
424 let authenticator_kind = listener.config.authenticator_kind();
425 let source = leak_listener_name(&name);
426 let tls = if listener.config.enable_tls() {
427 tls_reloading_context.clone()
428 } else {
429 None
430 };
431 let http_config = HttpConfig {
432 adapter_client_rx: adapter_client_rx.clone(),
433 active_connection_counter: active_connection_counter.clone(),
434 helm_chart_version: config.helm_chart_version.clone(),
435 http_host_name: config.http_host_name.clone(),
436 frontegg_oauth_issuer_url: config.frontegg_oauth_issuer_url.clone(),
437 source,
438 tls,
439 authenticator_kind,
440 frontegg: config.frontegg.clone(),
441 oidc_rx: authenticator_oidc_rx.clone(),
442 allowed_origin: config.cors_allowed_origin.clone(),
443 allowed_origin_list: config.cors_allowed_origin_list.clone(),
444 concurrent_webhook_req: webhook_concurrency_limit.semaphore(),
445 dyncfgs: Arc::clone(&config.system_dyncfgs),
446 metrics: metrics.clone(),
447 metrics_registry: metrics_registry.clone(),
448 mcp_metrics: mcp_metrics.clone(),
449 oauth_metadata_metrics: oauth_metadata_metrics.clone(),
450 internal_route_config: Arc::clone(&internal_route_config),
451 routes_enabled: listener.config.routes,
452 replica_http_locator: Arc::clone(&config.controller.replica_http_locator),
453 };
454 http_listener_handles.insert(name.clone(), listener.serve_http(http_config).await);
455 }
456
457 info!(
458 "startup: envd serve: preamble complete in {:?}",
459 serve_start.elapsed()
460 );
461
462 let catalog_init_start = Instant::now();
463 info!("startup: envd serve: catalog init beginning");
464
465 let boot_ts = (config.now)().into();
467
468 let persist_client = config
469 .catalog_config
470 .persist_clients
471 .open(config.controller.persist_location.clone())
472 .await
473 .context("opening persist client")?;
474 let mut openable_adapter_storage = mz_catalog::durable::persist_backed_catalog_state(
475 persist_client.clone(),
476 config.environment_id.organization_id(),
477 BUILD_INFO.semver_version(),
478 Some(config.controller.deploy_generation),
479 Arc::clone(&config.catalog_config.metrics),
480 )
481 .await?;
482
483 info!(
484 "startup: envd serve: catalog init complete in {:?}",
485 catalog_init_start.elapsed()
486 );
487
488 let system_param_sync_start = Instant::now();
489 info!("startup: envd serve: system parameter sync beginning");
490 let system_parameter_sync_config =
492 match (config.launchdarkly_sdk_key, config.config_sync_file_path) {
493 (None, None) => None,
494 (None, Some(f)) => {
495 info!("Using config file path {:?}", f);
496 Some(SystemParameterSyncConfig::new(
497 config.environment_id.clone(),
498 &BUILD_INFO,
499 &config.metrics_registry,
500 config.launchdarkly_key_map,
501 SystemParameterSyncClientConfig::File { path: f },
502 ))
503 }
504 (Some(key), None) => Some(SystemParameterSyncConfig::new(
505 config.environment_id.clone(),
506 &BUILD_INFO,
507 &config.metrics_registry,
508 config.launchdarkly_key_map,
509 SystemParameterSyncClientConfig::LaunchDarkly {
510 sdk_key: key,
511 base_uri: config.launchdarkly_base_uri,
512 now_fn: config.now.clone(),
513 },
514 )),
515
516 (Some(_), Some(_)) => {
517 panic!("Cannot configure both file and Launchdarkly based config syncing")
518 }
519 };
520
521 let remote_system_parameters = load_remote_system_parameters(
522 &mut openable_adapter_storage,
523 system_parameter_sync_config.clone(),
524 config.config_sync_timeout,
525 )
526 .await?;
527 info!(
528 "startup: envd serve: system parameter sync complete in {:?}",
529 system_param_sync_start.elapsed()
530 );
531
532 let preflight_checks_start = Instant::now();
533 info!("startup: envd serve: preflight checks beginning");
534
535 let with_0dt_deployment_max_wait = {
537 let cli_default = config
538 .system_parameter_defaults
539 .get(WITH_0DT_DEPLOYMENT_MAX_WAIT.name())
540 .map(|x| {
541 Duration::parse(VarInput::Flat(x)).map_err(|err| {
542 anyhow!(
543 "failed to parse default for {}: {:?}",
544 WITH_0DT_DEPLOYMENT_MAX_WAIT.name(),
545 err
546 )
547 })
548 })
549 .transpose()?;
550 let compiled_default = WITH_0DT_DEPLOYMENT_MAX_WAIT.default().clone();
551 let ld = get_ld_value(
552 WITH_0DT_DEPLOYMENT_MAX_WAIT.name(),
553 &remote_system_parameters,
554 |x| {
555 Duration::parse(VarInput::Flat(x)).map_err(|err| {
556 format!(
557 "failed to parse LD value {} for {}: {:?}",
558 x,
559 WITH_0DT_DEPLOYMENT_MAX_WAIT.name(),
560 err
561 )
562 })
563 },
564 )?;
565 let catalog = openable_adapter_storage
566 .get_0dt_deployment_max_wait()
567 .await?;
568 let computed = ld.or(catalog).or(cli_default).unwrap_or(compiled_default);
569 info!(
570 ?computed,
571 ?ld,
572 ?catalog,
573 ?cli_default,
574 ?compiled_default,
575 "determined value for {} system parameter",
576 WITH_0DT_DEPLOYMENT_MAX_WAIT.name()
577 );
578 computed
579 };
580 let with_0dt_deployment_ddl_check_interval = {
582 let cli_default = config
583 .system_parameter_defaults
584 .get(WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL.name())
585 .map(|x| {
586 Duration::parse(VarInput::Flat(x)).map_err(|err| {
587 anyhow!(
588 "failed to parse default for {}: {:?}",
589 WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL.name(),
590 err
591 )
592 })
593 })
594 .transpose()?;
595 let compiled_default = WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL.default().clone();
596 let ld = get_ld_value(
597 WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL.name(),
598 &remote_system_parameters,
599 |x| {
600 Duration::parse(VarInput::Flat(x)).map_err(|err| {
601 format!(
602 "failed to parse LD value {} for {}: {:?}",
603 x,
604 WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL.name(),
605 err
606 )
607 })
608 },
609 )?;
610 let catalog = openable_adapter_storage
611 .get_0dt_deployment_ddl_check_interval()
612 .await?;
613 let computed = ld.or(catalog).or(cli_default).unwrap_or(compiled_default);
614 info!(
615 ?computed,
616 ?ld,
617 ?catalog,
618 ?cli_default,
619 ?compiled_default,
620 "determined value for {} system parameter",
621 WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL.name()
622 );
623 computed
624 };
625
626 let enable_0dt_deployment_panic_after_timeout = {
629 let cli_default = config
630 .system_parameter_defaults
631 .get(ENABLE_0DT_DEPLOYMENT_PANIC_AFTER_TIMEOUT.name())
632 .map(|x| {
633 strconv::parse_bool(x).map_err(|err| {
634 anyhow!(
635 "failed to parse default for {}: {}",
636 ENABLE_0DT_DEPLOYMENT_PANIC_AFTER_TIMEOUT.name(),
637 err
638 )
639 })
640 })
641 .transpose()?;
642 let compiled_default = ENABLE_0DT_DEPLOYMENT_PANIC_AFTER_TIMEOUT.default().clone();
643 let ld = get_ld_value(
644 "enable_0dt_deployment_panic_after_timeout",
645 &remote_system_parameters,
646 |x| strconv::parse_bool(x).map_err(|x| x.to_string()),
647 )?;
648 let catalog = openable_adapter_storage
649 .get_enable_0dt_deployment_panic_after_timeout()
650 .await?;
651 let computed = ld.or(catalog).or(cli_default).unwrap_or(compiled_default);
652 info!(
653 %computed,
654 ?ld,
655 ?catalog,
656 ?cli_default,
657 ?compiled_default,
658 "determined value for enable_0dt_deployment_panic_after_timeout system parameter",
659 );
660 computed
661 };
662
663 let read_only = preflight::preflight_0dt(
667 openable_adapter_storage.as_mut(),
668 config.controller.deploy_generation,
669 )
670 .await?;
671
672 let bootstrap_args = BootstrapArgs {
673 default_cluster_replica_size: config.bootstrap_default_cluster_replica_size.clone(),
674 default_cluster_replication_factor: config.bootstrap_default_cluster_replication_factor,
675 bootstrap_role: config.bootstrap_role,
676 cluster_replica_size_map: config.cluster_replica_sizes.clone(),
677 };
678
679 let (caught_up_trigger, bootstrapped) = if read_only {
680 let (caught_up_trigger, caught_up_receiver) = mz_ore::channel::trigger::channel();
681 let (bootstrapped, bootstrapped_receiver) = tokio::sync::oneshot::channel();
682 let catchup_config = CatchupConfig {
683 boot_ts,
684 environment_id: config.environment_id.clone(),
685 persist_client,
686 deploy_generation: config.controller.deploy_generation,
687 deployment_state: deployment_state.clone(),
688 catalog_metrics: Arc::clone(&config.catalog_config.metrics),
689 caught_up_max_wait: with_0dt_deployment_max_wait,
690 panic_after_timeout: enable_0dt_deployment_panic_after_timeout,
691 bootstrap_args: bootstrap_args.clone(),
692 ddl_check_interval: with_0dt_deployment_ddl_check_interval,
693 };
694 preflight::spawn_catchup(catchup_config, caught_up_receiver, bootstrapped_receiver);
695 (Some(caught_up_trigger), Some(bootstrapped))
696 } else {
697 (None, None)
698 };
699
700 info!(
701 "startup: envd serve: preflight checks complete in {:?}",
702 preflight_checks_start.elapsed()
703 );
704
705 let catalog_open_start = Instant::now();
706 info!("startup: envd serve: durable catalog open beginning");
707
708 let mut adapter_storage = if read_only {
710 let adapter_storage = openable_adapter_storage
713 .open_savepoint(boot_ts, &bootstrap_args)
714 .await?;
715 adapter_storage
720 } else {
721 let adapter_storage = openable_adapter_storage
722 .open(boot_ts, &bootstrap_args)
723 .await?;
724
725 deployment_state.set_is_leader();
729
730 adapter_storage
731 };
732
733 let bootstrapped = match bootstrapped {
734 Some(tx) => Some((tx, preflight::get_next_ids(adapter_storage.as_mut()).await?)),
735 None => None,
736 };
737
738 if !read_only {
740 config.controller.persist_clients.cfg().enable_compaction();
741 }
742
743 info!(
744 "startup: envd serve: durable catalog open complete in {:?}",
745 catalog_open_start.elapsed()
746 );
747
748 let coord_init_start = Instant::now();
749 info!("startup: envd serve: coordinator init beginning");
750
751 if !config
752 .cluster_replica_sizes
753 .0
754 .contains_key(&config.bootstrap_default_cluster_replica_size)
755 {
756 return Err(anyhow!("bootstrap default cluster replica size is unknown").into());
757 }
758 let envd_epoch = adapter_storage.epoch();
759
760 let storage_usage_client = StorageUsageClient::open(
762 config
763 .controller
764 .persist_clients
765 .open(config.controller.persist_location.clone())
766 .await
767 .context("opening storage usage client")?,
768 );
769
770 let license_key = config.license_key.clone();
772 let segment_client = config.segment_api_key.map(|api_key| {
773 mz_segment::Client::new(mz_segment::Config {
774 api_key,
775 client_side: config.segment_client_side,
776 })
777 });
778 let connection_limiter = active_connection_counter.clone();
779 let connection_limit_callback = Box::new(move |limit, superuser_reserved| {
780 connection_limiter.update_limit(limit);
781 connection_limiter.update_superuser_reserved(superuser_reserved);
782 });
783
784 let (adapter_handle, adapter_client) = mz_adapter::serve(mz_adapter::Config {
785 connection_context: config.controller.connection_context.clone(),
786 connection_limit_callback,
787 controller_config: config.controller,
788 controller_envd_epoch: envd_epoch,
789 storage: adapter_storage,
790 timestamp_oracle_url: config.timestamp_oracle_url,
791 unsafe_mode: config.unsafe_mode,
792 all_features: config.all_features,
793 build_info: &BUILD_INFO,
794 environment_id: config.environment_id.clone(),
795 metrics_registry: config.metrics_registry.clone(),
796 now: config.now,
797 secrets_controller: config.secrets_controller,
798 cloud_resource_controller: config.cloud_resource_controller,
799 cluster_replica_sizes: config.cluster_replica_sizes,
800 builtin_system_cluster_config: config.bootstrap_builtin_system_cluster_config,
801 builtin_catalog_server_cluster_config: config
802 .bootstrap_builtin_catalog_server_cluster_config,
803 builtin_probe_cluster_config: config.bootstrap_builtin_probe_cluster_config,
804 builtin_support_cluster_config: config.bootstrap_builtin_support_cluster_config,
805 builtin_analytics_cluster_config: config.bootstrap_builtin_analytics_cluster_config,
806 availability_zones: config.availability_zones,
807 system_parameter_defaults: config.system_parameter_defaults,
808 storage_usage_client,
809 storage_usage_collection_interval: config.storage_usage_collection_interval,
810 storage_usage_retention_period: config.storage_usage_retention_period,
811 segment_client: segment_client.clone(),
812 egress_addresses: config.egress_addresses,
813 remote_system_parameters,
814 aws_account_id: config.aws_account_id,
815 aws_privatelink_availability_zones: config.aws_privatelink_availability_zones,
816 webhook_concurrency_limit: webhook_concurrency_limit.clone(),
817 http_host_name: config.http_host_name,
818 tracing_handle: config.tracing_handle,
819 read_only_controllers: read_only,
820 caught_up_trigger,
821 helm_chart_version: config.helm_chart_version.clone(),
822 license_key: config.license_key,
823 external_login_password_mz_system: config.external_login_password_mz_system,
824 force_builtin_schema_migration: config.force_builtin_schema_migration,
825 })
826 .instrument(info_span!("adapter::serve"))
827 .await?;
828
829 if let Some((tx, initial_ids)) = bootstrapped {
830 let _ = tx.send(initial_ids);
832 }
833
834 let oidc = GenericOidcAuthenticator::new(adapter_client.clone());
836
837 info!(
838 "startup: envd serve: coordinator init complete in {:?}",
839 coord_init_start.elapsed()
840 );
841
842 let serve_postamble_start = Instant::now();
843 info!("startup: envd serve: postamble beginning");
844
845 authenticator_oidc_tx
847 .send(oidc.clone())
848 .expect("rx known to be live");
849 adapter_client_tx
850 .send(adapter_client.clone())
851 .expect("internal HTTP server should not drop first");
852
853 let metrics = mz_pgwire::MetricsConfig::register_into(&config.metrics_registry);
854
855 let mut sql_listener_handles = BTreeMap::new();
857 for (name, listener) in self.sql {
858 sql_listener_handles.insert(
859 name.clone(),
860 listener
861 .serve_sql(
862 name,
863 active_connection_counter.clone(),
864 tls_reloading_context.clone(),
865 config.frontegg.clone(),
866 adapter_client.clone(),
867 oidc.clone(),
868 metrics.clone(),
869 config.helm_chart_version.clone(),
870 )
871 .await,
872 );
873 }
874
875 if let Some(segment_client) = segment_client {
877 telemetry::start_reporting(telemetry::Config {
878 segment_client,
879 adapter_client: adapter_client.clone(),
880 environment_id: config.environment_id,
881 license_key: license_key.clone(),
882 helm_chart_version: config.helm_chart_version.clone(),
883 report_interval: Duration::from_secs(3600),
884 });
885 } else if config.test_only_dummy_segment_client {
886 tracing::debug!("starting telemetry reporting with a dummy segment client");
892 let segment_client = mz_segment::Client::new_dummy_client();
893 telemetry::start_reporting(telemetry::Config {
894 segment_client,
895 adapter_client: adapter_client.clone(),
896 environment_id: config.environment_id,
897 license_key: license_key.clone(),
898 helm_chart_version: config.helm_chart_version.clone(),
899 report_interval: Duration::from_secs(180),
900 });
901 }
902
903 if let Some(system_parameter_sync_config) = system_parameter_sync_config {
906 task::spawn(
907 || "system_parameter_sync",
908 AssertUnwindSafe(system_parameter_sync(
909 system_parameter_sync_config,
910 adapter_client.clone(),
911 config.config_sync_loop_interval,
912 ))
913 .ore_catch_unwind(),
914 );
915 }
916
917 info!(
918 "startup: envd serve: postamble complete in {:?}",
919 serve_postamble_start.elapsed()
920 );
921 info!(
922 "startup: envd serve: complete in {:?}",
923 serve_start.elapsed()
924 );
925
926 Ok(Server {
927 sql_listener_handles,
928 http_listener_handles,
929 #[cfg(feature = "test")]
930 adapter_client,
931 _adapter_handle: adapter_handle,
932 })
933 }
934}
935
936fn get_ld_value<V>(
937 name: &str,
938 remote_system_parameters: &Option<BTreeMap<String, String>>,
939 parse: impl Fn(&str) -> Result<V, String>,
940) -> Result<Option<V>, anyhow::Error> {
941 remote_system_parameters
942 .as_ref()
943 .and_then(|params| params.get(name))
944 .map(|x| {
945 parse(x).map_err(|err| anyhow!("failed to parse remote value for {}: {}", name, err))
946 })
947 .transpose()
948}
949
950pub struct Server {
952 pub sql_listener_handles: BTreeMap<String, ListenerHandle>,
954 pub http_listener_handles: BTreeMap<String, ListenerHandle>,
955 #[cfg(feature = "test")]
956 adapter_client: AdapterClient,
957 _adapter_handle: mz_adapter::Handle,
958}
959
960impl Server {
961 #[cfg(feature = "test")]
963 pub fn adapter_client(&self) -> &AdapterClient {
964 &self.adapter_client
965 }
966}