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