1use std::sync::Arc;
16use std::time::Duration;
17
18use api::v1::meta::cluster_server::ClusterServer;
19use api::v1::meta::config_server::ConfigServer;
20use api::v1::meta::heartbeat_server::HeartbeatServer;
21use api::v1::meta::procedure_service_server::ProcedureServiceServer;
22use api::v1::meta::store_server::StoreServer;
23use client::client_manager::NodeClients;
24use common_grpc::channel_manager::{ChannelConfig, ChannelManager};
25use common_meta::key::TableMetadataManager;
26use common_meta::kv_backend::etcd::EtcdStore;
27use common_meta::kv_backend::memory::MemoryKvBackend;
28use common_meta::kv_backend::{KvBackendRef, ResettableKvBackendRef};
29use hyper_util::rt::TokioIo;
30use servers::grpc::GrpcOptions;
31use tonic::codec::CompressionEncoding;
32use tower::service_fn;
33
34use crate::add_compressed_service;
35use crate::metasrv::builder::MetasrvBuilder;
36use crate::metasrv::{Metasrv, MetasrvOptions, SelectorRef};
37
38#[derive(Clone)]
39pub struct MockInfo {
40 pub server_addr: String,
41 pub channel_manager: ChannelManager,
42 pub metasrv: Arc<Metasrv>,
43 pub kv_backend: KvBackendRef,
44 pub in_memory: Option<ResettableKvBackendRef>,
45}
46
47pub async fn mock_with_memstore() -> MockInfo {
48 let kv_backend = Arc::new(MemoryKvBackend::new());
49 let in_memory = Arc::new(MemoryKvBackend::new());
50 mock(
51 MetasrvOptions {
52 grpc: GrpcOptions {
53 server_addr: "127.0.0.1:3002".to_string(),
54 ..Default::default()
55 },
56 ..Default::default()
57 },
58 kv_backend,
59 None,
60 None,
61 Some(in_memory),
62 )
63 .await
64}
65
66pub async fn mock_with_etcdstore(addr: &str) -> MockInfo {
67 let kv_backend = EtcdStore::with_endpoints([addr], 128).await.unwrap();
68 mock(
69 MetasrvOptions {
70 grpc: GrpcOptions {
71 server_addr: "127.0.0.1:3002".to_string(),
72 ..Default::default()
73 },
74 ..Default::default()
75 },
76 kv_backend,
77 None,
78 None,
79 None,
80 )
81 .await
82}
83
84pub async fn mock(
85 opts: MetasrvOptions,
86 kv_backend: KvBackendRef,
87 selector: Option<SelectorRef>,
88 datanode_clients: Option<Arc<NodeClients>>,
89 in_memory: Option<ResettableKvBackendRef>,
90) -> MockInfo {
91 mock_inner(
92 opts,
93 kv_backend,
94 selector,
95 datanode_clients,
96 in_memory,
97 None,
98 )
99 .await
100}
101
102pub async fn mock_with_client_channel_config(
103 opts: MetasrvOptions,
104 kv_backend: KvBackendRef,
105 selector: Option<SelectorRef>,
106 datanode_clients: Option<Arc<NodeClients>>,
107 in_memory: Option<ResettableKvBackendRef>,
108 client_channel_config: ChannelConfig,
109) -> MockInfo {
110 mock_inner(
111 opts,
112 kv_backend,
113 selector,
114 datanode_clients,
115 in_memory,
116 Some(client_channel_config),
117 )
118 .await
119}
120
121async fn mock_inner(
122 opts: MetasrvOptions,
123 kv_backend: KvBackendRef,
124 selector: Option<SelectorRef>,
125 datanode_clients: Option<Arc<NodeClients>>,
126 in_memory: Option<ResettableKvBackendRef>,
127 client_channel_config: Option<ChannelConfig>,
128) -> MockInfo {
129 let server_addr = opts.grpc.server_addr.clone();
130 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
131
132 table_metadata_manager.init().await.unwrap();
133
134 let grpc_config = opts.grpc.as_config();
135 let grpc_options = opts.grpc.clone();
136 let builder = MetasrvBuilder::new()
137 .options(opts)
138 .kv_backend(kv_backend.clone());
139
140 let builder = match selector {
141 Some(s) => builder.selector(s),
142 None => builder,
143 };
144
145 let builder = match datanode_clients {
146 Some(clients) => builder.node_manager(clients),
147 None => builder,
148 };
149
150 let builder = match &in_memory {
151 Some(in_memory) => builder.in_memory(in_memory.clone()),
152 None => builder,
153 };
154
155 let metasrv = builder.build().await.unwrap();
156 metasrv.try_start().await.unwrap();
157
158 let (client, server) = tokio::io::duplex(1024);
159 let metasrv = Arc::new(metasrv);
160 let service = metasrv.clone();
161
162 let _handle = tokio::spawn(async move {
163 let mut router = tonic::transport::Server::builder();
164 let router = add_compressed_service!(
165 router,
166 HeartbeatServer::from_arc(service.clone()),
167 grpc_config
168 );
169 let router =
170 add_compressed_service!(router, StoreServer::from_arc(service.clone()), grpc_config);
171 let router = add_compressed_service!(
172 router,
173 ProcedureServiceServer::from_arc(service.clone()),
174 grpc_config
175 );
176 let router = add_compressed_service!(
177 router,
178 ClusterServer::from_arc(service.clone()),
179 grpc_config
180 );
181 let router =
182 add_compressed_service!(router, ConfigServer::from_arc(service.clone()), grpc_config);
183 router
184 .serve_with_incoming(futures::stream::iter(vec![Ok::<_, std::io::Error>(server)]))
185 .await
186 });
187
188 let config = client_channel_config.unwrap_or_else(|| {
190 let config = ChannelConfig::new()
191 .timeout(Some(Duration::from_secs(60)))
193 .connect_timeout(Duration::from_secs(10))
194 .tcp_nodelay(true);
195 ChannelConfig {
196 max_recv_message_size: grpc_options.max_recv_message_size,
197 max_send_message_size: grpc_options.max_send_message_size,
198 ..config
199 }
200 });
201 let channel_manager = ChannelManager::with_config(config, None);
202
203 let mut client = Some(client);
206 let res = channel_manager.reset_with_connector(
207 &server_addr,
208 service_fn(move |_| {
209 let client = client.take();
210
211 async move {
212 if let Some(client) = client {
213 Ok(TokioIo::new(client))
214 } else {
215 Err(std::io::Error::other("Client already taken"))
216 }
217 }
218 }),
219 );
220 let _ = res.unwrap();
221
222 MockInfo {
223 server_addr,
224 channel_manager,
225 metasrv,
226 kv_backend,
227 in_memory,
228 }
229}