Skip to main content

risingwave_meta_node/
server.rs

1// Copyright 2023 RisingWave Labs
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::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
103/// Used for standalone mode checking the status of the meta service.
104/// This can be easier and more accurate than checking the TCP connection.
105pub 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    /// Mark the meta service as started.
112    pub(crate) fn set() {
113        STARTED.store(true, Relaxed);
114    }
115
116    /// Check if the meta service has started.
117    pub fn get() -> bool {
118        STARTED.load(Relaxed)
119    }
120}
121
122/// A wrapper around [`rpc_serve_with_store`] that dispatches different store implementations.
123///
124/// For the timing of returning, see [`rpc_serve_with_store`].
125pub 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            // Use a dummy election client.
141            Arc::new(DummyElectionClient::new(
142                address_info.advertise_addr.clone(),
143            ))
144        }
145        MetaStoreBackend::Sql { .. } => {
146            // Init election client.
147            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
179/// Bootstraps the follower or leader service based on the election status.
180///
181/// Returns when the `shutdown` token is triggered, or when leader status is lost, or if the leader
182/// service fails to start.
183pub 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    // TODO(shutdown): directly use cancellation token
196    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            // Leader lost, shutdown the service.
210            shutdown.cancel();
211        }
212    });
213
214    // Spawn and run the follower service if not the leader.
215    // Watch the leader status and switch to the leader service when elected.
216    // TODO: the branch seems to be always hit since the default value of `is_leader` is false until
217    // the election is done (unless using `DummyElectionClient`).
218    if !election_client.is_leader() {
219        // The follower service can be shutdown separately if we're going to be the leader.
220        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        // Watch and wait until we become the leader.
229        let mut is_leader_watcher = election_client.subscribe();
230
231        while !*is_leader_watcher.borrow_and_update() {
232            tokio::select! {
233                // External shutdown signal. Directly return without switching to leader.
234                _ = 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    // Run the leader service.
250    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    // Leader service has stopped, shutdown the election service to gracefully resign.
264    election_shutdown_tx.send(()).ok();
265    let _ = election_handle.await;
266
267    result
268}
269
270/// Starts all services needed for the meta follower node.
271///
272/// Returns when the `shutdown` token is triggered.
273pub 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    // Wait for the shutdown signal.
304    shutdown.cancelled().await;
305    // Wait for the server to shutdown. This is necessary because we may be transitioning from follower
306    // to leader, and conflicts on the services must be avoided.
307    let _ = server_handle.await;
308}
309
310/// Starts all services needed for the meta leader node.
311///
312/// Returns when the `shutdown` token is triggered, or if the service initialization fails.
313pub 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    // Initialize services.
417    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    // TODO(shutdown): remove this as there's no need to gracefully shutdown some of these sub-tasks.
466    let mut sub_tasks = vec![shutdown_handle];
467
468    // Register before the barrier manager starts recovery. Dirty creating-job cleanup emits local
469    // serving mapping deletes, which must not be dropped during bootstrap.
470    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    // TODO: introduce compactor event stream handler to handle iceberg compaction events.
484    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 executed concurrently. Can be shutdown via shutdown_all
668    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    // May start telemetry reporting
773    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))] // `otlp-embedded` does not use madsim-patched tonic
838    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    // Wait for the shutdown signal.
853    shutdown.cancelled().await;
854    // TODO(shutdown): may warn user if there's any other node still running in the cluster.
855    // TODO(shutdown): do we have any other shutdown tasks?
856    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}