frontend/instance/
region_query.rs1use 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}