1use 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}