1#[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#[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 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 pub fn is_object_storage(&self) -> bool {
71 !matches!(self, Self::File(_))
72 }
73
74 pub fn config_name(&self) -> &str {
76 let name = match self {
77 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 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 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 pub enable_virtual_host_style: bool,
146 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#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
308#[serde(default)]
309#[cfg(feature = "hdfs-object-store")]
310pub struct HdfsConnection {
311 pub root: String,
313 pub name_node: String,
315 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#[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#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
377#[serde(default)]
378pub struct HttpClientConfig {
379 pub(crate) pool_max_idle_per_host: u32,
381
382 #[serde(with = "humantime_serde")]
384 pub(crate) connect_timeout: Duration,
385
386 #[serde(with = "humantime_serde")]
389 pub(crate) timeout: Duration,
390
391 #[serde(with = "humantime_serde")]
393 pub(crate) pool_idle_timeout: Duration,
394
395 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 pub enable_read_cache: bool,
416 pub cache_path: String,
418 pub cache_capacity: ReadableSize,
420}
421
422impl Default for ObjectStorageCacheConfig {
423 fn default() -> Self {
424 Self {
425 enable_read_cache: true,
426 cache_path: String::default(),
428 cache_capacity: DEFAULT_OBJECT_STORE_CACHE_SIZE,
429 }
430 }
431}
432
433impl ObjectStorageCacheConfig {
434 pub fn sanitize(&mut self, data_home: &str) {
436 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}