1mod 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 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 pub fn frontend_default_options() -> Self {
120 Self::new(0, Role::Frontend)
122 .enable_store()
123 .enable_heartbeat()
124 .enable_procedure()
125 .enable_access_cluster_info()
126 }
127
128 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 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 pub fn enable_store(self) -> Self {
156 Self {
157 enable_store: true,
158 ..self
159 }
160 }
161
162 #[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#[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#[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; 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 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 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#[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 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 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 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 pub async fn heartbeat(&self) -> Result<(HeartbeatSender, HeartbeatStream, HeartbeatConfig)> {
735 self.heartbeat_client()?.heartbeat().await
736 }
737
738 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 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 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 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 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 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 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 pub async fn query_procedure_state(&self, pid: &str) -> Result<ProcedureStateResponse> {
807 self.procedure_client()?.query_procedure_state(pid).await
808 }
809
810 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 pub async fn reconcile(&self, request: ReconcileRequest) -> Result<ReconcileResponse> {
829 self.procedure_client()?.reconcile(request).await
830 }
831
832 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 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 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 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 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 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 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 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 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 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}