1use std::borrow::Cow;
67use std::clone::Clone;
68use std::collections::BTreeMap;
69use std::fmt::Debug;
70use std::net::IpAddr;
71use std::num::NonZeroU32;
72use std::string::ToString;
73use std::sync::{Arc, LazyLock};
74use std::time::Duration;
75
76use chrono::{DateTime, Utc};
77use derivative::Derivative;
78use imbl::OrdMap;
79use mz_auth::user::ExternalUserMetadata;
80use mz_build_info::BuildInfo;
81use mz_dyncfg::{ConfigSet, ConfigType, ConfigUpdates, ConfigVal, ParameterScope};
82use mz_persist_client::cfg::{
83 CRDB_CONNECT_TIMEOUT, CRDB_KEEPALIVES_IDLE, CRDB_KEEPALIVES_INTERVAL, CRDB_KEEPALIVES_RETRIES,
84 CRDB_TCP_USER_TIMEOUT,
85};
86use mz_pgrepr::TextEncodeSettings;
87use mz_repr::adt::numeric::Numeric;
88use mz_repr::adt::timestamp::CheckedTimestamp;
89use mz_repr::bytes::ByteSize;
90use mz_repr::user::InternalUserMetadata;
91use mz_tracing::{CloneableEnvFilter, SerializableDirective};
92use serde::Serialize;
93use thiserror::Error;
94use uncased::UncasedStr;
95
96use crate::ast::Ident;
97use crate::session::user::User;
98
99pub(crate) mod constraints;
100pub(crate) mod definitions;
101pub(crate) mod errors;
102pub(crate) mod polyfill;
103pub(crate) mod value;
104
105pub use definitions::*;
106pub use errors::*;
107pub use value::*;
108
109#[derive(Debug, Clone, Copy, PartialEq, Eq)]
114pub enum EndTransactionAction {
115 Commit,
117 Rollback,
119}
120
121#[derive(Debug, Clone, Copy)]
127pub enum VarInput<'a> {
128 Flat(&'a str),
136 SqlSet(&'a [String]),
142}
143
144impl<'a> VarInput<'a> {
145 pub fn to_vec(&self) -> Vec<String> {
147 match self {
148 VarInput::Flat(v) => vec![v.to_string()],
149 VarInput::SqlSet(values) => values.into_iter().map(|v| v.to_string()).collect(),
150 }
151 }
152}
153
154#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize)]
156pub enum OwnedVarInput {
157 Flat(String),
164 SqlSet(Vec<String>),
168}
169
170impl OwnedVarInput {
171 pub fn borrow(&self) -> VarInput<'_> {
173 match self {
174 OwnedVarInput::Flat(v) => VarInput::Flat(v),
175 OwnedVarInput::SqlSet(v) => VarInput::SqlSet(v),
176 }
177 }
178}
179
180pub trait Var: Debug {
182 fn name(&self) -> &'static str;
184
185 fn value(&self) -> String;
191
192 fn description(&self) -> &'static str;
195
196 fn type_name(&self) -> Cow<'static, str>;
198
199 fn visible(&self, user: &User, system_vars: &SystemVars) -> Result<(), VarError>;
204
205 fn is_unsafe(&self) -> bool {
207 self.name().starts_with("unsafe_")
208 }
209
210 fn scope(&self) -> ParameterScope {
214 ParameterScope::Environment
215 }
216
217 fn as_var(&self) -> &dyn Var
220 where
221 Self: Sized,
222 {
223 self
224 }
225}
226
227#[derive(Debug)]
234pub struct SessionVar {
235 definition: VarDefinition,
236 default_value: Option<Box<dyn Value>>,
238 local_value: Option<Box<dyn Value>>,
240 staged_value: Option<Box<dyn Value>>,
242 session_value: Option<Box<dyn Value>>,
244}
245
246impl Clone for SessionVar {
247 fn clone(&self) -> Self {
248 SessionVar {
249 definition: self.definition.clone(),
250 default_value: self.default_value.as_ref().map(|v| v.box_clone()),
251 local_value: self.local_value.as_ref().map(|v| v.box_clone()),
252 staged_value: self.staged_value.as_ref().map(|v| v.box_clone()),
253 session_value: self.session_value.as_ref().map(|v| v.box_clone()),
254 }
255 }
256}
257
258impl SessionVar {
259 pub const fn new(var: VarDefinition) -> Self {
260 SessionVar {
261 definition: var,
262 default_value: None,
263 local_value: None,
264 staged_value: None,
265 session_value: None,
266 }
267 }
268
269 pub fn check(&self, input: VarInput) -> Result<String, VarError> {
272 let v = self.definition.parse(input)?;
273 self.validate_constraints(v.as_ref())?;
274
275 Ok(v.format())
276 }
277
278 pub fn set(&mut self, input: VarInput, local: bool) -> Result<(), VarError> {
280 let v = self.definition.parse(input)?;
281
282 self.validate_constraints(v.as_ref())?;
284
285 if local {
286 self.local_value = Some(v);
287 } else {
288 self.local_value = None;
289 self.staged_value = Some(v);
290 }
291 Ok(())
292 }
293
294 pub fn set_default(&mut self, input: VarInput) -> Result<(), VarError> {
296 let v = self.definition.parse(input)?;
297 self.validate_constraints(v.as_ref())?;
298 self.default_value = Some(v);
299 Ok(())
300 }
301
302 pub fn reset(&mut self, local: bool) {
304 let value = self
305 .default_value
306 .as_ref()
307 .map(|v| v.as_ref())
308 .unwrap_or_else(|| self.definition.value.value());
309 if local {
310 self.local_value = Some(value.box_clone());
311 } else {
312 self.local_value = None;
313 self.staged_value = Some(value.box_clone());
314 }
315 }
316
317 pub fn reset_durable(&mut self) {
326 self.local_value = None;
327 self.staged_value = None;
328 self.session_value = None;
329 }
330
331 #[must_use]
333 pub fn end_transaction(&self, action: EndTransactionAction) -> Option<Self> {
334 if !self.is_mutating() {
335 return None;
336 }
337 let mut next: Self = self.clone();
338 next.local_value = None;
339 match action {
340 EndTransactionAction::Commit if next.staged_value.is_some() => {
341 next.session_value = next.staged_value.take()
342 }
343 _ => next.staged_value = None,
344 }
345 Some(next)
346 }
347
348 pub fn is_mutating(&self) -> bool {
350 self.local_value.is_some() || self.staged_value.is_some()
351 }
352
353 pub fn value_dyn(&self) -> &dyn Value {
354 self.local_value
355 .as_deref()
356 .or(self.staged_value.as_deref())
357 .or(self.session_value.as_deref())
358 .or(self.default_value.as_deref())
359 .unwrap_or_else(|| self.definition.value.value())
360 }
361
362 pub fn inspect_session_value(&self) -> Option<&dyn Value> {
367 self.session_value.as_deref()
368 }
369
370 fn validate_constraints(&self, val: &dyn Value) -> Result<(), VarError> {
371 if let Some(constraint) = &self.definition.constraint {
372 constraint.check_constraint(self, self.value_dyn(), val)
373 } else {
374 Ok(())
375 }
376 }
377}
378
379impl Var for SessionVar {
380 fn name(&self) -> &'static str {
381 self.definition.name.as_str()
382 }
383
384 fn value(&self) -> String {
385 self.value_dyn().format()
386 }
387
388 fn description(&self) -> &'static str {
389 self.definition.description
390 }
391
392 fn type_name(&self) -> Cow<'static, str> {
393 self.definition.type_name()
394 }
395
396 fn visible(
397 &self,
398 user: &User,
399 system_vars: &super::vars::SystemVars,
400 ) -> Result<(), super::vars::VarError> {
401 self.definition.visible(user, system_vars)
402 }
403}
404
405#[derive(Debug, Clone, PartialEq, Eq)]
406pub struct MzVersion {
407 build_info: &'static BuildInfo,
409 helm_chart_version: Option<String>,
411}
412
413impl MzVersion {
414 pub fn new(build_info: &'static BuildInfo, helm_chart_version: Option<String>) -> Self {
415 MzVersion {
416 build_info,
417 helm_chart_version,
418 }
419 }
420}
421
422#[derive(Debug, Clone)]
427pub struct SessionVars {
428 vars: OrdMap<&'static UncasedStr, SessionVar>,
430 mz_version: MzVersion,
432 user: User,
434}
435
436impl SessionVars {
437 pub fn new_unchecked(
439 build_info: &'static BuildInfo,
440 user: User,
441 helm_chart_version: Option<String>,
442 ) -> SessionVars {
443 use definitions::*;
444
445 let vars = [
446 &FAILPOINTS,
447 &SERVER_VERSION,
448 &SERVER_VERSION_NUM,
449 &SQL_SAFE_UPDATES,
450 &REAL_TIME_RECENCY,
451 &EMIT_PLAN_INSIGHTS_NOTICE,
452 &EMIT_TIMESTAMP_NOTICE,
453 &EMIT_TRACE_ID_NOTICE,
454 &AUTO_ROUTE_CATALOG_QUERIES,
455 &ENABLE_SESSION_RBAC_CHECKS,
456 &RESTRICT_TO_USER_OBJECTS,
457 &ENABLE_SESSION_CARDINALITY_ESTIMATES,
458 &MAX_IDENTIFIER_LENGTH,
459 &STATEMENT_LOGGING_SAMPLE_RATE,
460 &EMIT_INTROSPECTION_QUERY_NOTICE,
461 &UNSAFE_NEW_TRANSACTION_WALL_TIME,
462 &WELCOME_MESSAGE,
463 ]
464 .into_iter()
465 .chain(SESSION_SYSTEM_VARS.iter().map(|(_name, var)| *var))
466 .map(|var| (var.name, SessionVar::new(var.clone())))
467 .collect();
468
469 SessionVars {
470 vars,
471 mz_version: MzVersion::new(build_info, helm_chart_version),
472 user,
473 }
474 }
475
476 fn expect_value<V: Value>(&self, var: &VarDefinition) -> &V {
477 let var = self
478 .vars
479 .get(var.name)
480 .expect("provided var should be in state");
481 let val = var.value_dyn();
482 val.as_any().downcast_ref::<V>().expect("success")
483 }
484
485 pub fn iter(&self) -> impl Iterator<Item = &dyn Var> {
492 #[allow(clippy::as_conversions)]
493 self.vars
494 .values()
495 .map(|v| v.as_var())
496 .chain([&self.mz_version as &dyn Var, &self.user])
497 }
498
499 pub fn notify_set(&self) -> impl Iterator<Item = &dyn Var> {
503 [
510 &APPLICATION_NAME,
511 &CLIENT_ENCODING,
512 &DATE_STYLE,
513 &INTEGER_DATETIMES,
514 &SERVER_VERSION,
515 &STANDARD_CONFORMING_STRINGS,
516 &TIMEZONE,
517 &INTERVAL_STYLE,
518 &CLUSTER,
523 &CLUSTER_REPLICA,
524 &DEFAULT_CLUSTER_REPLICATION_FACTOR,
525 &DATABASE,
526 &SEARCH_PATH,
527 ]
528 .into_iter()
529 .map(|v| self.vars[v.name].as_var())
530 .chain(std::iter::once(self.mz_version.as_var()))
537 }
538
539 pub fn reset_all(&mut self) -> BTreeMap<&'static str, String> {
548 let mut changed = BTreeMap::new();
549 let names: Vec<_> = self.vars.keys().copied().collect();
550 for name in names {
551 let var = &mut self.vars[name];
552 let before = var.value();
553 var.reset_durable();
554 let after = var.value();
555 if before != after {
556 changed.insert(var.name(), after);
557 }
558 }
559 changed
560 }
561
562 pub fn get(&self, system_vars: &SystemVars, name: &str) -> Result<&dyn Var, VarError> {
573 let name = compat_translate_name(name);
574
575 let name = UncasedStr::new(name);
576 if name == MZ_VERSION_NAME {
577 Ok(&self.mz_version)
578 } else if name == IS_SUPERUSER_NAME {
579 Ok(&self.user)
580 } else {
581 self.vars
582 .get(name)
583 .map(|v| {
584 v.visible(&self.user, system_vars)?;
585 Ok(v.as_var())
586 })
587 .transpose()?
588 .ok_or_else(|| VarError::UnknownParameter(name.to_string()))
589 }
590 }
591
592 pub fn inspect(&self, name: &str) -> Result<&SessionVar, VarError> {
597 let name = compat_translate_name(name);
598
599 self.vars
600 .get(UncasedStr::new(name))
601 .ok_or_else(|| VarError::UnknownParameter(name.to_string()))
602 }
603
604 pub fn set(
617 &mut self,
618 system_vars: &SystemVars,
619 name: &str,
620 input: VarInput,
621 local: bool,
622 ) -> Result<(), VarError> {
623 let (name, input) = compat_translate(name, input);
624
625 check_transaction_isolation_feature_flag(name, input, system_vars)?;
626
627 let name = UncasedStr::new(name);
628 self.check_read_only(name)?;
629
630 self.vars
631 .get_mut(name)
632 .map(|v| {
633 v.visible(&self.user, system_vars)?;
634 v.set(input, local)
635 })
636 .transpose()?
637 .ok_or_else(|| VarError::UnknownParameter(name.to_string()))
638 }
639
640 pub fn set_default(&mut self, name: &str, input: VarInput) -> Result<(), VarError> {
643 let (name, input) = compat_translate(name, input);
644
645 let name = UncasedStr::new(name);
646
647 if !Self::allow_role_default(name) {
651 self.check_read_only(name)?;
652 }
653
654 self.vars
655 .get_mut(name)
656 .map(|v| v.set_default(input))
658 .transpose()?
659 .ok_or_else(|| VarError::UnknownParameter(name.to_string()))
660 }
661
662 fn allow_role_default(name: &UncasedStr) -> bool {
670 name == RESTRICT_TO_USER_OBJECTS.name
671 }
672
673 pub fn reset(
687 &mut self,
688 system_vars: &SystemVars,
689 name: &str,
690 local: bool,
691 ) -> Result<(), VarError> {
692 let name = compat_translate_name(name);
693
694 let name = UncasedStr::new(name);
695 self.check_read_only(name)?;
696
697 self.vars
698 .get_mut(name)
699 .map(|v| {
700 v.visible(&self.user, system_vars)?;
701 v.reset(local);
702 Ok(())
703 })
704 .transpose()?
705 .ok_or_else(|| VarError::UnknownParameter(name.to_string()))
706 }
707
708 fn check_read_only(&self, name: &UncasedStr) -> Result<(), VarError> {
715 if name == MZ_VERSION_NAME {
716 Err(VarError::ReadOnlyParameter(MZ_VERSION_NAME.as_str()))
717 } else if name == IS_SUPERUSER_NAME {
718 Err(VarError::ReadOnlyParameter(IS_SUPERUSER_NAME.as_str()))
719 } else if name == MAX_IDENTIFIER_LENGTH.name {
720 Err(VarError::ReadOnlyParameter(
721 MAX_IDENTIFIER_LENGTH.name.as_str(),
722 ))
723 } else if name == RESTRICT_TO_USER_OBJECTS.name {
724 Err(VarError::ReadOnlyParameter(
728 RESTRICT_TO_USER_OBJECTS.name.as_str(),
729 ))
730 } else {
731 Ok(())
732 }
733 }
734
735 #[mz_ore::instrument(level = "debug")]
740 pub fn end_transaction(
741 &mut self,
742 action: EndTransactionAction,
743 ) -> BTreeMap<&'static str, String> {
744 let mut changed = BTreeMap::new();
745 let mut updates = Vec::new();
746 for (name, var) in self.vars.iter() {
747 if !var.is_mutating() {
748 continue;
749 }
750 let before = var.value();
751 let next = var.end_transaction(action).expect("must mutate");
752 let after = next.value();
753 updates.push((*name, next));
754
755 if before != after {
757 changed.insert(var.name(), after);
758 }
759 }
760 self.vars.extend(updates);
761 changed
762 }
763
764 pub fn application_name(&self) -> &str {
766 self.expect_value::<String>(&APPLICATION_NAME).as_str()
767 }
768
769 pub fn build_info(&self) -> &'static BuildInfo {
771 self.mz_version.build_info
772 }
773
774 pub fn client_encoding(&self) -> &ClientEncoding {
776 self.expect_value(&CLIENT_ENCODING)
777 }
778
779 pub fn client_min_messages(&self) -> &ClientSeverity {
781 self.expect_value(&CLIENT_MIN_MESSAGES)
782 }
783
784 pub fn cluster(&self) -> &str {
786 self.expect_value::<String>(&CLUSTER).as_str()
787 }
788
789 pub fn cluster_replica(&self) -> Option<&str> {
791 self.expect_value::<Option<String>>(&CLUSTER_REPLICA)
792 .as_deref()
793 }
794
795 pub fn current_object_missing_warnings(&self) -> bool {
798 *self.expect_value::<bool>(&CURRENT_OBJECT_MISSING_WARNINGS)
799 }
800
801 pub fn date_style(&self) -> &[&str] {
803 &self.expect_value::<DateStyle>(&DATE_STYLE).0
804 }
805
806 pub fn database(&self) -> &str {
808 self.expect_value::<String>(&DATABASE).as_str()
809 }
810
811 pub fn extra_float_digits(&self) -> i32 {
813 *self.expect_value(&EXTRA_FLOAT_DIGITS)
814 }
815
816 pub fn text_encode_settings(&self) -> TextEncodeSettings {
819 TextEncodeSettings {
820 extra_float_digits: self.extra_float_digits(),
821 }
822 }
823
824 pub fn integer_datetimes(&self) -> bool {
826 *self.expect_value(&INTEGER_DATETIMES)
827 }
828
829 pub fn intervalstyle(&self) -> &IntervalStyle {
831 self.expect_value(&INTERVAL_STYLE)
832 }
833
834 pub fn mz_version(&self) -> String {
836 self.mz_version.value()
837 }
838
839 pub fn search_path(&self) -> &[Ident] {
841 self.expect_value::<Vec<Ident>>(&SEARCH_PATH).as_slice()
842 }
843
844 pub fn server_version(&self) -> &str {
846 self.expect_value::<String>(&SERVER_VERSION).as_str()
847 }
848
849 pub fn server_version_num(&self) -> i32 {
851 *self.expect_value(&SERVER_VERSION_NUM)
852 }
853
854 pub fn sql_safe_updates(&self) -> bool {
856 *self.expect_value(&SQL_SAFE_UPDATES)
857 }
858
859 pub fn standard_conforming_strings(&self) -> bool {
862 *self.expect_value(&STANDARD_CONFORMING_STRINGS)
863 }
864
865 pub fn statement_timeout(&self) -> &Duration {
867 self.expect_value(&STATEMENT_TIMEOUT)
868 }
869
870 pub fn idle_in_transaction_session_timeout(&self) -> &Duration {
872 self.expect_value(&IDLE_IN_TRANSACTION_SESSION_TIMEOUT)
873 }
874
875 pub fn timezone(&self) -> &TimeZone {
877 self.expect_value(&TIMEZONE)
878 }
879
880 pub fn transaction_isolation(&self) -> &IsolationLevel {
883 self.expect_value(&TRANSACTION_ISOLATION)
884 }
885
886 pub fn real_time_recency(&self) -> bool {
888 *self.expect_value(&REAL_TIME_RECENCY)
889 }
890
891 pub fn real_time_recency_timeout(&self) -> &Duration {
893 self.expect_value(&REAL_TIME_RECENCY_TIMEOUT)
894 }
895
896 pub fn emit_plan_insights_notice(&self) -> bool {
898 *self.expect_value(&EMIT_PLAN_INSIGHTS_NOTICE)
899 }
900
901 pub fn emit_timestamp_notice(&self) -> bool {
903 *self.expect_value(&EMIT_TIMESTAMP_NOTICE)
904 }
905
906 pub fn emit_trace_id_notice(&self) -> bool {
908 *self.expect_value(&EMIT_TRACE_ID_NOTICE)
909 }
910
911 pub fn auto_route_catalog_queries(&self) -> bool {
913 *self.expect_value(&AUTO_ROUTE_CATALOG_QUERIES)
914 }
915
916 pub fn enable_session_rbac_checks(&self) -> bool {
918 *self.expect_value(&ENABLE_SESSION_RBAC_CHECKS)
919 }
920
921 pub fn restrict_to_user_objects(&self) -> bool {
923 *self.expect_value(&RESTRICT_TO_USER_OBJECTS)
924 }
925
926 pub fn enable_session_cardinality_estimates(&self) -> bool {
928 *self.expect_value(&ENABLE_SESSION_CARDINALITY_ESTIMATES)
929 }
930
931 pub fn is_superuser(&self) -> bool {
933 self.user.is_superuser()
934 }
935
936 pub fn user(&self) -> &User {
938 &self.user
939 }
940
941 pub fn max_query_result_size(&self) -> u64 {
943 self.expect_value::<ByteSize>(&MAX_QUERY_RESULT_SIZE)
944 .as_bytes()
945 }
946
947 pub fn set_internal_user_metadata(&mut self, metadata: InternalUserMetadata) {
949 self.user.internal_metadata = Some(metadata);
950 }
951
952 pub fn set_external_user_metadata(&mut self, metadata: ExternalUserMetadata) {
954 self.user.external_metadata = Some(metadata);
955 }
956
957 pub fn set_cluster(&mut self, cluster: String) {
958 let var = self
959 .vars
960 .get_mut(UncasedStr::new(CLUSTER.name()))
961 .expect("cluster variable must exist");
962 var.set(VarInput::Flat(&cluster), false)
963 .expect("setting cluster must succeed");
964 }
965
966 pub fn set_local_transaction_isolation(&mut self, transaction_isolation: IsolationLevel) {
967 let var = self
968 .vars
969 .get_mut(UncasedStr::new(TRANSACTION_ISOLATION.name()))
970 .expect("transaction_isolation variable must exist");
971 var.set(VarInput::Flat(&transaction_isolation.to_string()), true)
972 .expect("setting transaction isolation must succeed");
973 }
974
975 pub fn get_statement_logging_sample_rate(&self) -> Numeric {
976 *self.expect_value(&STATEMENT_LOGGING_SAMPLE_RATE)
977 }
978
979 pub fn emit_introspection_query_notice(&self) -> bool {
981 *self.expect_value(&EMIT_INTROSPECTION_QUERY_NOTICE)
982 }
983
984 pub fn unsafe_new_transaction_wall_time(&self) -> Option<CheckedTimestamp<DateTime<Utc>>> {
985 *self.expect_value(&UNSAFE_NEW_TRANSACTION_WALL_TIME)
986 }
987
988 pub fn welcome_message(&self) -> bool {
990 *self.expect_value(&WELCOME_MESSAGE)
991 }
992}
993
994pub const OLD_CATALOG_SERVER_CLUSTER: &str = "mz_introspection";
996pub const OLD_AUTO_ROUTE_CATALOG_QUERIES: &str = "auto_route_introspection_queries";
997
998fn compat_translate<'a, 'b>(name: &'a str, input: VarInput<'b>) -> (&'a str, VarInput<'b>) {
1006 if name == CLUSTER.name() {
1007 if let Ok(value) = CLUSTER.parse(input) {
1008 if value.format() == OLD_CATALOG_SERVER_CLUSTER {
1009 tracing::debug!(
1010 github_27285 = true,
1011 "encountered deprecated `cluster` variable value: {}",
1012 OLD_CATALOG_SERVER_CLUSTER,
1013 );
1014 return (name, VarInput::Flat("mz_catalog_server"));
1015 }
1016 }
1017 }
1018
1019 if name == OLD_AUTO_ROUTE_CATALOG_QUERIES {
1020 tracing::debug!(
1021 github_27285 = true,
1022 "encountered deprecated `{}` variable name",
1023 OLD_AUTO_ROUTE_CATALOG_QUERIES,
1024 );
1025 return (AUTO_ROUTE_CATALOG_QUERIES.name(), input);
1026 }
1027
1028 (name, input)
1029}
1030
1031fn compat_translate_name(name: &str) -> &str {
1032 let (name, _) = compat_translate(name, VarInput::Flat(""));
1033 name
1034}
1035
1036pub fn check_transaction_isolation_feature_flag(
1045 name: &str,
1046 input: VarInput,
1047 system_vars: &SystemVars,
1048) -> Result<(), VarError> {
1049 if UncasedStr::new(name) != UncasedStr::new(TRANSACTION_ISOLATION_VAR_NAME) {
1050 return Ok(());
1051 }
1052 let Ok(level) = IsolationLevel::parse(input) else {
1054 return Ok(());
1055 };
1056 match level {
1057 IsolationLevel::StrongSessionSerializable => ENABLE_SESSION_TIMELINES.require(system_vars),
1058 _ => Ok(()),
1059 }
1060}
1061
1062#[derive(Debug)]
1065pub struct SystemVar {
1066 definition: VarDefinition,
1067 persisted_value: Option<Box<dyn Value>>,
1069 dynamic_default: Option<Box<dyn Value>>,
1071}
1072
1073impl Clone for SystemVar {
1074 fn clone(&self) -> Self {
1075 SystemVar {
1076 definition: self.definition.clone(),
1077 persisted_value: self.persisted_value.as_ref().map(|v| v.box_clone()),
1078 dynamic_default: self.dynamic_default.as_ref().map(|v| v.box_clone()),
1079 }
1080 }
1081}
1082
1083impl SystemVar {
1084 pub fn new(definition: VarDefinition) -> Self {
1085 SystemVar {
1086 definition,
1087 persisted_value: None,
1088 dynamic_default: None,
1089 }
1090 }
1091
1092 fn is_default(&self, input: VarInput) -> Result<bool, VarError> {
1093 let v = self.definition.parse(input)?;
1094 Ok(self.definition.default_value() == v.as_ref())
1095 }
1096
1097 pub fn value_dyn(&self) -> &dyn Value {
1098 self.persisted_value
1099 .as_deref()
1100 .or(self.dynamic_default.as_deref())
1101 .unwrap_or_else(|| self.definition.default_value())
1102 }
1103
1104 pub fn value<V: 'static>(&self) -> &V {
1105 let val = self.value_dyn();
1106 val.as_any().downcast_ref::<V>().expect("success")
1107 }
1108
1109 fn parse(&self, input: VarInput) -> Result<Box<dyn Value>, VarError> {
1110 let v = self.definition.parse(input)?;
1111 self.validate_constraints(v.as_ref())?;
1113 Ok(v)
1114 }
1115
1116 fn set(&mut self, input: VarInput) -> Result<bool, VarError> {
1117 let v = self.parse(input)?;
1118
1119 if self.persisted_value.as_ref() != Some(&v) {
1120 self.persisted_value = Some(v);
1121 Ok(true)
1122 } else {
1123 Ok(false)
1124 }
1125 }
1126
1127 fn reset(&mut self) -> bool {
1128 if self.persisted_value.is_some() {
1129 self.persisted_value = None;
1130 true
1131 } else {
1132 false
1133 }
1134 }
1135
1136 fn set_default(&mut self, input: VarInput) -> Result<(), VarError> {
1137 let v = self.parse(input)?;
1138 self.dynamic_default = Some(v);
1139 Ok(())
1140 }
1141
1142 fn validate_constraints(&self, val: &dyn Value) -> Result<(), VarError> {
1143 if let Some(constraint) = &self.definition.constraint {
1144 constraint.check_constraint(self, self.value_dyn(), val)
1145 } else {
1146 Ok(())
1147 }
1148 }
1149}
1150
1151impl Var for SystemVar {
1152 fn name(&self) -> &'static str {
1153 self.definition.name.as_str()
1154 }
1155
1156 fn value(&self) -> String {
1157 self.value_dyn().format()
1158 }
1159
1160 fn description(&self) -> &'static str {
1161 self.definition.description
1162 }
1163
1164 fn type_name(&self) -> Cow<'static, str> {
1165 self.definition.type_name()
1166 }
1167
1168 fn scope(&self) -> ParameterScope {
1169 self.definition.scope()
1170 }
1171
1172 fn visible(&self, user: &User, system_vars: &SystemVars) -> Result<(), VarError> {
1173 self.definition.visible(user, system_vars)
1174 }
1175}
1176
1177#[derive(Debug, Error)]
1178pub enum NetworkPolicyError {
1179 #[error("Access denied for address {0}")]
1180 AddressDenied(IpAddr),
1181}
1182
1183#[derive(Derivative, Clone)]
1188#[derivative(Debug)]
1189pub struct SystemVars {
1190 allow_unsafe: bool,
1192 vars: BTreeMap<&'static UncasedStr, SystemVar>,
1194 #[derivative(Debug = "ignore")]
1196 callbacks: BTreeMap<String, Vec<Arc<dyn Fn(&SystemVars) + Send + Sync>>>,
1197
1198 dyncfgs: ConfigSet,
1202}
1203
1204impl Default for SystemVars {
1205 fn default() -> Self {
1206 Self::new()
1207 }
1208}
1209
1210impl SystemVars {
1211 pub fn new() -> Self {
1212 let system_vars = vec![
1213 &MAX_KAFKA_CONNECTIONS,
1214 &MAX_POSTGRES_CONNECTIONS,
1215 &MAX_MYSQL_CONNECTIONS,
1216 &MAX_SQL_SERVER_CONNECTIONS,
1217 &MAX_AWS_PRIVATELINK_CONNECTIONS,
1218 &MAX_TABLES,
1219 &MAX_SOURCES,
1220 &MAX_SINKS,
1221 &MAX_MATERIALIZED_VIEWS,
1222 &MAX_CLUSTERS,
1223 &MAX_REPLICAS_PER_CLUSTER,
1224 &MAX_CREDIT_CONSUMPTION_RATE,
1225 &MAX_DATABASES,
1226 &MAX_SCHEMAS_PER_DATABASE,
1227 &MAX_OBJECTS_PER_SCHEMA,
1228 &MAX_SECRETS,
1229 &MAX_ROLES,
1230 &MAX_NETWORK_POLICIES,
1231 &MAX_RULES_PER_NETWORK_POLICY,
1232 &MAX_RESULT_SIZE,
1233 &MAX_COPY_FROM_ROW_SIZE,
1234 &ALLOWED_CLUSTER_REPLICA_SIZES,
1235 &MAX_CONCURRENT_OCC_WRITES,
1236 &MAX_OCC_RETRIES,
1237 &upsert_rocksdb::UPSERT_ROCKSDB_COMPACTION_STYLE,
1238 &upsert_rocksdb::UPSERT_ROCKSDB_OPTIMIZE_COMPACTION_MEMTABLE_BUDGET,
1239 &upsert_rocksdb::UPSERT_ROCKSDB_LEVEL_COMPACTION_DYNAMIC_LEVEL_BYTES,
1240 &upsert_rocksdb::UPSERT_ROCKSDB_UNIVERSAL_COMPACTION_RATIO,
1241 &upsert_rocksdb::UPSERT_ROCKSDB_PARALLELISM,
1242 &upsert_rocksdb::UPSERT_ROCKSDB_COMPRESSION_TYPE,
1243 &upsert_rocksdb::UPSERT_ROCKSDB_BOTTOMMOST_COMPRESSION_TYPE,
1244 &upsert_rocksdb::UPSERT_ROCKSDB_BATCH_SIZE,
1245 &upsert_rocksdb::UPSERT_ROCKSDB_RETRY_DURATION,
1246 &upsert_rocksdb::UPSERT_ROCKSDB_STATS_LOG_INTERVAL_SECONDS,
1247 &upsert_rocksdb::UPSERT_ROCKSDB_STATS_PERSIST_INTERVAL_SECONDS,
1248 &upsert_rocksdb::UPSERT_ROCKSDB_POINT_LOOKUP_BLOCK_CACHE_SIZE_MB,
1249 &upsert_rocksdb::UPSERT_ROCKSDB_SHRINK_ALLOCATED_BUFFERS_BY_RATIO,
1250 &upsert_rocksdb::UPSERT_ROCKSDB_WRITE_BUFFER_MANAGER_CLUSTER_MEMORY_FRACTION,
1251 &upsert_rocksdb::UPSERT_ROCKSDB_WRITE_BUFFER_MANAGER_MEMORY_BYTES,
1252 &upsert_rocksdb::UPSERT_ROCKSDB_WRITE_BUFFER_MANAGER_ALLOW_STALL,
1253 &STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES,
1254 &STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES_TO_CLUSTER_SIZE_FRACTION,
1255 &STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES_DISK_ONLY,
1256 &STORAGE_STATISTICS_INTERVAL,
1257 &STORAGE_STATISTICS_COLLECTION_INTERVAL,
1258 &STORAGE_SHRINK_UPSERT_UNUSED_BUFFERS_BY_RATIO,
1259 &STORAGE_RECORD_SOURCE_SINK_NAMESPACED_ERRORS,
1260 &PERSIST_FAST_PATH_LIMIT,
1261 &METRICS_RETENTION,
1262 &DISABLED_METRIC_SINKS,
1263 &UNSAFE_MOCK_AUDIT_EVENT_TIMESTAMP,
1264 &ENABLE_RBAC_CHECKS,
1265 &PG_SOURCE_CONNECT_TIMEOUT,
1266 &PG_SOURCE_TCP_KEEPALIVES_IDLE,
1267 &PG_SOURCE_TCP_KEEPALIVES_INTERVAL,
1268 &PG_SOURCE_TCP_KEEPALIVES_RETRIES,
1269 &PG_SOURCE_TCP_USER_TIMEOUT,
1270 &PG_SOURCE_TCP_CONFIGURE_SERVER,
1271 &PG_SOURCE_SNAPSHOT_STATEMENT_TIMEOUT,
1272 &PG_SOURCE_WAL_SENDER_TIMEOUT,
1273 &PG_SOURCE_SNAPSHOT_COLLECT_STRICT_COUNT,
1274 &MYSQL_SOURCE_TCP_KEEPALIVE,
1275 &MYSQL_SOURCE_SNAPSHOT_MAX_EXECUTION_TIME,
1276 &MYSQL_SOURCE_SNAPSHOT_LOCK_WAIT_TIMEOUT,
1277 &MYSQL_SOURCE_SNAPSHOT_WAIT_TIMEOUT,
1278 &MYSQL_SOURCE_CONNECT_TIMEOUT,
1279 &SSH_CHECK_INTERVAL,
1280 &SSH_CONNECT_TIMEOUT,
1281 &SSH_KEEPALIVES_IDLE,
1282 &KAFKA_SOCKET_KEEPALIVE,
1283 &KAFKA_SOCKET_TIMEOUT,
1284 &KAFKA_TRANSACTION_TIMEOUT,
1285 &KAFKA_SOCKET_CONNECTION_SETUP_TIMEOUT,
1286 &KAFKA_FETCH_METADATA_TIMEOUT,
1287 &KAFKA_PROGRESS_RECORD_FETCH_TIMEOUT,
1288 &ENABLE_LAUNCHDARKLY,
1289 &MAX_CONNECTIONS,
1290 &NETWORK_POLICY,
1291 &SUPERUSER_RESERVED_CONNECTIONS,
1292 &KEEP_N_SOURCE_STATUS_HISTORY_ENTRIES,
1293 &KEEP_N_SINK_STATUS_HISTORY_ENTRIES,
1294 &KEEP_N_PRIVATELINK_STATUS_HISTORY_ENTRIES,
1295 &REPLICA_STATUS_HISTORY_RETENTION_WINDOW,
1296 &ENABLE_STORAGE_SHARD_FINALIZATION,
1297 &ENABLE_DEFAULT_CONNECTION_VALIDATION,
1298 &DEFAULT_TIMESTAMP_INTERVAL,
1299 &MIN_TIMESTAMP_INTERVAL,
1300 &MAX_TIMESTAMP_INTERVAL,
1301 &LOGGING_FILTER,
1302 &OPENTELEMETRY_FILTER,
1303 &LOGGING_FILTER_DEFAULTS,
1304 &OPENTELEMETRY_FILTER_DEFAULTS,
1305 &SENTRY_FILTERS,
1306 &WEBHOOKS_SECRETS_CACHING_TTL_SECS,
1307 &COORD_SLOW_MESSAGE_WARN_THRESHOLD,
1308 &grpc_client::CONNECT_TIMEOUT,
1309 &grpc_client::HTTP2_KEEP_ALIVE_INTERVAL,
1310 &grpc_client::HTTP2_KEEP_ALIVE_TIMEOUT,
1311 &cluster_scheduling::CLUSTER_MULTI_PROCESS_REPLICA_AZ_AFFINITY_WEIGHT,
1312 &cluster_scheduling::CLUSTER_SOFTEN_REPLICATION_ANTI_AFFINITY,
1313 &cluster_scheduling::CLUSTER_SOFTEN_REPLICATION_ANTI_AFFINITY_WEIGHT,
1314 &cluster_scheduling::CLUSTER_ENABLE_TOPOLOGY_SPREAD,
1315 &cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_IGNORE_NON_SINGULAR_SCALE,
1316 &cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_MAX_SKEW,
1317 &cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_MIN_DOMAINS,
1318 &cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_SOFT,
1319 &cluster_scheduling::CLUSTER_SOFTEN_AZ_AFFINITY,
1320 &cluster_scheduling::CLUSTER_SOFTEN_AZ_AFFINITY_WEIGHT,
1321 &cluster_scheduling::CLUSTER_ALTER_CHECK_READY_INTERVAL,
1322 &cluster_scheduling::CLUSTER_SECURITY_CONTEXT_ENABLED,
1323 &cluster_scheduling::CLUSTER_REFRESH_MV_COMPACTION_ESTIMATE,
1324 &grpc_client::HTTP2_KEEP_ALIVE_TIMEOUT,
1325 &STATEMENT_LOGGING_MAX_SAMPLE_RATE,
1326 &STATEMENT_LOGGING_DEFAULT_SAMPLE_RATE,
1327 &STATEMENT_LOGGING_TARGET_DATA_RATE,
1328 &STATEMENT_LOGGING_MAX_DATA_CREDIT,
1329 &ENABLE_INTERNAL_STATEMENT_LOGGING,
1330 &ENABLE_STATEMENT_ARRIVAL_LOGGING,
1331 &ENABLE_EXTENDED_PROTOCOL_IMPLICIT_TRANSACTION,
1332 &OPTIMIZER_STATS_TIMEOUT,
1333 &OPTIMIZER_ONESHOT_STATS_TIMEOUT,
1334 &PRIVATELINK_STATUS_UPDATE_QUOTA_PER_MINUTE,
1335 &WEBHOOK_CONCURRENT_REQUEST_LIMIT,
1336 &PG_TIMESTAMP_ORACLE_CONNECTION_POOL_MAX_SIZE,
1337 &PG_TIMESTAMP_ORACLE_CONNECTION_POOL_MAX_WAIT,
1338 &PG_TIMESTAMP_ORACLE_CONNECTION_POOL_TTL,
1339 &PG_TIMESTAMP_ORACLE_CONNECTION_POOL_TTL_STAGGER,
1340 &USER_STORAGE_MANAGED_COLLECTIONS_BATCH_DURATION,
1341 &FORCE_SOURCE_TABLE_SYNTAX,
1342 &OPTIMIZER_E2E_LATENCY_WARNING_THRESHOLD,
1343 &SCRAM_ITERATIONS,
1344 ];
1345
1346 let dyncfgs = mz_dyncfgs::all_dyncfgs();
1347 let dyncfg_vars: Vec<_> = dyncfgs
1348 .entries()
1349 .map(|cfg| {
1350 let var = match cfg.default() {
1351 ConfigVal::Bool(default) => {
1352 VarDefinition::new_runtime(cfg.name(), *default, cfg.desc(), false)
1353 }
1354 ConfigVal::U32(default) => {
1355 VarDefinition::new_runtime(cfg.name(), *default, cfg.desc(), false)
1356 }
1357 ConfigVal::Usize(default) => {
1358 VarDefinition::new_runtime(cfg.name(), *default, cfg.desc(), false)
1359 }
1360 ConfigVal::OptUsize(default) => {
1361 VarDefinition::new_runtime(cfg.name(), *default, cfg.desc(), false)
1362 }
1363 ConfigVal::F64(default) => {
1364 VarDefinition::new_runtime(cfg.name(), *default, cfg.desc(), false)
1365 }
1366 ConfigVal::String(default) => {
1367 VarDefinition::new_runtime(cfg.name(), default.clone(), cfg.desc(), false)
1368 }
1369 ConfigVal::OptString(default) => {
1370 VarDefinition::new_runtime(cfg.name(), default.clone(), cfg.desc(), false)
1371 }
1372 ConfigVal::Duration(default) => {
1373 VarDefinition::new_runtime(cfg.name(), default.clone(), cfg.desc(), false)
1374 }
1375 ConfigVal::Json(default) => {
1376 VarDefinition::new_runtime(cfg.name(), default.clone(), cfg.desc(), false)
1377 }
1378 };
1379 var.scoped(cfg.scope())
1382 })
1383 .collect();
1384
1385 let vars: BTreeMap<_, _> = system_vars
1386 .into_iter()
1387 .chain(definitions::FEATURE_FLAGS.iter().copied())
1389 .chain(SESSION_SYSTEM_VARS.values().copied())
1391 .cloned()
1392 .chain(dyncfg_vars)
1394 .map(|var| (var.name, SystemVar::new(var)))
1395 .collect();
1396
1397 let vars = SystemVars {
1398 vars,
1399 callbacks: BTreeMap::new(),
1400 allow_unsafe: false,
1401 dyncfgs,
1402 };
1403
1404 vars
1405 }
1406
1407 pub fn dyncfgs(&self) -> &ConfigSet {
1408 &self.dyncfgs
1409 }
1410
1411 pub fn set_unsafe(mut self, allow_unsafe: bool) -> Self {
1412 self.allow_unsafe = allow_unsafe;
1413 self
1414 }
1415
1416 pub fn allow_unsafe(&self) -> bool {
1417 self.allow_unsafe
1418 }
1419
1420 fn expect_value<V: 'static>(&self, var: &VarDefinition) -> &V {
1421 let val = self
1422 .vars
1423 .get(var.name)
1424 .expect("provided var should be in state");
1425
1426 val.value_dyn()
1427 .as_any()
1428 .downcast_ref::<V>()
1429 .expect("provided var type should matched stored var")
1430 }
1431
1432 fn expect_config_value<V: ConfigType + 'static>(&self, name: &UncasedStr) -> &V {
1433 let val = self
1434 .vars
1435 .get(name)
1436 .unwrap_or_else(|| panic!("provided var {name} should be in state"));
1437
1438 val.value_dyn()
1439 .as_any()
1440 .downcast_ref()
1441 .expect("provided var type should matched stored var")
1442 }
1443
1444 pub fn iter(&self) -> impl Iterator<Item = &dyn Var> {
1447 self.vars
1448 .values()
1449 .map(|v| v.as_var())
1450 .filter(|v| !SESSION_SYSTEM_VARS.contains_key(UncasedStr::new(v.name())))
1451 }
1452
1453 pub fn iter_synced(&self) -> impl Iterator<Item = &dyn Var> {
1457 self.iter().filter(|v| v.name() != ENABLE_LAUNCHDARKLY.name)
1458 }
1459
1460 pub fn iter_session(&self) -> impl Iterator<Item = &dyn Var> {
1462 self.vars
1463 .values()
1464 .map(|v| v.as_var())
1465 .filter(|v| SESSION_SYSTEM_VARS.contains_key(UncasedStr::new(v.name())))
1466 }
1467
1468 pub fn user_modifiable(&self, name: &str) -> bool {
1470 SESSION_SYSTEM_VARS.contains_key(UncasedStr::new(name))
1471 || name == ENABLE_RBAC_CHECKS.name()
1472 || name == NETWORK_POLICY.name()
1473 }
1474
1475 pub fn get(&self, name: &str) -> Result<&dyn Var, VarError> {
1496 self.vars
1497 .get(UncasedStr::new(name))
1498 .map(|v| v.as_var())
1499 .ok_or_else(|| VarError::UnknownParameter(name.into()))
1500 }
1501
1502 pub fn is_default(&self, name: &str, input: VarInput) -> Result<bool, VarError> {
1516 self.vars
1517 .get(UncasedStr::new(name))
1518 .ok_or_else(|| VarError::UnknownParameter(name.into()))
1519 .and_then(|v| v.is_default(input))
1520 }
1521
1522 pub fn set(&mut self, name: &str, input: VarInput) -> Result<bool, VarError> {
1545 let result = self
1546 .vars
1547 .get_mut(UncasedStr::new(name))
1548 .ok_or_else(|| VarError::UnknownParameter(name.into()))
1549 .and_then(|v| v.set(input))?;
1550 Ok(result)
1551 }
1552
1553 pub fn parse(&self, name: &str, input: VarInput) -> Result<Box<dyn Value>, VarError> {
1574 self.vars
1575 .get(UncasedStr::new(name))
1576 .ok_or_else(|| VarError::UnknownParameter(name.into()))
1577 .and_then(|v| v.parse(input))
1578 }
1579
1580 pub fn set_default(&mut self, name: &str, input: VarInput) -> Result<(), VarError> {
1588 self.vars
1589 .get_mut(UncasedStr::new(name))
1590 .ok_or_else(|| VarError::UnknownParameter(name.into()))
1591 .and_then(|v| v.set_default(input))?;
1592 Ok(())
1593 }
1594
1595 pub fn reset(&mut self, name: &str) -> Result<bool, VarError> {
1613 let result = self
1614 .vars
1615 .get_mut(UncasedStr::new(name))
1616 .ok_or_else(|| VarError::UnknownParameter(name.into()))
1617 .map(|v| v.reset())?;
1618 Ok(result)
1619 }
1620
1621 pub fn defaults(&self) -> BTreeMap<String, String> {
1623 self.vars
1624 .iter()
1625 .map(|(name, var)| {
1626 let default = var
1627 .dynamic_default
1628 .as_deref()
1629 .unwrap_or_else(|| var.definition.default_value());
1630 (name.as_str().to_owned(), default.format())
1631 })
1632 .collect()
1633 }
1634
1635 pub fn register_callback(
1654 &mut self,
1655 var: &VarDefinition,
1656 callback: Arc<dyn Fn(&SystemVars) + Send + Sync>,
1657 ) {
1658 self.callbacks
1659 .entry(var.name().to_string())
1660 .or_default()
1661 .push(callback);
1662 self.notify_callbacks(var.name());
1663 }
1664
1665 pub fn notify_all_callbacks(&self) {
1671 for callbacks in self.callbacks.values() {
1672 for callback in callbacks {
1673 (callback)(self);
1674 }
1675 }
1676 }
1677
1678 fn notify_callbacks(&self, name: &str) {
1680 if let Some(callbacks) = self.callbacks.get(name) {
1682 for callback in callbacks {
1683 (callback)(self);
1684 }
1685 }
1686 }
1687
1688 pub fn default_cluster(&self) -> String {
1691 self.expect_value::<String>(&CLUSTER).to_owned()
1692 }
1693
1694 pub fn max_kafka_connections(&self) -> u32 {
1696 *self.expect_value(&MAX_KAFKA_CONNECTIONS)
1697 }
1698
1699 pub fn max_postgres_connections(&self) -> u32 {
1701 *self.expect_value(&MAX_POSTGRES_CONNECTIONS)
1702 }
1703
1704 pub fn max_mysql_connections(&self) -> u32 {
1706 *self.expect_value(&MAX_MYSQL_CONNECTIONS)
1707 }
1708
1709 pub fn max_sql_server_connections(&self) -> u32 {
1711 *self.expect_value(&MAX_SQL_SERVER_CONNECTIONS)
1712 }
1713
1714 pub fn max_aws_privatelink_connections(&self) -> u32 {
1716 *self.expect_value(&MAX_AWS_PRIVATELINK_CONNECTIONS)
1717 }
1718
1719 pub fn max_tables(&self) -> u32 {
1721 *self.expect_value(&MAX_TABLES)
1722 }
1723
1724 pub fn max_sources(&self) -> u32 {
1726 *self.expect_value(&MAX_SOURCES)
1727 }
1728
1729 pub fn max_sinks(&self) -> u32 {
1731 *self.expect_value(&MAX_SINKS)
1732 }
1733
1734 pub fn max_materialized_views(&self) -> u32 {
1736 *self.expect_value(&MAX_MATERIALIZED_VIEWS)
1737 }
1738
1739 pub fn max_clusters(&self) -> u32 {
1741 *self.expect_value(&MAX_CLUSTERS)
1742 }
1743
1744 pub fn max_replicas_per_cluster(&self) -> u32 {
1746 *self.expect_value(&MAX_REPLICAS_PER_CLUSTER)
1747 }
1748
1749 pub fn max_credit_consumption_rate(&self) -> Numeric {
1751 *self.expect_value(&MAX_CREDIT_CONSUMPTION_RATE)
1752 }
1753
1754 pub fn max_databases(&self) -> u32 {
1756 *self.expect_value(&MAX_DATABASES)
1757 }
1758
1759 pub fn max_schemas_per_database(&self) -> u32 {
1761 *self.expect_value(&MAX_SCHEMAS_PER_DATABASE)
1762 }
1763
1764 pub fn max_objects_per_schema(&self) -> u32 {
1766 *self.expect_value(&MAX_OBJECTS_PER_SCHEMA)
1767 }
1768
1769 pub fn max_secrets(&self) -> u32 {
1771 *self.expect_value(&MAX_SECRETS)
1772 }
1773
1774 pub fn max_roles(&self) -> u32 {
1776 *self.expect_value(&MAX_ROLES)
1777 }
1778
1779 pub fn max_network_policies(&self) -> u32 {
1781 *self.expect_value(&MAX_NETWORK_POLICIES)
1782 }
1783
1784 pub fn max_rules_per_network_policy(&self) -> u32 {
1786 *self.expect_value(&MAX_RULES_PER_NETWORK_POLICY)
1787 }
1788
1789 pub fn max_result_size(&self) -> u64 {
1791 self.expect_value::<ByteSize>(&MAX_RESULT_SIZE).as_bytes()
1792 }
1793
1794 pub fn max_copy_from_row_size(&self) -> u64 {
1796 self.expect_value::<ByteSize>(&MAX_COPY_FROM_ROW_SIZE)
1797 .as_bytes()
1798 }
1799
1800 pub fn allowed_cluster_replica_sizes(&self) -> Vec<String> {
1802 self.expect_value::<Vec<Ident>>(&ALLOWED_CLUSTER_REPLICA_SIZES)
1803 .into_iter()
1804 .map(|s| s.as_str().into())
1805 .collect()
1806 }
1807
1808 pub fn max_concurrent_occ_writes(&self) -> u32 {
1810 *self.expect_value(&MAX_CONCURRENT_OCC_WRITES)
1811 }
1812
1813 pub fn max_occ_retries(&self) -> u32 {
1815 *self.expect_value(&MAX_OCC_RETRIES)
1816 }
1817
1818 pub fn default_cluster_replication_factor(&self) -> u32 {
1820 *self.expect_value::<u32>(&DEFAULT_CLUSTER_REPLICATION_FACTOR)
1821 }
1822
1823 pub fn upsert_rocksdb_compaction_style(&self) -> mz_rocksdb_types::config::CompactionStyle {
1824 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_COMPACTION_STYLE)
1825 }
1826
1827 pub fn upsert_rocksdb_optimize_compaction_memtable_budget(&self) -> usize {
1828 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_OPTIMIZE_COMPACTION_MEMTABLE_BUDGET)
1829 }
1830
1831 pub fn upsert_rocksdb_level_compaction_dynamic_level_bytes(&self) -> bool {
1832 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_LEVEL_COMPACTION_DYNAMIC_LEVEL_BYTES)
1833 }
1834
1835 pub fn upsert_rocksdb_universal_compaction_ratio(&self) -> i32 {
1836 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_UNIVERSAL_COMPACTION_RATIO)
1837 }
1838
1839 pub fn upsert_rocksdb_parallelism(&self) -> Option<i32> {
1840 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_PARALLELISM)
1841 }
1842
1843 pub fn upsert_rocksdb_compression_type(&self) -> mz_rocksdb_types::config::CompressionType {
1844 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_COMPRESSION_TYPE)
1845 }
1846
1847 pub fn upsert_rocksdb_bottommost_compression_type(
1848 &self,
1849 ) -> mz_rocksdb_types::config::CompressionType {
1850 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_BOTTOMMOST_COMPRESSION_TYPE)
1851 }
1852
1853 pub fn upsert_rocksdb_batch_size(&self) -> usize {
1854 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_BATCH_SIZE)
1855 }
1856
1857 pub fn upsert_rocksdb_retry_duration(&self) -> Duration {
1858 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_RETRY_DURATION)
1859 }
1860
1861 pub fn upsert_rocksdb_stats_log_interval_seconds(&self) -> u32 {
1862 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_STATS_LOG_INTERVAL_SECONDS)
1863 }
1864
1865 pub fn upsert_rocksdb_stats_persist_interval_seconds(&self) -> u32 {
1866 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_STATS_PERSIST_INTERVAL_SECONDS)
1867 }
1868
1869 pub fn upsert_rocksdb_point_lookup_block_cache_size_mb(&self) -> Option<u32> {
1870 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_POINT_LOOKUP_BLOCK_CACHE_SIZE_MB)
1871 }
1872
1873 pub fn upsert_rocksdb_shrink_allocated_buffers_by_ratio(&self) -> usize {
1874 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_SHRINK_ALLOCATED_BUFFERS_BY_RATIO)
1875 }
1876
1877 pub fn upsert_rocksdb_write_buffer_manager_cluster_memory_fraction(&self) -> Option<Numeric> {
1878 *self.expect_value(
1879 &upsert_rocksdb::UPSERT_ROCKSDB_WRITE_BUFFER_MANAGER_CLUSTER_MEMORY_FRACTION,
1880 )
1881 }
1882
1883 pub fn upsert_rocksdb_write_buffer_manager_memory_bytes(&self) -> Option<usize> {
1884 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_WRITE_BUFFER_MANAGER_MEMORY_BYTES)
1885 }
1886
1887 pub fn upsert_rocksdb_write_buffer_manager_allow_stall(&self) -> bool {
1888 *self.expect_value(&upsert_rocksdb::UPSERT_ROCKSDB_WRITE_BUFFER_MANAGER_ALLOW_STALL)
1889 }
1890
1891 pub fn persist_fast_path_limit(&self) -> usize {
1892 *self.expect_value(&PERSIST_FAST_PATH_LIMIT)
1893 }
1894
1895 pub fn pg_source_connect_timeout(&self) -> Duration {
1897 *self.expect_value(&PG_SOURCE_CONNECT_TIMEOUT)
1898 }
1899
1900 pub fn pg_source_tcp_keepalives_retries(&self) -> u32 {
1902 *self.expect_value(&PG_SOURCE_TCP_KEEPALIVES_RETRIES)
1903 }
1904
1905 pub fn pg_source_tcp_keepalives_idle(&self) -> Duration {
1907 *self.expect_value(&PG_SOURCE_TCP_KEEPALIVES_IDLE)
1908 }
1909
1910 pub fn pg_source_tcp_keepalives_interval(&self) -> Duration {
1912 *self.expect_value(&PG_SOURCE_TCP_KEEPALIVES_INTERVAL)
1913 }
1914
1915 pub fn pg_source_tcp_user_timeout(&self) -> Duration {
1917 *self.expect_value(&PG_SOURCE_TCP_USER_TIMEOUT)
1918 }
1919
1920 pub fn pg_source_tcp_configure_server(&self) -> bool {
1922 *self.expect_value(&PG_SOURCE_TCP_CONFIGURE_SERVER)
1923 }
1924
1925 pub fn pg_source_snapshot_statement_timeout(&self) -> Duration {
1927 *self.expect_value(&PG_SOURCE_SNAPSHOT_STATEMENT_TIMEOUT)
1928 }
1929
1930 pub fn pg_source_wal_sender_timeout(&self) -> Option<Duration> {
1932 *self.expect_value(&PG_SOURCE_WAL_SENDER_TIMEOUT)
1933 }
1934
1935 pub fn pg_source_snapshot_collect_strict_count(&self) -> bool {
1937 *self.expect_value(&PG_SOURCE_SNAPSHOT_COLLECT_STRICT_COUNT)
1938 }
1939
1940 pub fn mysql_source_tcp_keepalive(&self) -> Duration {
1942 *self.expect_value(&MYSQL_SOURCE_TCP_KEEPALIVE)
1943 }
1944
1945 pub fn mysql_source_snapshot_max_execution_time(&self) -> Duration {
1947 *self.expect_value(&MYSQL_SOURCE_SNAPSHOT_MAX_EXECUTION_TIME)
1948 }
1949
1950 pub fn mysql_source_snapshot_lock_wait_timeout(&self) -> Duration {
1952 *self.expect_value(&MYSQL_SOURCE_SNAPSHOT_LOCK_WAIT_TIMEOUT)
1953 }
1954
1955 pub fn mysql_source_snapshot_wait_timeout(&self) -> Duration {
1957 *self.expect_value(&MYSQL_SOURCE_SNAPSHOT_WAIT_TIMEOUT)
1958 }
1959
1960 pub fn mysql_source_connect_timeout(&self) -> Duration {
1962 *self.expect_value(&MYSQL_SOURCE_CONNECT_TIMEOUT)
1963 }
1964
1965 pub fn ssh_check_interval(&self) -> Duration {
1967 *self.expect_value(&SSH_CHECK_INTERVAL)
1968 }
1969
1970 pub fn ssh_connect_timeout(&self) -> Duration {
1972 *self.expect_value(&SSH_CONNECT_TIMEOUT)
1973 }
1974
1975 pub fn ssh_keepalives_idle(&self) -> Duration {
1977 *self.expect_value(&SSH_KEEPALIVES_IDLE)
1978 }
1979
1980 pub fn kafka_socket_keepalive(&self) -> bool {
1982 *self.expect_value(&KAFKA_SOCKET_KEEPALIVE)
1983 }
1984
1985 pub fn kafka_socket_timeout(&self) -> Option<Duration> {
1987 *self.expect_value(&KAFKA_SOCKET_TIMEOUT)
1988 }
1989
1990 pub fn kafka_transaction_timeout(&self) -> Duration {
1992 *self.expect_value(&KAFKA_TRANSACTION_TIMEOUT)
1993 }
1994
1995 pub fn kafka_socket_connection_setup_timeout(&self) -> Duration {
1997 *self.expect_value(&KAFKA_SOCKET_CONNECTION_SETUP_TIMEOUT)
1998 }
1999
2000 pub fn kafka_fetch_metadata_timeout(&self) -> Duration {
2002 *self.expect_value(&KAFKA_FETCH_METADATA_TIMEOUT)
2003 }
2004
2005 pub fn kafka_progress_record_fetch_timeout(&self) -> Option<Duration> {
2007 *self.expect_value(&KAFKA_PROGRESS_RECORD_FETCH_TIMEOUT)
2008 }
2009
2010 pub fn crdb_connect_timeout(&self) -> Duration {
2012 *self.expect_config_value(UncasedStr::new(
2013 mz_persist_client::cfg::CRDB_CONNECT_TIMEOUT.name(),
2014 ))
2015 }
2016
2017 pub fn crdb_tcp_user_timeout(&self) -> Duration {
2019 *self.expect_config_value(UncasedStr::new(
2020 mz_persist_client::cfg::CRDB_TCP_USER_TIMEOUT.name(),
2021 ))
2022 }
2023
2024 pub fn crdb_keepalives_idle(&self) -> Duration {
2026 *self.expect_config_value(UncasedStr::new(
2027 mz_persist_client::cfg::CRDB_KEEPALIVES_IDLE.name(),
2028 ))
2029 }
2030
2031 pub fn crdb_keepalives_interval(&self) -> Duration {
2033 *self.expect_config_value(UncasedStr::new(
2034 mz_persist_client::cfg::CRDB_KEEPALIVES_INTERVAL.name(),
2035 ))
2036 }
2037
2038 pub fn crdb_keepalives_retries(&self) -> u32 {
2040 *self.expect_config_value(UncasedStr::new(
2041 mz_persist_client::cfg::CRDB_KEEPALIVES_RETRIES.name(),
2042 ))
2043 }
2044
2045 pub fn storage_dataflow_max_inflight_bytes(&self) -> Option<usize> {
2047 *self.expect_value(&STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES)
2048 }
2049
2050 pub fn storage_dataflow_max_inflight_bytes_to_cluster_size_fraction(&self) -> Option<Numeric> {
2052 *self.expect_value(&STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES_TO_CLUSTER_SIZE_FRACTION)
2053 }
2054
2055 pub fn storage_shrink_upsert_unused_buffers_by_ratio(&self) -> usize {
2057 *self.expect_value(&STORAGE_SHRINK_UPSERT_UNUSED_BUFFERS_BY_RATIO)
2058 }
2059
2060 pub fn storage_dataflow_max_inflight_bytes_disk_only(&self) -> bool {
2062 *self.expect_value(&STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES_DISK_ONLY)
2063 }
2064
2065 pub fn storage_statistics_interval(&self) -> Duration {
2067 *self.expect_value(&STORAGE_STATISTICS_INTERVAL)
2068 }
2069
2070 pub fn storage_statistics_collection_interval(&self) -> Duration {
2072 *self.expect_value(&STORAGE_STATISTICS_COLLECTION_INTERVAL)
2073 }
2074
2075 pub fn storage_record_source_sink_namespaced_errors(&self) -> bool {
2077 *self.expect_value(&STORAGE_RECORD_SOURCE_SINK_NAMESPACED_ERRORS)
2078 }
2079
2080 pub fn persist_stats_filter_enabled(&self) -> bool {
2082 *self.expect_config_value(UncasedStr::new(
2083 mz_persist_client::stats::STATS_FILTER_ENABLED.name(),
2084 ))
2085 }
2086
2087 pub fn scram_iterations(&self) -> NonZeroU32 {
2088 *self.expect_value(&SCRAM_ITERATIONS)
2089 }
2090
2091 pub fn dyncfg_updates(&self) -> ConfigUpdates {
2096 let mut updates = ConfigUpdates::default();
2097 for entry in self.dyncfgs.entries() {
2098 let name = UncasedStr::new(entry.name());
2099 let val = match entry.val() {
2100 ConfigVal::Bool(_) => ConfigVal::from(*self.expect_config_value::<bool>(name)),
2101 ConfigVal::U32(_) => ConfigVal::from(*self.expect_config_value::<u32>(name)),
2102 ConfigVal::Usize(_) => ConfigVal::from(*self.expect_config_value::<usize>(name)),
2103 ConfigVal::OptUsize(_) => {
2104 ConfigVal::from(*self.expect_config_value::<Option<usize>>(name))
2105 }
2106 ConfigVal::F64(_) => ConfigVal::from(*self.expect_config_value::<f64>(name)),
2107 ConfigVal::String(_) => {
2108 ConfigVal::from(self.expect_config_value::<String>(name).clone())
2109 }
2110 ConfigVal::OptString(_) => {
2111 ConfigVal::from(self.expect_config_value::<Option<String>>(name).clone())
2112 }
2113 ConfigVal::Duration(_) => {
2114 ConfigVal::from(*self.expect_config_value::<Duration>(name))
2115 }
2116 ConfigVal::Json(_) => {
2117 ConfigVal::from(self.expect_config_value::<serde_json::Value>(name).clone())
2118 }
2119 };
2120 updates.add_dynamic(entry.name(), val);
2121 }
2122 updates
2123 }
2124
2125 pub fn sync_dyncfgs(&self) -> ConfigUpdates {
2133 let updates = self.dyncfg_updates();
2134 updates.apply(&self.dyncfgs);
2135 updates
2136 }
2137
2138 pub fn metrics_retention(&self) -> Duration {
2140 *self.expect_value(&METRICS_RETENTION)
2141 }
2142
2143 pub fn disabled_metric_sinks(&self) -> Vec<String> {
2145 self.expect_value::<Vec<Ident>>(&DISABLED_METRIC_SINKS)
2146 .into_iter()
2147 .map(|s| s.as_str().into())
2148 .collect()
2149 }
2150
2151 pub fn unsafe_mock_audit_event_timestamp(&self) -> Option<mz_repr::Timestamp> {
2153 *self.expect_value(&UNSAFE_MOCK_AUDIT_EVENT_TIMESTAMP)
2154 }
2155
2156 pub fn enable_rbac_checks(&self) -> bool {
2158 *self.expect_value(&ENABLE_RBAC_CHECKS)
2159 }
2160
2161 pub fn max_connections(&self) -> u32 {
2163 *self.expect_value(&MAX_CONNECTIONS)
2164 }
2165
2166 pub fn default_network_policy_name(&self) -> String {
2167 self.expect_value::<String>(&NETWORK_POLICY).clone()
2168 }
2169
2170 pub fn superuser_reserved_connections(&self) -> u32 {
2172 *self.expect_value(&SUPERUSER_RESERVED_CONNECTIONS)
2173 }
2174
2175 pub fn keep_n_source_status_history_entries(&self) -> usize {
2176 *self.expect_value(&KEEP_N_SOURCE_STATUS_HISTORY_ENTRIES)
2177 }
2178
2179 pub fn keep_n_sink_status_history_entries(&self) -> usize {
2180 *self.expect_value(&KEEP_N_SINK_STATUS_HISTORY_ENTRIES)
2181 }
2182
2183 pub fn keep_n_privatelink_status_history_entries(&self) -> usize {
2184 *self.expect_value(&KEEP_N_PRIVATELINK_STATUS_HISTORY_ENTRIES)
2185 }
2186
2187 pub fn replica_status_history_retention_window(&self) -> Duration {
2188 *self.expect_value(&REPLICA_STATUS_HISTORY_RETENTION_WINDOW)
2189 }
2190
2191 pub fn enable_storage_shard_finalization(&self) -> bool {
2193 *self.expect_value(&ENABLE_STORAGE_SHARD_FINALIZATION)
2194 }
2195
2196 pub fn enable_default_connection_validation(&self) -> bool {
2198 *self.expect_value(&ENABLE_DEFAULT_CONNECTION_VALIDATION)
2199 }
2200
2201 pub fn default_timestamp_interval(&self) -> Duration {
2203 *self.expect_value(&DEFAULT_TIMESTAMP_INTERVAL)
2204 }
2205
2206 pub fn min_timestamp_interval(&self) -> Duration {
2208 *self.expect_value(&MIN_TIMESTAMP_INTERVAL)
2209 }
2210 pub fn max_timestamp_interval(&self) -> Duration {
2212 *self.expect_value(&MAX_TIMESTAMP_INTERVAL)
2213 }
2214
2215 pub fn logging_filter(&self) -> CloneableEnvFilter {
2216 self.expect_value::<CloneableEnvFilter>(&LOGGING_FILTER)
2217 .clone()
2218 }
2219
2220 pub fn opentelemetry_filter(&self) -> CloneableEnvFilter {
2221 self.expect_value::<CloneableEnvFilter>(&OPENTELEMETRY_FILTER)
2222 .clone()
2223 }
2224
2225 pub fn logging_filter_defaults(&self) -> Vec<SerializableDirective> {
2226 self.expect_value::<Vec<SerializableDirective>>(&LOGGING_FILTER_DEFAULTS)
2227 .clone()
2228 }
2229
2230 pub fn opentelemetry_filter_defaults(&self) -> Vec<SerializableDirective> {
2231 self.expect_value::<Vec<SerializableDirective>>(&OPENTELEMETRY_FILTER_DEFAULTS)
2232 .clone()
2233 }
2234
2235 pub fn sentry_filters(&self) -> Vec<SerializableDirective> {
2236 self.expect_value::<Vec<SerializableDirective>>(&SENTRY_FILTERS)
2237 .clone()
2238 }
2239
2240 pub fn webhooks_secrets_caching_ttl_secs(&self) -> usize {
2241 *self.expect_value(&WEBHOOKS_SECRETS_CACHING_TTL_SECS)
2242 }
2243
2244 pub fn coord_slow_message_warn_threshold(&self) -> Duration {
2245 *self.expect_value(&COORD_SLOW_MESSAGE_WARN_THRESHOLD)
2246 }
2247
2248 pub fn grpc_client_http2_keep_alive_interval(&self) -> Duration {
2249 *self.expect_value(&grpc_client::HTTP2_KEEP_ALIVE_INTERVAL)
2250 }
2251
2252 pub fn grpc_client_http2_keep_alive_timeout(&self) -> Duration {
2253 *self.expect_value(&grpc_client::HTTP2_KEEP_ALIVE_TIMEOUT)
2254 }
2255
2256 pub fn grpc_connect_timeout(&self) -> Duration {
2257 *self.expect_value(&grpc_client::CONNECT_TIMEOUT)
2258 }
2259
2260 pub fn cluster_multi_process_replica_az_affinity_weight(&self) -> Option<i32> {
2261 *self.expect_value(&cluster_scheduling::CLUSTER_MULTI_PROCESS_REPLICA_AZ_AFFINITY_WEIGHT)
2262 }
2263
2264 pub fn cluster_soften_replication_anti_affinity(&self) -> bool {
2265 *self.expect_value(&cluster_scheduling::CLUSTER_SOFTEN_REPLICATION_ANTI_AFFINITY)
2266 }
2267
2268 pub fn cluster_soften_replication_anti_affinity_weight(&self) -> i32 {
2269 *self.expect_value(&cluster_scheduling::CLUSTER_SOFTEN_REPLICATION_ANTI_AFFINITY_WEIGHT)
2270 }
2271
2272 pub fn cluster_enable_topology_spread(&self) -> bool {
2273 *self.expect_value(&cluster_scheduling::CLUSTER_ENABLE_TOPOLOGY_SPREAD)
2274 }
2275
2276 pub fn cluster_topology_spread_ignore_non_singular_scale(&self) -> bool {
2277 *self.expect_value(&cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_IGNORE_NON_SINGULAR_SCALE)
2278 }
2279
2280 pub fn cluster_topology_spread_max_skew(&self) -> i32 {
2281 *self.expect_value(&cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_MAX_SKEW)
2282 }
2283
2284 pub fn cluster_topology_spread_set_min_domains(&self) -> Option<i32> {
2285 *self.expect_value(&cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_MIN_DOMAINS)
2286 }
2287
2288 pub fn cluster_topology_spread_soft(&self) -> bool {
2289 *self.expect_value(&cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_SOFT)
2290 }
2291
2292 pub fn cluster_soften_az_affinity(&self) -> bool {
2293 *self.expect_value(&cluster_scheduling::CLUSTER_SOFTEN_AZ_AFFINITY)
2294 }
2295
2296 pub fn cluster_soften_az_affinity_weight(&self) -> i32 {
2297 *self.expect_value(&cluster_scheduling::CLUSTER_SOFTEN_AZ_AFFINITY_WEIGHT)
2298 }
2299
2300 pub fn cluster_alter_check_ready_interval(&self) -> Duration {
2301 *self.expect_value(&cluster_scheduling::CLUSTER_ALTER_CHECK_READY_INTERVAL)
2302 }
2303
2304 pub fn cluster_security_context_enabled(&self) -> bool {
2305 *self.expect_value(&cluster_scheduling::CLUSTER_SECURITY_CONTEXT_ENABLED)
2306 }
2307
2308 pub fn cluster_refresh_mv_compaction_estimate(&self) -> Duration {
2309 *self.expect_value(&cluster_scheduling::CLUSTER_REFRESH_MV_COMPACTION_ESTIMATE)
2310 }
2311
2312 pub fn privatelink_status_update_quota_per_minute(&self) -> u32 {
2314 *self.expect_value(&PRIVATELINK_STATUS_UPDATE_QUOTA_PER_MINUTE)
2315 }
2316
2317 pub fn statement_logging_target_data_rate(&self) -> Option<usize> {
2318 *self.expect_value(&STATEMENT_LOGGING_TARGET_DATA_RATE)
2319 }
2320
2321 pub fn statement_logging_max_data_credit(&self) -> Option<usize> {
2322 *self.expect_value(&STATEMENT_LOGGING_MAX_DATA_CREDIT)
2323 }
2324
2325 pub fn statement_logging_max_sample_rate(&self) -> Numeric {
2327 *self.expect_value(&STATEMENT_LOGGING_MAX_SAMPLE_RATE)
2328 }
2329
2330 pub fn statement_logging_default_sample_rate(&self) -> Numeric {
2332 *self.expect_value(&STATEMENT_LOGGING_DEFAULT_SAMPLE_RATE)
2333 }
2334
2335 pub fn enable_internal_statement_logging(&self) -> bool {
2337 *self.expect_value(&ENABLE_INTERNAL_STATEMENT_LOGGING)
2338 }
2339
2340 pub fn enable_statement_arrival_logging(&self) -> bool {
2342 *self.expect_value(&ENABLE_STATEMENT_ARRIVAL_LOGGING)
2343 }
2344
2345 pub fn enable_extended_protocol_implicit_transaction(&self) -> bool {
2348 *self.expect_value(&ENABLE_EXTENDED_PROTOCOL_IMPLICIT_TRANSACTION)
2349 }
2350
2351 pub fn optimizer_stats_timeout(&self) -> Duration {
2353 *self.expect_value(&OPTIMIZER_STATS_TIMEOUT)
2354 }
2355
2356 pub fn optimizer_oneshot_stats_timeout(&self) -> Duration {
2358 *self.expect_value(&OPTIMIZER_ONESHOT_STATS_TIMEOUT)
2359 }
2360
2361 pub fn webhook_concurrent_request_limit(&self) -> usize {
2363 *self.expect_value(&WEBHOOK_CONCURRENT_REQUEST_LIMIT)
2364 }
2365
2366 pub fn pg_timestamp_oracle_connection_pool_max_size(&self) -> usize {
2368 *self.expect_value(&PG_TIMESTAMP_ORACLE_CONNECTION_POOL_MAX_SIZE)
2369 }
2370
2371 pub fn pg_timestamp_oracle_connection_pool_max_wait(&self) -> Option<Duration> {
2373 *self.expect_value(&PG_TIMESTAMP_ORACLE_CONNECTION_POOL_MAX_WAIT)
2374 }
2375
2376 pub fn pg_timestamp_oracle_connection_pool_ttl(&self) -> Duration {
2378 *self.expect_value(&PG_TIMESTAMP_ORACLE_CONNECTION_POOL_TTL)
2379 }
2380
2381 pub fn pg_timestamp_oracle_connection_pool_ttl_stagger(&self) -> Duration {
2383 *self.expect_value(&PG_TIMESTAMP_ORACLE_CONNECTION_POOL_TTL_STAGGER)
2384 }
2385
2386 pub fn user_storage_managed_collections_batch_duration(&self) -> Duration {
2388 *self.expect_value(&USER_STORAGE_MANAGED_COLLECTIONS_BATCH_DURATION)
2389 }
2390
2391 pub fn force_source_table_syntax(&self) -> bool {
2392 *self.expect_value(&FORCE_SOURCE_TABLE_SYNTAX)
2393 }
2394
2395 pub fn optimizer_e2e_latency_warning_threshold(&self) -> Duration {
2396 *self.expect_value(&OPTIMIZER_E2E_LATENCY_WARNING_THRESHOLD)
2397 }
2398
2399 pub fn is_controller_config_var(&self, name: &str) -> bool {
2401 self.is_dyncfg_var(name)
2402 }
2403
2404 pub fn is_compute_config_var(&self, name: &str) -> bool {
2408 name == MAX_RESULT_SIZE.name() || self.is_dyncfg_var(name) || is_tracing_var(name)
2409 }
2410
2411 pub fn is_metrics_config_var(&self, name: &str) -> bool {
2413 self.is_dyncfg_var(name)
2414 }
2415
2416 pub fn is_storage_config_var(&self, name: &str) -> bool {
2418 name == PG_SOURCE_CONNECT_TIMEOUT.name()
2419 || name == PG_SOURCE_TCP_KEEPALIVES_IDLE.name()
2420 || name == PG_SOURCE_TCP_KEEPALIVES_INTERVAL.name()
2421 || name == PG_SOURCE_TCP_KEEPALIVES_RETRIES.name()
2422 || name == PG_SOURCE_TCP_USER_TIMEOUT.name()
2423 || name == PG_SOURCE_TCP_CONFIGURE_SERVER.name()
2424 || name == PG_SOURCE_SNAPSHOT_STATEMENT_TIMEOUT.name()
2425 || name == PG_SOURCE_WAL_SENDER_TIMEOUT.name()
2426 || name == PG_SOURCE_SNAPSHOT_COLLECT_STRICT_COUNT.name()
2427 || name == MYSQL_SOURCE_TCP_KEEPALIVE.name()
2428 || name == MYSQL_SOURCE_SNAPSHOT_MAX_EXECUTION_TIME.name()
2429 || name == MYSQL_SOURCE_SNAPSHOT_LOCK_WAIT_TIMEOUT.name()
2430 || name == MYSQL_SOURCE_SNAPSHOT_WAIT_TIMEOUT.name()
2431 || name == MYSQL_SOURCE_CONNECT_TIMEOUT.name()
2432 || name == ENABLE_STORAGE_SHARD_FINALIZATION.name()
2433 || name == SSH_CHECK_INTERVAL.name()
2434 || name == SSH_CONNECT_TIMEOUT.name()
2435 || name == SSH_KEEPALIVES_IDLE.name()
2436 || name == KAFKA_SOCKET_KEEPALIVE.name()
2437 || name == KAFKA_SOCKET_TIMEOUT.name()
2438 || name == KAFKA_TRANSACTION_TIMEOUT.name()
2439 || name == KAFKA_SOCKET_CONNECTION_SETUP_TIMEOUT.name()
2440 || name == KAFKA_FETCH_METADATA_TIMEOUT.name()
2441 || name == KAFKA_PROGRESS_RECORD_FETCH_TIMEOUT.name()
2442 || name == STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES.name()
2443 || name == STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES_TO_CLUSTER_SIZE_FRACTION.name()
2444 || name == STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES_DISK_ONLY.name()
2445 || name == STORAGE_SHRINK_UPSERT_UNUSED_BUFFERS_BY_RATIO.name()
2446 || name == STORAGE_RECORD_SOURCE_SINK_NAMESPACED_ERRORS.name()
2447 || name == STORAGE_STATISTICS_INTERVAL.name()
2448 || name == STORAGE_STATISTICS_COLLECTION_INTERVAL.name()
2449 || name == USER_STORAGE_MANAGED_COLLECTIONS_BATCH_DURATION.name()
2450 || is_upsert_rocksdb_config_var(name)
2451 || self.is_dyncfg_var(name)
2452 || is_tracing_var(name)
2453 }
2454
2455 fn is_dyncfg_var(&self, name: &str) -> bool {
2457 self.dyncfgs.entries().any(|e| name == e.name())
2458 }
2459}
2460
2461pub fn is_tracing_var(name: &str) -> bool {
2462 name == LOGGING_FILTER.name()
2463 || name == LOGGING_FILTER_DEFAULTS.name()
2464 || name == OPENTELEMETRY_FILTER.name()
2465 || name == OPENTELEMETRY_FILTER_DEFAULTS.name()
2466 || name == SENTRY_FILTERS.name()
2467}
2468
2469pub fn is_secrets_caching_var(name: &str) -> bool {
2471 name == WEBHOOKS_SECRETS_CACHING_TTL_SECS.name()
2472}
2473
2474fn is_upsert_rocksdb_config_var(name: &str) -> bool {
2475 name == upsert_rocksdb::UPSERT_ROCKSDB_COMPACTION_STYLE.name()
2476 || name == upsert_rocksdb::UPSERT_ROCKSDB_OPTIMIZE_COMPACTION_MEMTABLE_BUDGET.name()
2477 || name == upsert_rocksdb::UPSERT_ROCKSDB_LEVEL_COMPACTION_DYNAMIC_LEVEL_BYTES.name()
2478 || name == upsert_rocksdb::UPSERT_ROCKSDB_UNIVERSAL_COMPACTION_RATIO.name()
2479 || name == upsert_rocksdb::UPSERT_ROCKSDB_PARALLELISM.name()
2480 || name == upsert_rocksdb::UPSERT_ROCKSDB_COMPRESSION_TYPE.name()
2481 || name == upsert_rocksdb::UPSERT_ROCKSDB_BOTTOMMOST_COMPRESSION_TYPE.name()
2482 || name == upsert_rocksdb::UPSERT_ROCKSDB_BATCH_SIZE.name()
2483 || name == upsert_rocksdb::UPSERT_ROCKSDB_STATS_LOG_INTERVAL_SECONDS.name()
2484 || name == upsert_rocksdb::UPSERT_ROCKSDB_STATS_PERSIST_INTERVAL_SECONDS.name()
2485 || name == upsert_rocksdb::UPSERT_ROCKSDB_POINT_LOOKUP_BLOCK_CACHE_SIZE_MB.name()
2486 || name == upsert_rocksdb::UPSERT_ROCKSDB_SHRINK_ALLOCATED_BUFFERS_BY_RATIO.name()
2487}
2488
2489pub fn is_timestamp_oracle_config_var(name: &str) -> bool {
2492 name == PG_TIMESTAMP_ORACLE_CONNECTION_POOL_MAX_SIZE.name()
2493 || name == PG_TIMESTAMP_ORACLE_CONNECTION_POOL_MAX_WAIT.name()
2494 || name == PG_TIMESTAMP_ORACLE_CONNECTION_POOL_TTL.name()
2495 || name == PG_TIMESTAMP_ORACLE_CONNECTION_POOL_TTL_STAGGER.name()
2496 || name == CRDB_CONNECT_TIMEOUT.name()
2497 || name == CRDB_TCP_USER_TIMEOUT.name()
2498 || name == CRDB_KEEPALIVES_IDLE.name()
2499 || name == CRDB_KEEPALIVES_INTERVAL.name()
2500 || name == CRDB_KEEPALIVES_RETRIES.name()
2501 || name == mz_adapter_types::dyncfgs::PG_TIMESTAMP_ORACLE_STATEMENT_TIMEOUT.name()
2502}
2503
2504pub fn is_cluster_scheduling_var(name: &str) -> bool {
2506 name == cluster_scheduling::CLUSTER_MULTI_PROCESS_REPLICA_AZ_AFFINITY_WEIGHT.name()
2507 || name == cluster_scheduling::CLUSTER_SOFTEN_REPLICATION_ANTI_AFFINITY.name()
2508 || name == cluster_scheduling::CLUSTER_SOFTEN_REPLICATION_ANTI_AFFINITY_WEIGHT.name()
2509 || name == cluster_scheduling::CLUSTER_ENABLE_TOPOLOGY_SPREAD.name()
2510 || name == cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_IGNORE_NON_SINGULAR_SCALE.name()
2511 || name == cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_MAX_SKEW.name()
2512 || name == cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_MIN_DOMAINS.name()
2513 || name == cluster_scheduling::CLUSTER_TOPOLOGY_SPREAD_SOFT.name()
2514 || name == cluster_scheduling::CLUSTER_SOFTEN_AZ_AFFINITY.name()
2515 || name == cluster_scheduling::CLUSTER_SOFTEN_AZ_AFFINITY_WEIGHT.name()
2516}
2517
2518pub fn is_http_config_var(name: &str) -> bool {
2520 name == WEBHOOK_CONCURRENT_REQUEST_LIMIT.name()
2521}
2522
2523static SESSION_SYSTEM_VARS: LazyLock<BTreeMap<&'static UncasedStr, &'static VarDefinition>> =
2527 LazyLock::new(|| {
2528 [
2529 &APPLICATION_NAME,
2530 &CLIENT_ENCODING,
2531 &CLIENT_MIN_MESSAGES,
2532 &CLUSTER,
2533 &CLUSTER_REPLICA,
2534 &DEFAULT_CLUSTER_REPLICATION_FACTOR,
2535 &CURRENT_OBJECT_MISSING_WARNINGS,
2536 &DATABASE,
2537 &DATE_STYLE,
2538 &EXTRA_FLOAT_DIGITS,
2539 &INTEGER_DATETIMES,
2540 &INTERVAL_STYLE,
2541 &REAL_TIME_RECENCY_TIMEOUT,
2542 &SEARCH_PATH,
2543 &STANDARD_CONFORMING_STRINGS,
2544 &STATEMENT_TIMEOUT,
2545 &IDLE_IN_TRANSACTION_SESSION_TIMEOUT,
2546 &TIMEZONE,
2547 &TRANSACTION_ISOLATION,
2548 &MAX_QUERY_RESULT_SIZE,
2549 ]
2550 .into_iter()
2551 .map(|var| (UncasedStr::new(var.name()), var))
2552 .collect()
2553 });
2554
2555#[derive(Debug)]
2558pub struct FeatureFlag {
2559 pub flag: &'static VarDefinition,
2560 pub feature_desc: &'static str,
2561}
2562
2563impl FeatureFlag {
2564 pub fn enabled(&'static self, system_vars: &SystemVars) -> bool {
2566 *system_vars.expect_value::<bool>(self.flag)
2567 }
2568
2569 pub fn require(&'static self, system_vars: &SystemVars) -> Result<(), VarError> {
2572 match self.enabled(system_vars) {
2573 true => Ok(()),
2574 false => Err(VarError::RequiresFeatureFlag { feature_flag: self }),
2575 }
2576 }
2577}
2578
2579impl PartialEq for FeatureFlag {
2580 fn eq(&self, other: &FeatureFlag) -> bool {
2581 self.flag.name() == other.flag.name()
2582 }
2583}
2584
2585impl Eq for FeatureFlag {}
2586
2587impl Var for MzVersion {
2588 fn name(&self) -> &'static str {
2589 MZ_VERSION_NAME.as_str()
2590 }
2591
2592 fn value(&self) -> String {
2593 self.build_info
2594 .human_version(self.helm_chart_version.clone())
2595 }
2596
2597 fn description(&self) -> &'static str {
2598 "Shows the Materialize server version (Materialize)."
2599 }
2600
2601 fn type_name(&self) -> Cow<'static, str> {
2602 String::type_name()
2603 }
2604
2605 fn visible(&self, _: &User, _: &SystemVars) -> Result<(), VarError> {
2606 Ok(())
2607 }
2608}
2609
2610impl Var for User {
2611 fn name(&self) -> &'static str {
2612 IS_SUPERUSER_NAME.as_str()
2613 }
2614
2615 fn value(&self) -> String {
2616 self.is_superuser().format()
2617 }
2618
2619 fn description(&self) -> &'static str {
2620 "Reports whether the current session is a superuser (PostgreSQL)."
2621 }
2622
2623 fn type_name(&self) -> Cow<'static, str> {
2624 bool::type_name()
2625 }
2626
2627 fn visible(&self, _: &User, _: &SystemVars) -> Result<(), VarError> {
2628 Ok(())
2629 }
2630}
2631
2632#[cfg(test)]
2633mod isolation_feature_flag_tests {
2634 use super::*;
2635
2636 #[mz_ore::test]
2637 fn gates_strong_session_serializable_value() {
2638 let mut system_vars = SystemVars::new();
2639
2640 for name in ["transaction_isolation", "TRANSACTION_ISOLATION"] {
2646 let err = check_transaction_isolation_feature_flag(
2647 name,
2648 VarInput::Flat("strong session serializable"),
2649 &system_vars,
2650 )
2651 .expect_err("flag off rejects strong session serializable");
2652 assert!(matches!(err, VarError::RequiresFeatureFlag { .. }));
2653 }
2654
2655 for level in ["serializable", "bounded staleness 5s"] {
2657 check_transaction_isolation_feature_flag(
2658 TRANSACTION_ISOLATION_VAR_NAME,
2659 VarInput::Flat(level),
2660 &system_vars,
2661 )
2662 .expect("ungated level always allowed");
2663 }
2664
2665 system_vars
2667 .set("enable_session_timelines", VarInput::Flat("on"))
2668 .expect("set flag");
2669 check_transaction_isolation_feature_flag(
2670 TRANSACTION_ISOLATION_VAR_NAME,
2671 VarInput::Flat("strong session serializable"),
2672 &system_vars,
2673 )
2674 .expect("flag on");
2675
2676 check_transaction_isolation_feature_flag(
2678 CLUSTER.name(),
2679 VarInput::Flat("strong session serializable"),
2680 &system_vars,
2681 )
2682 .expect("unrelated var ignored");
2683 }
2684}
2685
2686#[cfg(test)]
2687mod reset_all_tests {
2688 use super::*;
2689 use crate::session::user::SYSTEM_USER;
2690
2691 fn test_vars() -> SessionVars {
2692 SessionVars::new_unchecked(&mz_build_info::DUMMY_BUILD_INFO, SYSTEM_USER.clone(), None)
2693 }
2694
2695 #[mz_ore::test]
2699 fn reset_all_clears_committed_session_value() {
2700 let system_vars = SystemVars::new();
2701 let mut vars = test_vars();
2702 let default = vars.application_name().to_string();
2703
2704 vars.set(
2706 &system_vars,
2707 "application_name",
2708 VarInput::Flat("custom"),
2709 false,
2710 )
2711 .expect("set");
2712 vars.end_transaction(EndTransactionAction::Commit);
2713 assert_eq!(vars.application_name(), "custom");
2714 assert_eq!(
2715 vars.inspect("application_name")
2716 .unwrap()
2717 .inspect_session_value()
2718 .map(|v| v.format()),
2719 Some("custom".to_string())
2720 );
2721
2722 let changed = vars.reset_all();
2723 assert_eq!(
2724 changed,
2725 BTreeMap::from([("application_name", default.clone())])
2726 );
2727
2728 let var = vars.inspect("application_name").unwrap();
2731 assert_eq!(vars.application_name(), default);
2732 assert_eq!(var.inspect_session_value(), None);
2733 assert!(!var.is_mutating());
2734 }
2735
2736 #[mz_ore::test]
2740 fn reset_all_preserves_installed_default() {
2741 let system_vars = SystemVars::new();
2742 let mut vars = test_vars();
2743
2744 vars.set_default("application_name", VarInput::Flat("startup_default"))
2745 .expect("set_default");
2746 vars.set(
2747 &system_vars,
2748 "application_name",
2749 VarInput::Flat("custom"),
2750 false,
2751 )
2752 .expect("set");
2753 vars.end_transaction(EndTransactionAction::Commit);
2754 assert_eq!(vars.application_name(), "custom");
2755
2756 let changed = vars.reset_all();
2757 assert_eq!(
2758 changed,
2759 BTreeMap::from([("application_name", "startup_default".to_string())])
2760 );
2761
2762 assert_eq!(vars.application_name(), "startup_default");
2763 assert_eq!(
2764 vars.inspect("application_name")
2765 .unwrap()
2766 .inspect_session_value(),
2767 None
2768 );
2769 }
2770}