1use std::sync::Arc;
16use std::time::Duration;
17
18use otlp_embedded::TraceServiceServer;
19use regex::Regex;
20use risingwave_common::config::SessionInitConfig;
21use risingwave_common::monitor::{RouterExt, TcpConfig};
22use risingwave_common::secret::LocalSecretManager;
23use risingwave_common::system_param::reader::SystemParamsRead;
24use risingwave_common::telemetry::manager::TelemetryManager;
25use risingwave_common::telemetry::{report_scarf_enabled, report_to_scarf, telemetry_env_enabled};
26use risingwave_common::util::tokio_util::sync::CancellationToken;
27use risingwave_common_service::{AwaitTreeMiddlewareLayer, MetricsManager, TracingExtractLayer};
28use risingwave_meta::MetaStoreBackend;
29use risingwave_meta::barrier::GlobalBarrierManager;
30use risingwave_meta::controller::catalog::CatalogController;
31use risingwave_meta::controller::cluster::ClusterController;
32use risingwave_meta::hummock::IcebergCompactorManager;
33use risingwave_meta::manager::iceberg_compaction::IcebergCompactionManager;
34use risingwave_meta::manager::{META_NODE_ID, MetadataManager};
35use risingwave_meta::rpc::ElectionClientRef;
36use risingwave_meta::rpc::election::dummy::DummyElectionClient;
37use risingwave_meta::rpc::intercept::MetricsMiddlewareLayer;
38use risingwave_meta::stream::{GlobalRefreshManager, ScaleController};
39use risingwave_meta_service::AddressInfo;
40use risingwave_meta_service::backup_service::BackupServiceImpl;
41use risingwave_meta_service::cloud_service::CloudServiceImpl;
42use risingwave_meta_service::cluster_limit_service::ClusterLimitServiceImpl;
43use risingwave_meta_service::cluster_service::ClusterServiceImpl;
44use risingwave_meta_service::ddl_service::DdlServiceImpl;
45use risingwave_meta_service::event_log_service::EventLogServiceImpl;
46use risingwave_meta_service::health_service::HealthServiceImpl;
47use risingwave_meta_service::heartbeat_service::HeartbeatServiceImpl;
48use risingwave_meta_service::hosted_iceberg_catalog_service::HostedIcebergCatalogServiceImpl;
49use risingwave_meta_service::hummock_service::HummockServiceImpl;
50use risingwave_meta_service::meta_member_service::MetaMemberServiceImpl;
51use risingwave_meta_service::monitor_service::MonitorServiceImpl;
52use risingwave_meta_service::notification_service::NotificationServiceImpl;
53use risingwave_meta_service::scale_service::ScaleServiceImpl;
54use risingwave_meta_service::serving_service::ServingServiceImpl;
55use risingwave_meta_service::session_config::SessionParamsServiceImpl;
56use risingwave_meta_service::sink_coordination_service::SinkCoordinationServiceImpl;
57use risingwave_meta_service::stream_service::StreamServiceImpl;
58use risingwave_meta_service::system_params_service::SystemParamsServiceImpl;
59use risingwave_meta_service::telemetry_service::TelemetryInfoServiceImpl;
60use risingwave_meta_service::user_service::UserServiceImpl;
61use risingwave_pb::backup_service::backup_service_server::BackupServiceServer;
62use risingwave_pb::cloud_service::cloud_service_server::CloudServiceServer;
63use risingwave_pb::configured_monitor_service_server;
64use risingwave_pb::connector_service::sink_coordination_service_server::SinkCoordinationServiceServer;
65use risingwave_pb::ddl_service::ddl_service_server::DdlServiceServer;
66use risingwave_pb::health::health_server::HealthServer;
67use risingwave_pb::hummock::hummock_manager_service_server::HummockManagerServiceServer;
68use risingwave_pb::meta::SystemParams;
69use risingwave_pb::meta::cluster_limit_service_server::ClusterLimitServiceServer;
70use risingwave_pb::meta::cluster_service_server::ClusterServiceServer;
71use risingwave_pb::meta::event_log_service_server::EventLogServiceServer;
72use risingwave_pb::meta::heartbeat_service_server::HeartbeatServiceServer;
73use risingwave_pb::meta::hosted_iceberg_catalog_service_server::HostedIcebergCatalogServiceServer;
74use risingwave_pb::meta::meta_member_service_server::MetaMemberServiceServer;
75use risingwave_pb::meta::notification_service_server::NotificationServiceServer;
76use risingwave_pb::meta::scale_service_server::ScaleServiceServer;
77use risingwave_pb::meta::serving_service_server::ServingServiceServer;
78use risingwave_pb::meta::session_param_service_server::SessionParamServiceServer;
79use risingwave_pb::meta::stream_manager_service_server::StreamManagerServiceServer;
80use risingwave_pb::meta::system_params_service_server::SystemParamsServiceServer;
81use risingwave_pb::meta::telemetry_info_service_server::TelemetryInfoServiceServer;
82use risingwave_pb::monitor_service::monitor_service_server::MonitorServiceServer;
83use risingwave_pb::user::user_service_server::UserServiceServer;
84use sea_orm::{ConnectionTrait, DbBackend};
85use thiserror_ext::AsReport;
86use tokio::sync::watch;
87
88use crate::backup_restore::BackupManager;
89use crate::barrier::BarrierScheduler;
90use crate::controller::SqlMetaStore;
91use crate::controller::system_param::SystemParamsController;
92use crate::hummock::HummockManager;
93use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
94use crate::manager::sink_coordination::SinkCoordinatorManager;
95use crate::manager::{IdleManager, MetaOpts, MetaSrvEnv};
96use crate::rpc::election::sql::{MySqlDriver, PostgresDriver, SqlBackendElectionClient};
97use crate::rpc::metrics::{GLOBAL_META_METRICS, start_info_monitor, start_worker_info_monitor};
98use crate::serving::ServingVnodeMapping;
99use crate::stream::{GlobalStreamManager, SourceManager};
100use crate::telemetry::{MetaReportCreator, MetaTelemetryInfoFetcher};
101use crate::{MetaError, MetaResult, hummock, serving};
102
103pub mod started {
106 use std::sync::atomic::AtomicBool;
107 use std::sync::atomic::Ordering::Relaxed;
108
109 static STARTED: AtomicBool = AtomicBool::new(false);
110
111 pub(crate) fn set() {
113 STARTED.store(true, Relaxed);
114 }
115
116 pub fn get() -> bool {
118 STARTED.load(Relaxed)
119 }
120}
121
122pub async fn rpc_serve(
126 address_info: AddressInfo,
127 meta_store_backend: MetaStoreBackend,
128 max_cluster_heartbeat_interval: Duration,
129 lease_interval_secs: u64,
130 server_config: risingwave_common::config::ServerConfig,
131 opts: MetaOpts,
132 init_system_params: SystemParams,
133 session_init: SessionInitConfig,
134 shutdown: CancellationToken,
135) -> MetaResult<()> {
136 let meta_store_impl = SqlMetaStore::connect(meta_store_backend.clone()).await?;
137
138 let election_client = match meta_store_backend {
139 MetaStoreBackend::Mem => {
140 Arc::new(DummyElectionClient::new(
142 address_info.advertise_addr.clone(),
143 ))
144 }
145 MetaStoreBackend::Sql { .. } => {
146 let id = address_info.advertise_addr.clone();
148 let conn = meta_store_impl.conn.clone();
149 let election_client: ElectionClientRef = match conn.get_database_backend() {
150 DbBackend::Sqlite => Arc::new(DummyElectionClient::new(id)),
151 DbBackend::Postgres => {
152 Arc::new(SqlBackendElectionClient::new(id, PostgresDriver::new(conn)))
153 }
154 DbBackend::MySql => {
155 Arc::new(SqlBackendElectionClient::new(id, MySqlDriver::new(conn)))
156 }
157 };
158 election_client.init().await?;
159
160 election_client
161 }
162 };
163
164 Box::pin(rpc_serve_with_store(
165 meta_store_impl,
166 election_client,
167 address_info,
168 max_cluster_heartbeat_interval,
169 lease_interval_secs,
170 server_config,
171 opts,
172 init_system_params,
173 session_init,
174 shutdown,
175 ))
176 .await
177}
178
179pub async fn rpc_serve_with_store(
184 meta_store_impl: SqlMetaStore,
185 election_client: ElectionClientRef,
186 address_info: AddressInfo,
187 max_cluster_heartbeat_interval: Duration,
188 lease_interval_secs: u64,
189 server_config: risingwave_common::config::ServerConfig,
190 opts: MetaOpts,
191 init_system_params: SystemParams,
192 session_init: SessionInitConfig,
193 shutdown: CancellationToken,
194) -> MetaResult<()> {
195 let (election_shutdown_tx, election_shutdown_rx) = watch::channel(());
197
198 let election_handle = tokio::spawn({
199 let shutdown = shutdown.clone();
200 let election_client = election_client.clone();
201
202 async move {
203 while let Err(e) = election_client
204 .run_once(lease_interval_secs as i64, election_shutdown_rx.clone())
205 .await
206 {
207 tracing::error!(error = %e.as_report(), "an election error occurred");
208 }
209 shutdown.cancel();
211 }
212 });
213
214 if !election_client.is_leader() {
219 let follower_shutdown = shutdown.child_token();
221
222 let follower_handle = tokio::spawn(start_service_as_election_follower(
223 follower_shutdown.clone(),
224 address_info.clone(),
225 election_client.clone(),
226 ));
227
228 let mut is_leader_watcher = election_client.subscribe();
230
231 while !*is_leader_watcher.borrow_and_update() {
232 tokio::select! {
233 _ = shutdown.cancelled() => return Ok(()),
235
236 res = is_leader_watcher.changed() => {
237 if res.is_err() {
238 tracing::error!("failed to receive a leader watcher update");
239 }
240 }
241 }
242 }
243
244 tracing::info!("elected as leader, shutting down follower services");
245 follower_shutdown.cancel();
246 let _ = follower_handle.await;
247 }
248
249 let result = start_service_as_election_leader(
251 meta_store_impl,
252 address_info,
253 max_cluster_heartbeat_interval,
254 opts,
255 init_system_params,
256 session_init,
257 server_config,
258 election_client,
259 shutdown,
260 )
261 .await;
262
263 election_shutdown_tx.send(()).ok();
265 let _ = election_handle.await;
266
267 result
268}
269
270pub async fn start_service_as_election_follower(
274 shutdown: CancellationToken,
275 address_info: AddressInfo,
276 election_client: ElectionClientRef,
277) {
278 tracing::info!("starting follower services");
279
280 let meta_member_srv = MetaMemberServiceImpl::new(election_client);
281
282 let health_srv = HealthServiceImpl::new();
283
284 let server = tonic::transport::Server::builder()
285 .layer(MetricsMiddlewareLayer::new(Arc::new(
286 GLOBAL_META_METRICS.clone(),
287 )))
288 .layer(TracingExtractLayer::new())
289 .add_service(MetaMemberServiceServer::new(meta_member_srv))
290 .add_service(HealthServer::new(health_srv))
291 .monitored_serve_with_shutdown(
292 address_info.listen_addr,
293 "grpc-meta-follower-service",
294 TcpConfig {
295 tcp_nodelay: true,
296 keepalive_duration: None,
297 },
298 shutdown.clone().cancelled_owned(),
299 );
300 let server_handle = tokio::spawn(server);
301 started::set();
302
303 shutdown.cancelled().await;
305 let _ = server_handle.await;
308}
309
310pub async fn start_service_as_election_leader(
314 meta_store_impl: SqlMetaStore,
315 address_info: AddressInfo,
316 max_cluster_heartbeat_interval: Duration,
317 opts: MetaOpts,
318 init_system_params: SystemParams,
319 session_init: SessionInitConfig,
320 server_config: risingwave_common::config::ServerConfig,
321 election_client: ElectionClientRef,
322 shutdown: CancellationToken,
323) -> MetaResult<()> {
324 tracing::info!("starting leader services");
325
326 let env = MetaSrvEnv::new(
327 opts.clone(),
328 init_system_params,
329 session_init,
330 meta_store_impl,
331 )
332 .await?;
333 tracing::info!("MetaSrvEnv started");
334 let _ = env.may_start_watch_license_key_file()?;
335 let system_params_reader = env.system_params_reader().await;
336
337 let data_directory = system_params_reader.data_directory();
338 if !is_correct_data_directory(data_directory) {
339 return Err(MetaError::system_params(format!(
340 "The data directory {:?} is misconfigured.
341 Please use a combination of uppercase and lowercase letters and numbers, i.e. [a-z, A-Z, 0-9].
342 The string cannot start or end with '/', and consecutive '/' are not allowed.
343 The data directory cannot be empty and its length should not exceed 800 characters.",
344 data_directory
345 )));
346 }
347
348 let cluster_controller = Arc::new(
349 ClusterController::new(env.clone(), max_cluster_heartbeat_interval)
350 .await
351 .unwrap(),
352 );
353 let catalog_controller = Arc::new(CatalogController::new(env.clone()).await?);
354 let metadata_manager = MetadataManager::new(cluster_controller, catalog_controller);
355
356 let serving_vnode_mapping = Arc::new(ServingVnodeMapping::default());
357 let max_serving_parallelism = env
358 .session_params_manager_impl_ref()
359 .get_params()
360 .await
361 .batch_parallelism()
362 .map(|p| p.get());
363 serving::on_meta_start(
364 env.notification_manager_ref(),
365 &metadata_manager,
366 serving_vnode_mapping.clone(),
367 max_serving_parallelism,
368 )
369 .await;
370
371 let compactor_manager = Arc::new(
372 hummock::CompactorManager::with_meta(env.clone())
373 .await
374 .unwrap(),
375 );
376 tracing::info!("CompactorManager started");
377
378 let heartbeat_srv = HeartbeatServiceImpl::new(metadata_manager.clone());
379 tracing::info!("HeartbeatServiceImpl started");
380
381 let (compactor_streams_change_tx, compactor_streams_change_rx) =
382 tokio::sync::mpsc::unbounded_channel();
383
384 let meta_metrics = Arc::new(GLOBAL_META_METRICS.clone());
385
386 let hummock_manager = hummock::HummockManager::new(
387 env.clone(),
388 metadata_manager.clone(),
389 meta_metrics.clone(),
390 compactor_manager.clone(),
391 compactor_streams_change_tx,
392 )
393 .await
394 .unwrap();
395 tracing::info!("HummockManager started");
396 let object_store_media_type = hummock_manager.object_store_media_type();
397
398 let meta_member_srv = MetaMemberServiceImpl::new(election_client.clone());
399
400 let prometheus_client = opts.prometheus_endpoint.as_ref().map(|x| {
401 use std::str::FromStr;
402 prometheus_http_query::Client::from_str(x).unwrap()
403 });
404 let prometheus_selector = opts.prometheus_selector.unwrap_or_default();
405
406 let trace_state = otlp_embedded::State::new(otlp_embedded::Config {
407 max_length: opts.cached_traces_num,
408 max_memory_usage: opts.cached_traces_memory_limit_bytes,
409 });
410 let trace_srv = otlp_embedded::TraceServiceImpl::new(trace_state.clone());
411
412 let (barrier_scheduler, scheduled_barriers) =
413 BarrierScheduler::new_pair(hummock_manager.clone());
414 tracing::info!("BarrierScheduler started");
415
416 let backup_manager = BackupManager::new(
418 env.clone(),
419 hummock_manager.clone(),
420 meta_metrics.clone(),
421 system_params_reader.backup_storage_url(),
422 system_params_reader.backup_storage_directory(),
423 )
424 .await?;
425 tracing::info!("BackupManager started");
426
427 LocalSecretManager::init(
428 opts.temp_secret_file_dir,
429 env.cluster_id().to_string(),
430 META_NODE_ID,
431 );
432 tracing::info!("LocalSecretManager started");
433
434 let notification_srv = NotificationServiceImpl::new(
435 env.clone(),
436 metadata_manager.clone(),
437 hummock_manager.clone(),
438 backup_manager.clone(),
439 serving_vnode_mapping.clone(),
440 )
441 .await?;
442 tracing::info!("NotificationServiceImpl started");
443
444 let source_manager = Arc::new(
445 SourceManager::new(
446 barrier_scheduler.clone(),
447 metadata_manager.clone(),
448 meta_metrics.clone(),
449 env.clone(),
450 )
451 .await?,
452 );
453 tracing::info!("SourceManager started");
454
455 let (iceberg_compaction_stat_tx, iceberg_compaction_stat_rx) =
456 tokio::sync::mpsc::unbounded_channel();
457 let (sink_manager, shutdown_handle) = SinkCoordinatorManager::start_worker(
458 env.meta_store_ref().conn.clone(),
459 hummock_manager.clone(),
460 metadata_manager.clone(),
461 iceberg_compaction_stat_tx,
462 env.await_tree_reg().clone(),
463 );
464 tracing::info!("SinkCoordinatorManager started");
465 let mut sub_tasks = vec![shutdown_handle];
467
468 sub_tasks.push(serving::start_serving_vnode_mapping_worker(
471 env.notification_manager_ref(),
472 metadata_manager.clone(),
473 serving_vnode_mapping.clone(),
474 env.session_params_manager_impl_ref(),
475 ));
476
477 let iceberg_pk_index_sink_manager =
478 IcebergPkIndexSinkManager::new(env.meta_store_ref().conn.clone());
479 tracing::info!("IcebergPkIndexSinkManager started");
480
481 let iceberg_compactor_manager = Arc::new(IcebergCompactorManager::new());
482
483 let (iceberg_compaction_mgr, iceberg_compactor_event_rx) = IcebergCompactionManager::build(
485 env.clone(),
486 metadata_manager.clone(),
487 iceberg_compactor_manager.clone(),
488 meta_metrics.clone(),
489 );
490
491 sub_tasks.push(IcebergCompactionManager::compaction_stat_loop(
492 iceberg_compaction_mgr.clone(),
493 iceberg_compaction_stat_rx,
494 ));
495
496 sub_tasks.push(IcebergCompactionManager::gc_loop(
497 iceberg_compaction_mgr.clone(),
498 env.opts.iceberg_gc_interval_sec,
499 ));
500
501 let refresh_scheduler_interval = Duration::from_secs(env.opts.refresh_scheduler_interval_sec);
502 let (refresh_manager, refresh_handle, refresh_shutdown) = GlobalRefreshManager::start(
503 metadata_manager.clone(),
504 barrier_scheduler.clone(),
505 &env,
506 refresh_scheduler_interval,
507 )
508 .await?;
509 sub_tasks.push((refresh_handle, refresh_shutdown));
510
511 let scale_controller = Arc::new(ScaleController::new(
512 &metadata_manager,
513 source_manager.clone(),
514 env.clone(),
515 ));
516
517 let (barrier_manager, join_handle, shutdown_rx) = GlobalBarrierManager::start(
518 scheduled_barriers,
519 env.clone(),
520 metadata_manager.clone(),
521 hummock_manager.clone(),
522 serving_vnode_mapping.clone(),
523 source_manager.clone(),
524 sink_manager.clone(),
525 iceberg_pk_index_sink_manager.clone(),
526 iceberg_compaction_mgr.clone(),
527 scale_controller.clone(),
528 barrier_scheduler.clone(),
529 refresh_manager.clone(),
530 )
531 .await;
532 tracing::info!("GlobalBarrierManager started");
533 sub_tasks.push((join_handle, shutdown_rx));
534
535 {
536 let source_manager = source_manager.clone();
537 tokio::spawn(async move {
538 source_manager.run().await.unwrap();
539 });
540 }
541
542 let stream_manager = Arc::new(
543 GlobalStreamManager::new(
544 env.clone(),
545 metadata_manager.clone(),
546 barrier_scheduler.clone(),
547 hummock_manager.clone(),
548 source_manager.clone(),
549 refresh_manager.clone(),
550 iceberg_compaction_mgr.clone(),
551 scale_controller.clone(),
552 )
553 .unwrap(),
554 );
555
556 hummock_manager
557 .may_fill_backward_state_table_info()
558 .await
559 .unwrap();
560
561 let ddl_srv = DdlServiceImpl::new(
562 env.clone(),
563 metadata_manager.clone(),
564 stream_manager.clone(),
565 source_manager.clone(),
566 barrier_manager.clone(),
567 sink_manager.clone(),
568 meta_metrics.clone(),
569 iceberg_compaction_mgr.clone(),
570 barrier_scheduler.clone(),
571 iceberg_pk_index_sink_manager.clone(),
572 )
573 .await;
574
575 if env.opts.enable_legacy_table_migration {
576 sub_tasks.push(ddl_srv.start_migrate_table_fragments());
577 }
578
579 let user_srv = UserServiceImpl::new(metadata_manager.clone());
580
581 let scale_srv = ScaleServiceImpl::new(
582 metadata_manager.clone(),
583 stream_manager.clone(),
584 barrier_manager.clone(),
585 env.clone(),
586 );
587
588 let cluster_srv = ClusterServiceImpl::new(metadata_manager.clone(), barrier_manager.clone());
589 let stream_srv = StreamServiceImpl::new(
590 env.clone(),
591 barrier_scheduler.clone(),
592 barrier_manager.clone(),
593 stream_manager.clone(),
594 metadata_manager.clone(),
595 refresh_manager.clone(),
596 iceberg_compaction_mgr.clone(),
597 );
598 let sink_coordination_srv = SinkCoordinationServiceImpl::new(sink_manager);
599 let hummock_srv = HummockServiceImpl::new(
600 hummock_manager.clone(),
601 metadata_manager.clone(),
602 backup_manager.clone(),
603 iceberg_compaction_mgr.clone(),
604 );
605
606 let health_srv = HealthServiceImpl::new();
607 let backup_srv = BackupServiceImpl::new(backup_manager.clone());
608 let telemetry_srv = TelemetryInfoServiceImpl::new(env.meta_store());
609 let system_params_srv = SystemParamsServiceImpl::new(
610 env.system_params_manager_impl_ref(),
611 metadata_manager.clone(),
612 env.opts.license_key_path.is_some(),
613 );
614 let session_params_srv = SessionParamsServiceImpl::new(env.session_params_manager_impl_ref());
615 let serving_srv =
616 ServingServiceImpl::new(serving_vnode_mapping.clone(), metadata_manager.clone());
617 let cloud_srv = CloudServiceImpl::new();
618 let event_log_srv = EventLogServiceImpl::new(env.event_log_manager_ref());
619 let cluster_limit_srv = ClusterLimitServiceImpl::new(env.clone(), metadata_manager.clone());
620 let hosted_iceberg_catalog_srv = HostedIcebergCatalogServiceImpl::new(env.clone());
621 let monitor_srv = MonitorServiceImpl::new(
622 metadata_manager.clone(),
623 env.await_tree_reg().clone(),
624 server_config.clone(),
625 );
626 let diagnose_command = Arc::new(risingwave_meta::manager::diagnose::DiagnoseCommand::new(
627 metadata_manager.clone(),
628 env.await_tree_reg().clone(),
629 hummock_manager.clone(),
630 iceberg_compaction_mgr.clone(),
631 env.event_log_manager_ref(),
632 prometheus_client.clone(),
633 prometheus_selector.clone(),
634 opts.redact_sql_option_keywords.clone(),
635 env.system_params_manager_impl_ref(),
636 ));
637
638 #[cfg(not(madsim))]
639 let _dashboard_task = if let Some(ref dashboard_addr) = address_info.dashboard_addr {
640 use risingwave_common::config::RpcClientConfig;
641 use risingwave_rpc_client::MonitorClientPool;
642
643 let dashboard_service = crate::dashboard::DashboardService {
644 await_tree_reg: env.await_tree_reg().clone(),
645 dashboard_addr: *dashboard_addr,
646 prometheus_client,
647 prometheus_selector,
648 metadata_manager: metadata_manager.clone(),
649 hummock_manager: hummock_manager.clone(),
650 monitor_clients: MonitorClientPool::new(1, RpcClientConfig::default()),
651 diagnose_command,
652 profile_service: risingwave_common_heap_profiling::ProfileServiceImpl::new(
653 server_config.clone(),
654 ),
655 trace_state,
656 };
657 let task = tokio::spawn(dashboard_service.serve());
658 Some(task)
659 } else {
660 None
661 };
662
663 if let Some(prometheus_addr) = address_info.prometheus_addr {
664 MetricsManager::boot_metrics_service(prometheus_addr.to_string())
665 }
666
667 sub_tasks.extend(hummock::start_hummock_workers(
669 hummock_manager.clone(),
670 backup_manager.clone(),
671 &env.opts,
672 {
673 let catalog_controller = metadata_manager.catalog_controller.clone();
674 Box::new(move || {
675 let catalog_controller = catalog_controller.clone();
676 Box::pin(async move {
677 catalog_controller
678 .get_table_change_log_truncate_info()
679 .await
680 .map(Some)
681 .unwrap_or_else(|e| {
682 tracing::warn!(err = %e.as_report(), "failed to collect table change log retention metadata");
683 None
684 })
685 })
686 })
687 },
688 {
689 let catalog_controller = metadata_manager.catalog_controller.clone();
690 Box::new(move || {
691 let catalog_controller = catalog_controller.clone();
692 Box::pin(async move {
693 catalog_controller
694 .get_pinned_snapshot_epochs()
695 .await
696 .map(Some)
697 .unwrap_or_else(|e| {
698 tracing::warn!(
699 err = %e.as_report(),
700 "failed to collect pinned snapshot epochs; pausing time-travel vacuum",
701 );
702 None
703 })
704 })
705 })
706 },
707 ));
708 sub_tasks.push(start_worker_info_monitor(
709 metadata_manager.clone(),
710 election_client.clone(),
711 Duration::from_secs(env.opts.node_num_monitor_interval_sec),
712 meta_metrics.clone(),
713 ));
714 sub_tasks.push(start_info_monitor(
715 metadata_manager.clone(),
716 hummock_manager.clone(),
717 barrier_manager.clone(),
718 env.system_params_manager_impl_ref(),
719 meta_metrics.clone(),
720 ));
721 sub_tasks.push(SystemParamsController::start_params_notifier(
722 env.system_params_manager_impl_ref(),
723 ));
724 sub_tasks.push(HummockManager::hummock_timer_task(
725 hummock_manager.clone(),
726 Some(backup_manager),
727 ));
728 sub_tasks.extend(HummockManager::compaction_event_loop(
729 hummock_manager.clone(),
730 compactor_streams_change_rx,
731 ));
732
733 sub_tasks.extend(IcebergCompactionManager::iceberg_compaction_event_loop(
734 iceberg_compaction_mgr.clone(),
735 iceberg_compactor_event_rx,
736 ));
737
738 {
739 sub_tasks.push(ClusterController::start_heartbeat_checker(
740 metadata_manager.cluster_controller.clone(),
741 Duration::from_secs(1),
742 ));
743
744 if !env.opts.disable_automatic_parallelism_control {
745 sub_tasks.push(stream_manager.start_auto_parallelism_monitor());
746 }
747 }
748
749 let _idle_checker_handle = IdleManager::start_idle_checker(
750 env.idle_manager_ref(),
751 Duration::from_secs(30),
752 shutdown.clone(),
753 );
754
755 let (abort_sender, abort_recv) = tokio::sync::oneshot::channel();
756 let notification_mgr = env.notification_manager_ref();
757 let stream_abort_handler = tokio::spawn(async move {
758 let _ = abort_recv.await;
759 notification_mgr.abort_all();
760 compactor_manager.abort_all_compactors();
761 });
762 sub_tasks.push((stream_abort_handler, abort_sender));
763
764 let telemetry_manager = TelemetryManager::new(
765 Arc::new(MetaTelemetryInfoFetcher::new(env.cluster_id().clone())),
766 Arc::new(MetaReportCreator::new(
767 metadata_manager.clone(),
768 object_store_media_type,
769 )),
770 );
771
772 if env.opts.telemetry_enabled && telemetry_env_enabled() {
774 sub_tasks.push(telemetry_manager.start().await);
775 } else {
776 tracing::info!("Telemetry didn't start due to meta backend or config");
777 }
778 if !cfg!(madsim) && report_scarf_enabled() {
779 tokio::spawn(report_to_scarf());
780 } else {
781 tracing::info!("Scarf reporting is disabled");
782 };
783
784 if let Some(pair) = env.event_log_manager_ref().take_join_handle() {
785 sub_tasks.push(pair);
786 }
787
788 tracing::info!("Assigned cluster id {:?}", *env.cluster_id());
789 tracing::info!("Starting meta services");
790
791 let event = risingwave_pb::meta::event_log::EventMetaNodeStart {
792 advertise_addr: address_info.advertise_addr,
793 listen_addr: address_info.listen_addr.to_string(),
794 opts: serde_json::to_string(&env.opts).unwrap(),
795 };
796 env.event_log_manager_ref().add_event_logs(vec![
797 risingwave_pb::meta::event_log::Event::MetaNodeStart(event),
798 ]);
799
800 let server_builder = tonic::transport::Server::builder()
801 .layer(MetricsMiddlewareLayer::new(meta_metrics))
802 .layer(TracingExtractLayer::new())
803 .layer(AwaitTreeMiddlewareLayer::new(env.await_tree_reg().clone()))
804 .add_service(HeartbeatServiceServer::new(heartbeat_srv))
805 .add_service(ClusterServiceServer::new(cluster_srv))
806 .add_service(StreamManagerServiceServer::new(stream_srv))
807 .add_service(
808 HummockManagerServiceServer::new(hummock_srv).max_decoding_message_size(usize::MAX),
809 )
810 .add_service(NotificationServiceServer::new(notification_srv))
811 .add_service(MetaMemberServiceServer::new(meta_member_srv))
812 .add_service(DdlServiceServer::new(ddl_srv).max_decoding_message_size(usize::MAX))
813 .add_service(UserServiceServer::new(user_srv))
814 .add_service(CloudServiceServer::new(cloud_srv))
815 .add_service(ScaleServiceServer::new(scale_srv).max_decoding_message_size(usize::MAX))
816 .add_service(HealthServer::new(health_srv))
817 .add_service(BackupServiceServer::new(backup_srv))
818 .add_service(SystemParamsServiceServer::new(system_params_srv))
819 .add_service(SessionParamServiceServer::new(session_params_srv))
820 .add_service(TelemetryInfoServiceServer::new(telemetry_srv))
821 .add_service(ServingServiceServer::new(serving_srv))
822 .add_service(
823 SinkCoordinationServiceServer::new(sink_coordination_srv)
824 .max_decoding_message_size(usize::MAX),
825 )
826 .add_service(
827 EventLogServiceServer::new(event_log_srv).max_decoding_message_size(usize::MAX),
828 )
829 .add_service(ClusterLimitServiceServer::new(cluster_limit_srv))
830 .add_service(HostedIcebergCatalogServiceServer::new(
831 hosted_iceberg_catalog_srv,
832 ))
833 .add_service(configured_monitor_service_server(
834 MonitorServiceServer::new(monitor_srv),
835 ));
836
837 #[cfg(not(madsim))] let server_builder = server_builder.add_service(TraceServiceServer::new(trace_srv));
839
840 let server = server_builder.monitored_serve_with_shutdown(
841 address_info.listen_addr,
842 "grpc-meta-leader-service",
843 TcpConfig {
844 tcp_nodelay: true,
845 keepalive_duration: None,
846 },
847 shutdown.clone().cancelled_owned(),
848 );
849 started::set();
850 let _server_handle = tokio::spawn(server);
851
852 shutdown.cancelled().await;
854 Ok(())
857}
858
859fn is_correct_data_directory(data_directory: &str) -> bool {
860 let data_directory_regex = Regex::new(r"^[0-9a-zA-Z_/-]{1,}$").unwrap();
861 if data_directory.is_empty()
862 || !data_directory_regex.is_match(data_directory)
863 || data_directory.ends_with('/')
864 || data_directory.starts_with('/')
865 || data_directory.contains("//")
866 || data_directory.len() > 800
867 {
868 return false;
869 }
870 true
871}