Skip to main content

risingwave_meta_service/
cluster_service.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 risingwave_common::util::version::{current_rw_version, is_compatible_rw_version};
16use risingwave_meta::barrier::BarrierManagerRef;
17use risingwave_meta::manager::MetadataManager;
18use risingwave_pb::common::worker_node::State;
19use risingwave_pb::common::{HostAddress, WorkerType as PbWorkerType};
20use risingwave_pb::meta::cluster_service_server::ClusterService;
21use risingwave_pb::meta::{
22    ActivateWorkerNodeRequest, ActivateWorkerNodeResponse, AddWorkerNodeRequest,
23    AddWorkerNodeResponse, DeleteWorkerNodeRequest, DeleteWorkerNodeResponse,
24    GetClusterRecoveryStatusRequest, GetClusterRecoveryStatusResponse, GetMetaStoreInfoRequest,
25    GetMetaStoreInfoResponse, ListAllNodesRequest, ListAllNodesResponse,
26};
27use tonic::{Request, Response, Status};
28
29use crate::MetaError;
30
31#[derive(Clone)]
32pub struct ClusterServiceImpl {
33    metadata_manager: MetadataManager,
34    barrier_manager: BarrierManagerRef,
35}
36
37impl ClusterServiceImpl {
38    pub fn new(metadata_manager: MetadataManager, barrier_manager: BarrierManagerRef) -> Self {
39        ClusterServiceImpl {
40            metadata_manager,
41            barrier_manager,
42        }
43    }
44}
45
46#[async_trait::async_trait]
47impl ClusterService for ClusterServiceImpl {
48    async fn add_worker_node(
49        &self,
50        request: Request<AddWorkerNodeRequest>,
51    ) -> Result<Response<AddWorkerNodeResponse>, Status> {
52        let req = request.into_inner();
53        let worker_type = req.get_worker_type()?;
54        let host: HostAddress = req.get_host()?.clone();
55        let property = req
56            .property
57            .ok_or_else(|| MetaError::invalid_parameter("worker node property is not provided"))?;
58        let resource = req.resource.unwrap_or_default();
59        let current_rw_version = current_rw_version();
60        if matches!(
61            worker_type,
62            PbWorkerType::Frontend | PbWorkerType::ComputeNode | PbWorkerType::Compactor
63        ) && !is_compatible_rw_version(&resource.rw_version)
64        {
65            return Err(MetaError::invalid_parameter(format!(
66                "worker node version {} does not match meta node version {}",
67                resource.rw_version, current_rw_version,
68            ))
69            .into());
70        }
71        let start = std::time::Instant::now();
72        tracing::info!(
73            ?host,
74            ?worker_type,
75            "add_worker_node: received register request"
76        );
77        let worker_id = self
78            .metadata_manager
79            .add_worker_node(worker_type, host.clone(), property, resource)
80            .await?;
81        tracing::info!(
82            ?host,
83            ?worker_id,
84            ?worker_type,
85            elapsed_ms = start.elapsed().as_millis() as u64,
86            "add_worker_node: registered"
87        );
88        let cluster_id = self.metadata_manager.cluster_id().to_string();
89
90        Ok(Response::new(AddWorkerNodeResponse {
91            node_id: Some(worker_id),
92            cluster_id,
93        }))
94    }
95
96    async fn activate_worker_node(
97        &self,
98        request: Request<ActivateWorkerNodeRequest>,
99    ) -> Result<Response<ActivateWorkerNodeResponse>, Status> {
100        let req = request.into_inner();
101        let host = req.get_host()?.clone();
102        #[cfg(not(madsim))]
103        {
104            use risingwave_common::util::addr::try_resolve_dns;
105            use tracing::{error, info};
106            let socket_addr = try_resolve_dns(&host.host, host.port).await.map_err(|e| {
107                error!(e);
108                Status::internal(e)
109            })?;
110            info!(?socket_addr, ?host, "resolve host addr");
111        }
112        self.metadata_manager
113            .cluster_controller
114            .activate_worker(req.node_id)
115            .await?;
116        Ok(Response::new(ActivateWorkerNodeResponse { status: None }))
117    }
118
119    async fn delete_worker_node(
120        &self,
121        request: Request<DeleteWorkerNodeRequest>,
122    ) -> Result<Response<DeleteWorkerNodeResponse>, Status> {
123        let req = request.into_inner();
124        let host = req.get_host()?.clone();
125
126        let worker_node = self
127            .metadata_manager
128            .cluster_controller
129            .delete_worker(host)
130            .await?;
131        tracing::info!(
132            host = ?worker_node.host,
133            id = %worker_node.id,
134            r#type = ?worker_node.r#type(),
135            "deleted worker node",
136        );
137
138        Ok(Response::new(DeleteWorkerNodeResponse { status: None }))
139    }
140
141    async fn list_all_nodes(
142        &self,
143        request: Request<ListAllNodesRequest>,
144    ) -> Result<Response<ListAllNodesResponse>, Status> {
145        let req = request.into_inner();
146        let worker_type = req.worker_type.map(|wt| wt.try_into().unwrap());
147        let worker_states = if req.include_starting_nodes {
148            None
149        } else {
150            Some(State::Running)
151        };
152
153        let node_list = self
154            .metadata_manager
155            .list_worker_node(worker_type, worker_states)
156            .await?;
157        Ok(Response::new(ListAllNodesResponse {
158            status: None,
159            nodes: node_list,
160        }))
161    }
162
163    async fn get_cluster_recovery_status(
164        &self,
165        _request: Request<GetClusterRecoveryStatusRequest>,
166    ) -> Result<Response<GetClusterRecoveryStatusResponse>, Status> {
167        Ok(Response::new(GetClusterRecoveryStatusResponse {
168            status: self.barrier_manager.get_recovery_status() as _,
169        }))
170    }
171
172    async fn get_meta_store_info(
173        &self,
174        _request: Request<GetMetaStoreInfoRequest>,
175    ) -> Result<Response<GetMetaStoreInfoResponse>, Status> {
176        Ok(Response::new(GetMetaStoreInfoResponse {
177            meta_store_endpoint: self
178                .metadata_manager
179                .cluster_controller
180                .meta_store_endpoint(),
181        }))
182    }
183}