Skip to main content

object_store/
config.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15#[cfg(feature = "hdfs-object-store")]
16use std::collections::HashMap;
17use std::time::Duration;
18
19use common_base::readable_size::ReadableSize;
20use common_base::secrets::{ExposeSecret, SecretString};
21#[cfg(feature = "hdfs-object-store")]
22use opendal::services::HdfsNative;
23#[cfg(feature = "mysql-object-store")]
24use opendal::services::Mysql;
25use opendal::services::{Azblob, Gcs, Oss, S3};
26use serde::{Deserialize, Serialize};
27
28use crate::util;
29
30const DEFAULT_OBJECT_STORE_CACHE_SIZE: ReadableSize = ReadableSize::gb(5);
31
32/// Object storage config
33#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
34#[serde(tag = "type")]
35pub enum ObjectStoreConfig {
36    File(FileConfig),
37    S3(S3Config),
38    Oss(OssConfig),
39    Azblob(AzblobConfig),
40    Gcs(GcsConfig),
41    #[cfg(feature = "hdfs-object-store")]
42    Hdfs(HdfsConfig),
43    #[cfg(feature = "mysql-object-store")]
44    Mysql(MysqlConfig),
45}
46
47impl Default for ObjectStoreConfig {
48    fn default() -> Self {
49        ObjectStoreConfig::File(FileConfig {})
50    }
51}
52
53impl ObjectStoreConfig {
54    /// Returns the object storage type name, such as `S3`, `Oss` etc.
55    pub fn provider_name(&self) -> &'static str {
56        match self {
57            Self::File(_) => "File",
58            Self::S3(_) => "S3",
59            Self::Oss(_) => "Oss",
60            Self::Azblob(_) => "Azblob",
61            Self::Gcs(_) => "Gcs",
62            #[cfg(feature = "hdfs-object-store")]
63            Self::Hdfs(_) => "Hdfs",
64            #[cfg(feature = "mysql-object-store")]
65            Self::Mysql(_) => "Mysql",
66        }
67    }
68
69    /// Returns true when it's a remote object storage such as AWS s3 etc.
70    pub fn is_object_storage(&self) -> bool {
71        !matches!(self, Self::File(_))
72    }
73
74    /// Returns the object storage configuration name, return the provider name if it's empty.
75    pub fn config_name(&self) -> &str {
76        let name = match self {
77            // file storage doesn't support name
78            Self::File(_) => self.provider_name(),
79            Self::S3(s3) => &s3.name,
80            Self::Oss(oss) => &oss.name,
81            Self::Azblob(az) => &az.name,
82            Self::Gcs(gcs) => &gcs.name,
83            #[cfg(feature = "hdfs-object-store")]
84            Self::Hdfs(hdfs) => &hdfs.name,
85            #[cfg(feature = "mysql-object-store")]
86            Self::Mysql(mysql) => &mysql.name,
87        };
88
89        if name.trim().is_empty() {
90            return self.provider_name();
91        }
92
93        name
94    }
95
96    /// Returns the object storage cache configuration.
97    pub fn cache_config(&self) -> Option<&ObjectStorageCacheConfig> {
98        match self {
99            Self::File(_) => None,
100            Self::S3(s3) => Some(&s3.cache),
101            Self::Oss(oss) => Some(&oss.cache),
102            Self::Azblob(az) => Some(&az.cache),
103            Self::Gcs(gcs) => Some(&gcs.cache),
104            #[cfg(feature = "hdfs-object-store")]
105            Self::Hdfs(hdfs) => Some(&hdfs.cache),
106            #[cfg(feature = "mysql-object-store")]
107            Self::Mysql(mysql) => Some(&mysql.cache),
108        }
109    }
110
111    /// Returns the mutable object storage cache configuration.
112    pub fn cache_config_mut(&mut self) -> Option<&mut ObjectStorageCacheConfig> {
113        match self {
114            Self::File(_) => None,
115            Self::S3(s3) => Some(&mut s3.cache),
116            Self::Oss(oss) => Some(&mut oss.cache),
117            Self::Azblob(az) => Some(&mut az.cache),
118            Self::Gcs(gcs) => Some(&mut gcs.cache),
119            #[cfg(feature = "hdfs-object-store")]
120            Self::Hdfs(hdfs) => Some(&mut hdfs.cache),
121            #[cfg(feature = "mysql-object-store")]
122            Self::Mysql(mysql) => Some(&mut mysql.cache),
123        }
124    }
125}
126
127#[derive(Debug, Clone, Serialize, Default, Deserialize, Eq, PartialEq)]
128#[serde(default)]
129pub struct FileConfig {}
130
131#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
132#[serde(default)]
133pub struct S3Connection {
134    pub bucket: String,
135    pub root: String,
136    #[serde(skip_serializing)]
137    pub access_key_id: SecretString,
138    #[serde(skip_serializing)]
139    pub secret_access_key: SecretString,
140    pub endpoint: Option<String>,
141    pub region: Option<String>,
142    /// Enable virtual host style so that opendal will send API requests in virtual host style instead of path style.
143    /// By default, opendal will send API to https://s3.us-east-1.amazonaws.com/bucket_name
144    /// Enabled, opendal will send API to https://bucket_name.s3.us-east-1.amazonaws.com
145    pub enable_virtual_host_style: bool,
146    /// Disable EC2 metadata service.
147    /// By default, opendal will use EC2 metadata service to load credentials from the instance metadata,
148    /// when access key id and secret access key are not provided.
149    /// If enabled, opendal will *NOT* use EC2 metadata service.
150    pub disable_ec2_metadata: bool,
151}
152
153impl From<&S3Connection> for S3 {
154    fn from(connection: &S3Connection) -> Self {
155        let root = util::normalize_dir(&connection.root);
156
157        let mut builder = S3::default()
158            .root(&root)
159            .bucket(&connection.bucket)
160            .access_key_id(connection.access_key_id.expose_secret())
161            .secret_access_key(connection.secret_access_key.expose_secret());
162
163        if connection.disable_ec2_metadata {
164            builder = builder.disable_ec2_metadata();
165        }
166
167        if let Some(endpoint) = &connection.endpoint {
168            builder = builder.endpoint(endpoint);
169        }
170        if let Some(region) = &connection.region {
171            builder = builder.region(region);
172        }
173        if connection.enable_virtual_host_style {
174            builder = builder.enable_virtual_host_style();
175        }
176
177        builder
178    }
179}
180
181#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
182#[serde(default)]
183pub struct S3Config {
184    pub name: String,
185    #[serde(flatten)]
186    pub connection: S3Connection,
187    #[serde(flatten)]
188    pub cache: ObjectStorageCacheConfig,
189    pub http_client: HttpClientConfig,
190}
191
192#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
193#[serde(default)]
194pub struct OssConnection {
195    pub bucket: String,
196    pub root: String,
197    #[serde(skip_serializing)]
198    pub access_key_id: SecretString,
199    #[serde(skip_serializing)]
200    pub access_key_secret: SecretString,
201    pub endpoint: String,
202}
203
204impl From<&OssConnection> for Oss {
205    fn from(connection: &OssConnection) -> Self {
206        let root = util::normalize_dir(&connection.root);
207        Oss::default()
208            .root(&root)
209            .bucket(&connection.bucket)
210            .endpoint(&connection.endpoint)
211            .access_key_id(connection.access_key_id.expose_secret())
212            .access_key_secret(connection.access_key_secret.expose_secret())
213    }
214}
215
216#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
217#[serde(default)]
218pub struct OssConfig {
219    pub name: String,
220    #[serde(flatten)]
221    pub connection: OssConnection,
222    #[serde(flatten)]
223    pub cache: ObjectStorageCacheConfig,
224    pub http_client: HttpClientConfig,
225}
226
227#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
228#[serde(default)]
229pub struct AzblobConnection {
230    pub container: String,
231    pub root: String,
232    #[serde(skip_serializing)]
233    pub account_name: SecretString,
234    #[serde(skip_serializing)]
235    pub account_key: SecretString,
236    pub endpoint: String,
237    pub sas_token: Option<String>,
238}
239
240impl From<&AzblobConnection> for Azblob {
241    fn from(connection: &AzblobConnection) -> Self {
242        let root = util::normalize_dir(&connection.root);
243        let mut builder = Azblob::default()
244            .root(&root)
245            .container(&connection.container)
246            .endpoint(&connection.endpoint)
247            .account_name(connection.account_name.expose_secret())
248            .account_key(connection.account_key.expose_secret());
249
250        if let Some(token) = &connection.sas_token {
251            builder = builder.sas_token(token);
252        };
253
254        builder
255    }
256}
257
258#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
259#[serde(default)]
260pub struct AzblobConfig {
261    pub name: String,
262    #[serde(flatten)]
263    pub connection: AzblobConnection,
264    #[serde(flatten)]
265    pub cache: ObjectStorageCacheConfig,
266    pub http_client: HttpClientConfig,
267}
268
269#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
270#[serde(default)]
271pub struct GcsConnection {
272    pub root: String,
273    pub bucket: String,
274    pub scope: String,
275    #[serde(skip_serializing)]
276    pub credential_path: SecretString,
277    #[serde(skip_serializing)]
278    pub credential: SecretString,
279    pub endpoint: String,
280}
281
282#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
283#[serde(default)]
284pub struct GcsConfig {
285    pub name: String,
286    #[serde(flatten)]
287    pub connection: GcsConnection,
288    #[serde(flatten)]
289    pub cache: ObjectStorageCacheConfig,
290    pub http_client: HttpClientConfig,
291}
292
293impl From<&GcsConnection> for Gcs {
294    fn from(connection: &GcsConnection) -> Self {
295        let root = util::normalize_dir(&connection.root);
296        Gcs::default()
297            .root(&root)
298            .bucket(&connection.bucket)
299            .scope(&connection.scope)
300            .credential_path(connection.credential_path.expose_secret())
301            .credential(connection.credential.expose_secret())
302            .endpoint(&connection.endpoint)
303    }
304}
305
306/// Connection options for a Hadoop Distributed File System backend.
307#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
308#[serde(default)]
309#[cfg(feature = "hdfs-object-store")]
310pub struct HdfsConnection {
311    /// Working directory for all object-store operations.
312    pub root: String,
313    /// HDFS NameNode URI, for example `hdfs://127.0.0.1:9000`.
314    pub name_node: String,
315    /// Additional options passed to the native HDFS client.
316    pub options: HashMap<String, String>,
317}
318
319#[cfg(feature = "hdfs-object-store")]
320impl From<&HdfsConnection> for HdfsNative {
321    fn from(connection: &HdfsConnection) -> Self {
322        let root = util::normalize_dir(&connection.root);
323        HdfsNative::default()
324            .root(&root)
325            .name_node(&connection.name_node)
326            .options(connection.options.clone())
327    }
328}
329
330/// Hadoop Distributed File System object storage configuration.
331#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
332#[serde(default)]
333#[cfg(feature = "hdfs-object-store")]
334pub struct HdfsConfig {
335    pub name: String,
336    #[serde(flatten)]
337    pub connection: HdfsConnection,
338    #[serde(flatten)]
339    pub cache: ObjectStorageCacheConfig,
340}
341
342#[cfg(feature = "mysql-object-store")]
343#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
344#[serde(default)]
345pub struct MysqlConfig {
346    pub name: String,
347    pub root: String,
348    #[serde(skip_serializing)]
349    pub connection_string: SecretString,
350    pub table: Option<String>,
351    #[serde(flatten)]
352    pub cache: ObjectStorageCacheConfig,
353}
354
355#[cfg(feature = "mysql-object-store")]
356impl From<&MysqlConfig> for Mysql {
357    fn from(config: &MysqlConfig) -> Self {
358        let root = util::normalize_dir(&config.root);
359        let mut builder = Mysql::default()
360            .connection_string(config.connection_string.expose_secret())
361            .root(&root)
362            .key_field("key")
363            .value_field("value");
364
365        if let Some(table) = &config.table {
366            builder = builder.table(table);
367        } else {
368            builder = builder.table("greptime");
369        }
370
371        builder
372    }
373}
374
375/// The http client options to the storage.
376#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
377#[serde(default)]
378pub struct HttpClientConfig {
379    /// The maximum idle connection per host allowed in the pool.
380    pub(crate) pool_max_idle_per_host: u32,
381
382    /// The timeout for only the connect phase of a http client.
383    #[serde(with = "humantime_serde")]
384    pub(crate) connect_timeout: Duration,
385
386    /// The total request timeout, applied from when the request starts connecting until the response body has finished.
387    /// Also considered a total deadline.
388    #[serde(with = "humantime_serde")]
389    pub(crate) timeout: Duration,
390
391    /// The timeout for idle sockets being kept-alive.
392    #[serde(with = "humantime_serde")]
393    pub(crate) pool_idle_timeout: Duration,
394
395    /// Skip SSL certificate validation (insecure)
396    pub skip_ssl_validation: bool,
397}
398
399impl Default for HttpClientConfig {
400    fn default() -> Self {
401        Self {
402            pool_max_idle_per_host: 1024,
403            connect_timeout: Duration::from_secs(30),
404            timeout: Duration::from_secs(30),
405            pool_idle_timeout: Duration::from_secs(90),
406            skip_ssl_validation: false,
407        }
408    }
409}
410
411#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
412#[serde(default)]
413pub struct ObjectStorageCacheConfig {
414    /// Whether to enable read cache. If not set, the read cache will be enabled by default.
415    pub enable_read_cache: bool,
416    /// The local file cache directory
417    pub cache_path: String,
418    /// The cache capacity in bytes
419    pub cache_capacity: ReadableSize,
420}
421
422impl Default for ObjectStorageCacheConfig {
423    fn default() -> Self {
424        Self {
425            enable_read_cache: true,
426            // The cache directory is set to the value of data_home in the build_cache_layer process.
427            cache_path: String::default(),
428            cache_capacity: DEFAULT_OBJECT_STORE_CACHE_SIZE,
429        }
430    }
431}
432
433impl ObjectStorageCacheConfig {
434    /// Sanitize the `ObjectStorageCacheConfig` to ensure the config is valid.
435    pub fn sanitize(&mut self, data_home: &str) {
436        // If `cache_path` is unset, default to use `${data_home}` as the local read cache directory.
437        if self.cache_path.is_empty() {
438            self.cache_path = data_home.to_string();
439        }
440    }
441}
442
443#[cfg(test)]
444mod tests {
445    use super::*;
446    use crate::config::ObjectStoreConfig;
447
448    #[test]
449    fn test_config_name() {
450        let object_store_config = ObjectStoreConfig::default();
451        assert_eq!("File", object_store_config.config_name());
452
453        let s3_config = ObjectStoreConfig::S3(S3Config::default());
454        assert_eq!("S3", s3_config.config_name());
455        assert_eq!("S3", s3_config.provider_name());
456
457        let s3_config = ObjectStoreConfig::S3(S3Config {
458            name: "test".to_string(),
459            ..Default::default()
460        });
461        assert_eq!("test", s3_config.config_name());
462        assert_eq!("S3", s3_config.provider_name());
463
464        #[cfg(feature = "hdfs-object-store")]
465        {
466            let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig::default());
467            assert_eq!("Hdfs", hdfs_config.config_name());
468            assert_eq!("Hdfs", hdfs_config.provider_name());
469
470            let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig {
471                name: "test".to_string(),
472                ..Default::default()
473            });
474            assert_eq!("test", hdfs_config.config_name());
475            assert_eq!("Hdfs", hdfs_config.provider_name());
476        }
477
478        #[cfg(feature = "mysql-object-store")]
479        {
480            let mysql_config = ObjectStoreConfig::Mysql(MysqlConfig::default());
481            assert_eq!("Mysql", mysql_config.config_name());
482            assert_eq!("Mysql", mysql_config.provider_name());
483
484            let mysql_config = ObjectStoreConfig::Mysql(MysqlConfig {
485                name: "test".to_string(),
486                ..Default::default()
487            });
488            assert_eq!("test", mysql_config.config_name());
489            assert_eq!("Mysql", mysql_config.provider_name());
490        }
491    }
492
493    #[test]
494    fn test_is_object_storage() {
495        let store = ObjectStoreConfig::default();
496        assert!(!store.is_object_storage());
497        let s3_config = ObjectStoreConfig::S3(S3Config::default());
498        assert!(s3_config.is_object_storage());
499        let oss_config = ObjectStoreConfig::Oss(OssConfig::default());
500        assert!(oss_config.is_object_storage());
501        let gcs_config = ObjectStoreConfig::Gcs(GcsConfig::default());
502        assert!(gcs_config.is_object_storage());
503        let azblob_config = ObjectStoreConfig::Azblob(AzblobConfig::default());
504        assert!(azblob_config.is_object_storage());
505        #[cfg(feature = "hdfs-object-store")]
506        {
507            let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig::default());
508            assert!(hdfs_config.is_object_storage());
509        }
510        #[cfg(feature = "mysql-object-store")]
511        {
512            let mysql_config = ObjectStoreConfig::Mysql(MysqlConfig::default());
513            assert!(mysql_config.is_object_storage());
514        }
515    }
516
517    #[cfg(feature = "hdfs-object-store")]
518    #[test]
519    fn test_hdfs_config_serde() {
520        let config: ObjectStoreConfig = toml::from_str(
521            r#"
522type = "Hdfs"
523name = "hdfs-store"
524root = "/greptimedb"
525name_node = "hdfs://127.0.0.1:9000"
526
527[options]
528"dfs.client.block.write.replace-datanode-on-failure.enable" = "true"
529"#,
530        )
531        .unwrap();
532
533        let ObjectStoreConfig::Hdfs(hdfs_config) = config else {
534            unreachable!()
535        };
536
537        assert_eq!("hdfs-store", hdfs_config.name);
538        assert_eq!("/greptimedb", hdfs_config.connection.root);
539        assert_eq!("hdfs://127.0.0.1:9000", hdfs_config.connection.name_node);
540        assert_eq!(
541            Some(&"true".to_string()),
542            hdfs_config
543                .connection
544                .options
545                .get("dfs.client.block.write.replace-datanode-on-failure.enable")
546        );
547
548        let serialized = toml::to_string(&hdfs_config).unwrap();
549        assert!(serialized.contains("name_node = \"hdfs://127.0.0.1:9000\""));
550        assert!(
551            serialized.contains(
552                "\"dfs.client.block.write.replace-datanode-on-failure.enable\" = \"true\""
553            )
554        );
555    }
556
557    #[cfg(feature = "mysql-object-store")]
558    #[test]
559    fn test_mysql_config_connection_string_serde() {
560        let config: ObjectStoreConfig = toml::from_str(
561            r#"
562type = "Mysql"
563name = "mysql-store"
564root = "/greptimedb"
565connection_string = "mysql://user:password@127.0.0.1:3306/greptime"
566table = "object_store"
567"#,
568        )
569        .unwrap();
570
571        let ObjectStoreConfig::Mysql(mysql_config) = config else {
572            unreachable!()
573        };
574
575        assert_eq!("mysql-store", mysql_config.name);
576        assert_eq!("/greptimedb", mysql_config.root);
577        assert_eq!(
578            "mysql://user:password@127.0.0.1:3306/greptime",
579            mysql_config.connection_string.expose_secret()
580        );
581        assert_eq!(Some("object_store"), mysql_config.table.as_deref());
582
583        let serialized = toml::to_string(&mysql_config).unwrap();
584        assert!(!serialized.contains("connection_string"));
585        assert!(!serialized.contains("password"));
586    }
587}