meta_srv/utils/
database.rs1use std::sync::Arc;
16
17use client::error::{ExternalSnafu, Result as ClientResult};
18use client::{Client, Database, Output};
19use common_error::ext::BoxedError;
20use common_grpc::channel_manager::ChannelManager;
21use common_meta::peer::PeerDiscoveryRef;
22use common_telemetry::{debug, warn};
23use snafu::{ResultExt, ensure};
24use tokio::sync::RwLock;
25
26use crate::error::{ListActiveFrontendsSnafu, NoAvailableFrontendSnafu, Result};
27
28pub type DatabaseOperatorRef = Arc<DatabaseOperator>;
29
30#[derive(Debug, Clone, Copy)]
31pub struct DatabaseContext<'a> {
33 pub catalog: &'a str,
35 pub schema: &'a str,
37}
38
39impl<'a> DatabaseContext<'a> {
40 pub fn new(catalog: &'a str, schema: &'a str) -> Self {
42 Self { catalog, schema }
43 }
44}
45
46pub struct DatabaseOperator {
51 peer_discovery: PeerDiscoveryRef,
52 client: RwLock<Option<Client>>,
53}
54
55impl DatabaseOperator {
56 pub fn new(peer_discovery: PeerDiscoveryRef) -> Self {
58 Self {
59 peer_discovery,
60 client: RwLock::new(None),
61 }
62 }
63
64 pub async fn insert(
66 &self,
67 ctx: &DatabaseContext<'_>,
68 requests: api::v1::RowInsertRequests,
69 hints: &[(&str, &str)],
70 ) -> ClientResult<u32> {
71 let client = self.maybe_init_client().await?;
72 let database = Database::new(ctx.catalog, ctx.schema, client);
73
74 let result = database.row_inserts_with_hints(requests, hints).await;
75
76 if should_reset_client(&result) {
77 self.reset_client().await;
78 }
79
80 result
81 }
82
83 pub async fn logical_plan(
85 &self,
86 ctx: &DatabaseContext<'_>,
87 plan: Vec<u8>,
88 ) -> ClientResult<Output> {
89 let client = self.maybe_init_client().await?;
90 let database = Database::new(ctx.catalog, ctx.schema, client);
91
92 let result = database.logical_plan(plan).await;
93
94 if should_reset_client(&result) {
95 self.reset_client().await;
96 }
97
98 result
99 }
100
101 async fn build_client(&self) -> Result<Client> {
102 let frontends = self
103 .peer_discovery
104 .active_frontends()
105 .await
106 .context(ListActiveFrontendsSnafu)?;
107
108 ensure!(!frontends.is_empty(), NoAvailableFrontendSnafu);
109
110 let urls = frontends
111 .into_iter()
112 .map(|node| node.peer.addr)
113 .collect::<Vec<_>>();
114
115 debug!("Available frontend addresses: {:?}", urls);
116
117 Ok(Client::with_query_and_control_managers(
118 ChannelManager::new(),
119 ChannelManager::new(),
120 urls,
121 ))
122 }
123
124 async fn maybe_init_client(&self) -> ClientResult<Client> {
125 if let Some(client) = self.client.read().await.as_ref() {
126 return Ok(client.clone());
127 }
128
129 let client = self
130 .build_client()
131 .await
132 .map_err(BoxedError::new)
133 .context(ExternalSnafu)?;
134
135 let mut guard = self.client.write().await;
136 if let Some(client) = guard.as_ref() {
137 return Ok(client.clone());
138 }
139
140 *guard = Some(client.clone());
141 Ok(client)
142 }
143
144 async fn reset_client(&self) {
145 warn!("Resetting the client");
146 let mut guard = self.client.write().await;
147 guard.take();
148 }
149}
150
151fn should_reset_client<T>(result: &client::error::Result<T>) -> bool {
152 result
153 .as_ref()
154 .err()
155 .map(|err| err.is_connection_error())
156 .unwrap_or(false)
157}
158
159#[cfg(test)]
160mod tests {
161 use api::v1::meta::heartbeat_request::NodeWorkloads;
162 use common_meta::cluster::{FrontendStatus, NodeInfo, NodeStatus};
163 use common_meta::peer::{Peer, PeerDiscovery};
164
165 use super::*;
166
167 struct TestPeerDiscovery;
168
169 #[async_trait::async_trait]
170 impl PeerDiscovery for TestPeerDiscovery {
171 async fn active_frontends(&self) -> common_meta::error::Result<Vec<NodeInfo>> {
172 Ok(vec![NodeInfo {
173 peer: Peer::new(1, "127.0.0.1:3001".to_string()),
174 last_activity_ts: 0,
175 status: NodeStatus::Frontend(FrontendStatus::default()),
176 version: String::new(),
177 git_commit: String::new(),
178 start_time_ms: 0,
179 total_cpu_millicores: 0,
180 total_memory_bytes: 0,
181 cpu_usage_millicores: 0,
182 memory_usage_bytes: 0,
183 hostname: String::new(),
184 env_vars: Default::default(),
185 }])
186 }
187
188 async fn active_datanodes(
189 &self,
190 _filter: Option<for<'a> fn(&'a NodeWorkloads) -> bool>,
191 ) -> common_meta::error::Result<Vec<NodeInfo>> {
192 unreachable!()
193 }
194
195 async fn active_flownodes(
196 &self,
197 _filter: Option<for<'a> fn(&'a NodeWorkloads) -> bool>,
198 ) -> common_meta::error::Result<Vec<NodeInfo>> {
199 unreachable!()
200 }
201 }
202
203 #[tokio::test]
204 async fn test_build_client_uses_isolated_reused_channel_pools() {
205 let operator = DatabaseOperator::new(Arc::new(TestPeerDiscovery));
206 let client = operator.build_client().await.unwrap();
207
208 client.make_flight_client(false, false).unwrap();
209 client.make_flight_client(false, false).unwrap();
210 assert_eq!((1, 0), client.channel_pool_sizes());
211
212 client.find_channel().unwrap();
213 client.find_channel().unwrap();
214 assert_eq!((1, 1), client.channel_pool_sizes());
215 }
216}