Skip to main content

catalog/system_schema/
information_schema.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
15mod 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    // Memory tables in `information_schema`.
118    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
166/// A factory trait for making information schema tables.
167///
168/// This trait allows for extensibility of the information schema by providing
169/// a way to dynamically create custom information schema tables.
170pub trait InformationSchemaTableFactory {
171    fn make_information_table(&self, req: MakeInformationTableRequest) -> SystemTableRef;
172}
173
174pub type InformationSchemaTableFactoryRef = Arc<dyn InformationSchemaTableFactory + Send + Sync>;
175
176/// The `information_schema` tables info provider.
177pub 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        // SECURITY NOTE:
363        // Carefully consider the tables that may expose sensitive cluster configurations,
364        // authentication details, and other critical information.
365        // Only put these tables under `greptime` catalog to prevent info leak.
366        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        // Add memory tables
442        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
467// Provide compatibility for legacy `information_schema` code.
468impl<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/// The `InformationExtension` trait provides the extension methods for the `information_schema` tables.
500#[async_trait::async_trait]
501pub trait InformationExtension {
502    type Error: ErrorExt;
503
504    /// Gets the nodes information.
505    async fn nodes(&self) -> std::result::Result<Vec<NodeInfo>, Self::Error>;
506
507    /// Gets the procedures information.
508    async fn procedures(&self) -> std::result::Result<Vec<(String, ProcedureInfo)>, Self::Error>;
509
510    /// Gets the region statistics.
511    async fn region_stats(&self) -> std::result::Result<Vec<RegionStat>, Self::Error>;
512
513    /// Get the flow statistics. If no flownode is available, return `None`.
514    async fn flow_stats(&self) -> std::result::Result<Option<FlowStat>, Self::Error>;
515
516    /// Inspects the datanode.
517    async fn inspect_datanode(
518        &self,
519        request: DatanodeInspectRequest,
520    ) -> std::result::Result<SendableRecordBatchStream, Self::Error>;
521
522    /// Builds a physical plan for datanode inspect if the extension can expose
523    /// the distributed fan-in semantics to DataFusion.
524    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/// The request to inspect the datanode.
534#[derive(Debug, Clone, PartialEq)]
535pub struct DatanodeInspectRequest {
536    /// Kind to fetch from datanode.
537    pub kind: DatanodeInspectKind,
538
539    /// Pushdown scan configuration (projection/predicate/limit) for the returned stream.
540    /// This allows server-side filtering to reduce I/O and network costs.
541    pub scan: ScanRequest,
542}
543
544/// The kind of the datanode inspect request.
545#[derive(Debug, Clone, Copy, PartialEq, Eq)]
546pub enum DatanodeInspectKind {
547    /// List SST entries recorded in manifest
548    SstManifest,
549    /// List SST entries discovered in storage layer
550    SstStorage,
551    /// List index metadata collected from manifest
552    SstIndexMeta,
553    /// List region runtime and manifest info
554    RegionInfo,
555}
556
557impl DatanodeInspectRequest {
558    /// Builds a logical plan for the datanode inspect request.
559    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}