1use 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#[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
93pub 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 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}