1pub mod builder;
16#[allow(clippy::print_stdout)]
17pub(crate) mod objbench;
18#[cfg(feature = "dev-tools")]
19#[allow(clippy::print_stdout)]
20mod parquet_meta;
21#[cfg(feature = "dev-tools")]
22#[allow(clippy::print_stdout)]
23mod parquet_rewrite;
24#[allow(clippy::print_stdout)]
25pub mod parquetbench;
26#[allow(clippy::print_stdout)]
27pub mod scanbench;
28#[cfg(feature = "dev-tools")]
29#[allow(clippy::print_stdout)]
30mod sst_replace;
31mod tool_util;
32
33use std::path::Path;
34use std::time::Duration;
35
36use async_trait::async_trait;
37use clap::Parser;
38use common_config::Configurable;
39use common_telemetry::logging::{DEFAULT_LOGGING_DIR, TracingOptions};
40use common_telemetry::{info, warn};
41use common_wal::config::DatanodeWalConfig;
42use datanode::config::RegionEngineConfig;
43use datanode::datanode::Datanode;
44use meta_client::MetaClientOptions;
45use serde::{Deserialize, Serialize};
46use snafu::{ResultExt, ensure};
47use tracing_appender::non_blocking::WorkerGuard;
48
49use crate::App;
50use crate::datanode::builder::InstanceBuilder;
51use crate::datanode::objbench::ObjbenchCommand;
52#[cfg(feature = "dev-tools")]
53use crate::datanode::parquet_meta::ParquetMetaCommand;
54#[cfg(feature = "dev-tools")]
55use crate::datanode::parquet_rewrite::ParquetRewriteCommand;
56use crate::datanode::parquetbench::ParquetbenchCommand;
57use crate::datanode::scanbench::ScanbenchCommand;
58#[cfg(feature = "dev-tools")]
59use crate::datanode::sst_replace::SstReplaceCommand;
60use crate::error::{
61 LoadLayeredConfigSnafu, MissingConfigSnafu, Result, ShutdownDatanodeSnafu, StartDatanodeSnafu,
62};
63use crate::options::{GlobalOptions, GreptimeOptions};
64
65pub const APP_NAME: &str = "greptime-datanode";
66
67type DatanodeOptions = GreptimeOptions<datanode::config::DatanodeOptions>;
68
69pub struct Instance {
70 datanode: Datanode,
71
72 _guard: Vec<WorkerGuard>,
74}
75
76impl Instance {
77 pub fn new(datanode: Datanode, guard: Vec<WorkerGuard>) -> Self {
78 Self {
79 datanode,
80 _guard: guard,
81 }
82 }
83
84 pub fn datanode(&self) -> &Datanode {
85 &self.datanode
86 }
87
88 pub fn datanode_mut(&mut self) -> &mut Datanode {
90 &mut self.datanode
91 }
92}
93
94#[async_trait]
95impl App for Instance {
96 fn name(&self) -> &str {
97 APP_NAME
98 }
99
100 async fn start(&mut self) -> Result<()> {
101 plugins::start_datanode_plugins(&self.datanode)
102 .await
103 .context(StartDatanodeSnafu)?;
104
105 self.datanode.start().await.context(StartDatanodeSnafu)
106 }
107
108 async fn stop(&mut self) -> Result<()> {
109 self.datanode
110 .shutdown()
111 .await
112 .context(ShutdownDatanodeSnafu)
113 }
114}
115
116#[derive(Parser)]
117pub struct Command {
118 #[clap(subcommand)]
119 pub subcmd: SubCommand,
120}
121
122impl Command {
123 pub async fn build_with(&self, builder: InstanceBuilder) -> Result<Instance> {
124 self.subcmd.build_with(builder).await
125 }
126
127 pub fn load_options(&self, global_options: &GlobalOptions) -> Result<DatanodeOptions> {
128 match &self.subcmd {
129 SubCommand::Start(cmd) => cmd.load_options(global_options),
130 SubCommand::Objbench(_) | SubCommand::Scanbench(_) => Self::default_bench_options(),
132 SubCommand::Parquetbench(_) => Self::default_bench_options(),
133 #[cfg(feature = "dev-tools")]
134 SubCommand::ParquetMeta(_) => Self::default_bench_options(),
135 #[cfg(feature = "dev-tools")]
136 SubCommand::ParquetRewrite(_) => Self::default_bench_options(),
137 #[cfg(feature = "dev-tools")]
138 SubCommand::SstReplace(_) => Self::default_bench_options(),
139 }
140 }
141
142 fn default_bench_options() -> Result<DatanodeOptions> {
145 let mut opts = datanode::config::DatanodeOptions::default();
146 opts.sanitize();
147 Ok(DatanodeOptions {
148 runtime: Default::default(),
149 plugins: Default::default(),
150 component: opts,
151 })
152 }
153}
154
155#[derive(Parser)]
156pub enum SubCommand {
157 Start(StartCommand),
158 Objbench(ObjbenchCommand),
160 Scanbench(ScanbenchCommand),
162 Parquetbench(ParquetbenchCommand),
164 #[cfg(feature = "dev-tools")]
166 ParquetMeta(ParquetMetaCommand),
167 #[cfg(feature = "dev-tools")]
169 ParquetRewrite(ParquetRewriteCommand),
170 #[cfg(feature = "dev-tools")]
172 SstReplace(SstReplaceCommand),
173}
174
175impl SubCommand {
176 async fn build_with(&self, builder: InstanceBuilder) -> Result<Instance> {
177 match self {
178 SubCommand::Start(cmd) => {
179 info!("Building datanode with {:#?}", cmd);
180 builder.build().await
181 }
182 SubCommand::Objbench(cmd) => {
183 cmd.run().await?;
184 std::process::exit(0);
185 }
186 SubCommand::Scanbench(cmd) => {
187 cmd.run().await?;
188 std::process::exit(0);
189 }
190 SubCommand::Parquetbench(cmd) => {
191 cmd.run().await?;
192 std::process::exit(0);
193 }
194 #[cfg(feature = "dev-tools")]
195 SubCommand::ParquetMeta(cmd) => {
196 cmd.run().await?;
197 std::process::exit(0);
198 }
199 #[cfg(feature = "dev-tools")]
200 SubCommand::ParquetRewrite(cmd) => {
201 cmd.run().await?;
202 std::process::exit(0);
203 }
204 #[cfg(feature = "dev-tools")]
205 SubCommand::SstReplace(cmd) => {
206 cmd.run().await?;
207 std::process::exit(0);
208 }
209 }
210 }
211}
212
213#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
215#[serde(default)]
216pub struct StorageConfig {
217 pub data_home: String,
219 #[serde(flatten)]
220 pub store: object_store::config::ObjectStoreConfig,
221}
222
223#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
224#[serde(default)]
225struct StorageConfigWrapper {
226 storage: StorageConfig,
227 region_engine: Vec<RegionEngineConfig>,
228 #[serde(default, deserialize_with = "deserialize_wal_config")]
229 wal: DatanodeWalConfig,
230}
231
232fn deserialize_wal_config<'de, D>(
237 deserializer: D,
238) -> std::result::Result<DatanodeWalConfig, D::Error>
239where
240 D: serde::Deserializer<'de>,
241{
242 use serde::de::Error as _;
243
244 let mut table = <toml::value::Table as serde::Deserialize>::deserialize(deserializer)?;
245 if !table.contains_key("provider") {
246 table.insert(
247 "provider".to_string(),
248 toml::Value::String("raft_engine".to_string()),
249 );
250 }
251 DatanodeWalConfig::deserialize(toml::Value::Table(table)).map_err(D::Error::custom)
252}
253
254#[derive(Debug, Parser, Default)]
255pub struct StartCommand {
256 #[clap(long)]
257 node_id: Option<u64>,
258 #[clap(long = "grpc-bind-addr", alias = "rpc-bind-addr", alias = "rpc-addr")]
260 grpc_bind_addr: Option<String>,
261 #[clap(
265 long = "grpc-server-addr",
266 alias = "rpc-server-addr",
267 alias = "rpc-hostname"
268 )]
269 grpc_server_addr: Option<String>,
270 #[clap(long, value_delimiter = ',', num_args = 1..)]
271 metasrv_addrs: Option<Vec<String>>,
272 #[clap(short, long)]
273 config_file: Option<String>,
274 #[clap(long)]
275 data_home: Option<String>,
276 #[clap(long)]
277 wal_dir: Option<String>,
278 #[clap(long)]
279 http_addr: Option<String>,
280 #[clap(long)]
281 http_timeout: Option<u64>,
282 #[clap(long, default_value = "GREPTIMEDB_DATANODE")]
283 env_prefix: String,
284}
285
286impl StartCommand {
287 pub fn load_options(&self, global_options: &GlobalOptions) -> Result<DatanodeOptions> {
288 let mut opts = DatanodeOptions::load_layered_options(
289 self.config_file.as_deref(),
290 self.env_prefix.as_ref(),
291 )
292 .context(LoadLayeredConfigSnafu)?;
293
294 self.merge_with_cli_options(global_options, &mut opts)?;
295 opts.component.sanitize();
296
297 Ok(opts)
298 }
299
300 #[allow(deprecated)]
302 fn merge_with_cli_options(
303 &self,
304 global_options: &GlobalOptions,
305 opts: &mut DatanodeOptions,
306 ) -> Result<()> {
307 let opts = &mut opts.component;
308
309 if let Some(dir) = &global_options.log_dir {
310 opts.logging.dir.clone_from(dir);
311 }
312
313 if global_options.log_level.is_some() {
314 opts.logging.level.clone_from(&global_options.log_level);
315 }
316
317 opts.tracing = TracingOptions {
318 #[cfg(feature = "tokio-console")]
319 tokio_console_addr: global_options.tokio_console_addr.clone(),
320 };
321
322 if let Some(addr) = &self.grpc_bind_addr {
323 opts.grpc.bind_addr.clone_from(addr);
324 } else if let Some(addr) = &opts.rpc_addr {
325 warn!(
326 "Use the deprecated attribute `DatanodeOptions.rpc_addr`, please use `grpc.bind_addr` instead."
327 );
328 opts.grpc.bind_addr.clone_from(addr);
329 }
330
331 if let Some(server_addr) = &self.grpc_server_addr {
332 opts.grpc.server_addr.clone_from(server_addr);
333 } else if let Some(server_addr) = &opts.rpc_hostname {
334 warn!(
335 "Use the deprecated attribute `DatanodeOptions.rpc_hostname`, please use `grpc.server_addr` instead."
336 );
337 opts.grpc.server_addr.clone_from(server_addr);
338 }
339
340 if let Some(runtime_size) = opts.rpc_runtime_size {
341 warn!(
342 "Use the deprecated attribute `DatanodeOptions.rpc_runtime_size`, please use `grpc.runtime_size` instead."
343 );
344 opts.grpc.runtime_size = runtime_size;
345 }
346
347 if let Some(max_recv_message_size) = opts.rpc_max_recv_message_size {
348 warn!(
349 "Use the deprecated attribute `DatanodeOptions.rpc_max_recv_message_size`, please use `grpc.max_recv_message_size` instead."
350 );
351 opts.grpc.max_recv_message_size = max_recv_message_size;
352 }
353
354 if let Some(max_send_message_size) = opts.rpc_max_send_message_size {
355 warn!(
356 "Use the deprecated attribute `DatanodeOptions.rpc_max_send_message_size`, please use `grpc.max_send_message_size` instead."
357 );
358 opts.grpc.max_send_message_size = max_send_message_size;
359 }
360
361 if let Some(node_id) = self.node_id {
362 opts.node_id = Some(node_id);
363 }
364
365 if let Some(metasrv_addrs) = &self.metasrv_addrs {
366 opts.meta_client
367 .get_or_insert_with(MetaClientOptions::default)
368 .metasrv_addrs
369 .clone_from(metasrv_addrs);
370 }
371
372 ensure!(
373 opts.node_id.is_some(),
374 MissingConfigSnafu {
375 msg: "Missing node id option"
376 }
377 );
378
379 if let Some(data_home) = &self.data_home {
380 opts.storage.data_home.clone_from(data_home);
381 }
382
383 if let Some(wal_dir) = &self.wal_dir
385 && let DatanodeWalConfig::RaftEngine(raft_engine_config) = &mut opts.wal
386 {
387 if raft_engine_config
388 .dir
389 .as_ref()
390 .is_some_and(|original_dir| original_dir != wal_dir)
391 {
392 info!("The wal dir of raft-engine is altered to {wal_dir}");
393 }
394 raft_engine_config.dir.replace(wal_dir.clone());
395 }
396
397 if opts.logging.dir.is_empty() {
399 opts.logging.dir = Path::new(&opts.storage.data_home)
400 .join(DEFAULT_LOGGING_DIR)
401 .to_string_lossy()
402 .to_string();
403 }
404
405 if let Some(http_addr) = &self.http_addr {
406 opts.http.addr.clone_from(http_addr);
407 }
408
409 if let Some(http_timeout) = self.http_timeout {
410 opts.http.timeout = Duration::from_secs(http_timeout)
411 }
412
413 opts.http.disable_dashboard = true;
415
416 Ok(())
417 }
418}
419
420#[cfg(test)]
421mod tests {
422 use std::assert_matches;
423 use std::io::Write;
424 use std::time::Duration;
425
426 use clap::{CommandFactory, Parser};
427 use common_config::ENV_VAR_SEP;
428 use common_test_util::temp_dir::create_named_temp_file;
429 use object_store::config::{FileConfig, GcsConfig, ObjectStoreConfig, S3Config};
430
431 use super::*;
432 use crate::options::GlobalOptions;
433
434 #[test]
435 fn test_deprecated_cli_options() {
436 common_telemetry::init_default_ut_logging();
437 let mut file = create_named_temp_file();
438 let toml_str = r#"
439 enable_memory_catalog = false
440 node_id = 42
441
442 rpc_addr = "127.0.0.1:4001"
443 rpc_hostname = "192.168.0.1"
444 [grpc]
445 bind_addr = "127.0.0.1:3001"
446 server_addr = "127.0.0.1"
447 runtime_size = 8
448 "#;
449 write!(file, "{}", toml_str).unwrap();
450
451 let cmd = StartCommand {
452 config_file: Some(file.path().to_str().unwrap().to_string()),
453 ..Default::default()
454 };
455
456 let options = cmd.load_options(&Default::default()).unwrap().component;
457 assert_eq!("127.0.0.1:4001".to_string(), options.grpc.bind_addr);
458 assert_eq!("192.168.0.1".to_string(), options.grpc.server_addr);
459 }
460
461 #[test]
462 fn test_read_from_config_file() {
463 let mut file = create_named_temp_file();
464 let toml_str = r#"
465 enable_memory_catalog = false
466 node_id = 42
467
468 [grpc]
469 bind_addr = "127.0.0.1:3001"
470 server_addr = "127.0.0.1"
471 runtime_size = 8
472
473 [meta_client]
474 metasrv_addrs = ["127.0.0.1:3002"]
475 timeout = "3s"
476 connect_timeout = "5s"
477 ddl_timeout = "10s"
478 tcp_nodelay = true
479
480 [wal]
481 provider = "raft_engine"
482 dir = "/other/wal"
483 file_size = "1GB"
484 purge_threshold = "50GB"
485 purge_interval = "10m"
486 read_batch_size = 128
487 sync_write = false
488
489 [storage]
490 data_home = "./greptimedb_data/"
491 type = "File"
492
493 [[storage.providers]]
494 type = "Gcs"
495 bucket = "foo"
496 endpoint = "bar"
497
498 [[storage.providers]]
499 type = "S3"
500 bucket = "foo"
501
502 [logging]
503 level = "debug"
504 dir = "./greptimedb_data/test/logs"
505 "#;
506 write!(file, "{}", toml_str).unwrap();
507
508 let cmd = StartCommand {
509 config_file: Some(file.path().to_str().unwrap().to_string()),
510 ..Default::default()
511 };
512
513 let options = cmd.load_options(&Default::default()).unwrap().component;
514
515 assert_eq!("127.0.0.1:3001".to_string(), options.grpc.bind_addr);
516 assert_eq!("127.0.0.1".to_string(), options.grpc.server_addr);
517 assert_eq!(Some(42), options.node_id);
518
519 let DatanodeWalConfig::RaftEngine(raft_engine_config) = options.wal else {
520 unreachable!()
521 };
522 assert_eq!("/other/wal", raft_engine_config.dir.unwrap());
523 assert_eq!(Duration::from_secs(600), raft_engine_config.purge_interval);
524 assert_eq!(1024 * 1024 * 1024, raft_engine_config.file_size.0);
525 assert_eq!(
526 1024 * 1024 * 1024 * 50,
527 raft_engine_config.purge_threshold.0
528 );
529 assert!(!raft_engine_config.sync_write);
530
531 let MetaClientOptions {
532 metasrv_addrs: metasrv_addr,
533 timeout,
534 connect_timeout,
535 ddl_timeout,
536 tcp_nodelay,
537 ..
538 } = options.meta_client.unwrap();
539
540 assert_eq!(vec!["127.0.0.1:3002".to_string()], metasrv_addr);
541 assert_eq!(5000, connect_timeout.as_millis());
542 assert_eq!(10000, ddl_timeout.as_millis());
543 assert_eq!(3000, timeout.as_millis());
544 assert!(tcp_nodelay);
545 assert_eq!("./greptimedb_data/", options.storage.data_home);
546 assert!(matches!(
547 &options.storage.store,
548 ObjectStoreConfig::File(FileConfig { .. })
549 ));
550 assert_eq!(options.storage.providers.len(), 2);
551 assert!(matches!(
552 options.storage.providers[0],
553 ObjectStoreConfig::Gcs(GcsConfig { .. })
554 ));
555 assert!(matches!(
556 options.storage.providers[1],
557 ObjectStoreConfig::S3(S3Config { .. })
558 ));
559
560 assert_eq!("debug", options.logging.level.unwrap());
561 assert_eq!(
562 "./greptimedb_data/test/logs".to_string(),
563 options.logging.dir
564 );
565 }
566
567 #[test]
568 fn test_try_from_cmd() {
569 assert!(
570 (StartCommand {
571 metasrv_addrs: Some(vec!["127.0.0.1:3002".to_string()]),
572 ..Default::default()
573 })
574 .load_options(&GlobalOptions::default())
575 .is_err()
576 );
577
578 assert!(
580 (StartCommand {
581 node_id: Some(42),
582 ..Default::default()
583 })
584 .load_options(&GlobalOptions::default())
585 .is_ok()
586 );
587 }
588
589 #[test]
590 fn test_load_log_options_from_cli() {
591 let mut cmd = StartCommand::default();
592
593 let result = cmd.load_options(&GlobalOptions {
594 log_dir: Some("./greptimedb_data/test/logs".to_string()),
595 log_level: Some("debug".to_string()),
596
597 #[cfg(feature = "tokio-console")]
598 tokio_console_addr: None,
599 });
600 assert_matches!(result, Err(crate::error::Error::MissingConfig { .. }));
602
603 cmd.node_id = Some(42);
604
605 let options = cmd
606 .load_options(&GlobalOptions {
607 log_dir: Some("./greptimedb_data/test/logs".to_string()),
608 log_level: Some("debug".to_string()),
609
610 #[cfg(feature = "tokio-console")]
611 tokio_console_addr: None,
612 })
613 .unwrap()
614 .component;
615
616 let logging_opt = options.logging;
617 assert_eq!("./greptimedb_data/test/logs", logging_opt.dir);
618 assert_eq!("debug", logging_opt.level.as_ref().unwrap());
619 }
620
621 #[test]
622 fn test_config_precedence_order() {
623 let mut file = create_named_temp_file();
624 let toml_str = r#"
625 enable_memory_catalog = false
626 node_id = 42
627 rpc_addr = "127.0.0.1:3001"
628 rpc_runtime_size = 8
629 rpc_hostname = "10.103.174.219"
630
631 [meta_client]
632 timeout = "3s"
633 connect_timeout = "5s"
634 tcp_nodelay = true
635
636 [wal]
637 provider = "raft_engine"
638 file_size = "1GB"
639 purge_threshold = "50GB"
640 purge_interval = "5m"
641 sync_write = false
642
643 [storage]
644 type = "File"
645 data_home = "./greptimedb_data/"
646
647 [logging]
648 level = "debug"
649 dir = "./greptimedb_data/test/logs"
650 "#;
651 write!(file, "{}", toml_str).unwrap();
652
653 let env_prefix = "DATANODE_UT";
654 temp_env::with_vars(
655 [
656 (
657 [
659 env_prefix.to_string(),
660 "wal".to_uppercase(),
661 "purge_interval".to_uppercase(),
662 ]
663 .join(ENV_VAR_SEP),
664 Some("1m"),
665 ),
666 (
667 [
669 env_prefix.to_string(),
670 "wal".to_uppercase(),
671 "read_batch_size".to_uppercase(),
672 ]
673 .join(ENV_VAR_SEP),
674 Some("100"),
675 ),
676 (
677 [
679 env_prefix.to_string(),
680 "meta_client".to_uppercase(),
681 "metasrv_addrs".to_uppercase(),
682 ]
683 .join(ENV_VAR_SEP),
684 Some("127.0.0.1:3001,127.0.0.1:3002,127.0.0.1:3003"),
685 ),
686 ],
687 || {
688 let command = StartCommand {
689 config_file: Some(file.path().to_str().unwrap().to_string()),
690 wal_dir: Some("/other/wal/dir".to_string()),
691 env_prefix: env_prefix.to_string(),
692 ..Default::default()
693 };
694
695 let opts = command.load_options(&Default::default()).unwrap().component;
696
697 let DatanodeWalConfig::RaftEngine(raft_engine_config) = opts.wal else {
699 unreachable!()
700 };
701 assert_eq!(raft_engine_config.read_batch_size, 100);
702 assert_eq!(
703 opts.meta_client.unwrap().metasrv_addrs,
704 vec![
705 "127.0.0.1:3001".to_string(),
706 "127.0.0.1:3002".to_string(),
707 "127.0.0.1:3003".to_string()
708 ]
709 );
710
711 assert_eq!(
713 raft_engine_config.purge_interval,
714 Duration::from_secs(60 * 5)
715 );
716
717 assert_eq!(raft_engine_config.dir.unwrap(), "/other/wal/dir");
719
720 assert_eq!(
722 opts.http.addr,
723 DatanodeOptions::default().component.http.addr
724 );
725 assert_eq!(opts.grpc.server_addr, "10.103.174.219");
726 },
727 );
728 }
729
730 #[test]
731 fn test_parse_grpc_cli_aliases() {
732 let command = StartCommand::try_parse_from([
733 "datanode",
734 "--grpc-bind-addr",
735 "127.0.0.1:13001",
736 "--grpc-server-addr",
737 "10.0.0.1:13001",
738 ])
739 .unwrap();
740 assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:13001"));
741 assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.1:13001"));
742
743 let command = StartCommand::try_parse_from([
744 "datanode",
745 "--rpc-bind-addr",
746 "127.0.0.1:23001",
747 "--rpc-server-addr",
748 "10.0.0.2:23001",
749 ])
750 .unwrap();
751 assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:23001"));
752 assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.2:23001"));
753
754 let command = StartCommand::try_parse_from([
755 "datanode",
756 "--rpc-addr",
757 "127.0.0.1:33001",
758 "--rpc-hostname",
759 "10.0.0.3:33001",
760 ])
761 .unwrap();
762 assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:33001"));
763 assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.3:33001"));
764 }
765
766 #[test]
767 fn test_help_uses_grpc_option_names() {
768 let mut cmd = StartCommand::command();
769 let mut help = Vec::new();
770 cmd.write_long_help(&mut help).unwrap();
771 let help = String::from_utf8(help).unwrap();
772
773 assert!(help.contains("--grpc-bind-addr"));
774 assert!(help.contains("--grpc-server-addr"));
775 assert!(!help.contains("--rpc-bind-addr"));
776 assert!(!help.contains("--rpc-server-addr"));
777 assert!(!help.contains("--rpc-addr"));
778 assert!(!help.contains("--rpc-hostname"));
779 }
780}