1use 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 _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 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 #[clap(long)]
139 node_id: Option<u64>,
140 #[clap(long = "grpc-bind-addr", alias = "rpc-bind-addr", alias = "rpc-addr")]
142 grpc_bind_addr: Option<String>,
143 #[clap(
147 long = "grpc-server-addr",
148 alias = "rpc-server-addr",
149 alias = "rpc-hostname"
150 )]
151 grpc_server_addr: Option<String>,
152 #[clap(long, value_delimiter = ',', num_args = 1..)]
154 metasrv_addrs: Option<Vec<String>>,
155 #[clap(short, long)]
157 config_file: Option<String>,
158 #[clap(long, default_value = "GREPTIMEDB_FLOWNODE")]
160 env_prefix: String,
161 #[clap(long)]
162 http_addr: Option<String>,
163 #[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 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 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 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 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 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}