1mod cluster_info;
16pub mod columns;
17pub mod flow_statistics;
18pub mod flows;
19mod information_memory_table;
20pub mod key_column_usage;
21mod partitions;
22mod procedure_info;
23pub mod process_list;
24#[cfg(feature = "enterprise")]
25mod recycle_bin;
26mod region_info;
27pub mod region_peers;
28mod region_statistics;
29pub mod schemata;
30mod ssts;
31pub mod statistics;
32mod table_constraints;
33mod table_names;
34mod table_semantics;
35pub mod tables;
36mod views;
37
38#[cfg(all(test, feature = "enterprise"))]
39mod recycle_bin_test;
40
41use std::collections::HashMap;
42use std::sync::{Arc, Weak};
43
44use common_catalog::consts::{self, DEFAULT_CATALOG_NAME, INFORMATION_SCHEMA_NAME};
45use common_error::ext::ErrorExt;
46use common_meta::cluster::NodeInfo;
47use common_meta::datanode::RegionStat;
48use common_meta::key::flow::FlowMetadataManager;
49use common_meta::key::flow::flow_state::FlowStat;
50use common_meta::kv_backend::KvBackendRef;
51use common_procedure::ProcedureInfo;
52use common_recordbatch::SendableRecordBatchStream;
53use datafusion::error::DataFusionError;
54use datafusion::logical_expr::LogicalPlan;
55use datafusion::physical_plan::ExecutionPlan;
56use datatypes::schema::SchemaRef;
57use lazy_static::lazy_static;
58use paste::paste;
59use process_list::InformationSchemaProcessList;
60use region_info::InformationSchemaRegionInfo;
61use store_api::metric_engine_consts::{
62 MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING, PRIMARY_KEY_ENCODING,
63};
64use store_api::region_info::RegionInfoEntry;
65use store_api::sst_entry::{ManifestSstEntry, PuffinIndexMetaEntry, StorageSstEntry};
66use store_api::storage::{ScanRequest, TableId};
67use table::TableRef;
68use table::metadata::TableType;
69pub use table_names::*;
70use views::InformationSchemaViews;
71
72use self::columns::InformationSchemaColumns;
73use crate::CatalogManager;
74use crate::error::{Error, Result};
75use crate::process_manager::ProcessManagerRef;
76use crate::system_schema::information_schema::cluster_info::InformationSchemaClusterInfo;
77use crate::system_schema::information_schema::flow_statistics::InformationSchemaFlowStatistics;
78use crate::system_schema::information_schema::flows::InformationSchemaFlows;
79use crate::system_schema::information_schema::information_memory_table::get_schema_columns;
80use crate::system_schema::information_schema::key_column_usage::InformationSchemaKeyColumnUsage;
81use crate::system_schema::information_schema::partitions::InformationSchemaPartitions;
82#[cfg(feature = "enterprise")]
83use crate::system_schema::information_schema::recycle_bin::InformationSchemaRecycleBin;
84use crate::system_schema::information_schema::region_peers::InformationSchemaRegionPeers;
85use crate::system_schema::information_schema::schemata::InformationSchemaSchemata;
86use crate::system_schema::information_schema::ssts::{
87 InformationSchemaSstsIndexMeta, InformationSchemaSstsManifest, InformationSchemaSstsStorage,
88};
89use crate::system_schema::information_schema::statistics::InformationSchemaStatistics;
90use crate::system_schema::information_schema::table_constraints::InformationSchemaTableConstraints;
91use crate::system_schema::information_schema::table_semantics::InformationSchemaTableSemantics;
92use crate::system_schema::information_schema::tables::InformationSchemaTables;
93use crate::system_schema::memory_table::MemoryTable;
94pub(crate) use crate::system_schema::predicate::Predicates;
95use crate::system_schema::{
96 SystemSchemaProvider, SystemSchemaProviderInner, SystemTable, SystemTableRef,
97};
98
99const DENSE_PRIMARY_KEY_ENCODING: &str = "dense";
100const SPARSE_PRIMARY_KEY_ENCODING: &str = "sparse";
101
102pub(crate) fn primary_key_encoding_index_type(options: &HashMap<String, String>) -> &'static str {
103 options
104 .get(PRIMARY_KEY_ENCODING)
105 .or_else(|| options.get(MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING))
106 .map(|value| {
107 if value.eq_ignore_ascii_case(SPARSE_PRIMARY_KEY_ENCODING) {
108 SPARSE_PRIMARY_KEY_ENCODING
109 } else {
110 DENSE_PRIMARY_KEY_ENCODING
111 }
112 })
113 .unwrap_or(DENSE_PRIMARY_KEY_ENCODING)
114}
115
116lazy_static! {
117 static ref MEMORY_TABLES: &'static [&'static str] = &[
119 ENGINES,
120 COLUMN_PRIVILEGES,
121 COLUMN_STATISTICS,
122 CHARACTER_SETS,
123 COLLATIONS,
124 COLLATION_CHARACTER_SET_APPLICABILITY,
125 CHECK_CONSTRAINTS,
126 EVENTS,
127 FILES,
128 OPTIMIZER_TRACE,
129 PARAMETERS,
130 PROFILING,
131 REFERENTIAL_CONSTRAINTS,
132 ROUTINES,
133 SCHEMA_PRIVILEGES,
134 TABLE_PRIVILEGES,
135 GLOBAL_STATUS,
136 SESSION_STATUS,
137 PARTITIONS,
138 PLUGINS,
139 USER_PRIVILEGES,
140 PROCESSLIST,
141 ];
142}
143
144macro_rules! setup_memory_table {
145 ($name: expr) => {
146 paste! {
147 {
148 let (schema, columns) = get_schema_columns($name);
149 Some(Arc::new(MemoryTable::new(
150 consts::[<INFORMATION_SCHEMA_ $name _TABLE_ID>],
151 $name,
152 schema,
153 columns
154 )) as _)
155 }
156 }
157 };
158}
159
160pub struct MakeInformationTableRequest {
161 pub catalog_name: String,
162 pub catalog_manager: Weak<dyn CatalogManager>,
163 pub kv_backend: KvBackendRef,
164}
165
166pub trait InformationSchemaTableFactory {
171 fn make_information_table(&self, req: MakeInformationTableRequest) -> SystemTableRef;
172}
173
174pub type InformationSchemaTableFactoryRef = Arc<dyn InformationSchemaTableFactory + Send + Sync>;
175
176pub struct InformationSchemaProvider {
178 catalog_name: String,
179 catalog_manager: Weak<dyn CatalogManager>,
180 process_manager: Option<ProcessManagerRef>,
181 flow_metadata_manager: Arc<FlowMetadataManager>,
182 tables: HashMap<String, TableRef>,
183 kv_backend: KvBackendRef,
184 extra_table_factories: HashMap<String, InformationSchemaTableFactoryRef>,
185}
186
187impl SystemSchemaProvider for InformationSchemaProvider {
188 fn tables(&self) -> &HashMap<String, TableRef> {
189 assert!(!self.tables.is_empty());
190
191 &self.tables
192 }
193}
194
195impl SystemSchemaProviderInner for InformationSchemaProvider {
196 fn catalog_name(&self) -> &str {
197 &self.catalog_name
198 }
199 fn schema_name() -> &'static str {
200 INFORMATION_SCHEMA_NAME
201 }
202
203 fn system_table(&self, name: &str) -> Option<SystemTableRef> {
204 if let Some(factory) = self.extra_table_factories.get(name) {
205 let req = MakeInformationTableRequest {
206 catalog_name: self.catalog_name.clone(),
207 catalog_manager: self.catalog_manager.clone(),
208 kv_backend: self.kv_backend.clone(),
209 };
210 return Some(factory.make_information_table(req));
211 }
212
213 match name.to_ascii_lowercase().as_str() {
214 TABLES => Some(Arc::new(InformationSchemaTables::new(
215 self.catalog_name.clone(),
216 self.catalog_manager.clone(),
217 )) as _),
218 COLUMNS => Some(Arc::new(InformationSchemaColumns::new(
219 self.catalog_name.clone(),
220 self.catalog_manager.clone(),
221 )) as _),
222 ENGINES => setup_memory_table!(ENGINES),
223 COLUMN_PRIVILEGES => setup_memory_table!(COLUMN_PRIVILEGES),
224 COLUMN_STATISTICS => setup_memory_table!(COLUMN_STATISTICS),
225 BUILD_INFO => setup_memory_table!(BUILD_INFO),
226 CHARACTER_SETS => setup_memory_table!(CHARACTER_SETS),
227 COLLATIONS => setup_memory_table!(COLLATIONS),
228 COLLATION_CHARACTER_SET_APPLICABILITY => {
229 setup_memory_table!(COLLATION_CHARACTER_SET_APPLICABILITY)
230 }
231 CHECK_CONSTRAINTS => setup_memory_table!(CHECK_CONSTRAINTS),
232 EVENTS => setup_memory_table!(EVENTS),
233 FILES => setup_memory_table!(FILES),
234 OPTIMIZER_TRACE => setup_memory_table!(OPTIMIZER_TRACE),
235 PARAMETERS => setup_memory_table!(PARAMETERS),
236 PROFILING => setup_memory_table!(PROFILING),
237 REFERENTIAL_CONSTRAINTS => setup_memory_table!(REFERENTIAL_CONSTRAINTS),
238 ROUTINES => setup_memory_table!(ROUTINES),
239 SCHEMA_PRIVILEGES => setup_memory_table!(SCHEMA_PRIVILEGES),
240 TABLE_PRIVILEGES => setup_memory_table!(TABLE_PRIVILEGES),
241 GLOBAL_STATUS => setup_memory_table!(GLOBAL_STATUS),
242 SESSION_STATUS => setup_memory_table!(SESSION_STATUS),
243 PLUGINS => setup_memory_table!(PLUGINS),
244 USER_PRIVILEGES => setup_memory_table!(USER_PRIVILEGES),
245 PROCESSLIST => setup_memory_table!(PROCESSLIST),
246 KEY_COLUMN_USAGE => Some(Arc::new(InformationSchemaKeyColumnUsage::new(
247 self.catalog_name.clone(),
248 self.catalog_manager.clone(),
249 )) as _),
250 SCHEMATA => Some(Arc::new(InformationSchemaSchemata::new(
251 self.catalog_name.clone(),
252 self.catalog_manager.clone(),
253 )) as _),
254 PARTITIONS => Some(Arc::new(InformationSchemaPartitions::new(
255 self.catalog_name.clone(),
256 self.catalog_manager.clone(),
257 )) as _),
258 REGION_PEERS => Some(Arc::new(InformationSchemaRegionPeers::new(
259 self.catalog_name.clone(),
260 self.catalog_manager.clone(),
261 )) as _),
262 TABLE_CONSTRAINTS => Some(Arc::new(InformationSchemaTableConstraints::new(
263 self.catalog_name.clone(),
264 self.catalog_manager.clone(),
265 )) as _),
266 STATISTICS => Some(Arc::new(InformationSchemaStatistics::new(
267 self.catalog_name.clone(),
268 self.catalog_manager.clone(),
269 )) as _),
270 CLUSTER_INFO => Some(Arc::new(InformationSchemaClusterInfo::new(
271 self.catalog_manager.clone(),
272 )) as _),
273 VIEWS => Some(Arc::new(InformationSchemaViews::new(
274 self.catalog_name.clone(),
275 self.catalog_manager.clone(),
276 )) as _),
277 FLOWS => Some(Arc::new(InformationSchemaFlows::new(
278 self.catalog_name.clone(),
279 self.catalog_manager.clone(),
280 self.flow_metadata_manager.clone(),
281 )) as _),
282 FLOW_STATISTICS => Some(Arc::new(InformationSchemaFlowStatistics::new(
283 self.catalog_name.clone(),
284 self.catalog_manager.clone(),
285 self.flow_metadata_manager.clone(),
286 )) as _),
287 PROCEDURE_INFO => Some(
288 Arc::new(procedure_info::InformationSchemaProcedureInfo::new(
289 self.catalog_manager.clone(),
290 )) as _,
291 ),
292 #[cfg(feature = "enterprise")]
293 RECYCLE_BIN => Some(Arc::new(InformationSchemaRecycleBin::new(
294 self.catalog_name.clone(),
295 self.catalog_manager.clone(),
296 )) as _),
297 REGION_STATISTICS => Some(Arc::new(
298 region_statistics::InformationSchemaRegionStatistics::new(
299 self.catalog_manager.clone(),
300 ),
301 ) as _),
302 REGION_INFO => Some(Arc::new(InformationSchemaRegionInfo::new(
303 self.catalog_manager.clone(),
304 )) as _),
305 PROCESS_LIST => self
306 .process_manager
307 .as_ref()
308 .map(|p| Arc::new(InformationSchemaProcessList::new(p.clone())) as _),
309 SSTS_MANIFEST => Some(Arc::new(InformationSchemaSstsManifest::new(
310 self.catalog_manager.clone(),
311 )) as _),
312 SSTS_STORAGE => Some(Arc::new(InformationSchemaSstsStorage::new(
313 self.catalog_manager.clone(),
314 )) as _),
315 SSTS_INDEX_META => Some(Arc::new(InformationSchemaSstsIndexMeta::new(
316 self.catalog_manager.clone(),
317 )) as _),
318 TABLE_SEMANTICS => Some(Arc::new(InformationSchemaTableSemantics::new(
319 self.catalog_name.clone(),
320 self.catalog_manager.clone(),
321 )) as _),
322 _ => None,
323 }
324 }
325}
326
327impl InformationSchemaProvider {
328 pub fn new(
329 catalog_name: String,
330 catalog_manager: Weak<dyn CatalogManager>,
331 flow_metadata_manager: Arc<FlowMetadataManager>,
332 process_manager: Option<ProcessManagerRef>,
333 kv_backend: KvBackendRef,
334 ) -> Self {
335 let mut provider = Self {
336 catalog_name,
337 catalog_manager,
338 flow_metadata_manager,
339 process_manager,
340 tables: HashMap::new(),
341 kv_backend,
342 extra_table_factories: HashMap::new(),
343 };
344
345 provider.build_tables();
346
347 provider
348 }
349
350 pub(crate) fn with_extra_table_factories(
351 mut self,
352 factories: HashMap<String, InformationSchemaTableFactoryRef>,
353 ) -> Self {
354 self.extra_table_factories = factories;
355 self.build_tables();
356 self
357 }
358
359 fn build_tables(&mut self) {
360 let mut tables = HashMap::new();
361
362 if self.catalog_name == DEFAULT_CATALOG_NAME {
367 tables.insert(
368 BUILD_INFO.to_string(),
369 self.build_table(BUILD_INFO).unwrap(),
370 );
371 tables.insert(
372 REGION_PEERS.to_string(),
373 self.build_table(REGION_PEERS).unwrap(),
374 );
375 tables.insert(
376 CLUSTER_INFO.to_string(),
377 self.build_table(CLUSTER_INFO).unwrap(),
378 );
379 tables.insert(
380 PROCEDURE_INFO.to_string(),
381 self.build_table(PROCEDURE_INFO).unwrap(),
382 );
383 tables.insert(
384 REGION_STATISTICS.to_string(),
385 self.build_table(REGION_STATISTICS).unwrap(),
386 );
387 tables.insert(
388 REGION_INFO.to_string(),
389 self.build_table(REGION_INFO).unwrap(),
390 );
391 tables.insert(
392 SSTS_MANIFEST.to_string(),
393 self.build_table(SSTS_MANIFEST).unwrap(),
394 );
395 tables.insert(
396 SSTS_STORAGE.to_string(),
397 self.build_table(SSTS_STORAGE).unwrap(),
398 );
399 tables.insert(
400 SSTS_INDEX_META.to_string(),
401 self.build_table(SSTS_INDEX_META).unwrap(),
402 );
403 }
404
405 tables.insert(TABLES.to_string(), self.build_table(TABLES).unwrap());
406 tables.insert(VIEWS.to_string(), self.build_table(VIEWS).unwrap());
407 tables.insert(SCHEMATA.to_string(), self.build_table(SCHEMATA).unwrap());
408 tables.insert(COLUMNS.to_string(), self.build_table(COLUMNS).unwrap());
409 tables.insert(
410 KEY_COLUMN_USAGE.to_string(),
411 self.build_table(KEY_COLUMN_USAGE).unwrap(),
412 );
413 tables.insert(
414 TABLE_CONSTRAINTS.to_string(),
415 self.build_table(TABLE_CONSTRAINTS).unwrap(),
416 );
417 tables.insert(
418 STATISTICS.to_string(),
419 self.build_table(STATISTICS).unwrap(),
420 );
421 tables.insert(FLOWS.to_string(), self.build_table(FLOWS).unwrap());
422 tables.insert(
423 FLOW_STATISTICS.to_string(),
424 self.build_table(FLOW_STATISTICS).unwrap(),
425 );
426 #[cfg(feature = "enterprise")]
427 tables.insert(
428 RECYCLE_BIN.to_string(),
429 self.build_table(RECYCLE_BIN).unwrap(),
430 );
431 tables.insert(
432 TABLE_SEMANTICS.to_string(),
433 self.build_table(TABLE_SEMANTICS).unwrap(),
434 );
435 if let Some(process_list) = self.build_table(PROCESS_LIST) {
436 tables.insert(PROCESS_LIST.to_string(), process_list);
437 }
438 for name in self.extra_table_factories.keys() {
439 tables.insert(name.clone(), self.build_table(name).expect(name));
440 }
441 for name in MEMORY_TABLES.iter() {
443 tables.insert((*name).to_string(), self.build_table(name).expect(name));
444 }
445 self.tables = tables;
446 }
447}
448
449pub trait InformationTable {
450 fn table_id(&self) -> TableId;
451
452 fn table_name(&self) -> &'static str;
453
454 fn schema(&self) -> SchemaRef;
455
456 fn to_stream(&self, request: ScanRequest) -> Result<SendableRecordBatchStream>;
457
458 fn scan_plan(&self, _request: ScanRequest) -> Result<Option<Arc<dyn ExecutionPlan>>> {
459 Ok(None)
460 }
461
462 fn table_type(&self) -> TableType {
463 TableType::Temporary
464 }
465}
466
467impl<T> SystemTable for T
469where
470 T: InformationTable,
471{
472 fn table_id(&self) -> TableId {
473 InformationTable::table_id(self)
474 }
475
476 fn table_name(&self) -> &'static str {
477 InformationTable::table_name(self)
478 }
479
480 fn schema(&self) -> SchemaRef {
481 InformationTable::schema(self)
482 }
483
484 fn table_type(&self) -> TableType {
485 InformationTable::table_type(self)
486 }
487
488 fn to_stream(&self, request: ScanRequest) -> Result<SendableRecordBatchStream> {
489 InformationTable::to_stream(self, request)
490 }
491
492 fn scan_plan(&self, request: ScanRequest) -> Result<Option<Arc<dyn ExecutionPlan>>> {
493 InformationTable::scan_plan(self, request)
494 }
495}
496
497pub type InformationExtensionRef = Arc<dyn InformationExtension<Error = Error> + Send + Sync>;
498
499#[async_trait::async_trait]
501pub trait InformationExtension {
502 type Error: ErrorExt;
503
504 async fn nodes(&self) -> std::result::Result<Vec<NodeInfo>, Self::Error>;
506
507 async fn procedures(&self) -> std::result::Result<Vec<(String, ProcedureInfo)>, Self::Error>;
509
510 async fn region_stats(&self) -> std::result::Result<Vec<RegionStat>, Self::Error>;
512
513 async fn flow_stats(&self) -> std::result::Result<Option<FlowStat>, Self::Error>;
515
516 async fn inspect_datanode(
518 &self,
519 request: DatanodeInspectRequest,
520 ) -> std::result::Result<SendableRecordBatchStream, Self::Error>;
521
522 fn inspect_datanode_plan(
525 &self,
526 _request: DatanodeInspectRequest,
527 _schema: SchemaRef,
528 ) -> std::result::Result<Option<Arc<dyn ExecutionPlan>>, Self::Error> {
529 Ok(None)
530 }
531}
532
533#[derive(Debug, Clone, PartialEq)]
535pub struct DatanodeInspectRequest {
536 pub kind: DatanodeInspectKind,
538
539 pub scan: ScanRequest,
542}
543
544#[derive(Debug, Clone, Copy, PartialEq, Eq)]
546pub enum DatanodeInspectKind {
547 SstManifest,
549 SstStorage,
551 SstIndexMeta,
553 RegionInfo,
555}
556
557impl DatanodeInspectRequest {
558 pub fn build_plan(self) -> std::result::Result<LogicalPlan, DataFusionError> {
560 match self.kind {
561 DatanodeInspectKind::SstManifest => ManifestSstEntry::build_plan(self.scan),
562 DatanodeInspectKind::SstStorage => StorageSstEntry::build_plan(self.scan),
563 DatanodeInspectKind::SstIndexMeta => PuffinIndexMetaEntry::build_plan(self.scan),
564 DatanodeInspectKind::RegionInfo => RegionInfoEntry::build_plan(self.scan),
565 }
566 }
567}
568pub struct NoopInformationExtension;
569
570#[async_trait::async_trait]
571impl InformationExtension for NoopInformationExtension {
572 type Error = Error;
573
574 async fn nodes(&self) -> std::result::Result<Vec<NodeInfo>, Self::Error> {
575 Ok(vec![])
576 }
577
578 async fn procedures(&self) -> std::result::Result<Vec<(String, ProcedureInfo)>, Self::Error> {
579 Ok(vec![])
580 }
581
582 async fn region_stats(&self) -> std::result::Result<Vec<RegionStat>, Self::Error> {
583 Ok(vec![])
584 }
585
586 async fn flow_stats(&self) -> std::result::Result<Option<FlowStat>, Self::Error> {
587 Ok(None)
588 }
589
590 async fn inspect_datanode(
591 &self,
592 _request: DatanodeInspectRequest,
593 ) -> std::result::Result<SendableRecordBatchStream, Self::Error> {
594 Ok(common_recordbatch::RecordBatches::empty().as_stream())
595 }
596}
597
598#[cfg(test)]
599mod tests {
600 use store_api::region_info::RegionInfoEntry;
601
602 use super::*;
603
604 #[test]
605 fn test_datanode_inspect_region_info_build_plan() {
606 let plan = DatanodeInspectRequest {
607 kind: DatanodeInspectKind::RegionInfo,
608 scan: ScanRequest::default(),
609 }
610 .build_plan()
611 .unwrap();
612
613 let LogicalPlan::TableScan(scan) = plan else {
614 panic!("expected table scan");
615 };
616 assert_eq!(
617 scan.table_name.to_string(),
618 RegionInfoEntry::reserved_table_name_for_inspection()
619 );
620 }
621}