Skip to main content

meta_srv/
mocks.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::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    // Keep the mock client's codec limits aligned with the server by default.
189    let config = client_channel_config.unwrap_or_else(|| {
190        let config = ChannelConfig::new()
191            // Use an long timeout to prevent test failures due to slow operations (e.g., when testing with S3).
192            .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    // Move client to an option so we can _move_ the inner value
204    // on the first attempt to connect. All other attempts will fail.
205    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}