Skip to main content

cmd/
flownode.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::path::Path;
17use std::sync::Arc;
18use std::time::Duration;
19
20use cache::{build_fundamental_cache_registry, with_default_composite_cache_registry};
21use catalog::information_extension::DistributedInformationExtension;
22use catalog::kvbackend::{
23    CachedKvBackendBuilder, KvBackendCatalogManagerBuilder, new_read_only_meta_kv_backend,
24};
25use clap::Parser;
26use client::client_manager::NodeClients;
27use common_base::Plugins;
28use common_config::{Configurable, DEFAULT_DATA_HOME};
29use common_grpc::channel_manager::ChannelConfig;
30use common_meta::cache::{CacheRegistryBuilder, LayeredCacheRegistryBuilder};
31use common_meta::heartbeat::handler::HandlerGroupExecutor;
32use common_meta::heartbeat::handler::invalidate_table_cache::InvalidateCacheHandler;
33use common_meta::heartbeat::handler::parse_mailbox_message::ParseMailboxMessageHandler;
34use common_meta::key::TableMetadataManager;
35use common_meta::key::flow::FlowMetadataManager;
36use common_stat::ResourceStatImpl;
37use common_telemetry::info;
38use common_telemetry::logging::{DEFAULT_LOGGING_DIR, TracingOptions};
39use common_version::{short_version, verbose_version};
40use flow::{FlownodeBuilder, FlownodeInstance, FlownodeServiceBuilder, FrontendClient};
41use meta_client::{MetaClientOptions, MetaClientType};
42use plugins::flownode::context::GrpcConfigureContext;
43use servers::configurator::GrpcBuilderConfiguratorRef;
44use snafu::{OptionExt, ResultExt, ensure};
45use tracing_appender::non_blocking::WorkerGuard;
46
47use crate::error::{
48    BuildCacheRegistrySnafu, LoadLayeredConfigSnafu, MetaClientInitSnafu, MissingConfigSnafu,
49    OtherSnafu, Result, ShutdownFlownodeSnafu, StartFlownodeSnafu,
50};
51use crate::options::{GlobalOptions, GreptimeOptions};
52use crate::{App, create_resource_limit_metrics, log_versions, maybe_activate_heap_profile};
53
54pub const APP_NAME: &str = "greptime-flownode";
55
56type FlownodeOptions = GreptimeOptions<flow::FlownodeOptions>;
57
58pub struct Instance {
59    flownode: FlownodeInstance,
60    // Keep the logging guard to prevent the worker from being dropped.
61    _guard: Vec<WorkerGuard>,
62}
63
64impl Instance {
65    pub fn new(flownode: FlownodeInstance, guard: Vec<WorkerGuard>) -> Self {
66        Self {
67            flownode,
68            _guard: guard,
69        }
70    }
71
72    pub fn flownode(&self) -> &FlownodeInstance {
73        &self.flownode
74    }
75
76    /// allow customizing flownode for downstream projects
77    pub fn flownode_mut(&mut self) -> &mut FlownodeInstance {
78        &mut self.flownode
79    }
80}
81
82#[async_trait::async_trait]
83impl App for Instance {
84    fn name(&self) -> &str {
85        APP_NAME
86    }
87
88    async fn start(&mut self) -> Result<()> {
89        plugins::start_flownode_plugins(&self.flownode)
90            .await
91            .context(StartFlownodeSnafu)?;
92
93        self.flownode.start().await.context(StartFlownodeSnafu)
94    }
95
96    async fn stop(&mut self) -> Result<()> {
97        self.flownode
98            .shutdown()
99            .await
100            .context(ShutdownFlownodeSnafu)
101    }
102}
103
104#[derive(Parser)]
105pub struct Command {
106    #[clap(subcommand)]
107    subcmd: SubCommand,
108}
109
110impl Command {
111    pub async fn build(&self, opts: FlownodeOptions) -> Result<Instance> {
112        self.subcmd.build(opts).await
113    }
114
115    pub fn load_options(&self, global_options: &GlobalOptions) -> Result<FlownodeOptions> {
116        match &self.subcmd {
117            SubCommand::Start(cmd) => cmd.load_options(global_options),
118        }
119    }
120}
121
122#[derive(Parser)]
123enum SubCommand {
124    Start(StartCommand),
125}
126
127impl SubCommand {
128    async fn build(&self, opts: FlownodeOptions) -> Result<Instance> {
129        match self {
130            SubCommand::Start(cmd) => cmd.build(opts).await,
131        }
132    }
133}
134
135#[derive(Debug, Parser, Default)]
136struct StartCommand {
137    /// Flownode's id
138    #[clap(long)]
139    node_id: Option<u64>,
140    /// Bind address for the gRPC server.
141    #[clap(long = "grpc-bind-addr", alias = "rpc-bind-addr", alias = "rpc-addr")]
142    grpc_bind_addr: Option<String>,
143    /// The address advertised to the metasrv, and used for connections from outside the host.
144    /// If left empty or unset, the server will automatically use the IP address of the first network interface
145    /// on the host, with the same port number as the one specified in `grpc_bind_addr`.
146    #[clap(
147        long = "grpc-server-addr",
148        alias = "rpc-server-addr",
149        alias = "rpc-hostname"
150    )]
151    grpc_server_addr: Option<String>,
152    /// Metasrv address list;
153    #[clap(long, value_delimiter = ',', num_args = 1..)]
154    metasrv_addrs: Option<Vec<String>>,
155    /// The configuration file for flownode
156    #[clap(short, long)]
157    config_file: Option<String>,
158    /// The prefix of environment variables, default is `GREPTIMEDB_FLOWNODE`;
159    #[clap(long, default_value = "GREPTIMEDB_FLOWNODE")]
160    env_prefix: String,
161    #[clap(long)]
162    http_addr: Option<String>,
163    /// HTTP request timeout in seconds.
164    #[clap(long)]
165    http_timeout: Option<u64>,
166}
167
168impl StartCommand {
169    fn load_options(&self, global_options: &GlobalOptions) -> Result<FlownodeOptions> {
170        let mut opts = FlownodeOptions::load_layered_options(
171            self.config_file.as_deref(),
172            self.env_prefix.as_ref(),
173        )
174        .context(LoadLayeredConfigSnafu)?;
175
176        self.merge_with_cli_options(global_options, &mut opts)?;
177
178        Ok(opts)
179    }
180
181    // The precedence order is: cli > config file > environment variables > default values.
182    fn merge_with_cli_options(
183        &self,
184        global_options: &GlobalOptions,
185        opts: &mut FlownodeOptions,
186    ) -> Result<()> {
187        let opts = &mut opts.component;
188
189        if let Some(dir) = &global_options.log_dir {
190            opts.logging.dir.clone_from(dir);
191        }
192
193        // If the logging dir is not set, use the default logs dir in the data home.
194        if opts.logging.dir.is_empty() {
195            opts.logging.dir = Path::new(DEFAULT_DATA_HOME)
196                .join(DEFAULT_LOGGING_DIR)
197                .to_string_lossy()
198                .to_string();
199        }
200
201        if global_options.log_level.is_some() {
202            opts.logging.level.clone_from(&global_options.log_level);
203        }
204
205        opts.tracing = TracingOptions {
206            #[cfg(feature = "tokio-console")]
207            tokio_console_addr: global_options.tokio_console_addr.clone(),
208        };
209
210        if let Some(addr) = &self.grpc_bind_addr {
211            opts.grpc.bind_addr.clone_from(addr);
212        }
213
214        if let Some(server_addr) = &self.grpc_server_addr {
215            opts.grpc.server_addr.clone_from(server_addr);
216        }
217
218        if let Some(node_id) = self.node_id {
219            opts.node_id = Some(node_id);
220        }
221
222        if let Some(metasrv_addrs) = &self.metasrv_addrs {
223            opts.meta_client
224                .get_or_insert_with(MetaClientOptions::default)
225                .metasrv_addrs
226                .clone_from(metasrv_addrs);
227        }
228
229        if let Some(http_addr) = &self.http_addr {
230            opts.http.addr.clone_from(http_addr);
231        }
232
233        if let Some(http_timeout) = self.http_timeout {
234            opts.http.timeout = Duration::from_secs(http_timeout);
235        }
236
237        ensure!(
238            opts.node_id.is_some(),
239            MissingConfigSnafu {
240                msg: "Missing node id option"
241            }
242        );
243
244        Ok(())
245    }
246
247    async fn build(&self, opts: FlownodeOptions) -> Result<Instance> {
248        let guard = common_telemetry::init_global_logging(
249            APP_NAME,
250            &opts.component.logging,
251            &opts.component.tracing,
252            opts.component.node_id.map(|x| x.to_string()),
253            None,
254        );
255
256        common_runtime::init_global_runtimes(&opts.runtime);
257
258        crate::options::flush_dropped_plugin_warnings();
259        log_versions(verbose_version(), short_version(), APP_NAME);
260        maybe_activate_heap_profile(&opts.component.memory);
261        create_resource_limit_metrics(APP_NAME);
262
263        info!("Flownode start command: {:#?}", self);
264        info!("Flownode options: {:#?}", opts);
265
266        let plugin_opts = opts.plugins;
267        let mut opts = opts.component;
268        opts.grpc.detect_server_addr();
269
270        let mut plugins = Plugins::new();
271        plugins::setup_flownode_plugins_pre_build(&mut plugins, &plugin_opts, &opts)
272            .await
273            .context(StartFlownodeSnafu)?;
274
275        let member_id = opts
276            .node_id
277            .context(MissingConfigSnafu { msg: "'node_id'" })?;
278
279        let meta_config = opts.meta_client.as_ref().context(MissingConfigSnafu {
280            msg: "'meta_client_options'",
281        })?;
282
283        let meta_client = meta_client::create_meta_client(
284            MetaClientType::Flownode { member_id },
285            meta_config,
286            None,
287            None,
288        )
289        .await
290        .context(MetaClientInitSnafu)?;
291
292        let cache_max_capacity = meta_config.metadata_cache_max_capacity;
293        let cache_ttl = meta_config.metadata_cache_ttl;
294        let cache_tti = meta_config.metadata_cache_tti;
295
296        let readonly_meta_backend = new_read_only_meta_kv_backend(meta_client.clone());
297
298        // TODO(discord9): add helper function to ease the creation of cache registry&such
299        let cached_meta_backend = CachedKvBackendBuilder::new(readonly_meta_backend.clone())
300            .cache_max_capacity(cache_max_capacity)
301            .cache_ttl(cache_ttl)
302            .cache_tti(cache_tti)
303            .build();
304        let cached_meta_backend = Arc::new(cached_meta_backend);
305
306        // Builds cache registry
307        let layered_cache_builder = LayeredCacheRegistryBuilder::default().add_cache_registry(
308            CacheRegistryBuilder::default()
309                .add_cache(cached_meta_backend.clone())
310                .build(),
311        );
312        let fundamental_cache_registry =
313            build_fundamental_cache_registry(readonly_meta_backend.clone());
314        let layered_cache_registry = Arc::new(
315            with_default_composite_cache_registry(
316                layered_cache_builder.add_cache_registry(fundamental_cache_registry),
317            )
318            .context(BuildCacheRegistrySnafu)?
319            .build(),
320        );
321
322        // flownode's frontend to datanode need not timeout.
323        // Some queries are expected to take long time.
324        let channel_config = ChannelConfig {
325            timeout: None,
326            ..Default::default()
327        };
328        let client = Arc::new(NodeClients::new(channel_config));
329
330        let information_extension = Arc::new(DistributedInformationExtension::new(
331            meta_client.clone(),
332            client.clone(),
333        ));
334        let catalog_manager = KvBackendCatalogManagerBuilder::new(
335            information_extension,
336            cached_meta_backend.clone(),
337            layered_cache_registry.clone(),
338        )
339        .build();
340
341        let table_metadata_manager =
342            Arc::new(TableMetadataManager::new(cached_meta_backend.clone()));
343
344        let executor = HandlerGroupExecutor::new(vec![
345            Arc::new(ParseMailboxMessageHandler),
346            Arc::new(InvalidateCacheHandler::new(layered_cache_registry.clone())),
347        ]);
348
349        let mut resource_stat = ResourceStatImpl::default();
350        resource_stat.start_collect_cpu_usage();
351
352        let heartbeat_task = flow::heartbeat::HeartbeatTask::new(
353            &opts,
354            meta_client.clone(),
355            Arc::new(executor),
356            Arc::new(resource_stat),
357        );
358
359        let flow_metadata_manager = Arc::new(FlowMetadataManager::new(cached_meta_backend.clone()));
360        let frontend_client = FrontendClient::from_meta_client(
361            meta_client.clone(),
362            opts.query.clone(),
363            opts.flow.batching_mode.clone(),
364        )
365        .context(StartFlownodeSnafu)?;
366        let frontend_client = Arc::new(frontend_client);
367        let mut flownode_builder = FlownodeBuilder::new(
368            opts.clone(),
369            plugins.clone(),
370            table_metadata_manager,
371            catalog_manager.clone(),
372            flow_metadata_manager,
373            frontend_client.clone(),
374        )
375        .with_heartbeat_task(heartbeat_task);
376
377        plugins::setup_flownode_plugins_post_build(&mut plugins, &plugin_opts, &flownode_builder)
378            .await
379            .context(StartFlownodeSnafu)?;
380        flownode_builder.set_plugins(plugins.clone());
381
382        let mut flownode = flownode_builder.build().await.context(StartFlownodeSnafu)?;
383
384        let builder =
385            FlownodeServiceBuilder::grpc_server_builder(&opts, flownode.flownode_server());
386        let builder = if let Some(configurator) =
387            plugins.get::<GrpcBuilderConfiguratorRef<GrpcConfigureContext>>()
388        {
389            let context = GrpcConfigureContext {
390                kv_backend: cached_meta_backend.clone(),
391                fe_client: frontend_client.clone(),
392                flownode_id: member_id,
393                catalog_manager: catalog_manager.clone(),
394            };
395            configurator
396                .configure(builder, context)
397                .await
398                .context(OtherSnafu)?
399        } else {
400            builder
401        };
402        let grpc_server = builder.build();
403
404        let services = FlownodeServiceBuilder::new(&opts)
405            .with_grpc_server(grpc_server)
406            .enable_http_service()
407            .build()
408            .context(StartFlownodeSnafu)?;
409        flownode.setup_services(services);
410
411        Ok(Instance::new(flownode, guard))
412    }
413}
414
415#[cfg(test)]
416mod tests {
417    use clap::{CommandFactory, Parser};
418
419    use super::*;
420
421    #[test]
422    fn test_parse_grpc_cli_aliases() {
423        let command = StartCommand::try_parse_from([
424            "flownode",
425            "--grpc-bind-addr",
426            "127.0.0.1:14004",
427            "--grpc-server-addr",
428            "10.0.0.1:14004",
429        ])
430        .unwrap();
431        assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:14004"));
432        assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.1:14004"));
433
434        let command = StartCommand::try_parse_from([
435            "flownode",
436            "--rpc-bind-addr",
437            "127.0.0.1:24004",
438            "--rpc-server-addr",
439            "10.0.0.2:24004",
440        ])
441        .unwrap();
442        assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:24004"));
443        assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.2:24004"));
444
445        let command = StartCommand::try_parse_from([
446            "flownode",
447            "--rpc-addr",
448            "127.0.0.1:34004",
449            "--rpc-hostname",
450            "10.0.0.3:34004",
451        ])
452        .unwrap();
453        assert_eq!(command.grpc_bind_addr.as_deref(), Some("127.0.0.1:34004"));
454        assert_eq!(command.grpc_server_addr.as_deref(), Some("10.0.0.3:34004"));
455    }
456
457    #[test]
458    fn test_help_uses_grpc_option_names() {
459        let mut cmd = StartCommand::command();
460        let mut help = Vec::new();
461        cmd.write_long_help(&mut help).unwrap();
462        let help = String::from_utf8(help).unwrap();
463
464        assert!(help.contains("--grpc-bind-addr"));
465        assert!(help.contains("--grpc-server-addr"));
466        assert!(!help.contains("--rpc-bind-addr"));
467        assert!(!help.contains("--rpc-server-addr"));
468        assert!(!help.contains("--rpc-addr"));
469        assert!(!help.contains("--rpc-hostname"));
470        assert!(!help.contains("--user-provider"));
471    }
472
473    #[test]
474    fn test_user_provider_cli_option_is_removed() {
475        let command = StartCommand::try_parse_from([
476            "flownode",
477            "--node-id",
478            "14",
479            "--user-provider",
480            "static_user_provider:cmd:test=test",
481        ]);
482        assert!(command.is_err());
483    }
484}