Skip to main content

cmd/
frontend.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::future::Future;
17use std::path::Path;
18use std::pin::Pin;
19use std::sync::Arc;
20use std::time::Duration;
21
22use async_trait::async_trait;
23use cache::{build_fundamental_cache_registry, with_default_composite_cache_registry};
24use catalog::information_extension::DistributedInformationExtension;
25use catalog::information_schema::InformationExtensionRef;
26use catalog::kvbackend::{
27    CachedKvBackendBuilder, CatalogManagerConfiguratorRef, KvBackendCatalogManagerBuilder,
28    new_read_only_meta_kv_backend,
29};
30use catalog::process_manager::ProcessManager;
31use clap::Parser;
32use client::client_manager::NodeClients;
33use common_base::Plugins;
34use common_config::{Configurable, DEFAULT_DATA_HOME};
35use common_error::ext::BoxedError;
36use common_meta::cache::{CacheRegistryBuilder, LayeredCacheRegistryBuilder};
37use common_query::prelude::set_default_prefix;
38use common_stat::ResourceStatImpl;
39use common_telemetry::info;
40use common_telemetry::logging::{DEFAULT_LOGGING_DIR, TracingOptions};
41use common_time::timezone::set_default_timezone;
42use common_version::{short_version, verbose_version};
43use frontend::frontend::Frontend;
44use frontend::heartbeat::{
45    FrontendHeartbeatExtensions, HeartbeatTask, heartbeat_response_handler_executor,
46};
47use frontend::instance::builder::FrontendBuilder;
48use frontend::server::Services;
49use meta_client::{MetaClientOptions, MetaClientRef, MetaClientType};
50use plugins::PluginOptions;
51use plugins::frontend::context::{
52    CatalogManagerConfigureContext, DistributedCatalogManagerConfigureContext,
53};
54use plugins::options::PluginOptionsDeserializerImpl;
55use servers::addrs;
56use servers::grpc::GrpcOptions;
57use servers::tls::{TlsMode, TlsOption, merge_tls_option};
58use snafu::{OptionExt, ResultExt};
59use tracing_appender::non_blocking::WorkerGuard;
60
61use crate::error::{self, OtherSnafu, Result};
62use crate::options::{GlobalOptions, GreptimeOptions};
63use crate::{App, create_resource_limit_metrics, log_versions, maybe_activate_heap_profile};
64
65type FrontendOptions = GreptimeOptions<frontend::frontend::FrontendOptions>;
66type HeartbeatExtensionSetupFuture<'a> = Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
67
68pub struct Instance {
69    frontend: Frontend,
70    // Keep the logging guard to prevent the worker from being dropped.
71    _guard: Vec<WorkerGuard>,
72}
73
74pub const APP_NAME: &str = "greptime-frontend";
75
76impl Instance {
77    pub fn new(frontend: Frontend, _guard: Vec<WorkerGuard>) -> Self {
78        Self { frontend, _guard }
79    }
80
81    pub fn inner(&self) -> &Frontend {
82        &self.frontend
83    }
84
85    pub fn mut_inner(&mut self) -> &mut Frontend {
86        &mut self.frontend
87    }
88}
89
90#[async_trait]
91impl App for Instance {
92    fn name(&self) -> &str {
93        APP_NAME
94    }
95
96    async fn start(&mut self) -> Result<()> {
97        plugins::start_frontend_plugins(&self.frontend.instance)
98            .await
99            .context(error::StartFrontendSnafu)?;
100
101        self.frontend
102            .start()
103            .await
104            .context(error::StartFrontendSnafu)
105    }
106
107    async fn stop(&mut self) -> Result<()> {
108        self.frontend
109            .shutdown()
110            .await
111            .context(error::ShutdownFrontendSnafu)
112    }
113}
114
115#[derive(Parser)]
116pub struct Command {
117    #[clap(subcommand)]
118    pub subcmd: SubCommand,
119}
120
121impl Command {
122    pub async fn build(&self, opts: FrontendOptions) -> Result<Instance> {
123        self.build_with_heartbeat_extensions(opts, FrontendHeartbeatExtensions::default())
124            .await
125    }
126
127    /// Builds a frontend with pre-registered heartbeat extensions.
128    ///
129    /// Register extensions before calling this method. The supplied registry is installed before
130    /// normal plugin setup. After plugin heartbeat setup completes, the registry is frozen before
131    /// heartbeat handlers and the heartbeat task are built, so later registrations are rejected
132    /// and every heartbeat consumer observes the same extension membership.
133    pub async fn build_with_heartbeat_extensions(
134        &self,
135        opts: FrontendOptions,
136        heartbeat_extensions: FrontendHeartbeatExtensions,
137    ) -> Result<Instance> {
138        let plugins = plugins_with_heartbeat_extensions(heartbeat_extensions);
139        self.forward_build_with_plugins(opts, plugins, |command, opts, plugins| {
140            command.build_with_plugins(opts, plugins)
141        })
142        .await
143    }
144
145    fn forward_build_with_plugins<'a, T>(
146        &'a self,
147        opts: FrontendOptions,
148        plugins: Plugins,
149        build: impl FnOnce(&'a StartCommand, FrontendOptions, Plugins) -> T,
150    ) -> T {
151        self.subcmd.forward_build_with_plugins(opts, plugins, build)
152    }
153
154    pub fn load_options(&self, global_options: &GlobalOptions) -> Result<FrontendOptions> {
155        self.subcmd.load_options(global_options)
156    }
157}
158
159#[derive(Parser)]
160pub enum SubCommand {
161    Start(StartCommand),
162}
163
164impl SubCommand {
165    fn forward_build_with_plugins<'a, T>(
166        &'a self,
167        opts: FrontendOptions,
168        plugins: Plugins,
169        build: impl FnOnce(&'a StartCommand, FrontendOptions, Plugins) -> T,
170    ) -> T {
171        match self {
172            SubCommand::Start(cmd) => build(cmd, opts, plugins),
173        }
174    }
175
176    fn load_options(&self, global_options: &GlobalOptions) -> Result<FrontendOptions> {
177        match self {
178            SubCommand::Start(cmd) => cmd.load_options(global_options),
179        }
180    }
181}
182
183#[derive(Debug, Default, Parser)]
184pub struct StartCommand {
185    /// The address to bind the gRPC server.
186    #[clap(long = "grpc-bind-addr", alias = "rpc-bind-addr", alias = "rpc-addr")]
187    grpc_bind_addr: Option<String>,
188    /// The address advertised to the metasrv, and used for connections from outside the host.
189    /// If left empty or unset, the server will automatically use the IP address of the first network interface
190    /// on the host, with the same port number as the one specified in `grpc_bind_addr`.
191    #[clap(
192        long = "grpc-server-addr",
193        alias = "rpc-server-addr",
194        alias = "rpc-hostname"
195    )]
196    grpc_server_addr: Option<String>,
197    /// The address to bind the internal gRPC server.
198    #[clap(
199        long = "internal-grpc-bind-addr",
200        alias = "internal-rpc-bind-addr",
201        alias = "internal-rpc-addr"
202    )]
203    internal_grpc_bind_addr: Option<String>,
204    /// The address advertised to the metasrv, and used for connections from outside the host.
205    /// If left empty or unset, the server will automatically use the IP address of the first network interface
206    /// on the host, with the same port number as the one specified in `internal_grpc_bind_addr`.
207    #[clap(
208        long = "internal-grpc-server-addr",
209        alias = "internal-rpc-server-addr",
210        alias = "internal-rpc-hostname"
211    )]
212    internal_grpc_server_addr: Option<String>,
213    #[clap(long)]
214    http_addr: Option<String>,
215    #[clap(long)]
216    http_timeout: Option<u64>,
217    #[clap(long)]
218    mysql_addr: Option<String>,
219    #[clap(long)]
220    postgres_addr: Option<String>,
221    #[clap(short, long)]
222    pub config_file: Option<String>,
223    #[clap(short, long)]
224    influxdb_enable: Option<bool>,
225    #[clap(long, value_delimiter = ',', num_args = 1..)]
226    metasrv_addrs: Option<Vec<String>>,
227    #[clap(long)]
228    tls_mode: Option<TlsMode>,
229    #[clap(long)]
230    tls_cert_path: Option<String>,
231    #[clap(long)]
232    tls_key_path: Option<String>,
233    #[clap(long)]
234    tls_watch: bool,
235    #[clap(long)]
236    user_provider: Option<String>,
237    #[clap(long)]
238    disable_dashboard: Option<bool>,
239    #[clap(long, default_value = "GREPTIMEDB_FRONTEND")]
240    pub env_prefix: String,
241}
242
243impl StartCommand {
244    fn load_options(&self, global_options: &GlobalOptions) -> Result<FrontendOptions> {
245        let mut opts = FrontendOptions::load_layered_options(
246            self.config_file.as_deref(),
247            self.env_prefix.as_ref(),
248        )
249        .context(error::LoadLayeredConfigSnafu)?;
250
251        self.merge_with_cli_options(global_options, &mut opts)?;
252
253        Ok(opts)
254    }
255
256    // The precedence order is: cli > config file > environment variables > default values.
257    fn merge_with_cli_options(
258        &self,
259        global_options: &GlobalOptions,
260        opts: &mut FrontendOptions,
261    ) -> Result<()> {
262        let opts = &mut opts.component;
263
264        if let Some(dir) = &global_options.log_dir {
265            opts.logging.dir.clone_from(dir);
266        }
267
268        // If the logging dir is not set, use the default logs dir in the data home.
269        if opts.logging.dir.is_empty() {
270            opts.logging.dir = Path::new(DEFAULT_DATA_HOME)
271                .join(DEFAULT_LOGGING_DIR)
272                .to_string_lossy()
273                .to_string();
274        }
275
276        if global_options.log_level.is_some() {
277            opts.logging.level.clone_from(&global_options.log_level);
278        }
279
280        opts.tracing = TracingOptions {
281            #[cfg(feature = "tokio-console")]
282            tokio_console_addr: global_options.tokio_console_addr.clone(),
283        };
284
285        let tls_opts = TlsOption::new(
286            self.tls_mode,
287            self.tls_cert_path.clone(),
288            self.tls_key_path.clone(),
289            self.tls_watch,
290        );
291
292        if let Some(addr) = &self.http_addr {
293            opts.http.addr.clone_from(addr);
294        }
295
296        if let Some(http_timeout) = self.http_timeout {
297            opts.http.timeout = Duration::from_secs(http_timeout)
298        }
299
300        if let Some(disable_dashboard) = self.disable_dashboard {
301            opts.http.disable_dashboard = disable_dashboard;
302        }
303
304        if let Some(addr) = &self.grpc_bind_addr {
305            opts.grpc.bind_addr.clone_from(addr);
306            opts.grpc.tls = merge_tls_option(&opts.grpc.tls, tls_opts.clone());
307        }
308
309        if let Some(addr) = &self.grpc_server_addr {
310            opts.grpc.server_addr.clone_from(addr);
311        }
312
313        if let Some(addr) = &self.internal_grpc_bind_addr {
314            if let Some(internal_grpc) = &mut opts.internal_grpc {
315                internal_grpc.bind_addr = addr.clone();
316            } else {
317                let grpc_options = GrpcOptions {
318                    bind_addr: addr.clone(),
319                    ..Default::default()
320                };
321
322                opts.internal_grpc = Some(grpc_options);
323            }
324        }
325
326        if let Some(addr) = &self.internal_grpc_server_addr {
327            if let Some(internal_grpc) = &mut opts.internal_grpc {
328                internal_grpc.server_addr = addr.clone();
329            } else {
330                let grpc_options = GrpcOptions {
331                    server_addr: addr.clone(),
332                    ..Default::default()
333                };
334                opts.internal_grpc = Some(grpc_options);
335            }
336        }
337
338        if let Some(addr) = &self.mysql_addr {
339            opts.mysql.enable = true;
340            opts.mysql.addr.clone_from(addr);
341            opts.mysql.tls = merge_tls_option(&opts.mysql.tls, tls_opts.clone());
342        }
343
344        if let Some(addr) = &self.postgres_addr {
345            opts.postgres.enable = true;
346            opts.postgres.addr.clone_from(addr);
347            opts.postgres.tls = merge_tls_option(&opts.postgres.tls, tls_opts.clone());
348        }
349
350        if let Some(enable) = self.influxdb_enable {
351            opts.influxdb.enable = enable;
352        }
353
354        if let Some(metasrv_addrs) = &self.metasrv_addrs {
355            opts.meta_client
356                .get_or_insert_with(MetaClientOptions::default)
357                .metasrv_addrs
358                .clone_from(metasrv_addrs);
359        }
360
361        if let Some(user_provider) = &self.user_provider {
362            opts.user_provider = Some(user_provider.clone());
363        }
364
365        Ok(())
366    }
367
368    async fn build_with_plugins(
369        &self,
370        opts: FrontendOptions,
371        plugins: Plugins,
372    ) -> Result<Instance> {
373        let guard = common_telemetry::init_global_logging(
374            APP_NAME,
375            &opts.component.logging,
376            &opts.component.tracing,
377            opts.component.node_id.clone(),
378            Some(&opts.component.slow_query),
379        );
380
381        common_runtime::init_global_runtimes(&opts.runtime);
382
383        crate::options::flush_dropped_plugin_warnings();
384        log_versions(verbose_version(), short_version(), APP_NAME);
385        maybe_activate_heap_profile(&opts.component.memory);
386        create_resource_limit_metrics(APP_NAME);
387
388        info!("Frontend start command: {:#?}", self);
389        info!("Frontend options: {:#?}", opts);
390
391        let plugin_opts = opts.plugins;
392        let mut opts = opts.component;
393        opts.grpc.detect_server_addr();
394
395        set_default_timezone(opts.default_timezone.as_deref()).context(error::InitTimezoneSnafu)?;
396        set_default_prefix(opts.default_column_prefix.as_deref())
397            .map_err(BoxedError::new)
398            .context(error::BuildCliSnafu)?;
399
400        let meta_client_options = opts
401            .meta_client
402            .as_ref()
403            .context(error::MissingConfigSnafu {
404                msg: "'meta_client'",
405            })?;
406
407        let cache_max_capacity = meta_client_options.metadata_cache_max_capacity;
408        let cache_ttl = meta_client_options.metadata_cache_ttl;
409        let cache_tti = meta_client_options.metadata_cache_tti;
410
411        let meta_config: Vec<PluginOptions> = meta_client::create_meta_client(
412            MetaClientType::Frontend,
413            meta_client_options,
414            None,
415            None,
416        )
417        .await
418        .context(error::MetaClientInitSnafu)?
419        .pull_config(PluginOptionsDeserializerImpl)
420        .await
421        .context(error::MetaClientInitSnafu)?;
422
423        let mut plugins =
424            prepare_frontend_plugins(plugins, &plugin_opts, &opts, Some(&meta_config)).await?;
425
426        // now initialize the meta_client with plugins
427        let meta_client = meta_client::create_meta_client(
428            MetaClientType::Frontend,
429            meta_client_options,
430            Some(&plugins),
431            None,
432        )
433        .await
434        .context(error::MetaClientInitSnafu)?;
435
436        let readonly_meta_backend = new_read_only_meta_kv_backend(meta_client.clone());
437
438        // TODO(discord9): add helper function to ease the creation of cache registry&such
439        let cached_meta_backend = CachedKvBackendBuilder::new(readonly_meta_backend.clone())
440            .cache_max_capacity(cache_max_capacity)
441            .cache_ttl(cache_ttl)
442            .cache_tti(cache_tti)
443            .build();
444        let cached_meta_backend = Arc::new(cached_meta_backend);
445
446        // Builds cache registry
447        let layered_cache_builder = LayeredCacheRegistryBuilder::default().add_cache_registry(
448            CacheRegistryBuilder::default()
449                .add_cache(cached_meta_backend.clone())
450                .build(),
451        );
452        let fundamental_cache_registry =
453            build_fundamental_cache_registry(readonly_meta_backend.clone());
454        let mut layered_cache_builder = with_default_composite_cache_registry(
455            layered_cache_builder.add_cache_registry(fundamental_cache_registry),
456        )
457        .context(error::BuildCacheRegistrySnafu)?;
458
459        if let Some(plugin_cache_builder) = plugins::frontend::configure_cache_registry(&plugins) {
460            layered_cache_builder =
461                layered_cache_builder.add_cache_registry(plugin_cache_builder.build());
462        }
463
464        let layered_cache_registry = Arc::new(layered_cache_builder.build());
465
466        // frontend to datanode need not timeout.
467        // Some queries are expected to take long time.
468        let mut channel_config = opts.datanode.client.channel_config();
469        channel_config.timeout = None;
470        // Source Flight streams and sink unary responses share pooled connections.
471        channel_config.http2_adaptive_window = Some(true);
472        if opts.grpc.flight_compression.transport_compression() {
473            channel_config.accept_compression = true;
474            channel_config.send_compression = true;
475        }
476        let client = Arc::new(NodeClients::new(channel_config));
477
478        let information_extension = Arc::new(DistributedInformationExtension::new(
479            meta_client.clone(),
480            client.clone(),
481        ));
482        plugins.insert::<InformationExtensionRef>(information_extension.clone());
483
484        let process_manager = Arc::new(ProcessManager::new(
485            addrs::resolve_addr(&opts.grpc.bind_addr, Some(&opts.grpc.server_addr)),
486            Some(meta_client.clone()),
487        ));
488
489        let builder = KvBackendCatalogManagerBuilder::new(
490            information_extension,
491            cached_meta_backend.clone(),
492            layered_cache_registry.clone(),
493        )
494        .with_process_manager(process_manager.clone());
495        let builder = if let Some(configurator) =
496            plugins.get::<CatalogManagerConfiguratorRef<CatalogManagerConfigureContext>>()
497        {
498            let ctx = DistributedCatalogManagerConfigureContext {
499                meta_client: meta_client.clone(),
500            };
501            let ctx = CatalogManagerConfigureContext::Distributed(ctx);
502
503            configurator
504                .configure(builder, ctx)
505                .await
506                .context(OtherSnafu)?
507        } else {
508            builder
509        };
510        let catalog_manager = builder.build();
511
512        let builder = FrontendBuilder::new(
513            opts.clone(),
514            cached_meta_backend.clone(),
515            layered_cache_registry.clone(),
516            catalog_manager,
517            client,
518            meta_client.clone(),
519            process_manager,
520        );
521
522        plugins::setup_frontend_plugins_post_build(&mut plugins, &plugin_opts, &builder)
523            .await
524            .context(error::StartFrontendSnafu)?;
525
526        let instance = builder
527            .with_plugin(plugins.clone())
528            .with_local_cache_invalidator(layered_cache_registry)
529            .try_build()
530            .await
531            .context(error::StartFrontendSnafu)?;
532
533        let instance = Arc::new(instance);
534
535        let heartbeat_instance = instance.clone();
536        let heartbeat_extensions =
537            setup_and_freeze_frontend_heartbeat_extensions(&mut plugins, move |plugins| {
538                Box::pin(async move {
539                    plugins::setup_frontend_heartbeat_extensions(
540                        plugins,
541                        &plugin_opts,
542                        &heartbeat_instance,
543                    )
544                    .await
545                    .context(error::StartFrontendSnafu)
546                })
547            })
548            .await?;
549        let heartbeat_task = Some(create_heartbeat_task_with_extensions(
550            &opts,
551            meta_client,
552            &instance,
553            heartbeat_extensions,
554        ));
555
556        let servers = Services::new(opts, instance.clone(), plugins)
557            .build()
558            .context(error::StartFrontendSnafu)?;
559
560        let frontend = Frontend {
561            instance,
562            servers,
563            heartbeat_task,
564        };
565
566        Ok(Instance::new(frontend, guard))
567    }
568}
569
570async fn prepare_frontend_plugins(
571    mut plugins: Plugins,
572    plugin_opts: &[PluginOptions],
573    opts: &frontend::frontend::FrontendOptions,
574    meta_config: Option<&[PluginOptions]>,
575) -> Result<Plugins> {
576    plugins::setup_frontend_plugins_pre_build(&mut plugins, plugin_opts, opts, meta_config)
577        .await
578        .context(error::StartFrontendSnafu)?;
579    Ok(plugins)
580}
581
582fn plugins_with_heartbeat_extensions(heartbeat_extensions: FrontendHeartbeatExtensions) -> Plugins {
583    let plugins = Plugins::new();
584    plugins.insert(heartbeat_extensions);
585    plugins
586}
587
588async fn setup_and_freeze_frontend_heartbeat_extensions(
589    plugins: &mut Plugins,
590    setup: impl for<'a> FnOnce(&'a mut Plugins) -> HeartbeatExtensionSetupFuture<'a>,
591) -> Result<FrontendHeartbeatExtensions> {
592    setup(plugins).await?;
593    let extensions = plugins
594        .get::<FrontendHeartbeatExtensions>()
595        .unwrap_or_default();
596    extensions.freeze();
597    Ok(extensions)
598}
599
600pub fn create_heartbeat_task(
601    options: &frontend::frontend::FrontendOptions,
602    meta_client: MetaClientRef,
603    instance: &frontend::instance::Instance,
604) -> HeartbeatTask {
605    create_heartbeat_task_with_extensions(options, meta_client, instance, Default::default())
606}
607
608fn create_heartbeat_task_with_extensions(
609    options: &frontend::frontend::FrontendOptions,
610    meta_client: MetaClientRef,
611    instance: &frontend::instance::Instance,
612    extensions: FrontendHeartbeatExtensions,
613) -> HeartbeatTask {
614    let executor = heartbeat_response_handler_executor(
615        &extensions,
616        instance.suspend_state(),
617        instance.cache_invalidator().clone(),
618    );
619
620    let stat = {
621        let mut stat = ResourceStatImpl::default();
622        stat.start_collect_cpu_usage();
623        Arc::new(stat)
624    };
625
626    HeartbeatTask::new(
627        instance.frontend_peer_addr().to_string(),
628        options,
629        meta_client,
630        executor,
631        stat,
632    )
633    .with_extensions(extensions)
634}
635
636#[cfg(test)]
637mod tests {
638    use std::io::Write;
639    use std::sync::Arc;
640    use std::time::Duration;
641
642    use auth::{Identity, Password, UserProviderRef};
643    use clap::{CommandFactory, Parser};
644    use common_base::readable_size::ReadableSize;
645    use common_config::ENV_VAR_SEP;
646    use common_test_util::temp_dir::create_named_temp_file;
647    use servers::grpc::GrpcOptions;
648    use servers::http::HttpOptions;
649
650    use super::*;
651    use crate::options::GlobalOptions;
652
653    #[derive(Debug)]
654    struct TestExtension(&'static str);
655
656    #[async_trait]
657    impl frontend::heartbeat::FrontendHeartbeatExtension for TestExtension {
658        fn name(&self) -> &'static str {
659            self.0
660        }
661    }
662
663    #[test]
664    fn test_command_forwards_heartbeat_extensions_to_start_build() {
665        let command = Command {
666            subcmd: SubCommand::Start(StartCommand::default()),
667        };
668        let extensions = FrontendHeartbeatExtensions::default();
669        let shared_extensions = extensions.clone();
670
671        let forwarded_plugins = command.forward_build_with_plugins(
672            FrontendOptions::default(),
673            plugins_with_heartbeat_extensions(extensions),
674            |_, _, plugins| plugins,
675        );
676        let forwarded_extensions = forwarded_plugins
677            .get::<FrontendHeartbeatExtensions>()
678            .unwrap();
679
680        assert_registry_identity(&shared_extensions, &forwarded_extensions);
681    }
682
683    #[tokio::test]
684    async fn test_prefilled_heartbeat_extensions_survive_pre_build_setup() {
685        let extensions = FrontendHeartbeatExtensions::default();
686        let shared_extensions = extensions.clone();
687
688        let options = frontend::frontend::FrontendOptions {
689            meta_client: Some(MetaClientOptions::default()),
690            ..Default::default()
691        };
692        let plugins = prepare_frontend_plugins(
693            plugins_with_heartbeat_extensions(extensions),
694            &[],
695            &options,
696            None,
697        )
698        .await
699        .unwrap();
700
701        let setup_extensions = plugins.get_or_insert(FrontendHeartbeatExtensions::default);
702        assert_registry_identity(&shared_extensions, &setup_extensions);
703    }
704
705    #[tokio::test]
706    async fn test_heartbeat_extension_setup_completes_before_freeze() {
707        let extensions = FrontendHeartbeatExtensions::default();
708        assert_eq!(
709            extensions.try_register(Arc::new(TestExtension("caller"))),
710            Ok(())
711        );
712        let mut plugins = plugins_with_heartbeat_extensions(extensions.clone());
713
714        let finalized = setup_and_freeze_frontend_heartbeat_extensions(&mut plugins, |plugins| {
715            Box::pin(async move {
716                let setup_extensions = plugins.get_or_insert(FrontendHeartbeatExtensions::default);
717                assert_eq!(
718                    setup_extensions.try_register(Arc::new(TestExtension("plugin"))),
719                    Ok(())
720                );
721                Ok(())
722            })
723        })
724        .await
725        .unwrap();
726
727        assert_eq!(
728            finalized
729                .extensions()
730                .iter()
731                .map(|extension| extension.name())
732                .collect::<Vec<_>>(),
733            ["caller", "plugin"]
734        );
735        assert_eq!(
736            extensions.try_register(Arc::new(TestExtension("late"))),
737            Err(frontend::heartbeat::RegistrationError::Frozen)
738        );
739        assert_eq!(finalized.len(), 2);
740    }
741
742    fn assert_registry_identity(
743        expected: &FrontendHeartbeatExtensions,
744        actual: &FrontendHeartbeatExtensions,
745    ) {
746        let extension = Arc::new(TestExtension("cmd-test-extension"));
747        assert!(expected.register(extension.clone()));
748        let registered = actual.extensions();
749        assert_eq!(registered.len(), 1);
750        assert!(Arc::ptr_eq(
751            &(extension as Arc<dyn frontend::heartbeat::FrontendHeartbeatExtension>),
752            &registered[0]
753        ));
754    }
755
756    #[test]
757    fn test_try_from_start_command() {
758        let command = StartCommand {
759            http_addr: Some("127.0.0.1:1234".to_string()),
760            mysql_addr: Some("127.0.0.1:5678".to_string()),
761            postgres_addr: Some("127.0.0.1:5432".to_string()),
762            internal_grpc_bind_addr: Some("127.0.0.1:4010".to_string()),
763            internal_grpc_server_addr: Some("10.0.0.24:4010".to_string()),
764            influxdb_enable: Some(false),
765            disable_dashboard: Some(false),
766            ..Default::default()
767        };
768
769        let opts = command.load_options(&Default::default()).unwrap().component;
770
771        assert_eq!(opts.http.addr, "127.0.0.1:1234");
772        assert_eq!(ReadableSize::mb(64), opts.http.body_limit);
773        assert_eq!(opts.mysql.addr, "127.0.0.1:5678");
774        assert_eq!(opts.postgres.addr, "127.0.0.1:5432");
775
776        let internal_grpc = opts.internal_grpc.as_ref().unwrap();
777        assert_eq!(internal_grpc.bind_addr, "127.0.0.1:4010");
778        assert_eq!(internal_grpc.server_addr, "10.0.0.24:4010");
779
780        let default_opts = FrontendOptions::default().component;
781
782        assert_eq!(opts.grpc.bind_addr, default_opts.grpc.bind_addr);
783        assert!(opts.mysql.enable);
784        assert_eq!(opts.mysql.runtime_size, default_opts.mysql.runtime_size);
785        assert!(opts.postgres.enable);
786        assert_eq!(
787            opts.postgres.runtime_size,
788            default_opts.postgres.runtime_size
789        );
790        assert!(opts.opentsdb.enable);
791
792        assert!(!opts.influxdb.enable);
793    }
794
795    #[test]
796    fn test_read_from_config_file() {
797        let mut file = create_named_temp_file();
798        let toml_str = r#"
799            [http]
800            addr = "127.0.0.1:4000"
801            timeout = "0s"
802            body_limit = "2GB"
803
804            [opentsdb]
805            enable = false
806
807            [logging]
808            level = "debug"
809            dir = "./greptimedb_data/test/logs"
810        "#;
811        write!(file, "{}", toml_str).unwrap();
812
813        let command = StartCommand {
814            config_file: Some(file.path().to_str().unwrap().to_string()),
815            disable_dashboard: Some(false),
816            ..Default::default()
817        };
818
819        let fe_opts = command.load_options(&Default::default()).unwrap().component;
820
821        assert_eq!("127.0.0.1:4000".to_string(), fe_opts.http.addr);
822        assert_eq!(Duration::from_secs(0), fe_opts.http.timeout);
823
824        assert_eq!(ReadableSize::gb(2), fe_opts.http.body_limit);
825
826        assert_eq!("debug", fe_opts.logging.level.as_ref().unwrap());
827        assert_eq!(
828            "./greptimedb_data/test/logs".to_string(),
829            fe_opts.logging.dir
830        );
831        assert!(!fe_opts.opentsdb.enable);
832    }
833
834    #[tokio::test]
835    async fn test_try_from_start_command_to_anymap() {
836        let fe_opts = frontend::frontend::FrontendOptions {
837            http: HttpOptions {
838                disable_dashboard: false,
839                ..Default::default()
840            },
841            meta_client: Some(MetaClientOptions::default()),
842            user_provider: Some("static_user_provider:cmd:test=test".to_string()),
843            ..Default::default()
844        };
845
846        let mut plugins = Plugins::new();
847        plugins::setup_frontend_plugins_pre_build(&mut plugins, &[], &fe_opts, None)
848            .await
849            .unwrap();
850
851        let provider = plugins.get::<UserProviderRef>().unwrap();
852        let result = provider
853            .authenticate(
854                Identity::UserId("test", None),
855                Password::PlainText("test".to_string().into()),
856            )
857            .await;
858        let _ = result.unwrap();
859    }
860
861    #[test]
862    fn test_load_log_options_from_cli() {
863        let cmd = StartCommand {
864            disable_dashboard: Some(false),
865            ..Default::default()
866        };
867
868        let options = cmd
869            .load_options(&GlobalOptions {
870                log_dir: Some("./greptimedb_data/test/logs".to_string()),
871                log_level: Some("debug".to_string()),
872
873                #[cfg(feature = "tokio-console")]
874                tokio_console_addr: None,
875            })
876            .unwrap()
877            .component;
878
879        let logging_opt = options.logging;
880        assert_eq!("./greptimedb_data/test/logs", logging_opt.dir);
881        assert_eq!("debug", logging_opt.level.as_ref().unwrap());
882    }
883
884    #[test]
885    fn test_config_precedence_order() {
886        let mut file = create_named_temp_file();
887        let toml_str = r#"
888            [http]
889            addr = "127.0.0.1:4000"
890
891            [meta_client]
892            timeout = "3s"
893            connect_timeout = "5s"
894            tcp_nodelay = true
895
896            [mysql]
897            addr = "127.0.0.1:4002"
898        "#;
899        write!(file, "{}", toml_str).unwrap();
900
901        let env_prefix = "FRONTEND_UT";
902        temp_env::with_vars(
903            [
904                (
905                    // mysql.addr = 127.0.0.1:14002
906                    [
907                        env_prefix.to_string(),
908                        "mysql".to_uppercase(),
909                        "addr".to_uppercase(),
910                    ]
911                    .join(ENV_VAR_SEP),
912                    Some("127.0.0.1:14002"),
913                ),
914                (
915                    // mysql.runtime_size = 11
916                    [
917                        env_prefix.to_string(),
918                        "mysql".to_uppercase(),
919                        "runtime_size".to_uppercase(),
920                    ]
921                    .join(ENV_VAR_SEP),
922                    Some("11"),
923                ),
924                (
925                    // http.addr = 127.0.0.1:24000
926                    [
927                        env_prefix.to_string(),
928                        "http".to_uppercase(),
929                        "addr".to_uppercase(),
930                    ]
931                    .join(ENV_VAR_SEP),
932                    Some("127.0.0.1:24000"),
933                ),
934                (
935                    // meta_client.metasrv_addrs = 127.0.0.1:3001,127.0.0.1:3002,127.0.0.1:3003
936                    [
937                        env_prefix.to_string(),
938                        "meta_client".to_uppercase(),
939                        "metasrv_addrs".to_uppercase(),
940                    ]
941                    .join(ENV_VAR_SEP),
942                    Some("127.0.0.1:3001,127.0.0.1:3002,127.0.0.1:3003"),
943                ),
944            ],
945            || {
946                let command = StartCommand {
947                    config_file: Some(file.path().to_str().unwrap().to_string()),
948                    http_addr: Some("127.0.0.1:14000".to_string()),
949                    env_prefix: env_prefix.to_string(),
950                    ..Default::default()
951                };
952
953                let fe_opts = command.load_options(&Default::default()).unwrap().component;
954
955                // Should be read from env, env > default values.
956                assert_eq!(fe_opts.mysql.runtime_size, 11);
957                assert_eq!(
958                    fe_opts.meta_client.unwrap().metasrv_addrs,
959                    vec![
960                        "127.0.0.1:3001".to_string(),
961                        "127.0.0.1:3002".to_string(),
962                        "127.0.0.1:3003".to_string()
963                    ]
964                );
965
966                // Should be read from config file, config file > env > default values.
967                assert_eq!(fe_opts.mysql.addr, "127.0.0.1:4002");
968
969                // Should be read from cli, cli > config file > env > default values.
970                assert_eq!(fe_opts.http.addr, "127.0.0.1:14000");
971
972                // Should be default value.
973                assert_eq!(fe_opts.grpc.bind_addr, GrpcOptions::default().bind_addr);
974            },
975        );
976    }
977
978    #[test]
979    fn test_parse_grpc_cli_aliases() {
980        let command = StartCommand::try_parse_from([
981            "frontend",
982            "--grpc-bind-addr",
983            "127.0.0.1:14001",
984            "--grpc-server-addr",
985            "10.0.0.1:14001",
986            "--internal-grpc-bind-addr",
987            "127.0.0.1:14010",
988            "--internal-grpc-server-addr",
989            "10.0.0.1:14010",
990        ])
991        .unwrap();
992        assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:14001"));
993        assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.1:14001"));
994        assert_eq!(
995            command.internal_grpc_bind_addr.as_deref(),
996            Some("127.0.0.1:14010")
997        );
998        assert_eq!(
999            command.internal_grpc_server_addr.as_deref(),
1000            Some("10.0.0.1:14010")
1001        );
1002
1003        let command = StartCommand::try_parse_from([
1004            "frontend",
1005            "--rpc-bind-addr",
1006            "127.0.0.1:24001",
1007            "--rpc-server-addr",
1008            "10.0.0.2:24001",
1009            "--internal-rpc-bind-addr",
1010            "127.0.0.1:24010",
1011            "--internal-rpc-server-addr",
1012            "10.0.0.2:24010",
1013        ])
1014        .unwrap();
1015        assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:24001"));
1016        assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.2:24001"));
1017        assert_eq!(
1018            command.internal_grpc_bind_addr.as_deref(),
1019            Some("127.0.0.1:24010")
1020        );
1021        assert_eq!(
1022            command.internal_grpc_server_addr.as_deref(),
1023            Some("10.0.0.2:24010")
1024        );
1025
1026        let command = StartCommand::try_parse_from([
1027            "frontend",
1028            "--rpc-addr",
1029            "127.0.0.1:34001",
1030            "--rpc-hostname",
1031            "10.0.0.3:34001",
1032            "--internal-rpc-addr",
1033            "127.0.0.1:34010",
1034            "--internal-rpc-hostname",
1035            "10.0.0.3:34010",
1036        ])
1037        .unwrap();
1038        assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:34001"));
1039        assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.3:34001"));
1040        assert_eq!(
1041            command.internal_grpc_bind_addr.as_deref(),
1042            Some("127.0.0.1:34010")
1043        );
1044        assert_eq!(
1045            command.internal_grpc_server_addr.as_deref(),
1046            Some("10.0.0.3:34010")
1047        );
1048    }
1049
1050    #[test]
1051    fn test_help_uses_grpc_option_names() {
1052        let mut cmd = StartCommand::command();
1053        let mut help = Vec::new();
1054        cmd.write_long_help(&mut help).unwrap();
1055        let help = String::from_utf8(help).unwrap();
1056
1057        assert!(help.contains("--grpc-bind-addr"));
1058        assert!(help.contains("--grpc-server-addr"));
1059        assert!(help.contains("--internal-grpc-bind-addr"));
1060        assert!(help.contains("--internal-grpc-server-addr"));
1061        assert!(!help.contains("--rpc-bind-addr"));
1062        assert!(!help.contains("--rpc-server-addr"));
1063        assert!(!help.contains("--rpc-addr"));
1064        assert!(!help.contains("--rpc-hostname"));
1065        assert!(!help.contains("--internal-rpc-bind-addr"));
1066        assert!(!help.contains("--internal-rpc-server-addr"));
1067        assert!(!help.contains("--internal-rpc-addr"));
1068        assert!(!help.contains("--internal-rpc-hostname"));
1069    }
1070}