Skip to main content

object_store/
factory.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
15use std::{fs, path};
16
17use common_telemetry::info;
18#[cfg(feature = "hdfs-object-store")]
19use opendal::services::HdfsNative;
20#[cfg(feature = "mysql-object-store")]
21use opendal::services::Mysql;
22use opendal::services::{Fs, Gcs, Oss, S3};
23use snafu::prelude::*;
24
25#[cfg(feature = "hdfs-object-store")]
26use crate::config::HdfsConfig;
27#[cfg(feature = "mysql-object-store")]
28use crate::config::MysqlConfig;
29use crate::config::{AzblobConfig, FileConfig, GcsConfig, ObjectStoreConfig, OssConfig, S3Config};
30use crate::error::{self, Result};
31#[cfg(feature = "hdfs-object-store")]
32use crate::layers::HdfsCompatibilityLayer;
33use crate::services::Azblob;
34use crate::util::{build_http_context, clean_temp_dir, join_dir, normalize_dir};
35use crate::{ATOMIC_WRITE_DIR, OLD_ATOMIC_WRITE_DIR, ObjectStore, util};
36
37pub async fn new_raw_object_store(
38    store: &ObjectStoreConfig,
39    data_home: &str,
40) -> Result<ObjectStore> {
41    let data_home = normalize_dir(data_home);
42    match store {
43        ObjectStoreConfig::File(file_config) => new_fs_object_store(&data_home, file_config),
44        ObjectStoreConfig::S3(s3_config) => new_s3_object_store(s3_config).await,
45        ObjectStoreConfig::Oss(oss_config) => new_oss_object_store(oss_config).await,
46        ObjectStoreConfig::Azblob(azblob_config) => new_azblob_object_store(azblob_config).await,
47        ObjectStoreConfig::Gcs(gcs_config) => new_gcs_object_store(gcs_config).await,
48        #[cfg(feature = "hdfs-object-store")]
49        ObjectStoreConfig::Hdfs(hdfs_config) => new_hdfs_object_store(hdfs_config).await,
50        #[cfg(feature = "mysql-object-store")]
51        ObjectStoreConfig::Mysql(mysql_config) => new_mysql_object_store(mysql_config).await,
52    }
53}
54
55/// Creates an object store backed by a native HDFS client.
56#[cfg(feature = "hdfs-object-store")]
57pub async fn new_hdfs_object_store(hdfs_config: &HdfsConfig) -> Result<ObjectStore> {
58    let root = util::normalize_dir(&hdfs_config.connection.root);
59    info!(
60        "The HDFS NameNode is: {}, root is: {}",
61        hdfs_config.connection.name_node, root
62    );
63
64    let builder = HdfsNative::from(&hdfs_config.connection);
65    let compatibility_layer = HdfsCompatibilityLayer::new(
66        &hdfs_config.connection.name_node,
67        &hdfs_config.connection.root,
68        &hdfs_config.connection.options,
69    )
70    .context(error::InitBackendSnafu)?;
71    let operator = ObjectStore::new(builder)
72        .context(error::InitBackendSnafu)?
73        .layer(compatibility_layer);
74
75    Ok(operator)
76}
77
78#[cfg(feature = "mysql-object-store")]
79pub async fn new_mysql_object_store(mysql_config: &MysqlConfig) -> Result<ObjectStore> {
80    let root = util::normalize_dir(&mysql_config.root);
81    info!(
82        "The mysql object storage table is: {}, root is: {}",
83        mysql_config.table.as_deref().unwrap_or("greptime"),
84        root
85    );
86
87    let builder = Mysql::from(mysql_config);
88    let operator = ObjectStore::new(builder).context(error::InitBackendSnafu)?;
89
90    Ok(operator)
91}
92
93/// A helper function to create a file system object store.
94pub fn new_fs_object_store(data_home: &str, _file_config: &FileConfig) -> Result<ObjectStore> {
95    fs::create_dir_all(path::Path::new(&data_home))
96        .context(error::CreateDirSnafu { dir: data_home })?;
97    info!("The file storage home is: {}", data_home);
98
99    let atomic_write_dir = join_dir(data_home, ATOMIC_WRITE_DIR);
100    clean_temp_dir(&atomic_write_dir)?;
101
102    // Compatible code. Remove this after a major release.
103    let old_atomic_temp_dir = join_dir(data_home, OLD_ATOMIC_WRITE_DIR);
104    clean_temp_dir(&old_atomic_temp_dir)?;
105
106    let builder = Fs::default()
107        .root(data_home)
108        .atomic_write_dir(&atomic_write_dir);
109
110    let object_store = ObjectStore::new(builder).context(error::InitBackendSnafu)?;
111
112    Ok(object_store)
113}
114
115pub async fn new_azblob_object_store(azblob_config: &AzblobConfig) -> Result<ObjectStore> {
116    let root = util::normalize_dir(&azblob_config.connection.root);
117    info!(
118        "The azure storage container is: {}, root is: {}",
119        azblob_config.connection.container, &root
120    );
121
122    let ctx = build_http_context(&azblob_config.http_client)?;
123    let builder = Azblob::from(&azblob_config.connection);
124    let operator = ObjectStore::new(builder)
125        .context(error::InitBackendSnafu)?
126        .with_context(ctx);
127
128    Ok(operator)
129}
130
131pub async fn new_gcs_object_store(gcs_config: &GcsConfig) -> Result<ObjectStore> {
132    let root = util::normalize_dir(&gcs_config.connection.root);
133    info!(
134        "The gcs storage bucket is: {}, root is: {}",
135        gcs_config.connection.bucket, &root
136    );
137
138    let ctx = build_http_context(&gcs_config.http_client)?;
139    let builder = Gcs::from(&gcs_config.connection);
140    let operator = ObjectStore::new(builder)
141        .context(error::InitBackendSnafu)?
142        .with_context(ctx);
143
144    Ok(operator)
145}
146
147pub async fn new_oss_object_store(oss_config: &OssConfig) -> Result<ObjectStore> {
148    let root = util::normalize_dir(&oss_config.connection.root);
149    info!(
150        "The oss storage bucket is: {}, root is: {}",
151        oss_config.connection.bucket, &root
152    );
153
154    let ctx = build_http_context(&oss_config.http_client)?;
155    let builder = Oss::from(&oss_config.connection);
156    let operator = ObjectStore::new(builder)
157        .context(error::InitBackendSnafu)?
158        .with_context(ctx);
159
160    Ok(operator)
161}
162
163pub async fn new_s3_object_store(s3_config: &S3Config) -> Result<ObjectStore> {
164    let root = util::normalize_dir(&s3_config.connection.root);
165    info!(
166        "The s3 storage bucket is: {}, root is: {}",
167        s3_config.connection.bucket, &root
168    );
169
170    let ctx = build_http_context(&s3_config.http_client)?;
171    let builder = S3::from(&s3_config.connection);
172    let operator = ObjectStore::new(builder)
173        .context(error::InitBackendSnafu)?
174        .with_context(ctx);
175
176    Ok(operator)
177}
178
179#[cfg(all(test, feature = "hdfs-object-store"))]
180mod tests {
181    use opendal::services::HDFS_NATIVE_SCHEME;
182
183    use super::*;
184    use crate::config::HdfsConnection;
185
186    #[tokio::test]
187    async fn test_new_hdfs_object_store() {
188        let config = HdfsConfig {
189            connection: HdfsConnection {
190                root: "/greptimedb".to_string(),
191                name_node: "hdfs://127.0.0.1:9000".to_string(),
192                ..Default::default()
193            },
194            ..Default::default()
195        };
196
197        let store = new_hdfs_object_store(&config).await.unwrap();
198        assert_eq!(HDFS_NATIVE_SCHEME, store.info().scheme());
199        assert_eq!("/greptimedb/", store.info().root());
200    }
201
202    #[tokio::test]
203    async fn test_new_hdfs_object_store_requires_name_node() {
204        let result = new_hdfs_object_store(&HdfsConfig::default()).await;
205        assert!(result.is_err());
206    }
207}