Skip to main content

catalog/
information_extension.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::meta::ProcedureStatus;
18use common_error::ext::BoxedError;
19use common_meta::cluster::{ClusterInfo, NodeInfo, Role};
20use common_meta::datanode::RegionStat;
21use common_meta::key::flow::flow_state::FlowStat;
22use common_meta::node_manager::DatanodeManagerRef;
23use common_meta::procedure_executor::{ExecutorContext, ProcedureExecutor};
24use common_meta::rpc::procedure;
25use common_procedure::{ProcedureInfo, ProcedureState};
26use common_query::request::QueryRequest;
27use common_recordbatch::adapter::{AsyncRecordBatchStreamAdapter, DfRecordBatchStreamAdapter};
28use common_recordbatch::util::{ChainedRecordBatchStream, LimitedRecordBatchStream};
29use common_recordbatch::{DfSendableRecordBatchStream, SendableRecordBatchStream};
30use datafusion::common::Result;
31use datafusion::common::tree_node::TreeNodeRecursion;
32use datafusion::execution::TaskContext;
33use datafusion::physical_expr::{EquivalenceProperties, Partitioning, PhysicalExpr};
34use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
35use datafusion::physical_plan::{DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties};
36use datatypes::arrow::datatypes::SchemaRef as ArrowSchemaRef;
37use datatypes::schema::SchemaRef;
38use futures_util::future::try_join_all;
39use meta_client::MetaClientRef;
40use snafu::ResultExt;
41use store_api::storage::RegionId;
42
43use crate::error;
44use crate::information_schema::{DatanodeInspectRequest, InformationExtension};
45
46pub struct DistributedInformationExtension {
47    meta_client: MetaClientRef,
48    datanode_manager: DatanodeManagerRef,
49}
50
51impl DistributedInformationExtension {
52    pub fn new(meta_client: MetaClientRef, datanode_manager: DatanodeManagerRef) -> Self {
53        Self {
54            meta_client,
55            datanode_manager,
56        }
57    }
58
59    async fn inspect_datanode_stream(
60        meta_client: MetaClientRef,
61        datanode_manager: DatanodeManagerRef,
62        request: DatanodeInspectRequest,
63    ) -> std::result::Result<SendableRecordBatchStream, crate::error::Error> {
64        let nodes = meta_client
65            .list_nodes(Some(Role::Datanode))
66            .await
67            .map_err(BoxedError::new)
68            .context(crate::error::ListNodesSnafu)?;
69
70        let limit = request.scan.limit;
71        let plan = request
72            .build_plan()
73            .context(crate::error::DatafusionSnafu)?;
74
75        let streams = try_join_all(nodes.into_iter().map(|node| {
76            let datanode_manager = datanode_manager.clone();
77            let plan = plan.clone();
78            async move {
79                let client = datanode_manager.datanode(&node.peer).await;
80                client
81                    .handle_query(QueryRequest {
82                        plan,
83                        region_id: RegionId::default(),
84                        header: None,
85                    })
86                    .await
87                    .context(crate::error::HandleQuerySnafu)
88            }
89        }))
90        .await?;
91
92        let chained =
93            ChainedRecordBatchStream::new(streams).context(crate::error::CreateRecordBatchSnafu)?;
94        match limit {
95            Some(limit) => Ok(Box::pin(LimitedRecordBatchStream::new(
96                Box::pin(chained),
97                limit,
98            ))),
99            None => Ok(Box::pin(chained)),
100        }
101    }
102}
103
104struct DistributedInspectExec {
105    request: DatanodeInspectRequest,
106    meta_client: MetaClientRef,
107    datanode_manager: DatanodeManagerRef,
108    schema: SchemaRef,
109    arrow_schema: ArrowSchemaRef,
110    properties: Arc<PlanProperties>,
111}
112
113impl std::fmt::Debug for DistributedInspectExec {
114    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
115        f.debug_struct("DistributedInspectExec")
116            .field("request", &self.request)
117            .field("schema", &self.schema)
118            .finish_non_exhaustive()
119    }
120}
121
122impl DistributedInspectExec {
123    fn new(
124        request: DatanodeInspectRequest,
125        schema: SchemaRef,
126        meta_client: MetaClientRef,
127        datanode_manager: DatanodeManagerRef,
128    ) -> Self {
129        let arrow_schema = schema.arrow_schema().clone();
130        let properties = Arc::new(PlanProperties::new(
131            EquivalenceProperties::new(arrow_schema.clone()),
132            Partitioning::UnknownPartitioning(1),
133            EmissionType::Incremental,
134            Boundedness::Bounded,
135        ));
136        Self {
137            request,
138            meta_client,
139            datanode_manager,
140            schema,
141            arrow_schema,
142            properties,
143        }
144    }
145}
146
147impl DisplayAs for DistributedInspectExec {
148    fn fmt_as(&self, _t: DisplayFormatType, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
149        write!(
150            f,
151            "DistributedInspectExec: kind={:?}, scan={}, schema={:?}",
152            self.request.kind, self.request.scan, self.arrow_schema
153        )
154    }
155}
156
157impl ExecutionPlan for DistributedInspectExec {
158    fn name(&self) -> &str {
159        "DistributedInspectExec"
160    }
161
162    fn schema(&self) -> ArrowSchemaRef {
163        self.arrow_schema.clone()
164    }
165
166    fn properties(&self) -> &Arc<PlanProperties> {
167        &self.properties
168    }
169
170    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
171        vec![]
172    }
173
174    fn apply_expressions(
175        &self,
176        _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
177    ) -> Result<TreeNodeRecursion> {
178        Ok(TreeNodeRecursion::Continue)
179    }
180
181    fn with_new_children(
182        self: Arc<Self>,
183        children: Vec<Arc<dyn ExecutionPlan>>,
184    ) -> datafusion::error::Result<Arc<dyn ExecutionPlan>> {
185        if children.is_empty() {
186            Ok(self)
187        } else {
188            Err(datafusion::error::DataFusionError::Internal(
189                "DistributedInspectExec is a leaf execution plan and cannot accept children"
190                    .to_string(),
191            ))
192        }
193    }
194
195    fn execute(
196        &self,
197        partition: usize,
198        _context: Arc<TaskContext>,
199    ) -> datafusion::error::Result<DfSendableRecordBatchStream> {
200        if partition != 0 {
201            return Err(datafusion::error::DataFusionError::Execution(format!(
202                "DistributedInspectExec invalid partition. Expected 0, got {partition}"
203            )));
204        }
205
206        let meta_client = self.meta_client.clone();
207        let datanode_manager = self.datanode_manager.clone();
208        let request = self.request.clone();
209        let schema = self.schema.clone();
210        let future = async move {
211            DistributedInformationExtension::inspect_datanode_stream(
212                meta_client,
213                datanode_manager,
214                request,
215            )
216            .await
217            .map_err(BoxedError::new)
218            .context(common_recordbatch::error::ExternalSnafu)
219        };
220        let stream = AsyncRecordBatchStreamAdapter::new(schema, Box::pin(future));
221        Ok(Box::pin(DfRecordBatchStreamAdapter::new(Box::pin(stream))))
222    }
223}
224
225#[async_trait::async_trait]
226impl InformationExtension for DistributedInformationExtension {
227    type Error = crate::error::Error;
228
229    async fn nodes(&self) -> std::result::Result<Vec<NodeInfo>, Self::Error> {
230        self.meta_client
231            .list_nodes(None)
232            .await
233            .map_err(BoxedError::new)
234            .context(error::ListNodesSnafu)
235    }
236
237    async fn procedures(&self) -> std::result::Result<Vec<(String, ProcedureInfo)>, Self::Error> {
238        let procedures = self
239            .meta_client
240            .list_procedures(&ExecutorContext::default())
241            .await
242            .map_err(BoxedError::new)
243            .context(error::ListProceduresSnafu)?
244            .procedures;
245        let mut result = Vec::with_capacity(procedures.len());
246        for procedure in procedures {
247            let pid = match procedure.id {
248                Some(pid) => pid,
249                None => return error::ProcedureIdNotFoundSnafu {}.fail(),
250            };
251            let pid = procedure::pb_pid_to_pid(&pid)
252                .map_err(BoxedError::new)
253                .context(error::ConvertProtoDataSnafu)?;
254            let status = ProcedureStatus::try_from(procedure.status)
255                .map(|v| v.as_str_name())
256                .unwrap_or("Unknown")
257                .to_string();
258            let procedure_info = ProcedureInfo {
259                id: pid,
260                type_name: procedure.type_name,
261                start_time_ms: procedure.start_time_ms,
262                end_time_ms: procedure.end_time_ms,
263                state: ProcedureState::Running,
264                lock_keys: procedure.lock_keys,
265            };
266            result.push((status, procedure_info));
267        }
268
269        Ok(result)
270    }
271
272    async fn region_stats(&self) -> std::result::Result<Vec<RegionStat>, Self::Error> {
273        self.meta_client
274            .list_region_stats()
275            .await
276            .map_err(BoxedError::new)
277            .context(error::ListRegionStatsSnafu)
278    }
279
280    async fn flow_stats(&self) -> std::result::Result<Option<FlowStat>, Self::Error> {
281        self.meta_client
282            .list_flow_stats()
283            .await
284            .map_err(BoxedError::new)
285            .context(crate::error::ListFlowStatsSnafu)
286    }
287
288    async fn inspect_datanode(
289        &self,
290        request: DatanodeInspectRequest,
291    ) -> std::result::Result<SendableRecordBatchStream, Self::Error> {
292        Self::inspect_datanode_stream(
293            self.meta_client.clone(),
294            self.datanode_manager.clone(),
295            request,
296        )
297        .await
298    }
299
300    fn inspect_datanode_plan(
301        &self,
302        request: DatanodeInspectRequest,
303        schema: SchemaRef,
304    ) -> std::result::Result<Option<Arc<dyn ExecutionPlan>>, Self::Error> {
305        Ok(Some(Arc::new(DistributedInspectExec::new(
306            request,
307            schema,
308            self.meta_client.clone(),
309            self.datanode_manager.clone(),
310        ))))
311    }
312}