Skip to main content

mz_persist/
cfg.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//! Configuration for [crate::location] implementations.
11
12use std::collections::BTreeMap;
13use std::sync::Arc;
14use std::time::Duration;
15
16use anyhow::anyhow;
17use mz_dyncfg::ConfigSet;
18use mz_ore::error::ErrorExt;
19use mz_ore::url::SensitiveUrl;
20use tracing::warn;
21
22use mz_postgres_client::PostgresClientKnobs;
23use mz_postgres_client::metrics::PostgresClientMetrics;
24
25use crate::azure::{AzureBlob, AzureBlobConfig};
26use crate::file::{FileBlob, FileBlobConfig};
27use crate::hedge::HedgeSibling;
28use crate::location::{Blob, Consensus, Determinate, ExternalError};
29use crate::mem::{MemBlob, MemBlobConfig, MemConsensus};
30use crate::metrics::S3BlobMetrics;
31use crate::postgres::{PostgresConsensus, PostgresConsensusConfig};
32use crate::s3::{S3Blob, S3BlobConfig};
33
34/// Adds the full set of all mz_persist `Config`s.
35pub fn all_dyn_configs(configs: ConfigSet) -> ConfigSet {
36    configs
37        .add(&crate::postgres::PG_CONSENSUS_READ_COMMITTED)
38        .add(&crate::hedge::BLOB_HEDGED_GET_ENABLED)
39        .add(&crate::hedge::BLOB_HEDGED_GET_DELAY)
40        .add(&crate::hedge::BLOB_HEDGED_GET_MAX_CONCURRENT)
41        .add(&crate::hedge::BLOB_HEDGED_GET_BUDGET_RATIO)
42        .add(&crate::hedge::BLOB_HEDGED_GET_WARM_INTERVAL)
43}
44
45/// Opens the sibling handle that [crate::hedge::HedgedBlob] runs hedge
46/// requests on for `url`.
47///
48/// Contract:
49/// - An [HedgeSibling::Isolated] handle observes exactly the same durable
50///   store as a handle opened from the same `url`, but is built from a
51///   scratch client: it shares no HTTP connection pool, DNS state, or
52///   credential chain, so a hedge request on it can never be assigned a
53///   connection the primary's pool has already half-killed.
54/// - Backends where a second open would observe an independent store (mem,
55///   turmoil's simulated store), or that have no connection state to isolate
56///   (file), return [HedgeSibling::SharedWithPrimary] instead.
57/// - Callers must use the handle only for idempotent reads.
58///
59/// Errors opening the sibling degrade to [HedgeSibling::Unavailable] with a
60/// warning rather than failing: persist must come up even if hedging cannot.
61/// What a caller does with that is its own contract, see
62/// [crate::hedge::HedgedBlob::new_arming].
63pub async fn open_hedge_sibling(
64    url: &SensitiveUrl,
65    knobs: Box<dyn BlobKnobs>,
66    metrics: S3BlobMetrics,
67) -> HedgeSibling {
68    let config = match BlobConfig::try_from(url, knobs, metrics).await {
69        Ok(config) => config,
70        Err(err) => {
71            warn!(
72                "hedged blob gets unavailable, sibling config failed: {}",
73                err.display_with_causes()
74            );
75            return HedgeSibling::Unavailable;
76        }
77    };
78    match config {
79        // A second S3/Azure config builds its own SDK client and therefore
80        // its own connection pool, with DNS resolved per connect.
81        config @ (BlobConfig::S3(_) | BlobConfig::Azure(_)) => match config.open().await {
82            Ok(blob) => HedgeSibling::Isolated(blob),
83            Err(err) => {
84                warn!(
85                    "hedged blob gets unavailable, sibling open failed: {}",
86                    err.display_with_causes()
87                );
88                HedgeSibling::Unavailable
89            }
90        },
91        // File has no connection pool to isolate, so a second instance would
92        // buy nothing. A second open of Mem (or of turmoil's simulated
93        // store) would be actively wrong: it creates an INDEPENDENT store,
94        // and a hedged get against a different store can return `Ok(None)`
95        // for data that exists.
96        BlobConfig::File(_) | BlobConfig::Mem(_) => HedgeSibling::SharedWithPrimary,
97        #[cfg(feature = "turmoil")]
98        BlobConfig::Turmoil(_) => HedgeSibling::SharedWithPrimary,
99    }
100}
101
102/// Config for an implementation of [Blob].
103#[derive(Debug, Clone)]
104pub enum BlobConfig {
105    /// Config for [FileBlob].
106    File(FileBlobConfig),
107    /// Config for [S3Blob].
108    S3(S3BlobConfig),
109    /// Config for [MemBlob], only available in testing to prevent
110    /// footguns.
111    Mem(bool),
112    /// Config for [AzureBlob].
113    Azure(AzureBlobConfig),
114    #[cfg(feature = "turmoil")]
115    /// Config for [crate::turmoil::TurmoilBlob].
116    Turmoil(crate::turmoil::BlobConfig),
117}
118
119/// Configuration knobs for [Blob].
120pub trait BlobKnobs: std::fmt::Debug + Send + Sync {
121    /// Maximum time allowed for a network call, including retry attempts.
122    fn operation_timeout(&self) -> Duration;
123    /// Maximum time allowed for a single network call.
124    fn operation_attempt_timeout(&self) -> Duration;
125    /// Maximum time to wait for a socket connection to be made.
126    fn connect_timeout(&self) -> Duration;
127    /// Maximum time to wait to read the first byte of a response, including connection time.
128    fn read_timeout(&self) -> Duration;
129    /// Whether this is running in a "cc" sized cluster.
130    fn is_cc_active(&self) -> bool;
131}
132
133impl BlobConfig {
134    /// Opens the associated implementation of [Blob].
135    pub async fn open(self) -> Result<Arc<dyn Blob>, ExternalError> {
136        match self {
137            BlobConfig::File(config) => Ok(Arc::new(FileBlob::open(config).await?)),
138            BlobConfig::S3(config) => Ok(Arc::new(S3Blob::open(config).await?)),
139            BlobConfig::Azure(config) => Ok(Arc::new(AzureBlob::open(config).await?)),
140            BlobConfig::Mem(tombstone) => {
141                Ok(Arc::new(MemBlob::open(MemBlobConfig::new(tombstone))))
142            }
143            #[cfg(feature = "turmoil")]
144            BlobConfig::Turmoil(config) => Ok(Arc::new(crate::turmoil::TurmoilBlob::open(config))),
145        }
146    }
147
148    /// Parses a [Blob] config from a uri string.
149    pub async fn try_from(
150        url: &SensitiveUrl,
151        knobs: Box<dyn BlobKnobs>,
152        metrics: S3BlobMetrics,
153    ) -> Result<Self, ExternalError> {
154        let mut query_params = url.query_pairs().collect::<BTreeMap<_, _>>();
155
156        let config = match url.scheme() {
157            "file" => {
158                let mut config = FileBlobConfig::from(url.path());
159                if query_params.remove("tombstone").is_some() {
160                    config.tombstone = true;
161                }
162                Ok(BlobConfig::File(config))
163            }
164            "s3" => {
165                let bucket = url
166                    .host()
167                    .ok_or_else(|| anyhow!("missing bucket: {}", url))?
168                    .to_string();
169                let prefix = url
170                    .path()
171                    .strip_prefix('/')
172                    .unwrap_or_else(|| url.path())
173                    .to_string();
174                let role_arn = query_params.remove("role_arn").map(|x| x.into_owned());
175                let endpoint = query_params.remove("endpoint").map(|x| x.into_owned());
176                let region = query_params.remove("region").map(|x| x.into_owned());
177
178                let credentials = match url.password() {
179                    None => None,
180                    Some(password) => Some((
181                        String::from_utf8_lossy(&urlencoding::decode_binary(
182                            url.username().as_bytes(),
183                        ))
184                        .into_owned(),
185                        String::from_utf8_lossy(&urlencoding::decode_binary(password.as_bytes()))
186                            .into_owned(),
187                    )),
188                };
189
190                let config = S3BlobConfig::new(
191                    bucket,
192                    prefix,
193                    role_arn,
194                    endpoint,
195                    region,
196                    credentials,
197                    knobs,
198                    metrics,
199                )
200                .await?;
201
202                Ok(BlobConfig::S3(config))
203            }
204            "mem" => {
205                if !cfg!(debug_assertions) {
206                    warn!("persist unexpectedly using in-mem blob in a release binary");
207                }
208                let tombstone = match query_params.remove("tombstone").as_deref() {
209                    None | Some("true") => true,
210                    Some("false") => false,
211                    Some(other) => Err(Determinate::new(anyhow!(
212                        "invalid tombstone param value: {other}"
213                    )))?,
214                };
215                query_params.clear();
216                Ok(BlobConfig::Mem(tombstone))
217            }
218            "http" | "https" => match url
219                .host()
220                .ok_or_else(|| anyhow!("missing protocol: {}", url))?
221                .to_string()
222                .split_once('.')
223            {
224                // The Azurite emulator always uses the well-known account name devstoreaccount1
225                Some((account, root))
226                    if account == "devstoreaccount1" || root == "blob.core.windows.net" =>
227                {
228                    if let Some(container) = url
229                        .path_segments()
230                        .expect("azure blob storage container")
231                        .next()
232                    {
233                        query_params.clear();
234                        Ok(BlobConfig::Azure(AzureBlobConfig::new(
235                            account.to_string(),
236                            container.to_string(),
237                            // Azure doesn't support prefixes in the way S3 does.
238                            // This is always empty, but we leave the field for
239                            // compatibility with our existing test suite.
240                            "".to_string(),
241                            metrics,
242                            url.clone().into_redacted(),
243                            knobs,
244                        )?))
245                    } else {
246                        Err(anyhow!("unknown persist blob scheme: {}", url))
247                    }
248                }
249                _ => Err(anyhow!("unknown persist blob scheme: {}", url)),
250            },
251            #[cfg(feature = "turmoil")]
252            "turmoil" => {
253                let cfg = crate::turmoil::BlobConfig::new(url);
254                Ok(BlobConfig::Turmoil(cfg))
255            }
256            p => Err(anyhow!("unknown persist blob scheme {}: {}", p, url)),
257        }?;
258
259        if !query_params.is_empty() {
260            return Err(ExternalError::from(anyhow!(
261                "unknown blob location params {}: {}",
262                query_params
263                    .keys()
264                    .map(|x| x.as_ref())
265                    .collect::<Vec<_>>()
266                    .join(" "),
267                url,
268            )));
269        }
270
271        Ok(config)
272    }
273}
274
275/// Config for an implementation of [Consensus].
276#[derive(Debug, Clone)]
277pub enum ConsensusConfig {
278    /// Config for [PostgresConsensus].
279    Postgres(PostgresConsensusConfig),
280    /// Config for [MemConsensus], only available in testing.
281    Mem,
282    #[cfg(feature = "turmoil")]
283    /// Config for [crate::turmoil::TurmoilConsensus].
284    Turmoil(crate::turmoil::ConsensusConfig),
285}
286
287impl ConsensusConfig {
288    /// Opens the associated implementation of [Consensus].
289    pub async fn open(self) -> Result<Arc<dyn Consensus>, ExternalError> {
290        match self {
291            ConsensusConfig::Postgres(config) => {
292                Ok(Arc::new(PostgresConsensus::open(config).await?))
293            }
294            ConsensusConfig::Mem => Ok(Arc::new(MemConsensus::default())),
295            #[cfg(feature = "turmoil")]
296            ConsensusConfig::Turmoil(config) => {
297                Ok(Arc::new(crate::turmoil::TurmoilConsensus::open(config)))
298            }
299        }
300    }
301
302    /// Parses a [Consensus] config from a uri string.
303    pub fn try_from(
304        url: &SensitiveUrl,
305        knobs: Box<dyn PostgresClientKnobs>,
306        metrics: PostgresClientMetrics,
307        dyncfg: Arc<ConfigSet>,
308    ) -> Result<Self, ExternalError> {
309        let config = match url.scheme() {
310            "postgres" | "postgresql" => Ok(ConsensusConfig::Postgres(
311                PostgresConsensusConfig::new(url, knobs, metrics, dyncfg)?,
312            )),
313            "mem" => {
314                if !cfg!(debug_assertions) {
315                    warn!("persist unexpectedly using in-mem consensus in a release binary");
316                }
317                Ok(ConsensusConfig::Mem)
318            }
319            #[cfg(feature = "turmoil")]
320            "turmoil" => {
321                let cfg = crate::turmoil::ConsensusConfig::new(url);
322                Ok(ConsensusConfig::Turmoil(cfg))
323            }
324            p => Err(anyhow!("unknown persist consensus scheme {}: {}", p, url)),
325        }?;
326        Ok(config)
327    }
328}