Skip to main content

frontend/instance/
region_query.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 api::v1::region::{RemoteDynFilterUnregister, RemoteDynFilterUpdate};
18use async_trait::async_trait;
19use client::region::{
20    build_remote_dyn_filter_unregister_request, build_remote_dyn_filter_update_request,
21};
22use common_error::ext::BoxedError;
23use common_meta::node_manager::NodeManagerRef;
24use common_query::request::QueryRequest;
25use partition::manager::PartitionRuleManagerRef;
26use query::error::{RegionQuerySnafu, Result as QueryResult};
27use query::region_query::{RegionQueryHandler, RegionQueryTarget};
28use session::ReadPreference;
29use snafu::ResultExt;
30
31use crate::error::{FindRegionPeerSnafu, RequestQuerySnafu, Result};
32
33pub(crate) struct FrontendRegionQueryHandler {
34    partition_manager: PartitionRuleManagerRef,
35    node_manager: NodeManagerRef,
36}
37
38impl FrontendRegionQueryHandler {
39    pub fn arc(
40        partition_manager: PartitionRuleManagerRef,
41        node_manager: NodeManagerRef,
42    ) -> Arc<Self> {
43        Arc::new(Self {
44            partition_manager,
45            node_manager,
46        })
47    }
48}
49
50#[async_trait]
51impl RegionQueryHandler for FrontendRegionQueryHandler {
52    async fn select_target(
53        &self,
54        read_preference: ReadPreference,
55        region_id: store_api::storage::RegionId,
56    ) -> QueryResult<RegionQueryTarget> {
57        self.select_target_inner(read_preference, region_id)
58            .await
59            .map_err(BoxedError::new)
60            .context(RegionQuerySnafu)
61    }
62
63    async fn do_get(
64        &self,
65        target: &RegionQueryTarget,
66        request: QueryRequest,
67    ) -> QueryResult<common_recordbatch::SendableRecordBatchStream> {
68        self.do_get_inner(target, request)
69            .await
70            .map_err(BoxedError::new)
71            .context(RegionQuerySnafu)
72    }
73
74    async fn handle_remote_dyn_filter_update(
75        &self,
76        target: &RegionQueryTarget,
77        query_id: String,
78        update: RemoteDynFilterUpdate,
79    ) -> QueryResult<()> {
80        self.handle_remote_dyn_filter_update_inner(target, query_id, update)
81            .await
82            .map_err(BoxedError::new)
83            .context(RegionQuerySnafu)
84    }
85
86    async fn handle_remote_dyn_filter_unregister(
87        &self,
88        target: &RegionQueryTarget,
89        query_id: String,
90        unregister: RemoteDynFilterUnregister,
91    ) -> QueryResult<()> {
92        self.handle_remote_dyn_filter_unregister_inner(target, query_id, unregister)
93            .await
94            .map_err(BoxedError::new)
95            .context(RegionQuerySnafu)
96    }
97}
98
99impl FrontendRegionQueryHandler {
100    async fn select_target_inner(
101        &self,
102        read_preference: ReadPreference,
103        region_id: store_api::storage::RegionId,
104    ) -> Result<RegionQueryTarget> {
105        let peer = self
106            .partition_manager
107            .find_region_leader(region_id)
108            .await
109            .context(FindRegionPeerSnafu {
110                region_id,
111                read_preference,
112            })?;
113
114        Ok(RegionQueryTarget::new(peer))
115    }
116
117    async fn do_get_inner(
118        &self,
119        target: &RegionQueryTarget,
120        request: QueryRequest,
121    ) -> Result<common_recordbatch::SendableRecordBatchStream> {
122        self.node_manager
123            .datanode(target.peer())
124            .await
125            .handle_query(request)
126            .await
127            .context(RequestQuerySnafu)
128    }
129
130    async fn handle_remote_dyn_filter_update_inner(
131        &self,
132        target: &RegionQueryTarget,
133        query_id: String,
134        update: RemoteDynFilterUpdate,
135    ) -> Result<()> {
136        let client = self.node_manager.datanode(target.peer()).await;
137        client
138            .handle(build_remote_dyn_filter_update_request(query_id, update))
139            .await
140            .context(RequestQuerySnafu)?;
141        Ok(())
142    }
143
144    async fn handle_remote_dyn_filter_unregister_inner(
145        &self,
146        target: &RegionQueryTarget,
147        query_id: String,
148        unregister: RemoteDynFilterUnregister,
149    ) -> Result<()> {
150        let client = self.node_manager.datanode(target.peer()).await;
151        client
152            .handle(build_remote_dyn_filter_unregister_request(
153                query_id, unregister,
154            ))
155            .await
156            .context(RequestQuerySnafu)?;
157        Ok(())
158    }
159}