Skip to main content

meta_client/
client.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
15mod ask_leader;
16mod config;
17pub mod heartbeat;
18mod load_balance;
19mod procedure;
20
21mod cluster;
22mod store;
23mod util;
24
25use std::fmt::Debug;
26use std::sync::Arc;
27use std::time::Duration;
28
29use api::v1::meta::heartbeat_request::NodeWorkloads;
30use api::v1::meta::{
31    MetasrvNodeInfo, ProcedureDetailResponse, ReconcileRequest, ReconcileResponse, Role,
32};
33pub use ask_leader::{AskLeader, LeaderProvider, LeaderProviderRef};
34use cluster::Client as ClusterClient;
35pub use cluster::ClusterKvBackend;
36use common_error::ext::BoxedError;
37use common_grpc::channel_manager::{ChannelConfig, ChannelManager};
38use common_meta::cluster::{
39    ClusterInfo, MetasrvStatus, NodeInfo, NodeInfoKey, NodeStatus, Role as ClusterRole,
40};
41use common_meta::datanode::{DatanodeStatKey, DatanodeStatValue, RegionStat};
42use common_meta::distributed_time_constants::default_distributed_time_constants;
43use common_meta::error::{
44    self as meta_error, ExternalSnafu, Result as MetaResult, UnsupportedSnafu,
45};
46use common_meta::key::flow::flow_state::{FlowStat, FlowStateManager};
47use common_meta::kv_backend::KvBackendRef;
48use common_meta::peer::PeerDiscovery;
49use common_meta::procedure_executor::{ExecutorContext, ProcedureExecutor};
50use common_meta::range_stream::PaginationStream;
51use common_meta::rpc::KeyValue;
52use common_meta::rpc::ddl::{
53    CREATE_DATABASE_CREATOR_EXTENSION_KEY, CreateDatabaseTask, DdlTask, SubmitDdlTaskRequest,
54    SubmitDdlTaskResponse,
55};
56use common_meta::rpc::procedure::{
57    AddRegionFollowerRequest, AddTableFollowerRequest, GcRegionsRequest, GcResponse,
58    GcTableRequest, ManageRegionFollowerRequest, MigrateRegionRequest, MigrateRegionResponse,
59    ProcedureStateResponse, RemoveRegionFollowerRequest, RemoveTableFollowerRequest,
60};
61use common_meta::rpc::store::{
62    BatchDeleteRequest, BatchDeleteResponse, BatchGetRequest, BatchGetResponse, BatchPutRequest,
63    BatchPutResponse, CompareAndPutRequest, CompareAndPutResponse, DeleteRangeRequest,
64    DeleteRangeResponse, PutRequest, PutResponse, RangeRequest, RangeResponse,
65};
66use common_options::plugin_options::PluginOptionsDeserializer;
67use common_telemetry::info;
68use common_time::util::DefaultSystemTimer;
69use config::Client as ConfigClient;
70use futures::TryStreamExt;
71use heartbeat::{Client as HeartbeatClient, HeartbeatConfig};
72use procedure::{Client as ProcedureClient, procedure_actor, procedure_event_context};
73use serde::de::DeserializeOwned;
74use snafu::{OptionExt, ResultExt};
75use store::Client as StoreClient;
76
77pub use self::heartbeat::{HeartbeatSender, HeartbeatStream};
78use crate::client::ask_leader::{LeaderProviderFactoryImpl, LeaderProviderFactoryRef};
79use crate::error::{
80    ConvertMetaConfigSnafu, ConvertMetaRequestSnafu, ConvertMetaResponseSnafu, Error,
81    GetFlowStatSnafu, MissingQueryContextSnafu, NotStartedSnafu, Result,
82};
83
84pub type Id = u64;
85
86const DEFAULT_ASK_LEADER_MAX_RETRY: usize = 3;
87const DEFAULT_SUBMIT_DDL_MAX_RETRY: usize = 3;
88const DEFAULT_CLUSTER_CLIENT_MAX_RETRY: usize = 3;
89const DEFAULT_DDL_TIMEOUT: Duration = Duration::from_secs(10);
90
91#[derive(Clone, Debug, Default)]
92pub struct MetaClientBuilder {
93    id: Id,
94    role: Role,
95    enable_heartbeat: bool,
96    enable_store: bool,
97    #[cfg(test)]
98    enable_direct_store_writes: bool,
99    enable_procedure: bool,
100    enable_access_cluster_info: bool,
101    region_follower: Option<RegionFollowerClientRef>,
102    channel_manager: Option<ChannelManager>,
103    ddl_channel_manager: Option<ChannelManager>,
104    /// The default ddl timeout for each request.
105    ddl_timeout: Option<Duration>,
106    heartbeat_channel_manager: Option<ChannelManager>,
107}
108
109impl MetaClientBuilder {
110    pub fn new(member_id: u64, role: Role) -> Self {
111        Self {
112            id: member_id,
113            role,
114            ..Default::default()
115        }
116    }
117
118    /// Returns the role of Frontend's default options.
119    pub fn frontend_default_options() -> Self {
120        // Frontend does not need a member id.
121        Self::new(0, Role::Frontend)
122            .enable_store()
123            .enable_heartbeat()
124            .enable_procedure()
125            .enable_access_cluster_info()
126    }
127
128    /// Returns the role of Datanode's default options.
129    pub fn datanode_default_options(member_id: u64) -> Self {
130        Self::new(member_id, Role::Datanode)
131            .enable_store()
132            .enable_heartbeat()
133    }
134
135    /// Returns the role of Flownode's default options.
136    pub fn flownode_default_options(member_id: u64) -> Self {
137        Self::new(member_id, Role::Flownode)
138            .enable_store()
139            .enable_heartbeat()
140            .enable_procedure()
141            .enable_access_cluster_info()
142    }
143
144    pub fn enable_heartbeat(self) -> Self {
145        Self {
146            enable_heartbeat: true,
147            ..self
148        }
149    }
150
151    /// Enables the Store client in read-only mode.
152    ///
153    /// Store write methods fail fast by default. Metadata writes from production
154    /// frontend/datanode/flownode clients should go through metasrv procedures.
155    pub fn enable_store(self) -> Self {
156        Self {
157            enable_store: true,
158            ..self
159        }
160    }
161
162    /// Enables direct Store write RPCs for tests.
163    ///
164    /// Production metadata writes should use metasrv-owned write paths instead.
165    #[cfg(test)]
166    pub(super) fn enable_direct_store_writes_for_test(self) -> Self {
167        Self {
168            enable_store: true,
169            enable_direct_store_writes: true,
170            ..self
171        }
172    }
173
174    pub fn enable_procedure(self) -> Self {
175        Self {
176            enable_procedure: true,
177            ..self
178        }
179    }
180
181    pub fn enable_access_cluster_info(self) -> Self {
182        Self {
183            enable_access_cluster_info: true,
184            ..self
185        }
186    }
187
188    pub fn channel_manager(self, channel_manager: ChannelManager) -> Self {
189        Self {
190            channel_manager: Some(channel_manager),
191            ..self
192        }
193    }
194
195    pub fn ddl_channel_manager(self, channel_manager: ChannelManager) -> Self {
196        Self {
197            ddl_channel_manager: Some(channel_manager),
198            ..self
199        }
200    }
201
202    pub fn ddl_timeout(self, timeout: Duration) -> Self {
203        Self {
204            ddl_timeout: Some(timeout),
205            ..self
206        }
207    }
208
209    pub fn heartbeat_channel_manager(self, channel_manager: ChannelManager) -> Self {
210        Self {
211            heartbeat_channel_manager: Some(channel_manager),
212            ..self
213        }
214    }
215
216    pub fn with_region_follower(self, region_follower: RegionFollowerClientRef) -> Self {
217        Self {
218            region_follower: Some(region_follower),
219            ..self
220        }
221    }
222
223    pub fn build(self) -> MetaClient {
224        let mgr = self.channel_manager.unwrap_or_default();
225        let heartbeat_channel_manager = self
226            .heartbeat_channel_manager
227            .clone()
228            .unwrap_or_else(|| mgr.clone());
229
230        let heartbeat = self.enable_heartbeat.then(|| {
231            if self.heartbeat_channel_manager.is_some() {
232                info!("Enable heartbeat channel using the heartbeat channel manager.");
233            }
234
235            HeartbeatClient::new(self.id, self.role, heartbeat_channel_manager.clone())
236        });
237        let config = self
238            .enable_heartbeat
239            .then(|| ConfigClient::new(self.id, self.role, mgr.clone()));
240        let store = self.enable_store.then(|| {
241            #[cfg(test)]
242            {
243                if self.enable_direct_store_writes {
244                    return StoreClient::new_writable(self.id, self.role, mgr.clone());
245                }
246            }
247
248            StoreClient::new(self.id, self.role, mgr.clone())
249        });
250        let procedure = self.enable_procedure.then(|| {
251            let mgr = self.ddl_channel_manager.unwrap_or(mgr.clone());
252            ProcedureClient::new(
253                self.id,
254                self.role,
255                mgr,
256                DEFAULT_SUBMIT_DDL_MAX_RETRY,
257                self.ddl_timeout.unwrap_or(DEFAULT_DDL_TIMEOUT),
258            )
259        });
260        let cluster = self
261            .enable_access_cluster_info
262            .then(|| ClusterClient::new(mgr.clone(), DEFAULT_CLUSTER_CLIENT_MAX_RETRY));
263        let region_follower = self.region_follower.clone();
264
265        MetaClient {
266            id: self.id,
267            channel_manager: mgr.clone(),
268            leader_provider_factory: Arc::new(LeaderProviderFactoryImpl::new(
269                self.id,
270                self.role,
271                DEFAULT_ASK_LEADER_MAX_RETRY,
272                heartbeat_channel_manager,
273            )),
274            heartbeat,
275            config,
276            store,
277            procedure,
278            cluster,
279            region_follower,
280        }
281    }
282}
283
284#[derive(Debug)]
285pub struct MetaClient {
286    id: Id,
287    channel_manager: ChannelManager,
288    leader_provider_factory: LeaderProviderFactoryRef,
289    heartbeat: Option<HeartbeatClient>,
290    config: Option<ConfigClient>,
291    store: Option<StoreClient>,
292    procedure: Option<ProcedureClient>,
293    cluster: Option<ClusterClient>,
294    region_follower: Option<RegionFollowerClientRef>,
295}
296
297impl MetaClient {
298    pub fn new(id: Id, role: Role) -> Self {
299        Self {
300            id,
301            channel_manager: ChannelManager::default(),
302            leader_provider_factory: Arc::new(LeaderProviderFactoryImpl::new(
303                id,
304                role,
305                DEFAULT_ASK_LEADER_MAX_RETRY,
306                ChannelManager::default(),
307            )),
308            heartbeat: None,
309            config: None,
310            store: None,
311            procedure: None,
312            cluster: None,
313            region_follower: None,
314        }
315    }
316}
317
318pub type RegionFollowerClientRef = Arc<dyn RegionFollowerClient>;
319
320/// A trait for clients that can manage region followers.
321#[async_trait::async_trait]
322pub trait RegionFollowerClient: Sync + Send + Debug {
323    async fn add_region_follower(&self, request: AddRegionFollowerRequest) -> Result<()>;
324
325    async fn remove_region_follower(&self, request: RemoveRegionFollowerRequest) -> Result<()>;
326
327    async fn add_table_follower(&self, request: AddTableFollowerRequest) -> Result<()>;
328
329    async fn remove_table_follower(&self, request: RemoveTableFollowerRequest) -> Result<()>;
330
331    async fn start(&self, urls: &[&str]) -> Result<()>;
332
333    async fn start_with(&self, leader_provider: LeaderProviderRef) -> Result<()>;
334}
335
336#[async_trait::async_trait]
337impl ProcedureExecutor for MetaClient {
338    async fn submit_ddl_task(
339        &self,
340        ctx: ExecutorContext,
341        request: SubmitDdlTaskRequest,
342    ) -> MetaResult<SubmitDdlTaskResponse> {
343        MetaClient::submit_ddl_task(self, ctx, request)
344            .await
345            .map_err(BoxedError::new)
346            .context(meta_error::ExternalSnafu)
347    }
348
349    async fn migrate_region(
350        &self,
351        ctx: &ExecutorContext,
352        request: MigrateRegionRequest,
353    ) -> MetaResult<MigrateRegionResponse> {
354        self.migrate_region(ctx, request)
355            .await
356            .map_err(BoxedError::new)
357            .context(meta_error::ExternalSnafu)
358    }
359
360    async fn reconcile(
361        &self,
362        _ctx: &ExecutorContext,
363        request: ReconcileRequest,
364    ) -> MetaResult<ReconcileResponse> {
365        self.reconcile(request)
366            .await
367            .map_err(BoxedError::new)
368            .context(meta_error::ExternalSnafu)
369    }
370
371    async fn manage_region_follower(
372        &self,
373        _ctx: &ExecutorContext,
374        request: ManageRegionFollowerRequest,
375    ) -> MetaResult<()> {
376        if let Some(region_follower) = &self.region_follower {
377            match request {
378                ManageRegionFollowerRequest::AddRegionFollower(add_region_follower_request) => {
379                    region_follower
380                        .add_region_follower(add_region_follower_request)
381                        .await
382                }
383                ManageRegionFollowerRequest::RemoveRegionFollower(
384                    remove_region_follower_request,
385                ) => {
386                    region_follower
387                        .remove_region_follower(remove_region_follower_request)
388                        .await
389                }
390                ManageRegionFollowerRequest::AddTableFollower(add_table_follower_request) => {
391                    region_follower
392                        .add_table_follower(add_table_follower_request)
393                        .await
394                }
395                ManageRegionFollowerRequest::RemoveTableFollower(remove_table_follower_request) => {
396                    region_follower
397                        .remove_table_follower(remove_table_follower_request)
398                        .await
399                }
400            }
401            .map_err(BoxedError::new)
402            .context(meta_error::ExternalSnafu)
403        } else {
404            UnsupportedSnafu {
405                operation: "manage_region_follower",
406            }
407            .fail()
408        }
409    }
410
411    async fn query_procedure_state(
412        &self,
413        _ctx: &ExecutorContext,
414        pid: &str,
415    ) -> MetaResult<ProcedureStateResponse> {
416        self.query_procedure_state(pid)
417            .await
418            .map_err(BoxedError::new)
419            .context(meta_error::ExternalSnafu)
420    }
421
422    async fn gc_regions(
423        &self,
424        ctx: &ExecutorContext,
425        request: GcRegionsRequest,
426    ) -> MetaResult<GcResponse> {
427        self.gc_regions(ctx, request)
428            .await
429            .map_err(BoxedError::new)
430            .context(meta_error::ExternalSnafu)
431    }
432
433    async fn gc_table(
434        &self,
435        ctx: &ExecutorContext,
436        request: GcTableRequest,
437    ) -> MetaResult<GcResponse> {
438        self.gc_table(ctx, request)
439            .await
440            .map_err(BoxedError::new)
441            .context(meta_error::ExternalSnafu)
442    }
443
444    async fn list_procedures(&self, _ctx: &ExecutorContext) -> MetaResult<ProcedureDetailResponse> {
445        self.procedure_client()
446            .map_err(BoxedError::new)
447            .context(meta_error::ExternalSnafu)?
448            .list_procedures()
449            .await
450            .map_err(BoxedError::new)
451            .context(meta_error::ExternalSnafu)
452    }
453}
454
455// TODO(zyy17): Allow deprecated fields for backward compatibility. Remove this when the deprecated fields are removed from the proto.
456#[allow(deprecated)]
457#[async_trait::async_trait]
458impl ClusterInfo for MetaClient {
459    type Error = Error;
460
461    async fn list_nodes(&self, role: Option<ClusterRole>) -> Result<Vec<NodeInfo>> {
462        let cluster_client = self.cluster_client()?;
463
464        let (get_metasrv_nodes, nodes_key_prefix) = match role {
465            None => (true, Some(NodeInfoKey::key_prefix())),
466            Some(ClusterRole::Metasrv) => (true, None),
467            Some(role) => (false, Some(NodeInfoKey::key_prefix_with_role(role))),
468        };
469
470        let mut nodes = if get_metasrv_nodes {
471            let last_activity_ts = -1; // Metasrv does not provide this information.
472
473            let (leader, followers): (Option<MetasrvNodeInfo>, Vec<MetasrvNodeInfo>) =
474                cluster_client.get_metasrv_peers().await?;
475            followers
476                .into_iter()
477                .map(|node| {
478                    if let Some(node_info) = node.info {
479                        NodeInfo {
480                            peer: node.peer.unwrap_or_default(),
481                            last_activity_ts,
482                            status: NodeStatus::Metasrv(MetasrvStatus { is_leader: false }),
483                            version: node_info.version,
484                            git_commit: node_info.git_commit,
485                            start_time_ms: node_info.start_time_ms,
486                            total_cpu_millicores: node_info.total_cpu_millicores,
487                            total_memory_bytes: node_info.total_memory_bytes,
488                            cpu_usage_millicores: node_info.cpu_usage_millicores,
489                            memory_usage_bytes: node_info.memory_usage_bytes,
490                            hostname: node_info.hostname,
491                            env_vars: Default::default(),
492                        }
493                    } else {
494                        // TODO(zyy17): It's for backward compatibility. Remove this when the deprecated fields are removed from the proto.
495                        NodeInfo {
496                            peer: node.peer.unwrap_or_default(),
497                            last_activity_ts,
498                            status: NodeStatus::Metasrv(MetasrvStatus { is_leader: false }),
499                            version: node.version,
500                            git_commit: node.git_commit,
501                            start_time_ms: node.start_time_ms,
502                            total_cpu_millicores: node.cpus as i64,
503                            total_memory_bytes: node.memory_bytes as i64,
504                            cpu_usage_millicores: 0,
505                            memory_usage_bytes: 0,
506                            hostname: "".to_string(),
507                            env_vars: Default::default(),
508                        }
509                    }
510                })
511                .chain(leader.into_iter().map(|node| {
512                    if let Some(node_info) = node.info {
513                        NodeInfo {
514                            peer: node.peer.unwrap_or_default(),
515                            last_activity_ts,
516                            status: NodeStatus::Metasrv(MetasrvStatus { is_leader: true }),
517                            version: node_info.version,
518                            git_commit: node_info.git_commit,
519                            start_time_ms: node_info.start_time_ms,
520                            total_cpu_millicores: node_info.total_cpu_millicores,
521                            total_memory_bytes: node_info.total_memory_bytes,
522                            cpu_usage_millicores: node_info.cpu_usage_millicores,
523                            memory_usage_bytes: node_info.memory_usage_bytes,
524                            hostname: node_info.hostname,
525                            env_vars: Default::default(),
526                        }
527                    } else {
528                        // TODO(zyy17): It's for backward compatibility. Remove this when the deprecated fields are removed from the proto.
529                        NodeInfo {
530                            peer: node.peer.unwrap_or_default(),
531                            last_activity_ts,
532                            status: NodeStatus::Metasrv(MetasrvStatus { is_leader: true }),
533                            version: node.version,
534                            git_commit: node.git_commit,
535                            start_time_ms: node.start_time_ms,
536                            total_cpu_millicores: node.cpus as i64,
537                            total_memory_bytes: node.memory_bytes as i64,
538                            cpu_usage_millicores: 0,
539                            memory_usage_bytes: 0,
540                            hostname: "".to_string(),
541                            env_vars: Default::default(),
542                        }
543                    }
544                }))
545                .collect::<Vec<_>>()
546        } else {
547            Vec::new()
548        };
549
550        if let Some(prefix) = nodes_key_prefix {
551            let req = RangeRequest::new().with_prefix(prefix);
552            let res = cluster_client.range(req).await?;
553            for kv in res.kvs {
554                nodes.push(NodeInfo::try_from(kv.value).context(ConvertMetaResponseSnafu)?);
555            }
556        }
557
558        Ok(nodes)
559    }
560
561    async fn list_region_stats(&self) -> Result<Vec<RegionStat>> {
562        let cluster_kv_backend = Arc::new(self.cluster_client()?);
563        let range_prefix = DatanodeStatKey::prefix_key();
564        let req = RangeRequest::new().with_prefix(range_prefix);
565        let stream =
566            PaginationStream::new(cluster_kv_backend, req, 256, decode_stats).into_stream();
567        let mut datanode_stats = stream
568            .try_collect::<Vec<_>>()
569            .await
570            .context(ConvertMetaResponseSnafu)?;
571        let region_stats = datanode_stats
572            .iter_mut()
573            .flat_map(|datanode_stat| {
574                let last = datanode_stat.stats.pop();
575                last.map(|stat| stat.region_stats).unwrap_or_default()
576            })
577            .collect::<Vec<_>>();
578
579        Ok(region_stats)
580    }
581
582    async fn list_flow_stats(&self) -> Result<Option<FlowStat>> {
583        let cluster_backend = ClusterKvBackend::new(Arc::new(self.cluster_client()?));
584        let cluster_backend = Arc::new(cluster_backend) as KvBackendRef;
585        let flow_state_manager = FlowStateManager::new(cluster_backend);
586        let res = flow_state_manager.get().await.context(GetFlowStatSnafu)?;
587
588        Ok(res.map(|r| r.into()))
589    }
590}
591
592// TODO(weny): the discovery using client side timestamp may be inaccurate,
593// maybe we need to use the timestamp from metasrv in the future.
594#[async_trait::async_trait]
595impl PeerDiscovery for MetaClient {
596    async fn active_frontends(&self) -> MetaResult<Vec<NodeInfo>> {
597        let nodes = self
598            .list_nodes(Some(ClusterRole::Frontend))
599            .await
600            .map_err(BoxedError::new)
601            .context(ExternalSnafu)?;
602        Ok(util::alive_frontends(
603            &DefaultSystemTimer,
604            nodes,
605            // TODO(weny): the heartbeat interval should be received from metasrv
606            // instead of using the default value.
607            default_distributed_time_constants().frontend_heartbeat_interval,
608        ))
609    }
610
611    async fn active_datanodes(
612        &self,
613        filter: Option<for<'a> fn(&'a NodeWorkloads) -> bool>,
614    ) -> MetaResult<Vec<NodeInfo>> {
615        let nodes = self
616            .list_nodes(Some(ClusterRole::Datanode))
617            .await
618            .map_err(BoxedError::new)
619            .context(ExternalSnafu)?;
620        Ok(util::alive_datanodes(
621            &DefaultSystemTimer,
622            nodes,
623            default_distributed_time_constants().datanode_lease,
624            filter,
625        ))
626    }
627
628    async fn active_flownodes(
629        &self,
630        filter: Option<for<'a> fn(&'a NodeWorkloads) -> bool>,
631    ) -> MetaResult<Vec<NodeInfo>> {
632        let nodes = self
633            .list_nodes(Some(ClusterRole::Flownode))
634            .await
635            .map_err(BoxedError::new)
636            .context(ExternalSnafu)?;
637        Ok(util::alive_flownodes(
638            &DefaultSystemTimer,
639            nodes,
640            default_distributed_time_constants().flownode_lease,
641            filter,
642        ))
643    }
644}
645
646fn decode_stats(kv: KeyValue) -> MetaResult<DatanodeStatValue> {
647    DatanodeStatValue::try_from(kv.value)
648        .map_err(BoxedError::new)
649        .context(ExternalSnafu)
650}
651
652impl MetaClient {
653    pub async fn start<U, A>(&mut self, urls: A) -> Result<()>
654    where
655        U: AsRef<str>,
656        A: AsRef<[U]> + Clone,
657    {
658        info!("MetaClient channel config: {:?}", self.channel_config());
659
660        let urls = urls.as_ref().iter().map(|u| u.as_ref()).collect::<Vec<_>>();
661        let leader_provider = self.leader_provider_factory.create(&urls);
662
663        self.start_with(leader_provider, urls).await
664    }
665
666    /// Start the client with a [LeaderProvider] and other Metasrv peers' addresses.
667    pub(crate) async fn start_with<U, A>(
668        &mut self,
669        leader_provider: LeaderProviderRef,
670        peers: A,
671    ) -> Result<()>
672    where
673        U: AsRef<str>,
674        A: AsRef<[U]> + Clone,
675    {
676        if let Some(client) = &self.region_follower {
677            info!("Starting region follower client ...");
678            client.start_with(leader_provider.clone()).await?;
679        }
680
681        if let Some(client) = &self.heartbeat {
682            info!("Starting heartbeat client ...");
683            client.start_with(leader_provider.clone()).await?;
684        }
685
686        if let Some(client) = &self.config {
687            info!("Starting config client ...");
688            client.start_with(leader_provider.clone()).await?;
689        }
690
691        if let Some(client) = &mut self.store {
692            info!("Starting store client ...");
693            client.start(peers.clone()).await?;
694        }
695
696        if let Some(client) = &self.procedure {
697            info!("Starting procedure client ...");
698            client.start_with(leader_provider.clone()).await?;
699        }
700
701        if let Some(client) = &mut self.cluster {
702            info!("Starting cluster client ...");
703            client.start_with(leader_provider).await?;
704        }
705        Ok(())
706    }
707
708    /// Ask the leader address of `metasrv`, and the heartbeat component
709    /// needs to create a bidirectional streaming to the leader.
710    pub async fn ask_leader(&self) -> Result<String> {
711        self.heartbeat_client()?.ask_leader().await
712    }
713
714    pub async fn pull_config<T, U>(&self, deserializer: T) -> Result<U>
715    where
716        T: PluginOptionsDeserializer<U>,
717        U: DeserializeOwned,
718    {
719        let res = self.config_client()?.pull_config().await?;
720        let v = deserializer
721            .deserialize(&res.payload)
722            .context(ConvertMetaConfigSnafu)?;
723        Ok(v)
724    }
725
726    /// Returns a heartbeat bidirectional streaming: (sender, receiver), the
727    /// other end is the leader of `metasrv`.
728    ///
729    /// The `datanode` needs to use the sender to continuously send heartbeat
730    /// packets (some self-state data), and the receiver can receive a response
731    /// from "metasrv" (which may contain some scheduling instructions).
732    ///
733    /// Returns the heartbeat sender, stream, and configuration received from Metasrv.
734    pub async fn heartbeat(&self) -> Result<(HeartbeatSender, HeartbeatStream, HeartbeatConfig)> {
735        self.heartbeat_client()?.heartbeat().await
736    }
737
738    /// Range gets the keys in the range from the key-value store.
739    pub async fn range(&self, req: RangeRequest) -> Result<RangeResponse> {
740        self.store_client()?
741            .range(req.into())
742            .await?
743            .try_into()
744            .context(ConvertMetaResponseSnafu)
745    }
746
747    /// Put puts the given key into the key-value store.
748    pub async fn put(&self, req: PutRequest) -> Result<PutResponse> {
749        self.store_client()?
750            .put(req.into())
751            .await?
752            .try_into()
753            .context(ConvertMetaResponseSnafu)
754    }
755
756    /// BatchGet atomically get values by the given keys from the key-value store.
757    pub async fn batch_get(&self, req: BatchGetRequest) -> Result<BatchGetResponse> {
758        self.store_client()?
759            .batch_get(req.into())
760            .await?
761            .try_into()
762            .context(ConvertMetaResponseSnafu)
763    }
764
765    /// BatchPut atomically puts the given keys into the key-value store.
766    pub async fn batch_put(&self, req: BatchPutRequest) -> Result<BatchPutResponse> {
767        self.store_client()?
768            .batch_put(req.into())
769            .await?
770            .try_into()
771            .context(ConvertMetaResponseSnafu)
772    }
773
774    /// BatchDelete atomically deletes the given keys from the key-value store.
775    pub async fn batch_delete(&self, req: BatchDeleteRequest) -> Result<BatchDeleteResponse> {
776        self.store_client()?
777            .batch_delete(req.into())
778            .await?
779            .try_into()
780            .context(ConvertMetaResponseSnafu)
781    }
782
783    /// CompareAndPut atomically puts the value to the given updated
784    /// value if the current value == the expected value.
785    pub async fn compare_and_put(
786        &self,
787        req: CompareAndPutRequest,
788    ) -> Result<CompareAndPutResponse> {
789        self.store_client()?
790            .compare_and_put(req.into())
791            .await?
792            .try_into()
793            .context(ConvertMetaResponseSnafu)
794    }
795
796    /// DeleteRange deletes the given range from the key-value store.
797    pub async fn delete_range(&self, req: DeleteRangeRequest) -> Result<DeleteRangeResponse> {
798        self.store_client()?
799            .delete_range(req.into())
800            .await?
801            .try_into()
802            .context(ConvertMetaResponseSnafu)
803    }
804
805    /// Query the procedure state by its id.
806    pub async fn query_procedure_state(&self, pid: &str) -> Result<ProcedureStateResponse> {
807        self.procedure_client()?.query_procedure_state(pid).await
808    }
809
810    /// Submit a region migration task.
811    pub async fn migrate_region(
812        &self,
813        context: &ExecutorContext,
814        request: MigrateRegionRequest,
815    ) -> Result<MigrateRegionResponse> {
816        self.procedure_client()?
817            .migrate_region(
818                context,
819                request.region_id,
820                request.from_peer,
821                request.to_peer,
822                request.timeout,
823            )
824            .await
825    }
826
827    /// Reconcile the procedure state.
828    pub async fn reconcile(&self, request: ReconcileRequest) -> Result<ReconcileResponse> {
829        self.procedure_client()?.reconcile(request).await
830    }
831
832    /// Manually trigger GC for specific regions.
833    pub async fn gc_regions(
834        &self,
835        context: &ExecutorContext,
836        request: GcRegionsRequest,
837    ) -> Result<GcResponse> {
838        self.procedure_client()?.gc_regions(context, request).await
839    }
840
841    /// Manually trigger GC for a table (all its regions).
842    pub async fn gc_table(
843        &self,
844        context: &ExecutorContext,
845        request: GcTableRequest,
846    ) -> Result<GcResponse> {
847        self.procedure_client()?.gc_table(context, request).await
848    }
849
850    /// Submit a DDL task.
851    pub async fn submit_ddl_task(
852        &self,
853        context: ExecutorContext,
854        request: SubmitDdlTaskRequest,
855    ) -> Result<SubmitDdlTaskResponse> {
856        let event_context = procedure_event_context(&context);
857        let actor = procedure_actor(&context);
858        let mut query_context = context.query_context.context(MissingQueryContextSnafu)?;
859        query_context
860            .extensions
861            .remove(CREATE_DATABASE_CREATOR_EXTENSION_KEY);
862        if let DdlTask::CreateDatabase(CreateDatabaseTask {
863            creator: Some(creator),
864            ..
865        }) = &request.task
866        {
867            query_context.extensions.insert(
868                CREATE_DATABASE_CREATOR_EXTENSION_KEY.to_string(),
869                serde_json::to_string(creator).context(ConvertMetaConfigSnafu)?,
870            );
871        }
872
873        let mut request: api::v1::meta::DdlTaskRequest =
874            request.try_into().context(ConvertMetaRequestSnafu)?;
875        request.query_context = Some(api::v1::QueryContext::from(query_context));
876        request.event_context = event_context;
877        request.actor = actor;
878
879        self.procedure_client()?
880            .submit_ddl_task(request)
881            .await?
882            .try_into()
883            .context(ConvertMetaResponseSnafu)
884    }
885
886    pub fn heartbeat_client(&self) -> Result<HeartbeatClient> {
887        self.heartbeat.clone().context(NotStartedSnafu {
888            name: "heartbeat_client",
889        })
890    }
891
892    pub fn config_client(&self) -> Result<ConfigClient> {
893        self.config.clone().context(NotStartedSnafu {
894            name: "config_client",
895        })
896    }
897
898    pub fn store_client(&self) -> Result<StoreClient> {
899        self.store.clone().context(NotStartedSnafu {
900            name: "store_client",
901        })
902    }
903
904    pub fn procedure_client(&self) -> Result<ProcedureClient> {
905        self.procedure.clone().context(NotStartedSnafu {
906            name: "procedure_client",
907        })
908    }
909
910    pub fn cluster_client(&self) -> Result<ClusterClient> {
911        self.cluster.clone().context(NotStartedSnafu {
912            name: "cluster_client",
913        })
914    }
915
916    pub fn channel_config(&self) -> &ChannelConfig {
917        self.channel_manager.config()
918    }
919
920    pub fn id(&self) -> Id {
921        self.id
922    }
923}
924
925#[cfg(test)]
926mod tests {
927    use std::sync::atomic::{AtomicUsize, Ordering};
928    use std::sync::{Arc, Mutex};
929
930    use api::v1::meta::{HeartbeatRequest, Peer};
931    use common_base::readable_size::ReadableSize;
932    use common_meta::kv_backend::{KvBackend, KvBackendRef, ResettableKvBackendRef, TxnService};
933    use rand::Rng;
934
935    use super::*;
936    use crate::error;
937    use crate::mocks::{self, MockMetaContext};
938
939    const TEST_KEY_PREFIX: &str = "__unit_test__meta__";
940
941    struct TestClient {
942        ns: String,
943        client: MetaClient,
944        meta_ctx: MockMetaContext,
945    }
946
947    impl TestClient {
948        async fn new(ns: impl Into<String>) -> Self {
949            // can also test with etcd: mocks::mock_client_with_etcdstore("127.0.0.1:2379").await;
950            let (client, meta_ctx) = mocks::mock_client_with_memstore().await;
951            Self {
952                ns: ns.into(),
953                client,
954                meta_ctx,
955            }
956        }
957
958        async fn new_with_grpc_message_sizes(
959            ns: impl Into<String>,
960            server_max_recv_message_size: ReadableSize,
961            server_max_send_message_size: ReadableSize,
962            client_max_recv_message_size: ReadableSize,
963            client_max_send_message_size: ReadableSize,
964        ) -> Self {
965            let (client, meta_ctx) = mocks::mock_client_with_memstore_and_grpc_message_sizes(
966                server_max_recv_message_size,
967                server_max_send_message_size,
968                client_max_recv_message_size,
969                client_max_send_message_size,
970            )
971            .await;
972            Self {
973                ns: ns.into(),
974                client,
975                meta_ctx,
976            }
977        }
978
979        fn key(&self, name: &str) -> Vec<u8> {
980            format!("{}-{}-{}", TEST_KEY_PREFIX, self.ns, name).into_bytes()
981        }
982
983        async fn gen_data(&self) {
984            for i in 0..10 {
985                let req = PutRequest::new()
986                    .with_key(self.key(&format!("key-{i}")))
987                    .with_value(format!("{}-{}", "value", i).into_bytes())
988                    .with_prev_kv();
989                let res = self.client.put(req).await;
990                let _ = res.unwrap();
991            }
992        }
993
994        async fn clear_data(&self) {
995            let req =
996                DeleteRangeRequest::new().with_prefix(format!("{}-{}", TEST_KEY_PREFIX, self.ns));
997            let res = self.client.delete_range(req).await;
998            let _ = res.unwrap();
999        }
1000
1001        #[allow(dead_code)]
1002        fn kv_backend(&self) -> KvBackendRef {
1003            self.meta_ctx.kv_backend.clone()
1004        }
1005
1006        fn in_memory(&self) -> Option<ResettableKvBackendRef> {
1007            self.meta_ctx.in_memory.clone()
1008        }
1009    }
1010
1011    async fn new_client(ns: impl Into<String>) -> TestClient {
1012        let client = TestClient::new(ns).await;
1013        client.clear_data().await;
1014        client
1015    }
1016
1017    #[tokio::test]
1018    async fn test_meta_client_builder() {
1019        let urls = &["127.0.0.1:3001", "127.0.0.1:3002"];
1020
1021        let mut meta_client = MetaClientBuilder::new(0, Role::Datanode)
1022            .enable_heartbeat()
1023            .build();
1024        let _ = meta_client.heartbeat_client().unwrap();
1025        assert!(meta_client.store_client().is_err());
1026        meta_client.start(urls).await.unwrap();
1027
1028        let mut meta_client = MetaClientBuilder::new(0, Role::Datanode).build();
1029        assert!(meta_client.heartbeat_client().is_err());
1030        assert!(meta_client.store_client().is_err());
1031        meta_client.start(urls).await.unwrap();
1032
1033        let mut meta_client = MetaClientBuilder::new(0, Role::Datanode)
1034            .enable_store()
1035            .build();
1036        assert!(meta_client.heartbeat_client().is_err());
1037        let _ = meta_client.store_client().unwrap();
1038        meta_client.start(urls).await.unwrap();
1039
1040        let mut meta_client = MetaClientBuilder::new(2, Role::Datanode)
1041            .enable_heartbeat()
1042            .enable_store()
1043            .build();
1044        assert_eq!(2, meta_client.id());
1045        assert_eq!(2, meta_client.id());
1046        let _ = meta_client.heartbeat_client().unwrap();
1047        let _ = meta_client.store_client().unwrap();
1048        meta_client.start(urls).await.unwrap();
1049    }
1050
1051    #[tokio::test]
1052    async fn test_not_start_heartbeat_client() {
1053        let urls = &["127.0.0.1:3001", "127.0.0.1:3002"];
1054        let mut meta_client = MetaClientBuilder::new(0, Role::Datanode)
1055            .enable_store()
1056            .build();
1057        meta_client.start(urls).await.unwrap();
1058        let res = meta_client.ask_leader().await;
1059        assert!(matches!(res.err(), Some(error::Error::NotStarted { .. })));
1060    }
1061
1062    #[tokio::test]
1063    async fn test_not_start_store_client() {
1064        let urls = &["127.0.0.1:3001", "127.0.0.1:3002"];
1065        let mut meta_client = MetaClientBuilder::new(0, Role::Datanode)
1066            .enable_heartbeat()
1067            .build();
1068
1069        meta_client.start(urls).await.unwrap();
1070        let res = meta_client.put(PutRequest::default()).await;
1071        assert!(matches!(res.err(), Some(error::Error::NotStarted { .. })));
1072    }
1073
1074    #[tokio::test]
1075    async fn test_store_writes_are_read_only_by_default() {
1076        let meta_client = MetaClientBuilder::new(0, Role::Datanode)
1077            .enable_store()
1078            .build();
1079
1080        let res = meta_client.put(PutRequest::default()).await;
1081        assert!(matches!(
1082            res.err(),
1083            Some(error::Error::ReadOnlyKvBackend { .. })
1084        ));
1085    }
1086
1087    #[tokio::test]
1088    async fn test_ask_leader() {
1089        let tc = new_client("test_ask_leader").await;
1090        tc.client.ask_leader().await.unwrap();
1091    }
1092
1093    #[tokio::test]
1094    async fn test_heartbeat() {
1095        let tc = new_client("test_heartbeat").await;
1096        let (sender, mut receiver, _config) = tc.client.heartbeat().await.unwrap();
1097        // send heartbeats
1098
1099        let request_sent = Arc::new(AtomicUsize::new(0));
1100        let request_sent_clone = request_sent.clone();
1101        let _handle = tokio::spawn(async move {
1102            for _ in 0..5 {
1103                let req = HeartbeatRequest {
1104                    peer: Some(Peer {
1105                        id: 1,
1106                        addr: "meta_client_peer".to_string(),
1107                    }),
1108                    ..Default::default()
1109                };
1110                sender.send(req).await.unwrap();
1111                request_sent_clone.fetch_add(1, Ordering::Relaxed);
1112            }
1113        });
1114
1115        let heartbeat_count = Arc::new(AtomicUsize::new(0));
1116        let heartbeat_count_clone = heartbeat_count.clone();
1117        let handle = tokio::spawn(async move {
1118            while let Some(_resp) = receiver.message().await.unwrap() {
1119                heartbeat_count_clone.fetch_add(1, Ordering::Relaxed);
1120            }
1121        });
1122
1123        handle.await.unwrap();
1124        //+1 for the initial response
1125        assert_eq!(
1126            request_sent.load(Ordering::Relaxed) + 1,
1127            heartbeat_count.load(Ordering::Relaxed)
1128        );
1129    }
1130
1131    #[tokio::test]
1132    async fn test_range_get() {
1133        let tc = new_client("test_range_get").await;
1134        tc.gen_data().await;
1135
1136        let key = tc.key("key-0");
1137        let req = RangeRequest::new().with_key(key.as_slice());
1138        let res = tc.client.range(req).await;
1139        let mut kvs = res.unwrap().take_kvs();
1140        assert_eq!(1, kvs.len());
1141        let mut kv = kvs.pop().unwrap();
1142        assert_eq!(key, kv.take_key());
1143        assert_eq!(b"value-0".to_vec(), kv.take_value());
1144    }
1145
1146    #[tokio::test]
1147    async fn test_range_get_prefix() {
1148        let tc = new_client("test_range_get_prefix").await;
1149        tc.gen_data().await;
1150
1151        let req = RangeRequest::new().with_prefix(tc.key("key-"));
1152        let res = tc.client.range(req).await;
1153        let kvs = res.unwrap().take_kvs();
1154        assert_eq!(10, kvs.len());
1155        for (i, mut kv) in kvs.into_iter().enumerate() {
1156            assert_eq!(tc.key(&format!("key-{i}")), kv.take_key());
1157            assert_eq!(format!("{}-{}", "value", i).into_bytes(), kv.take_value());
1158        }
1159    }
1160
1161    #[tokio::test]
1162    async fn test_range() {
1163        let tc = new_client("test_range").await;
1164        tc.gen_data().await;
1165
1166        let req = RangeRequest::new().with_range(tc.key("key-5"), tc.key("key-8"));
1167        let res = tc.client.range(req).await;
1168        let kvs = res.unwrap().take_kvs();
1169        assert_eq!(3, kvs.len());
1170        for (i, mut kv) in kvs.into_iter().enumerate() {
1171            assert_eq!(tc.key(&format!("key-{}", i + 5)), kv.take_key());
1172            assert_eq!(
1173                format!("{}-{}", "value", i + 5).into_bytes(),
1174                kv.take_value()
1175            );
1176        }
1177    }
1178
1179    #[tokio::test]
1180    async fn test_range_keys_only() {
1181        let tc = new_client("test_range_keys_only").await;
1182        tc.gen_data().await;
1183
1184        let req = RangeRequest::new()
1185            .with_range(tc.key("key-5"), tc.key("key-8"))
1186            .with_keys_only();
1187        let res = tc.client.range(req).await;
1188        let kvs = res.unwrap().take_kvs();
1189        assert_eq!(3, kvs.len());
1190        for (i, mut kv) in kvs.into_iter().enumerate() {
1191            assert_eq!(tc.key(&format!("key-{}", i + 5)), kv.take_key());
1192            assert!(kv.take_value().is_empty());
1193        }
1194    }
1195
1196    #[tokio::test]
1197    async fn test_put() {
1198        let tc = new_client("test_put").await;
1199
1200        let req = PutRequest::new()
1201            .with_key(tc.key("key"))
1202            .with_value(b"value".to_vec());
1203        let res = tc.client.put(req).await;
1204        assert!(res.unwrap().prev_kv.is_none());
1205    }
1206
1207    #[tokio::test]
1208    async fn test_put_with_prev_kv() {
1209        let tc = new_client("test_put_with_prev_kv").await;
1210
1211        let key = tc.key("key");
1212        let req = PutRequest::new()
1213            .with_key(key.as_slice())
1214            .with_value(b"value".to_vec())
1215            .with_prev_kv();
1216        let res = tc.client.put(req).await;
1217        assert!(res.unwrap().prev_kv.is_none());
1218
1219        let req = PutRequest::new()
1220            .with_key(key.as_slice())
1221            .with_value(b"value1".to_vec())
1222            .with_prev_kv();
1223        let res = tc.client.put(req).await;
1224        let mut kv = res.unwrap().prev_kv.unwrap();
1225        assert_eq!(key, kv.take_key());
1226        assert_eq!(b"value".to_vec(), kv.take_value());
1227    }
1228
1229    #[tokio::test]
1230    async fn test_batch_put() {
1231        let tc = new_client("test_batch_put").await;
1232
1233        let mut req = BatchPutRequest::new();
1234        for i in 0..275 {
1235            req = req.add_kv(
1236                tc.key(&format!("key-{}", i)),
1237                format!("value-{}", i).into_bytes(),
1238            );
1239        }
1240
1241        let res = tc.client.batch_put(req).await;
1242        assert_eq!(0, res.unwrap().take_prev_kvs().len());
1243
1244        let req = RangeRequest::new().with_prefix(tc.key("key-"));
1245        let res = tc.client.range(req).await;
1246        let kvs = res.unwrap().take_kvs();
1247        assert_eq!(275, kvs.len());
1248    }
1249
1250    #[tokio::test]
1251    async fn test_batch_get() {
1252        let tc = new_client("test_batch_get").await;
1253        tc.gen_data().await;
1254
1255        let mut req = BatchGetRequest::default();
1256        for i in 0..256 {
1257            req = req.add_key(tc.key(&format!("key-{}", i)));
1258        }
1259        let res = tc.client.batch_get(req).await.unwrap();
1260        assert_eq!(10, res.kvs.len());
1261
1262        let req = BatchGetRequest::default()
1263            .add_key(tc.key("key-1"))
1264            .add_key(tc.key("key-999"));
1265        let res = tc.client.batch_get(req).await.unwrap();
1266        assert_eq!(1, res.kvs.len());
1267    }
1268
1269    #[tokio::test]
1270    async fn test_batch_put_with_prev_kv() {
1271        let tc = new_client("test_batch_put_with_prev_kv").await;
1272
1273        let key = tc.key("key");
1274        let key2 = tc.key("key2");
1275        let req = BatchPutRequest::new().add_kv(key.as_slice(), b"value".to_vec());
1276        let res = tc.client.batch_put(req).await;
1277        assert_eq!(0, res.unwrap().take_prev_kvs().len());
1278
1279        let req = BatchPutRequest::new()
1280            .add_kv(key.as_slice(), b"value-".to_vec())
1281            .add_kv(key2.as_slice(), b"value2-".to_vec())
1282            .with_prev_kv();
1283        let res = tc.client.batch_put(req).await;
1284        let mut kvs = res.unwrap().take_prev_kvs();
1285        assert_eq!(1, kvs.len());
1286        let mut kv = kvs.pop().unwrap();
1287        assert_eq!(key, kv.take_key());
1288        assert_eq!(b"value".to_vec(), kv.take_value());
1289    }
1290
1291    #[tokio::test]
1292    async fn test_compare_and_put() {
1293        let tc = new_client("test_compare_and_put").await;
1294
1295        let key = tc.key("key");
1296        let req = CompareAndPutRequest::new()
1297            .with_key(key.as_slice())
1298            .with_expect(b"expect".to_vec())
1299            .with_value(b"value".to_vec());
1300        let res = tc.client.compare_and_put(req).await;
1301        assert!(!res.unwrap().is_success());
1302
1303        // create if absent
1304        let req = CompareAndPutRequest::new()
1305            .with_key(key.as_slice())
1306            .with_value(b"value".to_vec());
1307        let res = tc.client.compare_and_put(req).await;
1308        let mut res = res.unwrap();
1309        assert!(res.is_success());
1310        assert!(res.take_prev_kv().is_none());
1311
1312        // compare and put fail
1313        let req = CompareAndPutRequest::new()
1314            .with_key(key.as_slice())
1315            .with_expect(b"not_eq".to_vec())
1316            .with_value(b"value2".to_vec());
1317        let res = tc.client.compare_and_put(req).await;
1318        let mut res = res.unwrap();
1319        assert!(!res.is_success());
1320        assert_eq!(b"value".to_vec(), res.take_prev_kv().unwrap().take_value());
1321
1322        // compare and put success
1323        let req = CompareAndPutRequest::new()
1324            .with_key(key.as_slice())
1325            .with_expect(b"value".to_vec())
1326            .with_value(b"value2".to_vec());
1327        let res = tc.client.compare_and_put(req).await;
1328        let mut res = res.unwrap();
1329        assert!(res.is_success());
1330
1331        // If compare-and-put is success, previous value doesn't need to be returned.
1332        assert!(res.take_prev_kv().is_none());
1333    }
1334
1335    #[tokio::test]
1336    async fn test_delete_with_key() {
1337        let tc = new_client("test_delete_with_key").await;
1338        tc.gen_data().await;
1339
1340        let req = DeleteRangeRequest::new()
1341            .with_key(tc.key("key-0"))
1342            .with_prev_kv();
1343        let res = tc.client.delete_range(req).await;
1344        let mut res = res.unwrap();
1345        assert_eq!(1, res.deleted());
1346        let mut kvs = res.take_prev_kvs();
1347        assert_eq!(1, kvs.len());
1348        let mut kv = kvs.pop().unwrap();
1349        assert_eq!(b"value-0".to_vec(), kv.take_value());
1350    }
1351
1352    #[tokio::test]
1353    async fn test_delete_with_prefix() {
1354        let tc = new_client("test_delete_with_prefix").await;
1355        tc.gen_data().await;
1356
1357        let req = DeleteRangeRequest::new()
1358            .with_prefix(tc.key("key-"))
1359            .with_prev_kv();
1360        let res = tc.client.delete_range(req).await;
1361        let mut res = res.unwrap();
1362        assert_eq!(10, res.deleted());
1363        let kvs = res.take_prev_kvs();
1364        assert_eq!(10, kvs.len());
1365        for (i, mut kv) in kvs.into_iter().enumerate() {
1366            assert_eq!(format!("{}-{}", "value", i).into_bytes(), kv.take_value());
1367        }
1368    }
1369
1370    #[tokio::test]
1371    async fn test_delete_with_range() {
1372        let tc = new_client("test_delete_with_range").await;
1373        tc.gen_data().await;
1374
1375        let req = DeleteRangeRequest::new()
1376            .with_range(tc.key("key-2"), tc.key("key-7"))
1377            .with_prev_kv();
1378        let res = tc.client.delete_range(req).await;
1379        let mut res = res.unwrap();
1380        assert_eq!(5, res.deleted());
1381        let kvs = res.take_prev_kvs();
1382        assert_eq!(5, kvs.len());
1383        for (i, mut kv) in kvs.into_iter().enumerate() {
1384            assert_eq!(
1385                format!("{}-{}", "value", i + 2).into_bytes(),
1386                kv.take_value()
1387            );
1388        }
1389    }
1390
1391    fn mock_decoder(kv: KeyValue) -> MetaResult<Vec<u8>> {
1392        Ok(kv.value)
1393    }
1394
1395    struct RecordingKvBackend {
1396        inner: KvBackendRef,
1397        limits: Arc<Mutex<Vec<i64>>>,
1398    }
1399
1400    #[async_trait::async_trait]
1401    impl TxnService for RecordingKvBackend {
1402        type Error = meta_error::Error;
1403    }
1404
1405    #[async_trait::async_trait]
1406    impl KvBackend for RecordingKvBackend {
1407        fn name(&self) -> &str {
1408            "RecordingKvBackend"
1409        }
1410
1411        fn as_any(&self) -> &dyn std::any::Any {
1412            self
1413        }
1414
1415        async fn range(&self, req: RangeRequest) -> MetaResult<RangeResponse> {
1416            self.limits.lock().unwrap().push(req.limit);
1417            self.inner.range(req).await
1418        }
1419
1420        async fn put(&self, req: PutRequest) -> MetaResult<PutResponse> {
1421            self.inner.put(req).await
1422        }
1423
1424        async fn batch_put(&self, req: BatchPutRequest) -> MetaResult<BatchPutResponse> {
1425            self.inner.batch_put(req).await
1426        }
1427
1428        async fn batch_get(&self, req: BatchGetRequest) -> MetaResult<BatchGetResponse> {
1429            self.inner.batch_get(req).await
1430        }
1431
1432        async fn delete_range(&self, req: DeleteRangeRequest) -> MetaResult<DeleteRangeResponse> {
1433            self.inner.delete_range(req).await
1434        }
1435
1436        async fn batch_delete(&self, req: BatchDeleteRequest) -> MetaResult<BatchDeleteResponse> {
1437            self.inner.batch_delete(req).await
1438        }
1439    }
1440
1441    async fn adaptive_range_with_message_sizes(
1442        server_max_recv_message_size: ReadableSize,
1443        server_max_send_message_size: ReadableSize,
1444        client_max_recv_message_size: ReadableSize,
1445        client_max_send_message_size: ReadableSize,
1446    ) -> (Vec<Vec<u8>>, Vec<Vec<u8>>, Vec<i64>) {
1447        let tx = TestClient::new_with_grpc_message_sizes(
1448            "test_cluster_client",
1449            server_max_recv_message_size,
1450            server_max_send_message_size,
1451            client_max_recv_message_size,
1452            client_max_send_message_size,
1453        )
1454        .await;
1455        let in_memory = tx.in_memory().unwrap();
1456        let cluster_client = tx.client.cluster_client().unwrap();
1457        let mut rng = rand::rng();
1458
1459        let mut expected = Vec::new();
1460        for i in 0..4 {
1461            let data: Vec<u8> = (0..256 * 1024).map(|_| rng.random::<u8>()).collect();
1462            in_memory
1463                .put(
1464                    PutRequest::new()
1465                        .with_key(format!("__prefix/{i}").as_bytes())
1466                        .with_value(data.clone()),
1467                )
1468                .await
1469                .unwrap();
1470            expected.push(data);
1471        }
1472
1473        let req = RangeRequest::new().with_prefix(b"__prefix/");
1474        let limits = Arc::new(Mutex::new(Vec::new()));
1475        let recording_backend = RecordingKvBackend {
1476            inner: Arc::new(cluster_client),
1477            limits: limits.clone(),
1478        };
1479        let stream =
1480            PaginationStream::new(Arc::new(recording_backend), req, 4, mock_decoder).into_stream();
1481
1482        let res = stream.try_collect::<Vec<_>>().await.unwrap();
1483        let limits = limits.lock().unwrap().clone();
1484        (expected, res, limits)
1485    }
1486
1487    #[tokio::test]
1488    async fn test_cluster_client_adaptive_range() {
1489        let (expected, res, limits) = adaptive_range_with_message_sizes(
1490            ReadableSize::mb(2),
1491            ReadableSize::mb(16),
1492            ReadableSize::mb(1),
1493            ReadableSize::mb(2),
1494        )
1495        .await;
1496
1497        assert_eq!(expected, res);
1498        assert_eq!(vec![4, 2, 2], limits);
1499    }
1500
1501    #[tokio::test]
1502    async fn test_cluster_client_adaptive_range_server_limit() {
1503        let (expected, res, limits) = adaptive_range_with_message_sizes(
1504            ReadableSize::mb(2),
1505            ReadableSize::mb(1),
1506            ReadableSize::mb(16),
1507            ReadableSize::mb(2),
1508        )
1509        .await;
1510
1511        assert_eq!(expected, res);
1512        assert_eq!(vec![4, 2, 2], limits);
1513    }
1514}