1use std::time::Duration;
11
12use mz_adapter_types::dyncfgs::PG_TIMESTAMP_ORACLE_STATEMENT_TIMEOUT;
13use mz_compute_client::protocol::command::ComputeParameters;
14use mz_orchestrator::scheduling_config::{ServiceSchedulingConfig, ServiceTopologySpreadConfig};
15use mz_ore::cast::CastFrom;
16use mz_ore::error::ErrorExt;
17use mz_service::params::GrpcClientParameters;
18use mz_sql::session::vars::SystemVars;
19use mz_storage_types::parameters::{
20 PgSourceSnapshotConfig, StorageMaxInflightBytesConfig, StorageParameters,
21};
22use mz_tracing::params::TracingParameters;
23
24use mz_timestamp_oracle::postgres_oracle::TimestampOracleParameters;
25
26pub fn compute_config(config: &SystemVars) -> ComputeParameters {
28 ComputeParameters {
29 workload_class: None,
30 max_result_size: Some(config.max_result_size()),
31 tracing: tracing_config(config),
32 grpc_client: grpc_client_config(config),
33 dyncfg_updates: config.dyncfg_updates(),
34 }
35}
36
37pub fn storage_config(config: &SystemVars) -> StorageParameters {
39 StorageParameters {
40 pg_source_connect_timeout: Some(config.pg_source_connect_timeout()),
41 pg_source_tcp_keepalives_retries: Some(config.pg_source_tcp_keepalives_retries()),
42 pg_source_tcp_keepalives_idle: Some(config.pg_source_tcp_keepalives_idle()),
43 pg_source_tcp_keepalives_interval: Some(config.pg_source_tcp_keepalives_interval()),
44 pg_source_tcp_user_timeout: Some(config.pg_source_tcp_user_timeout()),
45 pg_source_tcp_configure_server: config.pg_source_tcp_configure_server(),
46 pg_source_snapshot_statement_timeout: config.pg_source_snapshot_statement_timeout(),
47 pg_source_wal_sender_timeout: config.pg_source_wal_sender_timeout(),
48 mysql_source_timeouts: mz_mysql_util::TimeoutConfig::build(
49 config.mysql_source_snapshot_max_execution_time(),
50 config.mysql_source_snapshot_lock_wait_timeout(),
51 config.mysql_source_snapshot_wait_timeout(),
52 config.mysql_source_tcp_keepalive(),
53 config.mysql_source_connect_timeout(),
54 ),
55 keep_n_source_status_history_entries: config.keep_n_source_status_history_entries(),
56 keep_n_sink_status_history_entries: config.keep_n_sink_status_history_entries(),
57 keep_n_privatelink_status_history_entries: config
58 .keep_n_privatelink_status_history_entries(),
59 replica_status_history_retention_window: config.replica_status_history_retention_window(),
60 upsert_rocksdb_tuning_config: {
61 match mz_rocksdb_types::RocksDBTuningParameters::from_parameters(
62 config.upsert_rocksdb_compaction_style(),
63 config.upsert_rocksdb_optimize_compaction_memtable_budget(),
64 config.upsert_rocksdb_level_compaction_dynamic_level_bytes(),
65 config.upsert_rocksdb_universal_compaction_ratio(),
66 config.upsert_rocksdb_parallelism(),
67 config.upsert_rocksdb_compression_type(),
68 config.upsert_rocksdb_bottommost_compression_type(),
69 config.upsert_rocksdb_batch_size(),
70 config.upsert_rocksdb_retry_duration(),
71 config.upsert_rocksdb_stats_log_interval_seconds(),
72 config.upsert_rocksdb_stats_persist_interval_seconds(),
73 config.upsert_rocksdb_point_lookup_block_cache_size_mb(),
74 config.upsert_rocksdb_shrink_allocated_buffers_by_ratio(),
75 config.upsert_rocksdb_write_buffer_manager_memory_bytes(),
76 config
77 .upsert_rocksdb_write_buffer_manager_cluster_memory_fraction()
78 .and_then(|d| match d.try_into() {
79 Err(e) => {
80 tracing::error!(
81 "Couldn't convert upsert_rocksdb_write_buffer_manager_cluster_memory_fraction {:?} to f64, so defaulting to `None`: {e:?}",
82 config
83 .upsert_rocksdb_write_buffer_manager_cluster_memory_fraction()
84 );
85 None
86 }
87 Ok(o) => Some(o),
88 }),
89 config.upsert_rocksdb_write_buffer_manager_allow_stall(),
90 ) {
91 Ok(u) => u,
92 Err(e) => {
93 tracing::warn!(
94 "Failed to deserialize upsert_rocksdb parameters \
95 into a `RocksDBTuningParameters`, \
96 failing back to reasonable defaults: {}",
97 e.display_with_causes()
98 );
99 mz_rocksdb_types::RocksDBTuningParameters::default()
100 }
101 }
102 },
103 finalize_shards: config.enable_storage_shard_finalization(),
104 tracing: tracing_config(config),
105 storage_dataflow_max_inflight_bytes_config: StorageMaxInflightBytesConfig {
106 max_inflight_bytes_default: config.storage_dataflow_max_inflight_bytes(),
107 max_inflight_bytes_cluster_size_fraction: config
110 .storage_dataflow_max_inflight_bytes_to_cluster_size_fraction()
111 .and_then(|d| match d.try_into() {
112 Err(e) => {
113 tracing::error!(
114 "Couldn't convert {:?} to f64, so defaulting to `None`: {e:?}",
115 config.storage_dataflow_max_inflight_bytes_to_cluster_size_fraction()
116 );
117 None
118 }
119 Ok(o) => Some(o),
120 }),
121 disk_only: config.storage_dataflow_max_inflight_bytes_disk_only(),
122 },
123 grpc_client: grpc_client_config(config),
124 shrink_upsert_unused_buffers_by_ratio: config
125 .storage_shrink_upsert_unused_buffers_by_ratio(),
126 record_namespaced_errors: config.storage_record_source_sink_namespaced_errors(),
127 ssh_timeout_config: mz_ssh_util::tunnel::SshTimeoutConfig {
128 check_interval: config.ssh_check_interval(),
129 connect_timeout: config.ssh_connect_timeout(),
130 keepalives_idle: config.ssh_keepalives_idle(),
131 },
132 kafka_timeout_config: mz_kafka_util::client::TimeoutConfig::build(
133 config.kafka_socket_keepalive(),
134 config.kafka_socket_timeout(),
135 config.kafka_transaction_timeout(),
136 config.kafka_socket_connection_setup_timeout(),
137 config.kafka_fetch_metadata_timeout(),
138 config.kafka_progress_record_fetch_timeout(),
139 ),
140 statistics_interval: config.storage_statistics_interval(),
141 statistics_collection_interval: config.storage_statistics_collection_interval(),
142 pg_snapshot_config: PgSourceSnapshotConfig {
143 collect_strict_count: config.pg_source_snapshot_collect_strict_count(),
144 },
145 user_storage_managed_collections_batch_duration: config
146 .user_storage_managed_collections_batch_duration(),
147 dyncfg_updates: config.dyncfg_updates(),
148 }
149}
150
151pub fn tracing_config(config: &SystemVars) -> TracingParameters {
152 TracingParameters {
153 log_filter: Some(config.logging_filter()),
154 opentelemetry_filter: Some(config.opentelemetry_filter()),
155 log_filter_defaults: config.logging_filter_defaults(),
156 opentelemetry_filter_defaults: config.opentelemetry_filter_defaults(),
157 sentry_filters: config.sentry_filters(),
158 }
159}
160
161pub fn caching_config(config: &SystemVars) -> mz_secrets::CachingPolicy {
162 let ttl_secs = config.webhooks_secrets_caching_ttl_secs();
163 mz_secrets::CachingPolicy {
164 enabled: ttl_secs > 0,
165 ttl: Duration::from_secs(u64::cast_from(ttl_secs)),
166 }
167}
168
169pub fn timestamp_oracle_config(config: &SystemVars) -> TimestampOracleParameters {
170 TimestampOracleParameters {
171 pg_connection_pool_max_size: Some(config.pg_timestamp_oracle_connection_pool_max_size()),
172 pg_connection_pool_max_wait: Some(config.pg_timestamp_oracle_connection_pool_max_wait()),
173 pg_connection_pool_ttl: Some(config.pg_timestamp_oracle_connection_pool_ttl()),
174 pg_connection_pool_ttl_stagger: Some(
175 config.pg_timestamp_oracle_connection_pool_ttl_stagger(),
176 ),
177 pg_connection_pool_connect_timeout: Some(config.crdb_connect_timeout()),
181 pg_connection_pool_tcp_user_timeout: Some(config.crdb_tcp_user_timeout()),
182 pg_connection_pool_keepalives_idle: Some(config.crdb_keepalives_idle()),
183 pg_connection_pool_keepalives_interval: Some(config.crdb_keepalives_interval()),
184 pg_connection_pool_keepalives_retries: Some(config.crdb_keepalives_retries()),
185 pg_statement_timeout: Some(PG_TIMESTAMP_ORACLE_STATEMENT_TIMEOUT.get(config.dyncfgs())),
186 }
187}
188
189fn grpc_client_config(config: &SystemVars) -> GrpcClientParameters {
190 GrpcClientParameters {
191 connect_timeout: Some(config.grpc_connect_timeout()),
192 http2_keep_alive_interval: Some(config.grpc_client_http2_keep_alive_interval()),
193 http2_keep_alive_timeout: Some(config.grpc_client_http2_keep_alive_timeout()),
194 }
195}
196
197pub fn orchestrator_scheduling_config(config: &SystemVars) -> ServiceSchedulingConfig {
198 ServiceSchedulingConfig {
199 multi_pod_az_affinity_weight: config.cluster_multi_process_replica_az_affinity_weight(),
200 soften_replication_anti_affinity: config.cluster_soften_replication_anti_affinity(),
201 soften_replication_anti_affinity_weight: config
202 .cluster_soften_replication_anti_affinity_weight(),
203 topology_spread: ServiceTopologySpreadConfig {
204 enabled: config.cluster_enable_topology_spread(),
205 ignore_non_singular_scale: config.cluster_topology_spread_ignore_non_singular_scale(),
206 max_skew: config.cluster_topology_spread_max_skew(),
207 min_domains: config.cluster_topology_spread_set_min_domains(),
208 soft: config.cluster_topology_spread_soft(),
209 },
210 soften_az_affinity: config.cluster_soften_az_affinity(),
211 soften_az_affinity_weight: config.cluster_soften_az_affinity_weight(),
212 security_context_enabled: config.cluster_security_context_enabled(),
213 }
214}