1use std::any::Any;
16use std::future::Future;
17use std::sync::Arc;
18
19use api::greptime_proto::v1;
20use api::v1::meta::cluster_client::ClusterClient;
21use api::v1::meta::{MetasrvNodeInfo, MetasrvPeersRequest, ResponseHeader};
22use common_error::ext::BoxedError;
23use common_grpc::channel_manager::ChannelManager;
24use common_meta::error::{
25 Error as MetaError, ExternalSnafu, ResponseExceededSizeLimitSnafu, Result as MetaResult,
26};
27use common_meta::kv_backend::{KvBackend, TxnService};
28use common_meta::rpc::store::{
29 BatchDeleteRequest, BatchDeleteResponse, BatchGetRequest, BatchGetResponse, BatchPutRequest,
30 BatchPutResponse, DeleteRangeRequest, DeleteRangeResponse, PutRequest, PutResponse,
31 RangeRequest, RangeResponse,
32};
33use common_telemetry::{error, info, warn};
34use snafu::{ResultExt, ensure};
35use tokio::sync::RwLock;
36use tonic::Status;
37use tonic::transport::Channel;
38
39use crate::client::{LeaderProviderRef, util};
40use crate::error::{
41 ConvertMetaResponseSnafu, CreateChannelSnafu, Error, IllegalGrpcClientStateSnafu,
42 ReadOnlyKvBackendSnafu, Result, RetryTimesExceededSnafu,
43};
44
45#[derive(Clone, Debug)]
46pub struct Client {
47 inner: Arc<RwLock<Inner>>,
48}
49
50impl Client {
51 pub fn new(channel_manager: ChannelManager, max_retry: usize) -> Self {
52 let inner = Arc::new(RwLock::new(Inner {
53 channel_manager,
54 leader_provider: None,
55 max_retry,
56 }));
57
58 Self { inner }
59 }
60
61 pub(crate) async fn start_with(&self, leader_provider: LeaderProviderRef) -> Result<()> {
63 let mut inner = self.inner.write().await;
64 inner.start_with(leader_provider)
65 }
66
67 pub async fn range(&self, req: RangeRequest) -> Result<RangeResponse> {
68 let inner = self.inner.read().await;
69 inner.range(req).await
70 }
71
72 #[allow(dead_code)]
73 pub async fn batch_get(&self, req: BatchGetRequest) -> Result<BatchGetResponse> {
74 let inner = self.inner.read().await;
75 inner.batch_get(req).await
76 }
77
78 pub async fn get_metasrv_peers(
79 &self,
80 ) -> Result<(Option<MetasrvNodeInfo>, Vec<MetasrvNodeInfo>)> {
81 let inner = self.inner.read().await;
82 inner.get_metasrv_peers().await
83 }
84}
85
86impl TxnService for Client {
87 type Error = MetaError;
88}
89
90#[async_trait::async_trait]
91impl KvBackend for Client {
92 fn name(&self) -> &str {
93 "ClusterClientKvBackend"
94 }
95
96 fn as_any(&self) -> &dyn Any {
97 self
98 }
99
100 async fn range(&self, req: RangeRequest) -> MetaResult<RangeResponse> {
101 let resp = self.range(req).await;
102 match resp {
103 Ok(resp) => Ok(resp),
104 Err(err) if err.is_exceeded_size_limit() => {
105 Err(BoxedError::new(err)).context(ResponseExceededSizeLimitSnafu)
106 }
107 Err(err) => Err(BoxedError::new(err)).context(ExternalSnafu),
108 }
109 }
110
111 async fn put(&self, _: PutRequest) -> MetaResult<PutResponse> {
112 unimplemented!("`put` is not supported in cluster client kv backend")
113 }
114
115 async fn batch_put(&self, _: BatchPutRequest) -> MetaResult<BatchPutResponse> {
116 unimplemented!("`batch_put` is not supported in cluster client kv backend")
117 }
118
119 async fn batch_get(&self, req: BatchGetRequest) -> MetaResult<BatchGetResponse> {
120 self.batch_get(req)
121 .await
122 .map_err(BoxedError::new)
123 .context(ExternalSnafu)
124 }
125
126 async fn delete_range(&self, _: DeleteRangeRequest) -> MetaResult<DeleteRangeResponse> {
127 unimplemented!("`delete_range` is not supported in cluster client kv backend")
128 }
129
130 async fn batch_delete(&self, _: BatchDeleteRequest) -> MetaResult<BatchDeleteResponse> {
131 unimplemented!("`batch_delete` is not supported in cluster client kv backend")
132 }
133}
134
135#[derive(Debug)]
136struct Inner {
137 channel_manager: ChannelManager,
138 leader_provider: Option<LeaderProviderRef>,
139 max_retry: usize,
140}
141
142impl Inner {
143 fn start_with(&mut self, leader_provider: LeaderProviderRef) -> Result<()> {
144 ensure!(
145 !self.is_started(),
146 IllegalGrpcClientStateSnafu {
147 err_msg: "Cluster client already started",
148 }
149 );
150 self.leader_provider = Some(leader_provider);
151 Ok(())
152 }
153
154 fn make_client(&self, addr: impl AsRef<str>) -> Result<ClusterClient<Channel>> {
155 let channel = self.channel_manager.get(addr).context(CreateChannelSnafu)?;
156
157 Ok(common_grpc::configure_tonic_client!(
158 ClusterClient::new(channel),
159 self.channel_manager,
160 ))
161 }
162
163 #[inline]
164 fn is_started(&self) -> bool {
165 self.leader_provider.is_some()
166 }
167
168 async fn with_retry<T, F, R, H>(&self, task: &str, body_fn: F, get_header: H) -> Result<T>
169 where
170 R: Future<Output = std::result::Result<T, Status>>,
171 F: Fn(ClusterClient<Channel>) -> R,
172 H: Fn(&T) -> &Option<ResponseHeader>,
173 {
174 let Some(leader_provider) = self.leader_provider.as_ref() else {
175 return IllegalGrpcClientStateSnafu {
176 err_msg: "not started",
177 }
178 .fail();
179 };
180
181 let mut times = 0;
182 let mut last_error = None;
183
184 while times < self.max_retry {
185 if let Some(leader) = &leader_provider.leader() {
186 let client = self.make_client(leader)?;
187 match body_fn(client).await {
188 Ok(res) => {
189 if util::is_not_leader(get_header(&res)) {
190 last_error = Some(format!("{leader} is not a leader"));
191 warn!("Failed to {task} to {leader}, not a leader");
192 let leader = leader_provider.ask_leader().await?;
193 info!("Cluster client updated to new leader addr: {leader}");
194 times += 1;
195 continue;
196 }
197 return Ok(res);
198 }
199 Err(status) => {
200 if util::is_unreachable(&status) {
202 last_error = Some(status.to_string());
203 warn!("Failed to {task} to {leader}, source: {status}");
204 let leader = leader_provider.ask_leader().await?;
205 info!("Cluster client updated to new leader addr: {leader}");
206 times += 1;
207 continue;
208 } else {
209 error!("An error occurred in gRPC, status: {status}");
210 return Err(Error::from(status));
211 }
212 }
213 }
214 } else {
215 leader_provider.ask_leader().await?;
216 }
217 }
218
219 RetryTimesExceededSnafu {
220 msg: format!("Failed to {task}, last error: {:?}", last_error),
221 times: self.max_retry,
222 }
223 .fail()
224 }
225
226 async fn range(&self, request: RangeRequest) -> Result<RangeResponse> {
227 self.with_retry(
228 "range",
229 move |mut client| {
230 let inner_req = tonic::Request::new(v1::meta::RangeRequest::from(request.clone()));
231
232 async move { client.range(inner_req).await.map(|res| res.into_inner()) }
233 },
234 |res| &res.header,
235 )
236 .await?
237 .try_into()
238 .context(ConvertMetaResponseSnafu)
239 }
240
241 async fn batch_get(&self, request: BatchGetRequest) -> Result<BatchGetResponse> {
242 self.with_retry(
243 "batch_get",
244 move |mut client| {
245 let inner_req =
246 tonic::Request::new(v1::meta::BatchGetRequest::from(request.clone()));
247
248 async move {
249 client
250 .batch_get(inner_req)
251 .await
252 .map(|res| res.into_inner())
253 }
254 },
255 |res| &res.header,
256 )
257 .await?
258 .try_into()
259 .context(ConvertMetaResponseSnafu)
260 }
261
262 async fn get_metasrv_peers(&self) -> Result<(Option<MetasrvNodeInfo>, Vec<MetasrvNodeInfo>)> {
263 self.with_retry(
264 "get_metasrv_peers",
265 move |mut client| {
266 let inner_req = tonic::Request::new(MetasrvPeersRequest::default());
267
268 async move {
269 client
270 .metasrv_peers(inner_req)
271 .await
272 .map(|res| res.into_inner())
273 }
274 },
275 |res| &res.header,
276 )
277 .await
278 .map(|res| (res.leader, res.followers))
279 }
280}
281
282#[derive(Clone, Debug)]
285pub struct ClusterKvBackend {
286 inner: Arc<Client>,
287}
288
289impl ClusterKvBackend {
290 pub fn new(client: Arc<Client>) -> Self {
291 Self { inner: client }
292 }
293
294 fn unimpl(&self) -> common_meta::error::Error {
295 let ret: common_meta::error::Result<()> = ReadOnlyKvBackendSnafu {
296 name: self.name().to_string(),
297 }
298 .fail()
299 .map_err(BoxedError::new)
300 .context(common_meta::error::ExternalSnafu);
301 ret.unwrap_err()
302 }
303}
304
305impl TxnService for ClusterKvBackend {
306 type Error = common_meta::error::Error;
307}
308
309#[async_trait::async_trait]
310impl KvBackend for ClusterKvBackend {
311 fn name(&self) -> &str {
312 "ClusterKvBackend"
313 }
314
315 fn as_any(&self) -> &dyn Any {
316 self
317 }
318
319 async fn range(&self, req: RangeRequest) -> common_meta::error::Result<RangeResponse> {
320 self.inner
321 .range(req)
322 .await
323 .map_err(BoxedError::new)
324 .context(common_meta::error::ExternalSnafu)
325 }
326
327 async fn batch_get(&self, _: BatchGetRequest) -> common_meta::error::Result<BatchGetResponse> {
328 Err(self.unimpl())
329 }
330
331 async fn put(&self, _: PutRequest) -> common_meta::error::Result<PutResponse> {
332 Err(self.unimpl())
333 }
334
335 async fn batch_put(&self, _: BatchPutRequest) -> common_meta::error::Result<BatchPutResponse> {
336 Err(self.unimpl())
337 }
338
339 async fn delete_range(
340 &self,
341 _: DeleteRangeRequest,
342 ) -> common_meta::error::Result<DeleteRangeResponse> {
343 Err(self.unimpl())
344 }
345
346 async fn batch_delete(
347 &self,
348 _: BatchDeleteRequest,
349 ) -> common_meta::error::Result<BatchDeleteResponse> {
350 Err(self.unimpl())
351 }
352}