1use std::fmt::Debug;
16use std::net::SocketAddr;
17use std::path::Path;
18use std::sync::Arc;
19use std::{fs, path};
20
21use async_trait::async_trait;
22use cache::{build_fundamental_cache_registry, with_default_composite_cache_registry};
23use catalog::CatalogManagerRef;
24use catalog::information_schema::InformationExtensionRef;
25use catalog::kvbackend::{CatalogManagerConfiguratorRef, KvBackendCatalogManagerBuilder};
26use catalog::process_manager::ProcessManager;
27use clap::Parser;
28use common_base::Plugins;
29use common_catalog::consts::{MIN_USER_FLOW_ID, MIN_USER_TABLE_ID};
30use common_config::{Configurable, metadata_store_dir};
31use common_datasource::object_store::{LocalFileAccess, configured_local_path};
32use common_error::ext::BoxedError;
33use common_meta::DatanodeId;
34use common_meta::cache::{LayeredCacheRegistryBuilder, LayeredCacheRegistryRef};
35use common_meta::ddl::flow_meta::FlowMetadataAllocator;
36use common_meta::ddl::table_meta::TableMetadataAllocator;
37use common_meta::ddl::{DdlContext, NoopRegionFailureDetectorControl};
38use common_meta::ddl_manager::{DdlManager, DdlManagerConfiguratorRef, DdlManagerRef};
39use common_meta::key::flow::FlowMetadataManager;
40use common_meta::key::{TableMetadataManager, TableMetadataManagerRef};
41use common_meta::kv_backend::KvBackendRef;
42use common_meta::node_manager::{FlownodeRef, NodeManagerRef};
43use common_meta::procedure_executor::{LocalProcedureExecutor, ProcedureExecutorRef};
44use common_meta::region_keeper::MemoryRegionKeeper;
45use common_meta::region_registry::LeaderRegionRegistry;
46use common_meta::sequence::{Sequence, SequenceBuilder};
47use common_meta::wal_provider::{WalProviderRef, build_wal_provider};
48use common_options::plugin_options::StandaloneFlag;
49use common_procedure::ProcedureManagerRef;
50use common_query::prelude::set_default_prefix;
51use common_telemetry::info;
52use common_telemetry::logging::{DEFAULT_LOGGING_DIR, TracingOptions};
53use common_time::timezone::set_default_timezone;
54use common_version::{short_version, verbose_version};
55use datanode::config::{DatanodeOptions, StorageConfig};
56use datanode::datanode::{Datanode, DatanodeBuilder};
57use datanode::region_server::RegionServer;
58use flow::{
59 FlowDualEngineRef, FlownodeBuilder, FlownodeInstance, FlownodeOptions, FrontendClient,
60 FrontendInvoker, GrpcQueryHandlerWithBoxedError,
61};
62use frontend::frontend::Frontend;
63use frontend::instance::builder::FrontendBuilder;
64use frontend::server::Services;
65use meta_srv::metasrv::{FLOW_ID_SEQ, TABLE_ID_SEQ};
66use plugins::PluginOptions;
67use plugins::frontend::context::{
68 CatalogManagerConfigureContext, StandaloneCatalogManagerConfigureContext,
69};
70use plugins::standalone::context::DdlManagerConfigureContext;
71use servers::tls::{TlsMode, TlsOption, merge_tls_option};
72use snafu::{OptionExt, ResultExt};
73use standalone::options::StandaloneOptions;
74use standalone::{
75 StandaloneDatanodeManager, StandaloneInformationExtension,
76 StandaloneRepartitionProcedureFactory,
77};
78use tracing_appender::non_blocking::WorkerGuard;
79
80use crate::error::{OtherSnafu, Result, StartFlownodeSnafu};
81use crate::options::{GlobalOptions, GreptimeOptions};
82use crate::{App, create_resource_limit_metrics, error, log_versions, maybe_activate_heap_profile};
83
84pub const APP_NAME: &str = "greptime-standalone";
85
86fn standalone_local_file_access(
87 storage: &StorageConfig,
88) -> common_datasource::error::Result<LocalFileAccess> {
89 let data_home = configured_local_path(&storage.data_home)?;
90 let copy_root = match &storage.copy_root {
91 Some(root) => configured_local_path(root)?.with_context(|| {
92 common_datasource::error::InvalidLocalFileRootConfigSnafu {
93 root: root.clone(),
94 reason: "copy_root must be a local path or file URL".to_string(),
95 }
96 })?,
97 None => {
98 let Some(data_home) = &data_home else {
99 info!(
100 "SQL access to local files is disabled because storage.data_home is not a local path and storage.copy_root is unset"
101 );
102 return Ok(LocalFileAccess::Disabled);
103 };
104 data_home.join("copy")
105 }
106 };
107
108 let access = LocalFileAccess::sandboxed(©_root)?;
109 if let Some(data_home) = data_home {
110 let canonical_data_home = data_home.canonicalize().with_context(|_| {
111 common_datasource::error::InvalidLocalFileRootSnafu {
112 root: data_home.display().to_string(),
113 }
114 })?;
115 let canonical_copy_root = access.sandbox_root().with_context(|| {
116 common_datasource::error::InvalidLocalFileRootConfigSnafu {
117 root: copy_root.display().to_string(),
118 reason: "sandboxed local file access has no root".to_string(),
119 }
120 })?;
121 let default_copy_root = canonical_data_home.join("copy");
122 let exposes_internal_files = canonical_data_home.starts_with(canonical_copy_root)
123 || (canonical_copy_root.starts_with(&canonical_data_home)
124 && !canonical_copy_root.starts_with(default_copy_root));
125 if exposes_internal_files {
126 return common_datasource::error::InvalidLocalFileRootConfigSnafu {
127 root: copy_root.display().to_string(),
128 reason: "copy_root must not expose files in data_home outside data_home/copy"
129 .to_string(),
130 }
131 .fail();
132 }
133 }
134
135 Ok(access)
136}
137
138#[derive(Parser)]
139pub struct Command {
140 #[clap(subcommand)]
141 subcmd: SubCommand,
142}
143
144impl Command {
145 pub async fn build(&self, opts: GreptimeOptions<StandaloneOptions>) -> Result<Instance> {
146 self.subcmd.build(opts).await
147 }
148
149 pub fn load_options(
150 &self,
151 global_options: &GlobalOptions,
152 ) -> Result<GreptimeOptions<StandaloneOptions>> {
153 self.subcmd.load_options(global_options)
154 }
155
156 pub fn is_daemon(&self) -> bool {
158 self.subcmd.is_daemon()
159 }
160}
161
162#[derive(Parser)]
163enum SubCommand {
164 Start(StartCommand),
165}
166
167impl SubCommand {
168 async fn build(&self, opts: GreptimeOptions<StandaloneOptions>) -> Result<Instance> {
169 match self {
170 SubCommand::Start(cmd) => cmd.build(opts).await,
171 }
172 }
173
174 fn load_options(
175 &self,
176 global_options: &GlobalOptions,
177 ) -> Result<GreptimeOptions<StandaloneOptions>> {
178 match self {
179 SubCommand::Start(cmd) => cmd.load_options(global_options),
180 }
181 }
182
183 fn is_daemon(&self) -> bool {
184 match self {
185 SubCommand::Start(cmd) => cmd.is_daemon(),
186 }
187 }
188}
189
190pub struct Instance {
191 datanode: Datanode,
192 frontend: Frontend,
193 flownode: FlownodeInstance,
194 procedure_manager: ProcedureManagerRef,
195 leader_services_controller: Box<dyn StandaloneLeaderServicesController>,
196 leader_services_context: LeaderServicesContext,
197 _guard: Vec<WorkerGuard>,
199}
200
201impl Instance {
202 pub fn server_addr(&self, name: &str) -> Option<SocketAddr> {
204 self.frontend.server_handlers().addr(name)
205 }
206
207 pub fn mut_frontend(&mut self) -> &mut Frontend {
210 &mut self.frontend
211 }
212
213 pub fn datanode(&self) -> &Datanode {
216 &self.datanode
217 }
218}
219
220#[async_trait]
221impl App for Instance {
222 fn name(&self) -> &str {
223 APP_NAME
224 }
225
226 async fn start(&mut self) -> Result<()> {
227 self.datanode.start_telemetry();
228
229 self.leader_services_controller
230 .start(self.leader_services_context.clone())
231 .await?;
232
233 plugins::start_frontend_plugins(&self.frontend.instance)
234 .await
235 .context(error::StartFrontendSnafu)?;
236
237 self.frontend
238 .start()
239 .await
240 .context(error::StartFrontendSnafu)?;
241
242 self.flownode.start().await.context(StartFlownodeSnafu)?;
243
244 Ok(())
245 }
246
247 async fn stop(&mut self) -> Result<()> {
248 self.frontend
249 .shutdown()
250 .await
251 .context(error::ShutdownFrontendSnafu)?;
252
253 self.leader_services_controller
254 .stop(
255 self.procedure_manager.clone(),
256 self.datanode.region_server(),
257 )
258 .await?;
259
260 self.datanode
261 .shutdown()
262 .await
263 .context(error::ShutdownDatanodeSnafu)?;
264
265 self.flownode
266 .shutdown()
267 .await
268 .context(error::ShutdownFlownodeSnafu)?;
269
270 info!("Datanode instance stopped.");
271
272 Ok(())
273 }
274}
275
276#[derive(Debug, Default, Parser)]
277pub struct StartCommand {
278 #[clap(long)]
279 http_addr: Option<String>,
280 #[clap(long = "grpc-bind-addr", alias = "rpc-bind-addr", alias = "rpc-addr")]
281 grpc_bind_addr: Option<String>,
282 #[clap(long)]
283 mysql_addr: Option<String>,
284 #[clap(long)]
285 postgres_addr: Option<String>,
286 #[clap(short, long)]
287 influxdb_enable: bool,
288 #[clap(short, long)]
289 pub config_file: Option<String>,
290 #[clap(long)]
291 tls_mode: Option<TlsMode>,
292 #[clap(long)]
293 tls_cert_path: Option<String>,
294 #[clap(long)]
295 tls_key_path: Option<String>,
296 #[clap(long)]
297 tls_watch: bool,
298 #[clap(long)]
299 user_provider: Option<String>,
300 #[clap(long, default_value = "GREPTIMEDB_STANDALONE")]
301 pub env_prefix: String,
302 #[clap(long)]
304 data_home: Option<String>,
305 #[cfg(unix)]
307 #[clap(short, long)]
308 daemon: bool,
309}
310
311impl StartCommand {
312 #[cfg(unix)]
314 pub(crate) fn is_daemon(&self) -> bool {
315 self.daemon
316 }
317
318 #[cfg(not(unix))]
319 pub(crate) fn is_daemon(&self) -> bool {
320 false
321 }
322
323 pub fn load_options(
325 &self,
326 global_options: &GlobalOptions,
327 ) -> Result<GreptimeOptions<StandaloneOptions>> {
328 let mut opts = GreptimeOptions::<StandaloneOptions>::load_layered_options(
329 self.config_file.as_deref(),
330 self.env_prefix.as_ref(),
331 )
332 .context(error::LoadLayeredConfigSnafu)?;
333
334 self.merge_with_cli_options(global_options, &mut opts.component)?;
335 opts.component.sanitize();
336
337 Ok(opts)
338 }
339
340 pub fn merge_with_cli_options(
342 &self,
343 global_options: &GlobalOptions,
344 opts: &mut StandaloneOptions,
345 ) -> Result<()> {
346 if let Some(dir) = &global_options.log_dir {
347 opts.logging.dir.clone_from(dir);
348 }
349
350 if global_options.log_level.is_some() {
351 opts.logging.level.clone_from(&global_options.log_level);
352 }
353
354 opts.tracing = TracingOptions {
355 #[cfg(feature = "tokio-console")]
356 tokio_console_addr: global_options.tokio_console_addr.clone(),
357 };
358
359 let tls_opts = TlsOption::new(
360 self.tls_mode,
361 self.tls_cert_path.clone(),
362 self.tls_key_path.clone(),
363 self.tls_watch,
364 );
365
366 if let Some(addr) = &self.http_addr {
367 opts.http.addr.clone_from(addr);
368 }
369
370 if let Some(data_home) = &self.data_home {
371 opts.storage.data_home.clone_from(data_home);
372 }
373
374 if opts.logging.dir.is_empty() {
376 opts.logging.dir = Path::new(&opts.storage.data_home)
377 .join(DEFAULT_LOGGING_DIR)
378 .to_string_lossy()
379 .to_string();
380 }
381
382 if let Some(addr) = &self.grpc_bind_addr {
383 let datanode_grpc_addr = DatanodeOptions::default().grpc.bind_addr;
385 if addr.eq(&datanode_grpc_addr) {
386 return error::IllegalConfigSnafu {
387 msg: format!(
388 "gRPC listen address conflicts with datanode reserved gRPC addr: {datanode_grpc_addr}",
389 ),
390 }.fail();
391 }
392 opts.grpc.bind_addr.clone_from(addr);
393 opts.grpc.tls = merge_tls_option(&opts.grpc.tls, tls_opts.clone());
394 }
395
396 if let Some(addr) = &self.mysql_addr {
397 opts.mysql.enable = true;
398 opts.mysql.addr.clone_from(addr);
399 opts.mysql.tls = merge_tls_option(&opts.mysql.tls, tls_opts.clone());
400 }
401
402 if let Some(addr) = &self.postgres_addr {
403 opts.postgres.enable = true;
404 opts.postgres.addr.clone_from(addr);
405 opts.postgres.tls = merge_tls_option(&opts.postgres.tls, tls_opts.clone());
406 }
407
408 if self.influxdb_enable {
409 opts.influxdb.enable = self.influxdb_enable;
410 }
411
412 if let Some(user_provider) = &self.user_provider {
413 opts.user_provider = Some(user_provider.clone());
414 }
415
416 Ok(())
417 }
418
419 #[allow(unreachable_code)]
420 #[allow(unused_variables)]
421 #[allow(clippy::diverging_sub_expression)]
422 pub async fn build(&self, opts: GreptimeOptions<StandaloneOptions>) -> Result<Instance> {
424 let guard = common_telemetry::init_global_logging(
425 APP_NAME,
426 &opts.component.logging,
427 &opts.component.tracing,
428 None,
429 Some(&opts.component.slow_query),
430 );
431
432 common_runtime::init_standalone_runtimes(&opts.runtime);
433
434 crate::options::flush_dropped_plugin_warnings();
435 log_versions(verbose_version(), short_version(), APP_NAME);
436 maybe_activate_heap_profile(&opts.component.memory);
437 create_resource_limit_metrics(APP_NAME);
438
439 info!("Standalone start command: {:#?}", self);
440 info!("Standalone options: {opts:#?}");
441
442 let (mut instance, _) =
443 Self::build_with(opts.component, opts.plugins, InstanceCreator::default()).await?;
444 instance._guard.extend(guard);
445 Ok(instance)
446 }
447
448 pub async fn build_with(
449 mut opts: StandaloneOptions,
450 plugin_opts: Vec<PluginOptions>,
451 creator: InstanceCreator,
452 ) -> Result<(Instance, InstanceCreatorResult)> {
453 let mut plugins = Plugins::new();
454 plugins.insert(StandaloneFlag);
455 set_default_prefix(opts.default_column_prefix.as_deref())
456 .map_err(BoxedError::new)
457 .context(error::BuildCliSnafu)?;
458
459 opts.grpc.detect_server_addr();
460 let fe_opts = opts.frontend_options();
461 let dn_opts = opts.datanode_options();
462 let node_id = dn_opts.node_id;
463 let init_regions_parallelism = dn_opts.init_regions_parallelism;
464
465 plugins::setup_frontend_plugins_pre_build(&mut plugins, &plugin_opts, &fe_opts, None)
466 .await
467 .context(error::StartFrontendSnafu)?;
468
469 plugins::setup_datanode_plugins_pre_build(&mut plugins, &plugin_opts, &dn_opts)
470 .await
471 .context(error::StartDatanodeSnafu)?;
472
473 set_default_timezone(fe_opts.default_timezone.as_deref())
474 .context(error::InitTimezoneSnafu)?;
475
476 let data_home = &dn_opts.storage.data_home;
477 fs::create_dir_all(path::Path::new(data_home))
479 .context(error::CreateDirSnafu { dir: data_home })?;
480 let local_file_access = standalone_local_file_access(&dn_opts.storage)
481 .map_err(BoxedError::new)
482 .context(OtherSnafu)?;
483
484 let metadata_dir = metadata_store_dir(data_home);
485 let kv_backend = creator
486 .metadata_kv_backend_creator
487 .create(metadata_dir, &opts)
488 .await?;
489 let (procedure_manager, event_recorder_handle) =
490 standalone::build_procedure_manager(kv_backend.clone(), opts.procedure);
491
492 plugins::setup_standalone_plugins(&mut plugins, &plugin_opts, &opts, kv_backend.clone())
493 .await
494 .context(error::SetupStandalonePluginsSnafu)?;
495
496 let layered_cache_builder = LayeredCacheRegistryBuilder::default();
498 let fundamental_cache_registry = build_fundamental_cache_registry(kv_backend.clone());
499 let mut layered_cache_builder = with_default_composite_cache_registry(
500 layered_cache_builder.add_cache_registry(fundamental_cache_registry),
501 )
502 .context(error::BuildCacheRegistrySnafu)?;
503
504 if let Some(plugin_cache_builder) = plugins::standalone::configure_cache_registry(&plugins)
505 {
506 layered_cache_builder =
507 layered_cache_builder.add_cache_registry(plugin_cache_builder.build());
508 }
509
510 let layered_cache_registry = Arc::new(layered_cache_builder.build());
511
512 let mut builder = DatanodeBuilder::new(dn_opts, plugins.clone(), kv_backend.clone());
513 builder.with_cache_registry(layered_cache_registry.clone());
514 builder.with_local_file_access(local_file_access.clone());
515 if let Some(writable) = creator.open_regions_writable_override {
516 builder.with_open_regions_writable_override(writable);
517 }
518
519 plugins::setup_datanode_plugins_post_build(&mut plugins, &plugin_opts, &builder)
520 .await
521 .context(error::StartDatanodeSnafu)?;
522 builder.set_plugins(plugins.clone());
523
524 let datanode = builder.build().await.context(error::StartDatanodeSnafu)?;
525
526 let information_extension = Arc::new(StandaloneInformationExtension::new(
527 datanode.region_server(),
528 procedure_manager.clone(),
529 ));
530
531 plugins.insert::<InformationExtensionRef>(information_extension.clone());
532
533 let process_manager = Arc::new(ProcessManager::new(opts.grpc.server_addr.clone(), None));
534
535 let (frontend_client, frontend_instance_handler) =
538 FrontendClient::from_empty_grpc_handler(opts.query.clone());
539 let frontend_client = Arc::new(frontend_client);
540
541 let builder = KvBackendCatalogManagerBuilder::new(
542 information_extension.clone(),
543 kv_backend.clone(),
544 layered_cache_registry.clone(),
545 )
546 .with_procedure_manager(procedure_manager.clone())
547 .with_process_manager(process_manager.clone());
548 let builder = if let Some(configurator) =
549 plugins.get::<CatalogManagerConfiguratorRef<CatalogManagerConfigureContext>>()
550 {
551 let ctx = StandaloneCatalogManagerConfigureContext {
552 fe_client: frontend_client.clone(),
553 };
554 let ctx = CatalogManagerConfigureContext::Standalone(ctx);
555 configurator
556 .configure(builder, ctx)
557 .await
558 .context(OtherSnafu)?
559 } else {
560 builder
561 };
562 let catalog_manager = builder.build();
563
564 let table_metadata_manager =
565 Self::create_table_metadata_manager(kv_backend.clone()).await?;
566
567 let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
568 let flownode_options = FlownodeOptions {
569 flow: opts.flow.clone(),
570 ..Default::default()
571 };
572
573 let mut flow_builder = FlownodeBuilder::new(
574 flownode_options,
575 plugins.clone(),
576 table_metadata_manager.clone(),
577 catalog_manager.clone(),
578 flow_metadata_manager.clone(),
579 frontend_client.clone(),
580 );
581
582 plugins::setup_flownode_plugins_post_build(&mut plugins, &plugin_opts, &flow_builder)
583 .await
584 .context(error::StartFlownodeSnafu)?;
585 flow_builder.set_plugins(plugins.clone());
586
587 let flownode = flow_builder
588 .build()
589 .await
590 .map_err(BoxedError::new)
591 .context(error::OtherSnafu)?;
592 let flow_engine = flownode.flow_engine();
593
594 {
596 information_extension
597 .set_flow_engine(flow_engine.clone())
598 .await;
599 }
600
601 let node_manager = creator
602 .node_manager_creator
603 .create(&kv_backend, datanode.region_server(), flow_engine.clone())
604 .await?;
605
606 let table_id_allocator = creator.table_id_allocator_creator.create(&kv_backend);
607 let flow_id_sequence = Arc::new(
608 SequenceBuilder::new(FLOW_ID_SEQ, kv_backend.clone())
609 .initial(MIN_USER_FLOW_ID as u64)
610 .step(10)
611 .build(),
612 );
613 let kafka_options = opts
614 .wal
615 .clone()
616 .try_into()
617 .context(error::InvalidWalProviderSnafu)?;
618 let wal_provider = build_wal_provider(&kafka_options, kv_backend.clone())
619 .await
620 .context(error::BuildWalProviderSnafu)?;
621 let wal_provider = Arc::new(wal_provider);
622 let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
623 table_id_allocator.clone(),
624 wal_provider.clone(),
625 ));
626 let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
627 flow_id_sequence,
628 ));
629
630 let ddl_context = DdlContext {
631 node_manager: node_manager.clone(),
632 cache_invalidator: layered_cache_registry.clone(),
633 memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
634 leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
635 table_metadata_manager: table_metadata_manager.clone(),
636 table_metadata_allocator: table_metadata_allocator.clone(),
637 flow_metadata_manager: flow_metadata_manager.clone(),
638 flow_metadata_allocator: flow_metadata_allocator.clone(),
639 region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
640 soft_drop_enabled: false,
641 soft_drop_retention: None,
642 create_database_metadata_committer: None,
643 };
644
645 let ddl_manager = DdlManager::new(
646 ddl_context,
647 procedure_manager.clone(),
648 Arc::new(StandaloneRepartitionProcedureFactory),
649 );
650
651 let ddl_manager = if let Some(configurator) =
652 plugins.get::<DdlManagerConfiguratorRef<DdlManagerConfigureContext>>()
653 {
654 let ctx = DdlManagerConfigureContext {
655 kv_backend: kv_backend.clone(),
656 fe_client: frontend_client.clone(),
657 catalog_manager: catalog_manager.clone(),
658 };
659 configurator
660 .configure(ddl_manager, ctx)
661 .await
662 .context(OtherSnafu)?
663 } else {
664 ddl_manager
665 };
666 ddl_manager
667 .register_loaders()
668 .context(error::InitDdlManagerSnafu)?;
669
670 let procedure_executor = creator
671 .procedure_executor_creator
672 .create(Arc::new(ddl_manager), procedure_manager.clone())
673 .await?;
674
675 let fe_instance = FrontendBuilder::new(
676 fe_opts.clone(),
677 kv_backend.clone(),
678 layered_cache_registry.clone(),
679 catalog_manager.clone(),
680 node_manager.clone(),
681 procedure_executor.clone(),
682 process_manager,
683 )
684 .with_local_file_access(local_file_access);
685
686 plugins::setup_frontend_plugins_post_build(&mut plugins, &plugin_opts, &fe_instance)
687 .await
688 .context(error::StartFrontendSnafu)?;
689
690 let fe_instance = fe_instance
691 .with_plugin(plugins.clone())
692 .try_build()
693 .await
694 .context(error::StartFrontendSnafu)?;
695 let fe_instance = Arc::new(fe_instance);
696
697 event_recorder_handle.install(fe_instance.event_recorder());
698
699 let grpc_handler = fe_instance.clone() as Arc<dyn GrpcQueryHandlerWithBoxedError>;
701 let weak_grpc_handler = Arc::downgrade(&grpc_handler);
702 frontend_instance_handler
703 .set_handler(weak_grpc_handler)
704 .await;
705
706 let flow_streaming_engine = flow_engine.streaming_engine();
708 let invoker = FrontendInvoker::build_from(
710 flow_streaming_engine.clone(),
711 catalog_manager.clone(),
712 kv_backend.clone(),
713 layered_cache_registry.clone(),
714 procedure_executor,
715 node_manager.clone(),
716 fe_instance.frontend_peer_addr().to_string(),
717 )
718 .await
719 .context(StartFlownodeSnafu)?;
720 flow_streaming_engine.set_frontend_invoker(invoker).await;
721
722 let servers = Services::new(opts, fe_instance.clone(), plugins.clone())
723 .build()
724 .context(error::StartFrontendSnafu)?;
725
726 let frontend = Frontend {
727 instance: fe_instance,
728 servers,
729 heartbeat_task: None,
730 };
731 let leader_services_context = LeaderServicesContext {
732 procedure_manager: procedure_manager.clone(),
733 wal_provider: wal_provider.clone(),
734 region_server: datanode.region_server(),
735 kv_backend: kv_backend.clone(),
736 cache_registry: layered_cache_registry,
737 catalog_manager,
738 flow_engine,
739 frontend_client,
740 node_id,
741 init_regions_parallelism,
742 plugin_options: plugin_opts,
743 };
744
745 let instance = Instance {
746 datanode,
747 frontend,
748 flownode,
749 procedure_manager,
750 leader_services_controller: creator.leader_services_controller,
751 leader_services_context,
752 _guard: vec![],
753 };
754 let result = InstanceCreatorResult {
755 kv_backend,
756 node_manager,
757 table_id_allocator,
758 };
759 Ok((instance, result))
760 }
761
762 pub async fn create_table_metadata_manager(
763 kv_backend: KvBackendRef,
764 ) -> Result<TableMetadataManagerRef> {
765 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend));
766
767 table_metadata_manager
768 .init()
769 .await
770 .context(error::InitMetadataSnafu)?;
771
772 Ok(table_metadata_manager)
773 }
774}
775
776#[async_trait]
777pub trait NodeManagerCreator: Send + Sync {
778 async fn create(
779 &self,
780 kv_backend: &KvBackendRef,
781 region_server: RegionServer,
782 flow_server: FlownodeRef,
783 ) -> Result<NodeManagerRef>;
784}
785
786pub struct DefaultNodeManagerCreator;
787
788#[async_trait]
789impl NodeManagerCreator for DefaultNodeManagerCreator {
790 async fn create(
791 &self,
792 _: &KvBackendRef,
793 region_server: RegionServer,
794 flow_server: FlownodeRef,
795 ) -> Result<NodeManagerRef> {
796 Ok(Arc::new(StandaloneDatanodeManager {
797 region_server,
798 flow_server,
799 }))
800 }
801}
802
803#[async_trait]
809pub trait MetadataKvBackendCreator: Send + Sync {
810 async fn create(&self, metadata_dir: String, opts: &StandaloneOptions) -> Result<KvBackendRef>;
811}
812
813pub struct DefaultMetadataKvBackendCreator;
814
815#[async_trait]
816impl MetadataKvBackendCreator for DefaultMetadataKvBackendCreator {
817 async fn create(&self, metadata_dir: String, opts: &StandaloneOptions) -> Result<KvBackendRef> {
818 standalone::build_metadata_kvbackend(metadata_dir, opts.metadata_store)
819 .context(error::BuildMetadataKvbackendSnafu)
820 }
821}
822
823pub trait TableIdAllocatorCreator: Send + Sync {
824 fn create(&self, kv_backend: &KvBackendRef) -> Arc<Sequence>;
825}
826
827struct DefaultTableIdAllocatorCreator;
828
829impl TableIdAllocatorCreator for DefaultTableIdAllocatorCreator {
830 fn create(&self, kv_backend: &KvBackendRef) -> Arc<Sequence> {
831 Arc::new(
832 SequenceBuilder::new(TABLE_ID_SEQ, kv_backend.clone())
833 .initial(MIN_USER_TABLE_ID as u64)
834 .step(10)
835 .build(),
836 )
837 }
838}
839
840#[async_trait]
841pub trait ProcedureExecutorCreator: Send + Sync {
842 async fn create(
843 &self,
844 ddl_manager: DdlManagerRef,
845 procedure_manager: ProcedureManagerRef,
846 ) -> Result<ProcedureExecutorRef>;
847}
848
849pub struct DefaultProcedureExecutorCreator;
850
851#[async_trait]
852impl ProcedureExecutorCreator for DefaultProcedureExecutorCreator {
853 async fn create(
854 &self,
855 ddl_manager: DdlManagerRef,
856 procedure_manager: ProcedureManagerRef,
857 ) -> Result<ProcedureExecutorRef> {
858 Ok(Arc::new(LocalProcedureExecutor::new(
859 ddl_manager,
860 procedure_manager,
861 )))
862 }
863}
864
865#[async_trait]
866pub trait StandaloneLeaderServicesController: Send + Sync {
867 async fn start(&self, context: LeaderServicesContext) -> Result<()>;
872
873 async fn stop(
875 &self,
876 procedure_manager: ProcedureManagerRef,
877 region_server: RegionServer,
878 ) -> Result<()>;
879}
880
881#[derive(Clone)]
882pub struct LeaderServicesContext {
892 pub procedure_manager: ProcedureManagerRef,
893 pub wal_provider: WalProviderRef,
894 pub region_server: RegionServer,
895 pub kv_backend: KvBackendRef,
896 pub cache_registry: LayeredCacheRegistryRef,
897 pub catalog_manager: CatalogManagerRef,
898 pub flow_engine: FlowDualEngineRef,
899 pub frontend_client: Arc<FrontendClient>,
900 pub node_id: Option<DatanodeId>,
901 pub init_regions_parallelism: usize,
902 pub plugin_options: Vec<PluginOptions>,
903}
904
905pub struct DefaultStandaloneLeaderServicesController;
906
907#[async_trait]
908impl StandaloneLeaderServicesController for DefaultStandaloneLeaderServicesController {
909 async fn start(&self, context: LeaderServicesContext) -> Result<()> {
910 context
911 .procedure_manager
912 .start()
913 .await
914 .context(error::StartProcedureManagerSnafu)?;
915 context
916 .wal_provider
917 .start()
918 .await
919 .context(error::StartWalProviderSnafu)
920 }
921
922 async fn stop(
923 &self,
924 procedure_manager: ProcedureManagerRef,
925 _region_server: RegionServer,
926 ) -> Result<()> {
927 procedure_manager
928 .stop()
929 .await
930 .context(error::StopProcedureManagerSnafu)
931 }
932}
933
934pub struct InstanceCreator {
937 metadata_kv_backend_creator: Box<dyn MetadataKvBackendCreator>,
940 node_manager_creator: Box<dyn NodeManagerCreator>,
941 table_id_allocator_creator: Box<dyn TableIdAllocatorCreator>,
942 procedure_executor_creator: Box<dyn ProcedureExecutorCreator>,
943 leader_services_controller: Box<dyn StandaloneLeaderServicesController>,
944 open_regions_writable_override: Option<bool>,
945}
946
947impl InstanceCreator {
948 pub fn new(
949 node_manager_creator: Box<dyn NodeManagerCreator>,
950 table_id_allocator_creator: Box<dyn TableIdAllocatorCreator>,
951 procedure_executor_creator: Box<dyn ProcedureExecutorCreator>,
952 ) -> Self {
953 Self {
954 metadata_kv_backend_creator: Box::new(DefaultMetadataKvBackendCreator),
955 node_manager_creator,
956 table_id_allocator_creator,
957 procedure_executor_creator,
958 leader_services_controller: Box::new(DefaultStandaloneLeaderServicesController),
959 open_regions_writable_override: None,
960 }
961 }
962
963 pub fn with_metadata_kv_backend_creator(
964 mut self,
965 metadata_kv_backend_creator: Box<dyn MetadataKvBackendCreator>,
966 ) -> Self {
967 self.metadata_kv_backend_creator = metadata_kv_backend_creator;
968 self
969 }
970
971 pub fn map_metadata_kv_backend_creator<F>(mut self, f: F) -> Self
976 where
977 F: FnOnce(Box<dyn MetadataKvBackendCreator>) -> Box<dyn MetadataKvBackendCreator>,
978 {
979 self.metadata_kv_backend_creator = f(self.metadata_kv_backend_creator);
980 self
981 }
982
983 pub fn map_node_manager_creator<F>(mut self, f: F) -> Self
985 where
986 F: FnOnce(Box<dyn NodeManagerCreator>) -> Box<dyn NodeManagerCreator>,
987 {
988 self.node_manager_creator = f(self.node_manager_creator);
989 self
990 }
991
992 pub fn map_procedure_executor_creator<F>(mut self, f: F) -> Self
994 where
995 F: FnOnce(Box<dyn ProcedureExecutorCreator>) -> Box<dyn ProcedureExecutorCreator>,
996 {
997 self.procedure_executor_creator = f(self.procedure_executor_creator);
998 self
999 }
1000
1001 pub fn with_leader_services_controller(
1003 mut self,
1004 leader_services_controller: Box<dyn StandaloneLeaderServicesController>,
1005 ) -> Self {
1006 self.leader_services_controller = leader_services_controller;
1007 self
1008 }
1009
1010 pub fn with_open_regions_writable_override(mut self, writable: bool) -> Self {
1018 self.open_regions_writable_override = Some(writable);
1019 self
1020 }
1021}
1022
1023impl Default for InstanceCreator {
1024 fn default() -> Self {
1025 Self {
1026 metadata_kv_backend_creator: Box::new(DefaultMetadataKvBackendCreator),
1027 node_manager_creator: Box::new(DefaultNodeManagerCreator),
1028 table_id_allocator_creator: Box::new(DefaultTableIdAllocatorCreator),
1029 procedure_executor_creator: Box::new(DefaultProcedureExecutorCreator),
1030 leader_services_controller: Box::new(DefaultStandaloneLeaderServicesController),
1031 open_regions_writable_override: None,
1032 }
1033 }
1034}
1035
1036pub struct InstanceCreatorResult {
1039 pub kv_backend: KvBackendRef,
1040 pub node_manager: NodeManagerRef,
1041 pub table_id_allocator: Arc<Sequence>,
1042}
1043
1044#[cfg(test)]
1045mod tests {
1046 use std::default::Default;
1047 use std::io::Write;
1048 use std::time::Duration;
1049
1050 use auth::{Identity, Password, UserProviderRef};
1051 use clap::{CommandFactory, Parser};
1052 use common_base::readable_size::ReadableSize;
1053 use common_config::ENV_VAR_SEP;
1054 use common_options::plugin_options::StandaloneFlag;
1055 use common_test_util::temp_dir::{create_named_temp_file, create_temp_dir};
1056 use common_wal::config::DatanodeWalConfig;
1057 use frontend::frontend::FrontendOptions;
1058 use object_store::config::{FileConfig, GcsConfig};
1059 use servers::grpc::GrpcOptions;
1060
1061 use super::*;
1062 use crate::options::GlobalOptions;
1063
1064 #[test]
1065 fn test_standalone_local_file_access_config() {
1066 let data_home = create_temp_dir("standalone_copy_root");
1067 let storage = StorageConfig {
1068 data_home: data_home.path().display().to_string(),
1069 ..Default::default()
1070 };
1071 let access = standalone_local_file_access(&storage).unwrap();
1072 assert_eq!(
1073 access.sandbox_root().unwrap(),
1074 data_home.path().join("copy").canonicalize().unwrap()
1075 );
1076
1077 let remote_data_home = StorageConfig {
1078 data_home: "s3://bucket/data".to_string(),
1079 ..Default::default()
1080 };
1081 assert!(matches!(
1082 standalone_local_file_access(&remote_data_home).unwrap(),
1083 LocalFileAccess::Disabled
1084 ));
1085
1086 let explicit_root = create_temp_dir("standalone_explicit_copy_root");
1087 let remote_with_explicit_root = StorageConfig {
1088 data_home: "s3://bucket/data".to_string(),
1089 copy_root: Some(explicit_root.path().display().to_string()),
1090 ..Default::default()
1091 };
1092 assert_eq!(
1093 standalone_local_file_access(&remote_with_explicit_root)
1094 .unwrap()
1095 .sandbox_root()
1096 .unwrap(),
1097 explicit_root.path().canonicalize().unwrap()
1098 );
1099
1100 let remote_copy_root = StorageConfig {
1101 data_home: data_home.path().display().to_string(),
1102 copy_root: Some("s3://bucket/copy".to_string()),
1103 ..Default::default()
1104 };
1105 assert!(matches!(
1106 standalone_local_file_access(&remote_copy_root),
1107 Err(common_datasource::error::Error::InvalidLocalFileRootConfig { .. })
1108 ));
1109
1110 let exposes_internal = StorageConfig {
1111 data_home: data_home.path().display().to_string(),
1112 copy_root: Some(data_home.path().join("data").display().to_string()),
1113 ..Default::default()
1114 };
1115 assert!(matches!(
1116 standalone_local_file_access(&exposes_internal),
1117 Err(common_datasource::error::Error::InvalidLocalFileRootConfig { .. })
1118 ));
1119
1120 let exposes_data_home = StorageConfig {
1121 data_home: data_home.path().display().to_string(),
1122 copy_root: Some(data_home.path().parent().unwrap().display().to_string()),
1123 ..Default::default()
1124 };
1125 assert!(matches!(
1126 standalone_local_file_access(&exposes_data_home),
1127 Err(common_datasource::error::Error::InvalidLocalFileRootConfig { .. })
1128 ));
1129 }
1130
1131 #[tokio::test]
1132 async fn test_try_from_start_command_to_anymap() {
1133 let fe_opts = FrontendOptions {
1134 user_provider: Some("static_user_provider:cmd:test=test".to_string()),
1135 ..Default::default()
1136 };
1137
1138 let mut plugins = Plugins::new();
1139 plugins.insert(StandaloneFlag);
1140 plugins::setup_frontend_plugins_pre_build(&mut plugins, &[], &fe_opts, None)
1141 .await
1142 .unwrap();
1143
1144 let provider = plugins.get::<UserProviderRef>().unwrap();
1145 let result = provider
1146 .authenticate(
1147 Identity::UserId("test", None),
1148 Password::PlainText("test".to_string().into()),
1149 )
1150 .await;
1151 let _ = result.unwrap();
1152 }
1153
1154 #[test]
1155 fn test_toml() {
1156 let opts = StandaloneOptions::default();
1157 let toml_string = toml::to_string(&opts).unwrap();
1158 assert!(toml_string.contains("experimental_enable_exponential_histogram = false"));
1159 let parsed: StandaloneOptions = toml::from_str(&toml_string).unwrap();
1160 assert_eq!(parsed.otlp, opts.otlp);
1161 }
1162
1163 #[test]
1164 fn test_read_from_config_file() {
1165 let mut file = create_named_temp_file();
1166 let toml_str = r#"
1167 enable_memory_catalog = true
1168
1169 [wal]
1170 provider = "raft_engine"
1171 dir = "./greptimedb_data/test/wal"
1172 file_size = "1GB"
1173 purge_threshold = "50GB"
1174 purge_interval = "10m"
1175 read_batch_size = 128
1176 sync_write = false
1177
1178 [storage]
1179 data_home = "./greptimedb_data/"
1180 type = "File"
1181
1182 [[storage.providers]]
1183 type = "Gcs"
1184 bucket = "foo"
1185 endpoint = "bar"
1186
1187 [[storage.providers]]
1188 type = "S3"
1189 access_key_id = "access_key_id"
1190 secret_access_key = "secret_access_key"
1191
1192 [storage.compaction]
1193 max_inflight_tasks = 3
1194 max_files_in_level0 = 7
1195 max_purge_tasks = 32
1196
1197 [storage.manifest]
1198 checkpoint_margin = 9
1199 gc_duration = '7s'
1200
1201 [http]
1202 addr = "127.0.0.1:4000"
1203 timeout = "33s"
1204 body_limit = "128MB"
1205
1206 [opentsdb]
1207 enable = true
1208
1209 [logging]
1210 level = "debug"
1211 dir = "./greptimedb_data/test/logs"
1212 "#;
1213 write!(file, "{}", toml_str).unwrap();
1214 let cmd = StartCommand {
1215 config_file: Some(file.path().to_str().unwrap().to_string()),
1216 user_provider: Some("static_user_provider:cmd:test=test".to_string()),
1217 ..Default::default()
1218 };
1219
1220 let options = cmd
1221 .load_options(&GlobalOptions::default())
1222 .unwrap()
1223 .component;
1224 let fe_opts = options.frontend_options();
1225 let dn_opts = options.datanode_options();
1226 let logging_opts = options.logging;
1227 assert_eq!("127.0.0.1:4000".to_string(), fe_opts.http.addr);
1228 assert_eq!(Duration::from_secs(33), fe_opts.http.timeout);
1229 assert_eq!(ReadableSize::mb(128), fe_opts.http.body_limit);
1230 assert_eq!("127.0.0.1:4001".to_string(), fe_opts.grpc.bind_addr);
1231 assert!(fe_opts.mysql.enable);
1232 assert_eq!("127.0.0.1:4002", fe_opts.mysql.addr);
1233 assert_eq!(2, fe_opts.mysql.runtime_size);
1234 assert_eq!(None, fe_opts.mysql.reject_no_database);
1235 assert!(fe_opts.influxdb.enable);
1236 assert!(fe_opts.opentsdb.enable);
1237
1238 let DatanodeWalConfig::RaftEngine(raft_engine_config) = dn_opts.wal else {
1239 unreachable!()
1240 };
1241 assert_eq!(
1242 "./greptimedb_data/test/wal",
1243 raft_engine_config.dir.unwrap()
1244 );
1245
1246 assert!(matches!(
1247 &dn_opts.storage.store,
1248 object_store::config::ObjectStoreConfig::File(FileConfig { .. })
1249 ));
1250 assert_eq!(dn_opts.storage.providers.len(), 2);
1251 assert!(matches!(
1252 dn_opts.storage.providers[0],
1253 object_store::config::ObjectStoreConfig::Gcs(GcsConfig { .. })
1254 ));
1255 match &dn_opts.storage.providers[1] {
1256 object_store::config::ObjectStoreConfig::S3(s3_config) => {
1257 assert_eq!(
1258 "SecretBox<alloc::string::String>([REDACTED])".to_string(),
1259 format!("{:?}", s3_config.connection.access_key_id)
1260 );
1261 }
1262 _ => {
1263 unreachable!()
1264 }
1265 }
1266
1267 assert_eq!("debug", logging_opts.level.as_ref().unwrap());
1268 assert_eq!("./greptimedb_data/test/logs".to_string(), logging_opts.dir);
1269 }
1270
1271 #[test]
1272 fn test_load_log_options_from_cli() {
1273 let cmd = StartCommand {
1274 user_provider: Some("static_user_provider:cmd:test=test".to_string()),
1275 mysql_addr: Some("127.0.0.1:4002".to_string()),
1276 postgres_addr: Some("127.0.0.1:4003".to_string()),
1277 ..Default::default()
1278 };
1279
1280 let opts = cmd
1281 .load_options(&GlobalOptions {
1282 log_dir: Some("./greptimedb_data/test/logs".to_string()),
1283 log_level: Some("debug".to_string()),
1284
1285 #[cfg(feature = "tokio-console")]
1286 tokio_console_addr: None,
1287 })
1288 .unwrap()
1289 .component;
1290
1291 assert_eq!("./greptimedb_data/test/logs", opts.logging.dir);
1292 assert_eq!("debug", opts.logging.level.unwrap());
1293 }
1294
1295 #[test]
1296 fn test_config_precedence_order() {
1297 let mut file = create_named_temp_file();
1298 let toml_str = r#"
1299 [http]
1300 addr = "127.0.0.1:4000"
1301
1302 [logging]
1303 level = "debug"
1304 "#;
1305 write!(file, "{}", toml_str).unwrap();
1306
1307 let env_prefix = "STANDALONE_UT";
1308 temp_env::with_vars(
1309 [
1310 (
1311 [
1313 env_prefix.to_string(),
1314 "logging".to_uppercase(),
1315 "dir".to_uppercase(),
1316 ]
1317 .join(ENV_VAR_SEP),
1318 Some("/other/log/dir"),
1319 ),
1320 (
1321 [
1323 env_prefix.to_string(),
1324 "logging".to_uppercase(),
1325 "level".to_uppercase(),
1326 ]
1327 .join(ENV_VAR_SEP),
1328 Some("info"),
1329 ),
1330 (
1331 [
1333 env_prefix.to_string(),
1334 "http".to_uppercase(),
1335 "addr".to_uppercase(),
1336 ]
1337 .join(ENV_VAR_SEP),
1338 Some("127.0.0.1:24000"),
1339 ),
1340 ],
1341 || {
1342 let command = StartCommand {
1343 config_file: Some(file.path().to_str().unwrap().to_string()),
1344 http_addr: Some("127.0.0.1:14000".to_string()),
1345 env_prefix: env_prefix.to_string(),
1346 ..Default::default()
1347 };
1348
1349 let opts = command.load_options(&Default::default()).unwrap().component;
1350
1351 assert_eq!(opts.logging.dir, "/other/log/dir");
1353
1354 assert_eq!(opts.logging.level.as_ref().unwrap(), "debug");
1356
1357 let fe_opts = opts.frontend_options();
1359 assert_eq!(fe_opts.http.addr, "127.0.0.1:14000");
1360 assert_eq!(ReadableSize::mb(64), fe_opts.http.body_limit);
1361
1362 assert_eq!(fe_opts.grpc.bind_addr, GrpcOptions::default().bind_addr);
1364 },
1365 );
1366 }
1367
1368 #[test]
1369 fn test_parse_grpc_bind_addr_aliases() {
1370 let command =
1371 StartCommand::try_parse_from(["standalone", "--grpc-bind-addr", "127.0.0.1:14001"])
1372 .unwrap();
1373 assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:14001"));
1374
1375 let command =
1376 StartCommand::try_parse_from(["standalone", "--rpc-bind-addr", "127.0.0.1:24001"])
1377 .unwrap();
1378 assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:24001"));
1379
1380 let command =
1381 StartCommand::try_parse_from(["standalone", "--rpc-addr", "127.0.0.1:34001"]).unwrap();
1382 assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:34001"));
1383 }
1384
1385 #[cfg(unix)]
1386 #[test]
1387 fn test_parse_daemon_flag() {
1388 let command = StartCommand::try_parse_from(["standalone", "--daemon"]).unwrap();
1389 assert!(command.is_daemon());
1390
1391 let command = StartCommand::try_parse_from(["standalone", "-d"]).unwrap();
1392 assert!(command.is_daemon());
1393
1394 let command = StartCommand::try_parse_from(["standalone"]).unwrap();
1395 assert!(!command.is_daemon());
1396 }
1397
1398 #[test]
1399 fn test_help_uses_grpc_option_names() {
1400 let mut cmd = StartCommand::command();
1401 let mut help = Vec::new();
1402 cmd.write_long_help(&mut help).unwrap();
1403 let help = String::from_utf8(help).unwrap();
1404
1405 assert!(help.contains("--grpc-bind-addr"));
1406 assert!(!help.contains("--rpc-bind-addr"));
1407 assert!(!help.contains("--rpc-addr"));
1408 }
1409
1410 #[test]
1411 fn test_load_default_standalone_options() {
1412 let options =
1413 StandaloneOptions::load_layered_options(None, "GREPTIMEDB_STANDALONE").unwrap();
1414 let default_options = StandaloneOptions::default();
1415 assert_eq!(options.enable_telemetry, default_options.enable_telemetry);
1416 assert_eq!(options.http, default_options.http);
1417 assert_eq!(options.grpc, default_options.grpc);
1418 assert_eq!(options.mysql, default_options.mysql);
1419 assert_eq!(options.postgres, default_options.postgres);
1420 assert_eq!(options.opentsdb, default_options.opentsdb);
1421 assert_eq!(options.influxdb, default_options.influxdb);
1422 assert_eq!(options.prom_store, default_options.prom_store);
1423 assert_eq!(options.wal, default_options.wal);
1424 assert_eq!(options.metadata_store, default_options.metadata_store);
1425 assert_eq!(options.procedure, default_options.procedure);
1426 assert_eq!(options.logging, default_options.logging);
1427 assert_eq!(options.region_engine, default_options.region_engine);
1428 }
1429
1430 #[test]
1431 fn test_cache_config() {
1432 let toml_str = r#"
1433 [storage]
1434 data_home = "test_data_home"
1435 type = "S3"
1436 [storage.cache_config]
1437 enable_read_cache = true
1438 "#;
1439 let mut opts: StandaloneOptions = toml::from_str(toml_str).unwrap();
1440 opts.sanitize();
1441 assert!(opts.storage.store.cache_config().unwrap().enable_read_cache);
1442 assert_eq!(
1443 opts.storage.store.cache_config().unwrap().cache_path,
1444 "test_data_home"
1445 );
1446 }
1447
1448 #[test]
1449 #[cfg(not(feature = "enterprise"))]
1450 fn test_load_options_ignores_unknown_plugin_options() {
1451 let mut file = create_named_temp_file();
1455 write!(
1456 file,
1457 r#"
1458[[plugins]]
1459SomeUnknownPlugin = {{ feature = "foo", count = 5 }}
1460
1461[[plugins]]
1462AnotherUnknownPlugin = {{}}
1463"#
1464 )
1465 .unwrap();
1466
1467 let opts = GreptimeOptions::<StandaloneOptions>::load_layered_options(
1468 Some(file.path().to_str().unwrap()),
1469 "GREPTIMEDB_STANDALONE_UT",
1470 )
1471 .expect(
1472 "loading a config with unrecognized plugin options should succeed, \
1473 ignoring the unknown ones",
1474 );
1475 assert!(opts.plugins.is_empty());
1478 }
1479
1480 #[test]
1481 fn test_load_options_errors_on_malformed_known_plugin_option() {
1482 let mut file = create_named_temp_file();
1487 write!(
1488 file,
1489 r#"
1490[[plugins]]
1491Dummy = {{ unexpected = "payload" }}
1492"#
1493 )
1494 .unwrap();
1495
1496 let result = GreptimeOptions::<StandaloneOptions>::load_layered_options(
1497 Some(file.path().to_str().unwrap()),
1498 "GREPTIMEDB_STANDALONE_UT",
1499 );
1500 assert!(
1501 result.is_err(),
1502 "a malformed payload for a known plugin variant must fail config loading"
1503 );
1504 }
1505
1506 #[test]
1507 #[cfg(not(feature = "enterprise"))]
1508 fn test_load_options_ignores_multi_key_unknown_plugin_entry() {
1509 let mut file = create_named_temp_file();
1513 write!(
1514 file,
1515 r#"
1516[[plugins]]
1517FirstUnknownPlugin = {{ a = 1 }}
1518SecondUnknownPlugin = {{ b = 2 }}
1519"#
1520 )
1521 .unwrap();
1522
1523 let opts = GreptimeOptions::<StandaloneOptions>::load_layered_options(
1524 Some(file.path().to_str().unwrap()),
1525 "GREPTIMEDB_STANDALONE_UT",
1526 )
1527 .expect(
1528 "loading a config with a multi-key unknown plugin entry should succeed, \
1529 ignoring the unknown ones",
1530 );
1531 assert!(opts.plugins.is_empty());
1532 }
1533}