Skip to main content

meta_srv/utils/
database.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;
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)]
31/// Database-level request context used by metasrv forwarding.
32pub struct DatabaseContext<'a> {
33    /// Catalog name carried in forwarded requests.
34    pub catalog: &'a str,
35    /// Schema name carried in forwarded requests.
36    pub schema: &'a str,
37}
38
39impl<'a> DatabaseContext<'a> {
40    /// Creates a new database context from catalog and schema.
41    pub fn new(catalog: &'a str, schema: &'a str) -> Self {
42        Self { catalog, schema }
43    }
44}
45
46/// A cached frontend database operator used by metasrv.
47///
48/// Its cached clients use independent query and control channel-manager pools;
49/// the legacy single-manager client constructors intentionally remain shared.
50pub struct DatabaseOperator {
51    peer_discovery: PeerDiscoveryRef,
52    client: RwLock<Option<Client>>,
53}
54
55impl DatabaseOperator {
56    /// Creates a database operator backed by discovered frontend peers.
57    pub fn new(peer_discovery: PeerDiscoveryRef) -> Self {
58        Self {
59            peer_discovery,
60            client: RwLock::new(None),
61        }
62    }
63
64    /// Forwards row inserts to an available frontend database client.
65    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    /// Executes a serialized logical plan on an available frontend.
84    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}