Skip to main content

risingwave_meta_service/
system_params_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 async_trait::async_trait;
16use futures::future::try_join_all;
17use risingwave_common::config::RpcClientConfig;
18use risingwave_common::system_param::LICENSE_KEY_KEY;
19use risingwave_meta::controller::system_param::SystemParamsControllerRef;
20use risingwave_meta::manager::MetadataManager;
21use risingwave_pb::common::WorkerType;
22use risingwave_pb::meta::system_params_service_server::SystemParamsService;
23use risingwave_pb::meta::{
24    ClearFileCacheRequest, ClearFileCacheResponse, GetSystemParamsRequest, GetSystemParamsResponse,
25    SetSystemParamRequest, SetSystemParamResponse,
26};
27use risingwave_rpc_client::ComputeClient;
28use thiserror_ext::AsReport;
29use tonic::{Request, Response, Status};
30
31pub struct SystemParamsServiceImpl {
32    system_params_manager: SystemParamsControllerRef,
33    metadata_manager: MetadataManager,
34
35    /// Whether the license key is managed by license key file, i.e., `--license-key-path` is set.
36    managed_license_key: bool,
37}
38
39impl SystemParamsServiceImpl {
40    pub fn new(
41        system_params_manager: SystemParamsControllerRef,
42        metadata_manager: MetadataManager,
43        managed_license_key: bool,
44    ) -> Self {
45        Self {
46            system_params_manager,
47            metadata_manager,
48            managed_license_key,
49        }
50    }
51}
52
53#[async_trait]
54impl SystemParamsService for SystemParamsServiceImpl {
55    async fn get_system_params(
56        &self,
57        _request: Request<GetSystemParamsRequest>,
58    ) -> Result<Response<GetSystemParamsResponse>, Status> {
59        let params = self.system_params_manager.get_pb_params().await;
60
61        Ok(Response::new(GetSystemParamsResponse {
62            params: Some(params),
63        }))
64    }
65
66    async fn set_system_param(
67        &self,
68        request: Request<SetSystemParamRequest>,
69    ) -> Result<Response<SetSystemParamResponse>, Status> {
70        let req = request.into_inner();
71
72        // When license key path is specified, license key from system parameters can be easily
73        // overwritten. So we simply reject this case.
74        if self.managed_license_key && req.param == LICENSE_KEY_KEY {
75            return Err(Status::permission_denied(
76                "cannot alter license key manually when \
77                argument `--license-key-path` (or env var `RW_LICENSE_KEY_PATH`) is set, \
78                please update the license key file instead",
79            ));
80        }
81
82        let params = self
83            .system_params_manager
84            .set_param(&req.param, req.value)
85            .await?;
86
87        Ok(Response::new(SetSystemParamResponse {
88            params: Some(params),
89        }))
90    }
91
92    async fn clear_file_cache(
93        &self,
94        request: Request<ClearFileCacheRequest>,
95    ) -> Result<Response<ClearFileCacheResponse>, Status> {
96        let request = request.into_inner();
97        let worker_nodes = self
98            .metadata_manager
99            .list_worker_node(Some(WorkerType::ComputeNode), None)
100            .await
101            .map_err(|e| Status::internal(e.to_report_string()))?;
102
103        let clear_futures = worker_nodes.into_iter().map(|worker| async move {
104            let worker_id = worker.id;
105            let host = worker
106                .get_host()
107                .map_err(|e| Status::internal(e.to_report_string()))?
108                .clone();
109            let client = ComputeClient::new((&host).into(), &RpcClientConfig::default())
110                .await
111                .map_err(|e| {
112                    Status::internal(format!(
113                        "connect to compute node {worker_id}: {}",
114                        e.as_report()
115                    ))
116                })?;
117            client
118                .resize_cache(risingwave_pb::compute::ResizeCacheRequest {
119                    meta_cache_capacity: 0,
120                    data_cache_capacity: 0,
121                    clear_meta_cache: request.clear_meta_cache,
122                    clear_data_cache: request.clear_data_cache,
123                })
124                .await
125                .map_err(|e| {
126                    Status::internal(format!(
127                        "clear file cache on {worker_id}: {}",
128                        e.as_report()
129                    ))
130                })?;
131            Ok::<_, Status>(())
132        });
133
134        try_join_all(clear_futures).await?;
135        Ok(Response::new(ClearFileCacheResponse {}))
136    }
137}