Skip to main content

meta_client/client/
cluster.rs

1// Copyright 2023 Greptime Team
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 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    /// Start the client with a [LeaderProvider].
62    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                        // The leader may be unreachable.
201                        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/// A client for the cluster info. Read only and corresponding to
283/// `in_memory` kvbackend in the meta-srv.
284#[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}