risingwave_meta_service/
cluster_service.rs1use 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}