Skip to main content

cmd/
standalone.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::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(&copy_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    /// Whether the `standalone start` command requested daemonization.
157    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    // Keep the logging guard to prevent the worker from being dropped.
198    _guard: Vec<WorkerGuard>,
199}
200
201impl Instance {
202    /// Find the socket addr of a server by its `name`.
203    pub fn server_addr(&self, name: &str) -> Option<SocketAddr> {
204        self.frontend.server_handlers().addr(name)
205    }
206
207    /// Get the mutable Frontend component of this Standalone instance for externally modification
208    /// by others (might not be in this code base, so don't delete this function).
209    pub fn mut_frontend(&mut self) -> &mut Frontend {
210        &mut self.frontend
211    }
212
213    /// Get the Datanode component of this Standalone instance for externally usage
214    /// by others (might not be in this code base, so don't delete this function).
215    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    /// The working home directory of this standalone instance.
303    #[clap(long)]
304    data_home: Option<String>,
305    /// Run in the background as a daemon.
306    #[cfg(unix)]
307    #[clap(short, long)]
308    daemon: bool,
309}
310
311impl StartCommand {
312    /// Whether the `standalone start` command requested daemonization.
313    #[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    /// Load the GreptimeDB options from various sources (command line, config file or env).
324    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    // The precedence order is: cli > config file > environment variables > default values.
341    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 the logging dir is not set, use the default logs dir in the data home.
375        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            // frontend grpc addr conflict with datanode default grpc addr
384            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    /// Build GreptimeDB instance with the loaded options.
423    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        // Ensure the data_home directory exists.
478        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        // Builds cache registry
497        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        // for standalone not use grpc, but get a handler to frontend grpc client without
536        // actually make a connection
537        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        // set the ref to query for the local flow state
595        {
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        // set the frontend client for flownode
700        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        // set the frontend invoker for flownode
707        let flow_streaming_engine = flow_engine.streaming_engine();
708        // flow server need to be able to use frontend to write insert requests back
709        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/// Customizes how standalone opens its metadata KV backend.
804///
805/// The default implementation preserves the built-in raft-engine path. Other
806/// callers can provide a custom implementation without changing standalone
807/// configuration types.
808#[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    /// Starts leader services that manage standalone metadata or WAL state.
868    ///
869    /// The default implementation starts the procedure manager and WAL provider
870    /// during instance startup.
871    async fn start(&self, context: LeaderServicesContext) -> Result<()>;
872
873    /// Stops services started by [`StandaloneLeaderServicesController::start`].
874    async fn stop(
875        &self,
876        procedure_manager: ProcedureManagerRef,
877        region_server: RegionServer,
878    ) -> Result<()>;
879}
880
881#[derive(Clone)]
882/// Additional runtime handles for custom leader-service controllers.
883///
884/// The default standalone startup only needs to start/stop the procedure
885/// manager and WAL provider. Some embedders need to do more work around
886/// leader-service startup, for example reconciling metadata-backed runtime
887/// state before publishing writable leadership. Grouping those handles here
888/// keeps `Instance` small and avoids expanding
889/// [`StandaloneLeaderServicesController::start`] every time a custom lifecycle
890/// needs one more standalone component.
891pub 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
934/// `InstanceCreator` is used for grouping various component creators for building the
935/// Standalone instance, suitable for customizing how the instance can be built.
936pub struct InstanceCreator {
937    /// Hook for replacing metadata KV construction while reusing the rest of the
938    /// standalone build flow.
939    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    /// Wraps the metadata backend creator while retaining the default creator.
972    ///
973    /// This is useful for callers that need to add runtime behavior around
974    /// metadata access without reimplementing backend selection.
975    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    /// Wraps node-manager creation while preserving the selected standalone node manager.
984    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    /// Wraps procedure-executor creation while preserving the current setup.
993    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    /// Replaces startup/shutdown ownership for procedure manager and WAL provider.
1002    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    /// Overrides whether regions opened during startup should become writable.
1011    ///
1012    /// `None` keeps the default startup behavior (regions open writable).
1013    ///
1014    /// Warning: setting this to `false` in standalone mode will leave reopened regions
1015    /// permanently read-only. Standalone has no metasrv heartbeat or region-role
1016    /// reconciliation, so there is no path to promote regions to Leader after startup.
1017    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
1036/// `InstanceCreatorResult` is expected to be used paired with [InstanceCreator].
1037/// It stores the created and other important components for further reusing.
1038pub 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                    // logging.dir = /other/log/dir
1312                    [
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                    // logging.level = info
1322                    [
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                    // http.addr = 127.0.0.1:24000
1332                    [
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                // Should be read from env, env > default values.
1352                assert_eq!(opts.logging.dir, "/other/log/dir");
1353
1354                // Should be read from config file, config file > env > default values.
1355                assert_eq!(opts.logging.level.as_ref().unwrap(), "debug");
1356
1357                // Should be read from cli, cli > config file > env > default values.
1358                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                // Should be default value.
1363                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        // Plugin options that are not recognized by the current build (for example,
1452        // an enterprise plugin option seen by an open-source build) must not abort
1453        // startup. They should be dropped with a warning instead.
1454        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        // Unknown plugin options are dropped; the recognized list is empty in the
1476        // open-source build.
1477        assert!(opts.plugins.is_empty());
1478    }
1479
1480    #[test]
1481    fn test_load_options_errors_on_malformed_known_plugin_option() {
1482        // A *known* variant with a malformed payload must NOT be silently
1483        // dropped; config loading must fail so a misconfigured plugin is not
1484        // disabled without notice. Here `Dummy` (a unit variant) is given a map
1485        // payload, which is invalid.
1486        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        // A single plugin table carrying several *unknown* keys must not abort
1510        // startup with serde's "expected map with a single key" error; it is
1511        // dropped like any other unrecognized plugin option.
1512        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}