1use 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
34pub 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
45pub 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 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 BlobConfig::File(_) | BlobConfig::Mem(_) => HedgeSibling::SharedWithPrimary,
97 #[cfg(feature = "turmoil")]
98 BlobConfig::Turmoil(_) => HedgeSibling::SharedWithPrimary,
99 }
100}
101
102#[derive(Debug, Clone)]
104pub enum BlobConfig {
105 File(FileBlobConfig),
107 S3(S3BlobConfig),
109 Mem(bool),
112 Azure(AzureBlobConfig),
114 #[cfg(feature = "turmoil")]
115 Turmoil(crate::turmoil::BlobConfig),
117}
118
119pub trait BlobKnobs: std::fmt::Debug + Send + Sync {
121 fn operation_timeout(&self) -> Duration;
123 fn operation_attempt_timeout(&self) -> Duration;
125 fn connect_timeout(&self) -> Duration;
127 fn read_timeout(&self) -> Duration;
129 fn is_cc_active(&self) -> bool;
131}
132
133impl BlobConfig {
134 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 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 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 "".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#[derive(Debug, Clone)]
277pub enum ConsensusConfig {
278 Postgres(PostgresConsensusConfig),
280 Mem,
282 #[cfg(feature = "turmoil")]
283 Turmoil(crate::turmoil::ConsensusConfig),
285}
286
287impl ConsensusConfig {
288 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 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}