Skip to main content

mz_sql/session/
vars.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Run-time configuration parameters
11//!
12//! ## Overview
13//! Materialize roughly follows the PostgreSQL configuration model, which works
14//! as follows. There is a global set of named configuration parameters, like
15//! `DateStyle` and `client_encoding`. These parameters can be set in several
16//! places: in an on-disk configuration file (in Postgres, named
17//! postgresql.conf), in command line arguments when the server is started, or
18//! at runtime via the `ALTER SYSTEM` or `SET` statements. Parameters that are
19//! set in a session take precedence over database defaults, which in turn take
20//! precedence over command line arguments, which in turn take precedence over
21//! settings in the on-disk configuration. Note that changing the value of
22//! parameters obeys transaction semantics: if a transaction fails to commit,
23//! any parameters that were changed in that transaction (i.e., via `SET`) will
24//! be rolled back to their previous value.
25//!
26//! The Materialize configuration hierarchy at the moment is much simpler.
27//! Global defaults are hardcoded into the binary, and a select few parameters
28//! can be overridden per session. A select few parameters can be overridden on
29//! disk.
30//!
31//! The set of variables that can be overridden per session and the set of
32//! variables that can be overridden on disk are currently disjoint. The
33//! infrastructure has been designed with an eye towards merging these two sets
34//! and supporting additional layers to the hierarchy, however, should the need
35//! arise.
36//!
37//! The configuration parameters that exist are driven by compatibility with
38//! PostgreSQL drivers that expect them, not because they are particularly
39//! important.
40//!
41//! ## Structure
42//! The most meaningful exports from this module are:
43//!
44//! - [`SessionVars`] represent per-session parameters, which each user can
45//!   access independently of one another, and are accessed via `SET`.
46//!
47//!   The fields of [`SessionVars`] are either;
48//!     - `SessionVar`, which is preferable and simply requires full support of
49//!       the `SessionVar` impl for its embedded value type.
50//!     - `ServerVar` for types that do not currently support everything
51//!       required by `SessionVar`, e.g. they are fixed-value parameters.
52//!
53//!   In the fullness of time, all fields in [`SessionVars`] should be
54//!   `SessionVar`.
55//!
56//! - [`SystemVars`] represent system-wide configuration settings and are
57//!   accessed via `ALTER SYSTEM SET`.
58//!
59//!   All elements of [`SystemVars`] are `SystemVar`.
60//!
61//! Some [`VarDefinition`] are also marked as a [`FeatureFlag`]; this is just a
62//! wrapper to make working with a set of [`VarDefinition`] easier, primarily from
63//! within SQL planning, where we might want to check if a feature is enabled
64//! before planning it.
65
66use 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/// The action to take during end_transaction.
110///
111/// This enum lives here because of convenience: it's more of an adapter
112/// concept but [`SessionVars::end_transaction`] takes it.
113#[derive(Debug, Clone, Copy, PartialEq, Eq)]
114pub enum EndTransactionAction {
115    /// Commit the transaction.
116    Commit,
117    /// Rollback the transaction.
118    Rollback,
119}
120
121/// Represents the input to a variable.
122///
123/// Each variable has different rules for how it handles each style of input.
124/// This type allows us to defer interpretation of the input until the
125/// variable-specific interpretation can be applied.
126#[derive(Debug, Clone, Copy)]
127pub enum VarInput<'a> {
128    /// The input has been flattened into a single string.
129    ///
130    /// NOTE: when adding a new variant here (or in [`OwnedVarInput`]), extend
131    /// the `mz_catalog.mz_role_parameters` materialized view in
132    /// `src/catalog/src/builtin/mz_catalog.rs`. That MV discriminates on the
133    /// externally-tagged JSON shape of [`OwnedVarInput`] to format
134    /// `parameter_value`.
135    Flat(&'a str),
136    /// The input comes from a SQL `SET` statement and is jumbled across
137    /// multiple components.
138    ///
139    /// NOTE: see the doc-comment on [`VarInput::Flat`] — adding a new variant
140    /// requires extending `mz_catalog.mz_role_parameters`.
141    SqlSet(&'a [String]),
142}
143
144impl<'a> VarInput<'a> {
145    /// Converts the variable input to an owned vector of strings.
146    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/// An owned version of [`VarInput`].
155#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize)]
156pub enum OwnedVarInput {
157    /// See [`VarInput::Flat`].
158    ///
159    /// NOTE: adding a new variant requires extending the
160    /// `mz_catalog.mz_role_parameters` materialized view in
161    /// `src/catalog/src/builtin/mz_catalog.rs`, which discriminates on the
162    /// externally-tagged JSON shape of this enum.
163    Flat(String),
164    /// See [`VarInput::SqlSet`].
165    ///
166    /// NOTE: see the doc-comment on [`OwnedVarInput::Flat`].
167    SqlSet(Vec<String>),
168}
169
170impl OwnedVarInput {
171    /// Converts this owned variable input as a [`VarInput`].
172    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
180/// A `Var` represents a configuration parameter of an arbitrary type.
181pub trait Var: Debug {
182    /// Returns the name of the configuration parameter.
183    fn name(&self) -> &'static str;
184
185    /// Constructs a flattened string representation of the current value of the
186    /// configuration parameter.
187    ///
188    /// The resulting string is guaranteed to be parsable if provided to
189    /// `Value::parse` as a [`VarInput::Flat`].
190    fn value(&self) -> String;
191
192    /// Returns a short sentence describing the purpose of the configuration
193    /// parameter.
194    fn description(&self) -> &'static str;
195
196    /// Returns the name of the type of this variable.
197    fn type_name(&self) -> Cow<'static, str>;
198
199    /// Indicates wither the [`Var`] is visible as a function of the `user` and `system_vars`.
200    /// "Invisible" parameters return `VarErrors`.
201    ///
202    /// Variables marked as `internal` are only visible for the system user.
203    fn visible(&self, user: &User, system_vars: &SystemVars) -> Result<(), VarError>;
204
205    /// Reports whether the variable is only visible in unsafe mode.
206    fn is_unsafe(&self) -> bool {
207        self.name().starts_with("unsafe_")
208    }
209
210    /// Returns the [`ParameterScope`] at which this variable's value may be
211    /// overridden by the LaunchDarkly sync loop. Defaults to
212    /// [`ParameterScope::Environment`].
213    fn scope(&self) -> ParameterScope {
214        ParameterScope::Environment
215    }
216
217    /// Upcast `self` to a `dyn Var`, useful when working with multiple different implementors of
218    /// [`Var`].
219    fn as_var(&self) -> &dyn Var
220    where
221        Self: Sized,
222    {
223        self
224    }
225}
226
227/// A `SessionVar` is the session value for a configuration parameter. If unset,
228/// the server default is used instead.
229///
230/// Note: even though all of the different `*_value` fields are `Box<dyn Value>` they are enforced
231/// to be the same type because we use the `definition`s `parse(...)` method. This is guaranteed to
232/// return the same type as the compiled in default.
233#[derive(Debug)]
234pub struct SessionVar {
235    definition: VarDefinition,
236    /// System or Role default value.
237    default_value: Option<Box<dyn Value>>,
238    /// Value `LOCAL` to a transaction, will be unset at the completion of the transaction.
239    local_value: Option<Box<dyn Value>>,
240    /// Value set during a transaction, will be set if the transaction is committed.
241    staged_value: Option<Box<dyn Value>>,
242    /// Value that overrides the default.
243    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    /// Checks if the provided [`VarInput`] is valid for the current session variable, returning
270    /// the formatted output if it's valid.
271    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    /// Parse the input and update the stored value to match.
279    pub fn set(&mut self, input: VarInput, local: bool) -> Result<(), VarError> {
280        let v = self.definition.parse(input)?;
281
282        // Validate our parsed value.
283        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    /// Sets the default value for the variable.
295    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    /// Reset the stored value to the default.
303    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    /// Resets the variable to its default, discarding any session, staged, or
318    /// local override. Unlike [`SessionVar::reset`], which stages the reset for
319    /// the current transaction, this takes effect immediately and does not
320    /// depend on a later [`SessionVar::end_transaction`] commit.
321    ///
322    /// `DISCARD ALL` requires this: it ends the transaction before resetting, so
323    /// there is no commit left to promote a staged value. `default_value` (the
324    /// system/role/startup default) is deliberately preserved.
325    pub fn reset_durable(&mut self) {
326        self.local_value = None;
327        self.staged_value = None;
328        self.session_value = None;
329    }
330
331    /// Returns a possibly new SessionVar if this needs to mutate at transaction end.
332    #[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    /// Whether this Var needs to mutate at the end of a transaction.
349    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    /// Returns the [`Value`] that is currently stored as the `session_value`.
363    ///
364    /// Note: This should __only__ be used for inspection, if you want to determine the current
365    /// value of this [`SessionVar`] you should use [`SessionVar::value`].
366    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    /// Inputs to computed variables.
408    build_info: &'static BuildInfo,
409    /// Helm chart version
410    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/// Session variables.
423///
424/// See the [`crate::session::vars`] module documentation for more details on the
425/// Materialize configuration model.
426#[derive(Debug, Clone)]
427pub struct SessionVars {
428    /// The set of all session variables.
429    vars: OrdMap<&'static UncasedStr, SessionVar>,
430    /// Inputs to computed variables.
431    mz_version: MzVersion,
432    /// Information about the user associated with this Session.
433    user: User,
434}
435
436impl SessionVars {
437    /// Creates a new [`SessionVars`] without considering the System or Role defaults.
438    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    /// Returns an iterator over the configuration parameters and their current
486    /// values for this session.
487    ///
488    /// Note that this function does not check that the access variable should
489    /// be visible because of other settings or users. Before or after accessing
490    /// this method, you should call `Var::visible`.
491    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    /// Returns an iterator over configuration parameters (and their current
500    /// values for this session) that are expected to be sent to the client when
501    /// a new connection is established or when their value changes.
502    pub fn notify_set(&self) -> impl Iterator<Item = &dyn Var> {
503        // WARNING: variables in this set are not checked for visibility, and
504        // are assumed to be visible for all sessions.
505        //
506        // This is fixible with some elbow grease, but at the moment it seems
507        // unlikely that we'll have a variable in the notify set that shouldn't
508        // be visible to all sessions.
509        [
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            // Including `cluster`, `cluster_replica`, `database`, and `search_path` in the notify
519            // set is a Materialize extension. Doing so allows users to more easily identify where
520            // their queries will be executing, which is important to know when you consider the
521            // size of a cluster, what indexes are present, etc.
522            &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        // Including `mz_version` in the notify set is a Materialize
531        // extension. Doing so allows applications to detect whether they
532        // are talking to Materialize or PostgreSQL without an additional
533        // network roundtrip. This is known to be safe because CockroachDB
534        // has an analogous extension [0].
535        // [0]: https://github.com/cockroachdb/cockroach/blob/369c4057a/pkg/sql/pgwire/conn.go#L1840
536        .chain(std::iter::once(self.mz_version.as_var()))
537    }
538
539    /// Durably resets all variables to their default value.
540    ///
541    /// Unlike a staged [`SessionVar::reset`], this takes effect immediately and
542    /// does not depend on a later transaction commit. Used by `DISCARD ALL`,
543    /// which ends the transaction before resetting, so there is no commit left
544    /// to promote a staged value. System/role/startup defaults are preserved.
545    ///
546    /// Returns the parameters whose value changed, with their new values.
547    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    /// Returns a [`Var`] representing the configuration parameter with the
563    /// specified name.
564    ///
565    /// Configuration parameters are matched case insensitively. If no such
566    /// configuration parameter exists, `get` returns an error.
567    ///
568    /// Note that if `name` is known at compile time, you should instead use the
569    /// named accessor to access the variable with its true Rust type. For
570    /// example, `self.get("sql_safe_updates").value()` returns the string
571    /// `"true"` or `"false"`, while `self.sql_safe_updates()` returns a bool.
572    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    /// Returns a [`SessionVar`] for inspection.
593    ///
594    /// Note: If you're trying to determine the value of the variable with `name` you should
595    /// instead use the named accessor, or [`SessionVars::get`].
596    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    /// Sets the configuration parameter named `name` to the value represented
605    /// by `value`.
606    ///
607    /// The new value may be either committed or rolled back by the next call to
608    /// [`SessionVars::end_transaction`]. If `local` is true, the new value is always
609    /// discarded by the next call to [`SessionVars::end_transaction`], even if the
610    /// transaction is marked to commit.
611    ///
612    /// Like with [`SessionVars::get`], configuration parameters are matched case
613    /// insensitively. If `value` is not valid, as determined by the underlying
614    /// configuration parameter, or if the named configuration parameter does
615    /// not exist, an error is returned.
616    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    /// Sets the default value for the parameter named `name` to the value
641    /// represented by `value`.
642    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        // Check if this variable is allowed to be set as a role default.
648        // Most read-only variables are blocked, but some (like restrict_to_user_objects)
649        // are specifically designed to be set via ALTER ROLE by superusers.
650        if !Self::allow_role_default(name) {
651            self.check_read_only(name)?;
652        }
653
654        self.vars
655            .get_mut(name)
656            // Note: visibility is checked when persisting a role default.
657            .map(|v| v.set_default(input))
658            .transpose()?
659            .ok_or_else(|| VarError::UnknownParameter(name.to_string()))
660    }
661
662    /// Returns true if the variable can be set as a role default even if it's
663    /// otherwise read-only from direct SET commands.
664    ///
665    /// SECURITY: Any variable listed here must also have a corresponding
666    /// superuser RBAC check in `generate_rbac_requirements` in `rbac.rs`
667    /// (see the `PlannedAlterRoleOption::Variable` match arm). Without that
668    /// check, any role could set the variable on themselves via ALTER ROLE.
669    fn allow_role_default(name: &UncasedStr) -> bool {
670        name == RESTRICT_TO_USER_OBJECTS.name
671    }
672
673    /// Sets the configuration parameter named `name` to its default value.
674    ///
675    /// The new value may be either committed or rolled back by the next call to
676    /// [`SessionVars::end_transaction`]. If `local` is true, the new value is
677    /// always discarded by the next call to [`SessionVars::end_transaction`],
678    /// even if the transaction is marked to commit.
679    ///
680    /// Like with [`SessionVars::get`], configuration parameters are matched
681    /// case insensitively. If the named configuration parameter does not exist,
682    /// an error is returned.
683    ///
684    /// If the variable does not exist or the user does not have the visibility
685    /// requires, this function returns an error.
686    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    /// Returns an error if the variable corresponding to `name` is read only.
709    ///
710    /// Note: This is called by `set()` (for SQL SET commands) but NOT by
711    /// `set_default()` (for role defaults). This allows variables like
712    /// `restrict_to_user_objects` to be set via `ALTER ROLE ... SET` by
713    /// superusers while blocking direct `SET` commands from regular users.
714    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            // This variable can only be set via ALTER ROLE ... SET by superusers,
725            // not via direct SET commands. This prevents malicious queries from
726            // bypassing the restriction.
727            Err(VarError::ReadOnlyParameter(
728                RESTRICT_TO_USER_OBJECTS.name.as_str(),
729            ))
730        } else {
731            Ok(())
732        }
733    }
734
735    /// Commits or rolls back configuration parameter updates made via
736    /// [`SessionVars::set`] since the last call to `end_transaction`.
737    ///
738    /// Returns any session parameters that changed because the transaction ended.
739    #[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            // Report the new value of the parameter.
756            if before != after {
757                changed.insert(var.name(), after);
758            }
759        }
760        self.vars.extend(updates);
761        changed
762    }
763
764    /// Returns the value of the `application_name` configuration parameter.
765    pub fn application_name(&self) -> &str {
766        self.expect_value::<String>(&APPLICATION_NAME).as_str()
767    }
768
769    /// Returns the build info.
770    pub fn build_info(&self) -> &'static BuildInfo {
771        self.mz_version.build_info
772    }
773
774    /// Returns the value of the `client_encoding` configuration parameter.
775    pub fn client_encoding(&self) -> &ClientEncoding {
776        self.expect_value(&CLIENT_ENCODING)
777    }
778
779    /// Returns the value of the `client_min_messages` configuration parameter.
780    pub fn client_min_messages(&self) -> &ClientSeverity {
781        self.expect_value(&CLIENT_MIN_MESSAGES)
782    }
783
784    /// Returns the value of the `cluster` configuration parameter.
785    pub fn cluster(&self) -> &str {
786        self.expect_value::<String>(&CLUSTER).as_str()
787    }
788
789    /// Returns the value of the `cluster_replica` configuration parameter.
790    pub fn cluster_replica(&self) -> Option<&str> {
791        self.expect_value::<Option<String>>(&CLUSTER_REPLICA)
792            .as_deref()
793    }
794
795    /// Returns the value of the `current_object_missing_warnings` configuration
796    /// parameter.
797    pub fn current_object_missing_warnings(&self) -> bool {
798        *self.expect_value::<bool>(&CURRENT_OBJECT_MISSING_WARNINGS)
799    }
800
801    /// Returns the value of the `DateStyle` configuration parameter.
802    pub fn date_style(&self) -> &[&str] {
803        &self.expect_value::<DateStyle>(&DATE_STYLE).0
804    }
805
806    /// Returns the value of the `database` configuration parameter.
807    pub fn database(&self) -> &str {
808        self.expect_value::<String>(&DATABASE).as_str()
809    }
810
811    /// Returns the value of the `extra_float_digits` configuration parameter.
812    pub fn extra_float_digits(&self) -> i32 {
813        *self.expect_value(&EXTRA_FLOAT_DIGITS)
814    }
815
816    /// Returns the settings that govern how this session encodes values as
817    /// text.
818    pub fn text_encode_settings(&self) -> TextEncodeSettings {
819        TextEncodeSettings {
820            extra_float_digits: self.extra_float_digits(),
821        }
822    }
823
824    /// Returns the value of the `integer_datetimes` configuration parameter.
825    pub fn integer_datetimes(&self) -> bool {
826        *self.expect_value(&INTEGER_DATETIMES)
827    }
828
829    /// Returns the value of the `intervalstyle` configuration parameter.
830    pub fn intervalstyle(&self) -> &IntervalStyle {
831        self.expect_value(&INTERVAL_STYLE)
832    }
833
834    /// Returns the value of the `mz_version` configuration parameter.
835    pub fn mz_version(&self) -> String {
836        self.mz_version.value()
837    }
838
839    /// Returns the value of the `search_path` configuration parameter.
840    pub fn search_path(&self) -> &[Ident] {
841        self.expect_value::<Vec<Ident>>(&SEARCH_PATH).as_slice()
842    }
843
844    /// Returns the value of the `server_version` configuration parameter.
845    pub fn server_version(&self) -> &str {
846        self.expect_value::<String>(&SERVER_VERSION).as_str()
847    }
848
849    /// Returns the value of the `server_version_num` configuration parameter.
850    pub fn server_version_num(&self) -> i32 {
851        *self.expect_value(&SERVER_VERSION_NUM)
852    }
853
854    /// Returns the value of the `sql_safe_updates` configuration parameter.
855    pub fn sql_safe_updates(&self) -> bool {
856        *self.expect_value(&SQL_SAFE_UPDATES)
857    }
858
859    /// Returns the value of the `standard_conforming_strings` configuration
860    /// parameter.
861    pub fn standard_conforming_strings(&self) -> bool {
862        *self.expect_value(&STANDARD_CONFORMING_STRINGS)
863    }
864
865    /// Returns the value of the `statement_timeout` configuration parameter.
866    pub fn statement_timeout(&self) -> &Duration {
867        self.expect_value(&STATEMENT_TIMEOUT)
868    }
869
870    /// Returns the value of the `idle_in_transaction_session_timeout` configuration parameter.
871    pub fn idle_in_transaction_session_timeout(&self) -> &Duration {
872        self.expect_value(&IDLE_IN_TRANSACTION_SESSION_TIMEOUT)
873    }
874
875    /// Returns the value of the `timezone` configuration parameter.
876    pub fn timezone(&self) -> &TimeZone {
877        self.expect_value(&TIMEZONE)
878    }
879
880    /// Returns the value of the `transaction_isolation` configuration
881    /// parameter.
882    pub fn transaction_isolation(&self) -> &IsolationLevel {
883        self.expect_value(&TRANSACTION_ISOLATION)
884    }
885
886    /// Returns the value of `real_time_recency` configuration parameter.
887    pub fn real_time_recency(&self) -> bool {
888        *self.expect_value(&REAL_TIME_RECENCY)
889    }
890
891    /// Returns the value of the `real_time_recency_timeout` configuration parameter.
892    pub fn real_time_recency_timeout(&self) -> &Duration {
893        self.expect_value(&REAL_TIME_RECENCY_TIMEOUT)
894    }
895
896    /// Returns the value of `emit_plan_insights_notice` configuration parameter.
897    pub fn emit_plan_insights_notice(&self) -> bool {
898        *self.expect_value(&EMIT_PLAN_INSIGHTS_NOTICE)
899    }
900
901    /// Returns the value of `emit_timestamp_notice` configuration parameter.
902    pub fn emit_timestamp_notice(&self) -> bool {
903        *self.expect_value(&EMIT_TIMESTAMP_NOTICE)
904    }
905
906    /// Returns the value of `emit_trace_id_notice` configuration parameter.
907    pub fn emit_trace_id_notice(&self) -> bool {
908        *self.expect_value(&EMIT_TRACE_ID_NOTICE)
909    }
910
911    /// Returns the value of `auto_route_catalog_queries` configuration parameter.
912    pub fn auto_route_catalog_queries(&self) -> bool {
913        *self.expect_value(&AUTO_ROUTE_CATALOG_QUERIES)
914    }
915
916    /// Returns the value of `enable_session_rbac_checks` configuration parameter.
917    pub fn enable_session_rbac_checks(&self) -> bool {
918        *self.expect_value(&ENABLE_SESSION_RBAC_CHECKS)
919    }
920
921    /// Returns the value of `restrict_to_user_objects` configuration parameter.
922    pub fn restrict_to_user_objects(&self) -> bool {
923        *self.expect_value(&RESTRICT_TO_USER_OBJECTS)
924    }
925
926    /// Returns the value of `enable_session_cardinality_estimates` configuration parameter.
927    pub fn enable_session_cardinality_estimates(&self) -> bool {
928        *self.expect_value(&ENABLE_SESSION_CARDINALITY_ESTIMATES)
929    }
930
931    /// Returns the value of `is_superuser` configuration parameter.
932    pub fn is_superuser(&self) -> bool {
933        self.user.is_superuser()
934    }
935
936    /// Returns the user associated with this `SessionVars` instance.
937    pub fn user(&self) -> &User {
938        &self.user
939    }
940
941    /// Returns the value of the `max_query_result_size` configuration parameter.
942    pub fn max_query_result_size(&self) -> u64 {
943        self.expect_value::<ByteSize>(&MAX_QUERY_RESULT_SIZE)
944            .as_bytes()
945    }
946
947    /// Sets the internal metadata associated with the user.
948    pub fn set_internal_user_metadata(&mut self, metadata: InternalUserMetadata) {
949        self.user.internal_metadata = Some(metadata);
950    }
951
952    /// Sets the external metadata associated with the user.
953    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    /// Returns the value of the `emit_introspection_query_notice` configuration parameter.
980    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    /// Returns the value of the `welcome_message` configuration parameter.
989    pub fn welcome_message(&self) -> bool {
990        *self.expect_value(&WELCOME_MESSAGE)
991    }
992}
993
994// TODO(database-issues#8069) remove together with `compat_translate`
995pub const OLD_CATALOG_SERVER_CLUSTER: &str = "mz_introspection";
996pub const OLD_AUTO_ROUTE_CATALOG_QUERIES: &str = "auto_route_introspection_queries";
997
998/// If the given variable name and/or input is deprecated, return a corresponding updated value,
999/// otherwise return the original.
1000///
1001/// This method was introduced to gracefully handle the rename of the `mz_introspection` cluster to
1002/// `mz_cluster_server`. The plan is to remove it once all users have migrated to the new name. The
1003/// debug logs will be helpful for checking this in production.
1004// TODO(database-issues#8069) remove this after sufficient time has passed
1005fn 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
1036/// Enforces feature-flag gating for `transaction_isolation` levels that sit
1037/// behind a flag (`strong session serializable`).
1038///
1039/// Returns `Ok(())` for any other variable, and for an unparseable value
1040/// (parse errors surface on the actual set). This is shared by every path that
1041/// assigns `transaction_isolation` — `SET`, `SET TRANSACTION`,
1042/// `ALTER ROLE ... SET`, and connection options — so that the gate cannot be
1043/// bypassed by choosing a different syntax or letter case.
1044pub 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    // Ignore parse failures here; the actual set surfaces them.
1053    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/// A `SystemVar` is persisted on disk value for a configuration parameter. If unset,
1063/// the server default is used instead.
1064#[derive(Debug)]
1065pub struct SystemVar {
1066    definition: VarDefinition,
1067    /// Value currently persisted to disk.
1068    persisted_value: Option<Box<dyn Value>>,
1069    /// Current default, not persisted to disk.
1070    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        // Validate our parsed value.
1112        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/// On disk variables.
1184///
1185/// See the [`crate::session::vars`] module documentation for more details on the
1186/// Materialize configuration model.
1187#[derive(Derivative, Clone)]
1188#[derivative(Debug)]
1189pub struct SystemVars {
1190    /// Allows "unsafe" parameters to be set.
1191    allow_unsafe: bool,
1192    /// Set of all [`SystemVar`]s.
1193    vars: BTreeMap<&'static UncasedStr, SystemVar>,
1194    /// External components interested in when a [`SystemVar`] gets updated.
1195    #[derivative(Debug = "ignore")]
1196    callbacks: BTreeMap<String, Vec<Arc<dyn Fn(&SystemVars) + Send + Sync>>>,
1197
1198    /// NB: This is intentionally disconnected from the one that is plumbed around to persist and
1199    /// the controllers. This is so we can explicitly control and reason about when changes to config
1200    /// values are propagated to the rest of the system.
1201    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                // Carry the dyncfg's declared scope through to the system var,
1380                // so scoped resolution and introspection see it.
1381                var.scoped(cfg.scope())
1382            })
1383            .collect();
1384
1385        let vars: BTreeMap<_, _> = system_vars
1386            .into_iter()
1387            // Include all of our feature flags.
1388            .chain(definitions::FEATURE_FLAGS.iter().copied())
1389            // Include the subset of Session variables we allow system defaults for.
1390            .chain(SESSION_SYSTEM_VARS.values().copied())
1391            .cloned()
1392            // Include Persist configs.
1393            .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    /// Returns an iterator over the configuration parameters and their current
1445    /// values on disk.
1446    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    /// Returns an iterator over the configuration parameters and their current
1454    /// values on disk. Compared to [`SystemVars::iter`], this should omit vars
1455    /// that shouldn't be synced by SystemParameterFrontend.
1456    pub fn iter_synced(&self) -> impl Iterator<Item = &dyn Var> {
1457        self.iter().filter(|v| v.name() != ENABLE_LAUNCHDARKLY.name)
1458    }
1459
1460    /// Returns an iterator over the configuration parameters that can be overriden per-Session.
1461    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    /// Returns whether or not this parameter can be modified by a superuser.
1469    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    /// Returns a [`Var`] representing the configuration parameter with the
1476    /// specified name.
1477    ///
1478    /// Configuration parameters are matched case insensitively. If no such
1479    /// configuration parameter exists, `get` returns an error.
1480    ///
1481    /// Note that:
1482    /// - If `name` is known at compile time, you should instead use the named
1483    /// accessor to access the variable with its true Rust type. For example,
1484    /// `self.get("max_tables").value()` returns the string `"25"` or the
1485    /// current value, while `self.max_tables()` returns an i32.
1486    ///
1487    /// - This function does not check that the access variable should be
1488    /// visible because of other settings or users. Before or after accessing
1489    /// this method, you should call `Var::visible`.
1490    ///
1491    /// # Errors
1492    ///
1493    /// The call will return an error:
1494    /// 1. If `name` does not refer to a valid [`SystemVars`] field.
1495    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    /// Check if the given `values` is the default value for the [`Var`]
1503    /// identified by `name`.
1504    ///
1505    /// Note that this function does not check that the access variable should
1506    /// be visible because of other settings or users. Before or after accessing
1507    /// this method, you should call `Var::visible`.
1508    ///
1509    /// # Errors
1510    ///
1511    /// The call will return an error:
1512    /// 1. If `name` does not refer to a valid [`SystemVars`] field.
1513    /// 2. If `values` does not represent a valid [`SystemVars`] value for
1514    ///    `name`.
1515    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    /// Sets the configuration parameter named `name` to the value represented
1523    /// by `input`.
1524    ///
1525    /// Like with [`SystemVars::get`], configuration parameters are matched case
1526    /// insensitively. If `input` is not valid, as determined by the underlying
1527    /// configuration parameter, or if the named configuration parameter does
1528    /// not exist, an error is returned.
1529    ///
1530    /// Return a `bool` value indicating whether the [`Var`] identified by
1531    /// `name` was modified by this call (it won't be if it already had the
1532    /// given `input`).
1533    ///
1534    /// Note that this function does not check that the access variable should
1535    /// be visible because of other settings or users. Before or after accessing
1536    /// this method, you should call `Var::visible`.
1537    ///
1538    /// # Errors
1539    ///
1540    /// The call will return an error:
1541    /// 1. If `name` does not refer to a valid [`SystemVars`] field.
1542    /// 2. If `input` does not represent a valid [`SystemVars`] value for
1543    ///    `name`.
1544    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    /// Parses the configuration parameter value represented by `input` named
1554    /// `name`.
1555    ///
1556    /// Like with [`SystemVars::get`], configuration parameters are matched case
1557    /// insensitively. If `input` is not valid, as determined by the underlying
1558    /// configuration parameter, or if the named configuration parameter does
1559    /// not exist, an error is returned.
1560    ///
1561    /// Return a `Box<dyn Value>` that is the result of parsing `input`.
1562    ///
1563    /// Note that this function does not check that the access variable should
1564    /// be visible because of other settings or users. Before or after accessing
1565    /// this method, you should call `Var::visible`.
1566    ///
1567    /// # Errors
1568    ///
1569    /// The call will return an error:
1570    /// 1. If `name` does not refer to a valid [`SystemVars`] field.
1571    /// 2. If `input` does not represent a valid [`SystemVars`] value for
1572    ///    `name`.
1573    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    /// Set the default for this variable. This is the value this
1581    /// variable will be be `reset` to. If no default is set, the static default in the
1582    /// variable definition is used instead.
1583    ///
1584    /// Note that this function does not check that the access variable should
1585    /// be visible because of other settings or users. Before or after accessing
1586    /// this method, you should call `Var::visible`.
1587    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    /// Sets the configuration parameter named `name` to its default value.
1596    ///
1597    /// Like with [`SystemVars::get`], configuration parameters are matched case
1598    /// insensitively. If the named configuration parameter does not exist, an
1599    /// error is returned.
1600    ///
1601    /// Return a `bool` value indicating whether the [`Var`] identified by
1602    /// `name` was modified by this call (it won't be if was already reset).
1603    ///
1604    /// Note that this function does not check that the access variable should
1605    /// be visible because of other settings or users. Before or after accessing
1606    /// this method, you should call `Var::visible`.
1607    ///
1608    /// # Errors
1609    ///
1610    /// The call will return an error:
1611    /// 1. If `name` does not refer to a valid [`SystemVars`] field.
1612    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    /// Returns a map from each system parameter's name to its default value.
1622    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    /// Registers a closure that mirrors the value of the given
1636    /// [`VarDefinition`] into out-of-band state.
1637    ///
1638    /// The callback has to be an idempotent read of the passed [`SystemVars`],
1639    /// because we don't promise to only call it when its var actually changed.
1640    /// It runs once right now against the current values, and then again at
1641    /// every catalog commit boundary whose transaction touched a system var
1642    /// (see `Coordinator::apply_catalog_implications` and
1643    /// [`SystemVars::notify_all_callbacks`]). Speculative mutations never
1644    /// trigger it, so an aborted or dry-run transaction leaves the mirror
1645    /// untouched.
1646    ///
1647    /// NOTE: a callback on a `feature_flags!` var won't observe the transient
1648    /// flip that `CatalogState::with_enable_for_item_parsing` performs during
1649    /// item parsing. That flip mutates the value and then restores the prior
1650    /// `Arc` wholesale without re-notifying, so the mirror keeps tracking
1651    /// committed state throughout, which is the contract here. Committed changes
1652    /// to a feature flag (via `ALTER SYSTEM`) still notify like any other var.
1653    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    /// Re-runs every registered callback against the current values.
1666    ///
1667    /// This fires all of them, even ones whose var didn't change, which is why
1668    /// callbacks have to be idempotent reads of the passed [`SystemVars`]. See
1669    /// [`SystemVars::register_callback`].
1670    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    /// Notify any external components interested in this variable.
1679    fn notify_callbacks(&self, name: &str) {
1680        // Get the callbacks interested in this variable.
1681        if let Some(callbacks) = self.callbacks.get(name) {
1682            for callback in callbacks {
1683                (callback)(self);
1684            }
1685        }
1686    }
1687
1688    /// Returns the system default for the [`CLUSTER`] session variable. To know the active cluster
1689    /// for the current session, you must check the [`SessionVars`].
1690    pub fn default_cluster(&self) -> String {
1691        self.expect_value::<String>(&CLUSTER).to_owned()
1692    }
1693
1694    /// Returns the value of the `max_kafka_connections` configuration parameter.
1695    pub fn max_kafka_connections(&self) -> u32 {
1696        *self.expect_value(&MAX_KAFKA_CONNECTIONS)
1697    }
1698
1699    /// Returns the value of the `max_postgres_connections` configuration parameter.
1700    pub fn max_postgres_connections(&self) -> u32 {
1701        *self.expect_value(&MAX_POSTGRES_CONNECTIONS)
1702    }
1703
1704    /// Returns the value of the `max_mysql_connections` configuration parameter.
1705    pub fn max_mysql_connections(&self) -> u32 {
1706        *self.expect_value(&MAX_MYSQL_CONNECTIONS)
1707    }
1708
1709    /// Returns the value of the `max_sql_server_connections` configuration parameter.
1710    pub fn max_sql_server_connections(&self) -> u32 {
1711        *self.expect_value(&MAX_SQL_SERVER_CONNECTIONS)
1712    }
1713
1714    /// Returns the value of the `max_aws_privatelink_connections` configuration parameter.
1715    pub fn max_aws_privatelink_connections(&self) -> u32 {
1716        *self.expect_value(&MAX_AWS_PRIVATELINK_CONNECTIONS)
1717    }
1718
1719    /// Returns the value of the `max_tables` configuration parameter.
1720    pub fn max_tables(&self) -> u32 {
1721        *self.expect_value(&MAX_TABLES)
1722    }
1723
1724    /// Returns the value of the `max_sources` configuration parameter.
1725    pub fn max_sources(&self) -> u32 {
1726        *self.expect_value(&MAX_SOURCES)
1727    }
1728
1729    /// Returns the value of the `max_sinks` configuration parameter.
1730    pub fn max_sinks(&self) -> u32 {
1731        *self.expect_value(&MAX_SINKS)
1732    }
1733
1734    /// Returns the value of the `max_materialized_views` configuration parameter.
1735    pub fn max_materialized_views(&self) -> u32 {
1736        *self.expect_value(&MAX_MATERIALIZED_VIEWS)
1737    }
1738
1739    /// Returns the value of the `max_clusters` configuration parameter.
1740    pub fn max_clusters(&self) -> u32 {
1741        *self.expect_value(&MAX_CLUSTERS)
1742    }
1743
1744    /// Returns the value of the `max_replicas_per_cluster` configuration parameter.
1745    pub fn max_replicas_per_cluster(&self) -> u32 {
1746        *self.expect_value(&MAX_REPLICAS_PER_CLUSTER)
1747    }
1748
1749    /// Returns the value of the `max_credit_consumption_rate` configuration parameter.
1750    pub fn max_credit_consumption_rate(&self) -> Numeric {
1751        *self.expect_value(&MAX_CREDIT_CONSUMPTION_RATE)
1752    }
1753
1754    /// Returns the value of the `max_databases` configuration parameter.
1755    pub fn max_databases(&self) -> u32 {
1756        *self.expect_value(&MAX_DATABASES)
1757    }
1758
1759    /// Returns the value of the `max_schemas_per_database` configuration parameter.
1760    pub fn max_schemas_per_database(&self) -> u32 {
1761        *self.expect_value(&MAX_SCHEMAS_PER_DATABASE)
1762    }
1763
1764    /// Returns the value of the `max_objects_per_schema` configuration parameter.
1765    pub fn max_objects_per_schema(&self) -> u32 {
1766        *self.expect_value(&MAX_OBJECTS_PER_SCHEMA)
1767    }
1768
1769    /// Returns the value of the `max_secrets` configuration parameter.
1770    pub fn max_secrets(&self) -> u32 {
1771        *self.expect_value(&MAX_SECRETS)
1772    }
1773
1774    /// Returns the value of the `max_roles` configuration parameter.
1775    pub fn max_roles(&self) -> u32 {
1776        *self.expect_value(&MAX_ROLES)
1777    }
1778
1779    /// Returns the value of the `max_network_policies` configuration parameter.
1780    pub fn max_network_policies(&self) -> u32 {
1781        *self.expect_value(&MAX_NETWORK_POLICIES)
1782    }
1783
1784    /// Returns the value of the `max_network_policies` configuration parameter.
1785    pub fn max_rules_per_network_policy(&self) -> u32 {
1786        *self.expect_value(&MAX_RULES_PER_NETWORK_POLICY)
1787    }
1788
1789    /// Returns the value of the `max_result_size` configuration parameter.
1790    pub fn max_result_size(&self) -> u64 {
1791        self.expect_value::<ByteSize>(&MAX_RESULT_SIZE).as_bytes()
1792    }
1793
1794    /// Returns the value of the `max_copy_from_row_size` configuration parameter.
1795    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    /// Returns the value of the `allowed_cluster_replica_sizes` configuration parameter.
1801    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    /// Returns the value of the `max_concurrent_occ_writes` configuration parameter.
1809    pub fn max_concurrent_occ_writes(&self) -> u32 {
1810        *self.expect_value(&MAX_CONCURRENT_OCC_WRITES)
1811    }
1812
1813    /// Returns the value of the `max_occ_retries` configuration parameter.
1814    pub fn max_occ_retries(&self) -> u32 {
1815        *self.expect_value(&MAX_OCC_RETRIES)
1816    }
1817
1818    /// Returns the value of the `default_cluster_replication_factor` configuration parameter.
1819    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    /// Returns the `pg_source_connect_timeout` configuration parameter.
1896    pub fn pg_source_connect_timeout(&self) -> Duration {
1897        *self.expect_value(&PG_SOURCE_CONNECT_TIMEOUT)
1898    }
1899
1900    /// Returns the `pg_source_tcp_keepalives_retries` configuration parameter.
1901    pub fn pg_source_tcp_keepalives_retries(&self) -> u32 {
1902        *self.expect_value(&PG_SOURCE_TCP_KEEPALIVES_RETRIES)
1903    }
1904
1905    /// Returns the `pg_source_tcp_keepalives_idle` configuration parameter.
1906    pub fn pg_source_tcp_keepalives_idle(&self) -> Duration {
1907        *self.expect_value(&PG_SOURCE_TCP_KEEPALIVES_IDLE)
1908    }
1909
1910    /// Returns the `pg_source_tcp_keepalives_interval` configuration parameter.
1911    pub fn pg_source_tcp_keepalives_interval(&self) -> Duration {
1912        *self.expect_value(&PG_SOURCE_TCP_KEEPALIVES_INTERVAL)
1913    }
1914
1915    /// Returns the `pg_source_tcp_user_timeout` configuration parameter.
1916    pub fn pg_source_tcp_user_timeout(&self) -> Duration {
1917        *self.expect_value(&PG_SOURCE_TCP_USER_TIMEOUT)
1918    }
1919
1920    /// Returns the `pg_source_tcp_configure_server` configuration parameter.
1921    pub fn pg_source_tcp_configure_server(&self) -> bool {
1922        *self.expect_value(&PG_SOURCE_TCP_CONFIGURE_SERVER)
1923    }
1924
1925    /// Returns the `pg_source_snapshot_statement_timeout` configuration parameter.
1926    pub fn pg_source_snapshot_statement_timeout(&self) -> Duration {
1927        *self.expect_value(&PG_SOURCE_SNAPSHOT_STATEMENT_TIMEOUT)
1928    }
1929
1930    /// Returns the `pg_source_wal_sender_timeout` configuration parameter.
1931    pub fn pg_source_wal_sender_timeout(&self) -> Option<Duration> {
1932        *self.expect_value(&PG_SOURCE_WAL_SENDER_TIMEOUT)
1933    }
1934
1935    /// Returns the `pg_source_snapshot_collect_strict_count` configuration parameter.
1936    pub fn pg_source_snapshot_collect_strict_count(&self) -> bool {
1937        *self.expect_value(&PG_SOURCE_SNAPSHOT_COLLECT_STRICT_COUNT)
1938    }
1939
1940    /// Returns the `mysql_source_tcp_keepalive` configuration parameter.
1941    pub fn mysql_source_tcp_keepalive(&self) -> Duration {
1942        *self.expect_value(&MYSQL_SOURCE_TCP_KEEPALIVE)
1943    }
1944
1945    /// Returns the `mysql_source_snapshot_max_execution_time` configuration parameter.
1946    pub fn mysql_source_snapshot_max_execution_time(&self) -> Duration {
1947        *self.expect_value(&MYSQL_SOURCE_SNAPSHOT_MAX_EXECUTION_TIME)
1948    }
1949
1950    /// Returns the `mysql_source_snapshot_lock_wait_timeout` configuration parameter.
1951    pub fn mysql_source_snapshot_lock_wait_timeout(&self) -> Duration {
1952        *self.expect_value(&MYSQL_SOURCE_SNAPSHOT_LOCK_WAIT_TIMEOUT)
1953    }
1954
1955    /// Returns the `mysql_source_snapshot_wait_timeout` configuration parameter.
1956    pub fn mysql_source_snapshot_wait_timeout(&self) -> Duration {
1957        *self.expect_value(&MYSQL_SOURCE_SNAPSHOT_WAIT_TIMEOUT)
1958    }
1959
1960    /// Returns the `mysql_source_connect_timeout` configuration parameter.
1961    pub fn mysql_source_connect_timeout(&self) -> Duration {
1962        *self.expect_value(&MYSQL_SOURCE_CONNECT_TIMEOUT)
1963    }
1964
1965    /// Returns the `ssh_check_interval` configuration parameter.
1966    pub fn ssh_check_interval(&self) -> Duration {
1967        *self.expect_value(&SSH_CHECK_INTERVAL)
1968    }
1969
1970    /// Returns the `ssh_connect_timeout` configuration parameter.
1971    pub fn ssh_connect_timeout(&self) -> Duration {
1972        *self.expect_value(&SSH_CONNECT_TIMEOUT)
1973    }
1974
1975    /// Returns the `ssh_keepalives_idle` configuration parameter.
1976    pub fn ssh_keepalives_idle(&self) -> Duration {
1977        *self.expect_value(&SSH_KEEPALIVES_IDLE)
1978    }
1979
1980    /// Returns the `kafka_socket_keepalive` configuration parameter.
1981    pub fn kafka_socket_keepalive(&self) -> bool {
1982        *self.expect_value(&KAFKA_SOCKET_KEEPALIVE)
1983    }
1984
1985    /// Returns the `kafka_socket_timeout` configuration parameter.
1986    pub fn kafka_socket_timeout(&self) -> Option<Duration> {
1987        *self.expect_value(&KAFKA_SOCKET_TIMEOUT)
1988    }
1989
1990    /// Returns the `kafka_transaction_timeout` configuration parameter.
1991    pub fn kafka_transaction_timeout(&self) -> Duration {
1992        *self.expect_value(&KAFKA_TRANSACTION_TIMEOUT)
1993    }
1994
1995    /// Returns the `kafka_socket_connection_setup_timeout` configuration parameter.
1996    pub fn kafka_socket_connection_setup_timeout(&self) -> Duration {
1997        *self.expect_value(&KAFKA_SOCKET_CONNECTION_SETUP_TIMEOUT)
1998    }
1999
2000    /// Returns the `kafka_fetch_metadata_timeout` configuration parameter.
2001    pub fn kafka_fetch_metadata_timeout(&self) -> Duration {
2002        *self.expect_value(&KAFKA_FETCH_METADATA_TIMEOUT)
2003    }
2004
2005    /// Returns the `kafka_progress_record_fetch_timeout` configuration parameter.
2006    pub fn kafka_progress_record_fetch_timeout(&self) -> Option<Duration> {
2007        *self.expect_value(&KAFKA_PROGRESS_RECORD_FETCH_TIMEOUT)
2008    }
2009
2010    /// Returns the `crdb_connect_timeout` configuration parameter.
2011    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    /// Returns the `crdb_tcp_user_timeout` configuration parameter.
2018    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    /// Returns the `crdb_keepalives_idle` configuration parameter.
2025    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    /// Returns the `crdb_keepalives_interval` configuration parameter.
2032    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    /// Returns the `crdb_keepalives_retries` configuration parameter.
2039    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    /// Returns the `storage_dataflow_max_inflight_bytes` configuration parameter.
2046    pub fn storage_dataflow_max_inflight_bytes(&self) -> Option<usize> {
2047        *self.expect_value(&STORAGE_DATAFLOW_MAX_INFLIGHT_BYTES)
2048    }
2049
2050    /// Returns the `storage_dataflow_max_inflight_bytes_to_cluster_size_fraction` configuration parameter.
2051    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    /// Returns the `storage_shrink_upsert_unused_buffers_by_ratio` configuration parameter.
2056    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    /// Returns the `storage_dataflow_max_inflight_bytes_disk_only` configuration parameter.
2061    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    /// Returns the `storage_statistics_interval` configuration parameter.
2066    pub fn storage_statistics_interval(&self) -> Duration {
2067        *self.expect_value(&STORAGE_STATISTICS_INTERVAL)
2068    }
2069
2070    /// Returns the `storage_statistics_collection_interval` configuration parameter.
2071    pub fn storage_statistics_collection_interval(&self) -> Duration {
2072        *self.expect_value(&STORAGE_STATISTICS_COLLECTION_INTERVAL)
2073    }
2074
2075    /// Returns the `storage_record_source_sink_namespaced_errors` configuration parameter.
2076    pub fn storage_record_source_sink_namespaced_errors(&self) -> bool {
2077        *self.expect_value(&STORAGE_RECORD_SOURCE_SINK_NAMESPACED_ERRORS)
2078    }
2079
2080    /// Returns the `persist_stats_filter_enabled` configuration parameter.
2081    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    /// Computes an update for every dyncfg from its configured value.
2092    ///
2093    /// This does not touch our own [`ConfigSet`]. Callers that need this
2094    /// process's dyncfgs to reflect the result want [`Self::sync_dyncfgs`].
2095    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    /// Applies the configured dyncfg values to our own [`ConfigSet`], and
2126    /// returns them.
2127    ///
2128    /// Two callers own keeping this process's dyncfgs in step with the
2129    /// catalog: catalog open, and every durable system-config change. Everyone
2130    /// else only forwards the updates elsewhere and wants
2131    /// [`Self::dyncfg_updates`] instead.
2132    pub fn sync_dyncfgs(&self) -> ConfigUpdates {
2133        let updates = self.dyncfg_updates();
2134        updates.apply(&self.dyncfgs);
2135        updates
2136    }
2137
2138    /// Returns the `metrics_retention` configuration parameter.
2139    pub fn metrics_retention(&self) -> Duration {
2140        *self.expect_value(&METRICS_RETENTION)
2141    }
2142
2143    /// Returns the `disabled_metric_sinks` configuration parameter.
2144    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    /// Returns the `unsafe_mock_audit_event_timestamp` configuration parameter.
2152    pub fn unsafe_mock_audit_event_timestamp(&self) -> Option<mz_repr::Timestamp> {
2153        *self.expect_value(&UNSAFE_MOCK_AUDIT_EVENT_TIMESTAMP)
2154    }
2155
2156    /// Returns the `enable_rbac_checks` configuration parameter.
2157    pub fn enable_rbac_checks(&self) -> bool {
2158        *self.expect_value(&ENABLE_RBAC_CHECKS)
2159    }
2160
2161    /// Returns the `max_connections` configuration parameter.
2162    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    /// Returns the `superuser_reserved_connections` configuration parameter.
2171    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    /// Returns the `enable_storage_shard_finalization` configuration parameter.
2192    pub fn enable_storage_shard_finalization(&self) -> bool {
2193        *self.expect_value(&ENABLE_STORAGE_SHARD_FINALIZATION)
2194    }
2195
2196    /// Returns the `enable_default_connection_validation` configuration parameter.
2197    pub fn enable_default_connection_validation(&self) -> bool {
2198        *self.expect_value(&ENABLE_DEFAULT_CONNECTION_VALIDATION)
2199    }
2200
2201    /// Returns the `default_timestamp_interval` configuration parameter.
2202    pub fn default_timestamp_interval(&self) -> Duration {
2203        *self.expect_value(&DEFAULT_TIMESTAMP_INTERVAL)
2204    }
2205
2206    /// Returns the `min_timestamp_interval` configuration parameter.
2207    pub fn min_timestamp_interval(&self) -> Duration {
2208        *self.expect_value(&MIN_TIMESTAMP_INTERVAL)
2209    }
2210    /// Returns the `max_timestamp_interval` configuration parameter.
2211    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    /// Returns the `privatelink_status_update_quota_per_minute` configuration parameter.
2313    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    /// Returns the `statement_logging_max_sample_rate` configuration parameter.
2326    pub fn statement_logging_max_sample_rate(&self) -> Numeric {
2327        *self.expect_value(&STATEMENT_LOGGING_MAX_SAMPLE_RATE)
2328    }
2329
2330    /// Returns the `statement_logging_default_sample_rate` configuration parameter.
2331    pub fn statement_logging_default_sample_rate(&self) -> Numeric {
2332        *self.expect_value(&STATEMENT_LOGGING_DEFAULT_SAMPLE_RATE)
2333    }
2334
2335    /// Returns the `enable_internal_statement_logging` configuration parameter.
2336    pub fn enable_internal_statement_logging(&self) -> bool {
2337        *self.expect_value(&ENABLE_INTERNAL_STATEMENT_LOGGING)
2338    }
2339
2340    /// Returns the `enable_statement_arrival_logging` configuration parameter.
2341    pub fn enable_statement_arrival_logging(&self) -> bool {
2342        *self.expect_value(&ENABLE_STATEMENT_ARRIVAL_LOGGING)
2343    }
2344
2345    /// Returns the `enable_extended_protocol_implicit_transaction` configuration
2346    /// parameter.
2347    pub fn enable_extended_protocol_implicit_transaction(&self) -> bool {
2348        *self.expect_value(&ENABLE_EXTENDED_PROTOCOL_IMPLICIT_TRANSACTION)
2349    }
2350
2351    /// Returns the `optimizer_stats_timeout` configuration parameter.
2352    pub fn optimizer_stats_timeout(&self) -> Duration {
2353        *self.expect_value(&OPTIMIZER_STATS_TIMEOUT)
2354    }
2355
2356    /// Returns the `optimizer_oneshot_stats_timeout` configuration parameter.
2357    pub fn optimizer_oneshot_stats_timeout(&self) -> Duration {
2358        *self.expect_value(&OPTIMIZER_ONESHOT_STATS_TIMEOUT)
2359    }
2360
2361    /// Returns the `webhook_concurrent_request_limit` configuration parameter.
2362    pub fn webhook_concurrent_request_limit(&self) -> usize {
2363        *self.expect_value(&WEBHOOK_CONCURRENT_REQUEST_LIMIT)
2364    }
2365
2366    /// Returns the `pg_timestamp_oracle_connection_pool_max_size` configuration parameter.
2367    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    /// Returns the `pg_timestamp_oracle_connection_pool_max_wait` configuration parameter.
2372    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    /// Returns the `pg_timestamp_oracle_connection_pool_ttl` configuration parameter.
2377    pub fn pg_timestamp_oracle_connection_pool_ttl(&self) -> Duration {
2378        *self.expect_value(&PG_TIMESTAMP_ORACLE_CONNECTION_POOL_TTL)
2379    }
2380
2381    /// Returns the `pg_timestamp_oracle_connection_pool_ttl_stagger` configuration parameter.
2382    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    /// Returns the `user_storage_managed_collections_batch_duration` configuration parameter.
2387    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    /// Returns whether the named variable is a controller configuration parameter.
2400    pub fn is_controller_config_var(&self, name: &str) -> bool {
2401        self.is_dyncfg_var(name)
2402    }
2403
2404    /// Returns whether the named variable is a compute configuration parameter
2405    /// (things that go in `ComputeParameters` and are sent to replicas via `UpdateConfiguration`
2406    /// commands).
2407    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    /// Returns whether the named variable is a metrics configuration parameter
2412    pub fn is_metrics_config_var(&self, name: &str) -> bool {
2413        self.is_dyncfg_var(name)
2414    }
2415
2416    /// Returns whether the named variable is a storage configuration parameter.
2417    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    /// Returns whether the named variable is a dyncfg configuration parameter.
2456    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
2469/// Returns whether the named variable is a caching configuration parameter.
2470pub 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
2489/// Returns whether the named variable is a (Postgres/CRDB) timestamp oracle
2490/// configuration parameter.
2491pub 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
2504/// Returns whether the named variable is a cluster scheduling config
2505pub 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
2518/// Returns whether the named variable is an HTTP server related config var.
2519pub fn is_http_config_var(name: &str) -> bool {
2520    name == WEBHOOK_CONCURRENT_REQUEST_LIMIT.name()
2521}
2522
2523/// Set of [`SystemVar`]s that can also get set at a per-Session level.
2524///
2525/// TODO(parkmycar): Instead of a separate list, make this a field on VarDefinition.
2526static 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// Provides a wrapper to express that a particular `ServerVar` is meant to be used as a feature
2556/// flag.
2557#[derive(Debug)]
2558pub struct FeatureFlag {
2559    pub flag: &'static VarDefinition,
2560    pub feature_desc: &'static str,
2561}
2562
2563impl FeatureFlag {
2564    /// Returns whether the feature flag is enabled in the provided `system_vars`.
2565    pub fn enabled(&'static self, system_vars: &SystemVars) -> bool {
2566        *system_vars.expect_value::<bool>(self.flag)
2567    }
2568
2569    /// Returns an error unless the feature flag is enabled in the provided
2570    /// `system_vars`.
2571    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        // The flag defaults off: the value is rejected regardless of the letter
2641        // case of the variable name. This covers `SET`,
2642        // `SET "TRANSACTION_ISOLATION"`, `ALTER ROLE ... SET`, and connection
2643        // options, which all route through `SessionVars::set` and this shared
2644        // check.
2645        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        // Ungated levels pass regardless of the flag.
2656        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        // With the flag on, the gated value passes too.
2666        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        // Unrelated variables are ignored, even with a gated-looking value.
2677        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    // `reset_all` (used by `DISCARD ALL`) must clear a committed session
2696    // override durably, without depending on a later transaction commit to
2697    // promote the reset. Regression coverage for SQL-529.
2698    #[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        // Set non-locally and commit, so the override lives in `session_value`.
2705        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        // The value falls back to the default, the var is unset, and it will
2729        // not mutate at a later transaction end.
2730        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    // `reset_all` must fall back to a system/role/startup default installed via
2737    // `set_default`, not the compiled-in default. Guards the "startup/role
2738    // defaults survive DISCARD ALL" contract.
2739    #[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}