Skip to main content

risingwave_common/catalog/
mod.rs

1// Copyright 2022 RisingWave Labs
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
15pub mod cdc_type_compatibility;
16mod column;
17mod external_table;
18mod internal_table;
19mod physical_table;
20mod schema;
21pub mod test_utils;
22
23use std::fmt::Binary;
24use std::sync::Arc;
25
26pub use column::*;
27pub use external_table::*;
28use futures::stream::BoxStream;
29pub use internal_table::*;
30use parse_display::Display;
31pub use physical_table::*;
32use risingwave_pb::catalog::table::PbEngine;
33use risingwave_pb::catalog::{
34    CreateType as PbCreateType, HandleConflictBehavior as PbHandleConflictBehavior,
35    StreamJobStatus as PbStreamJobStatus,
36};
37use risingwave_pb::plan_common::ColumnDescVersion;
38pub use schema::{Field, FieldDisplay, FieldLike, Schema, test_utils as schema_test_utils};
39
40use crate::array::DataChunk;
41pub use crate::constants::hummock;
42use crate::error::BoxedError;
43pub use crate::id::*;
44
45/// The global version of the catalog.
46pub type CatalogVersion = u64;
47
48/// The version number of the per-table catalog.
49pub type TableVersionId = u64;
50/// The default version ID for a new table.
51pub const INITIAL_TABLE_VERSION_ID: u64 = 0;
52/// The version number of the per-source catalog.
53pub type SourceVersionId = u64;
54/// The default version ID for a new source.
55pub const INITIAL_SOURCE_VERSION_ID: u64 = 0;
56
57pub const DEFAULT_DATABASE_NAME: &str = "dev";
58pub const DEFAULT_SCHEMA_NAME: &str = "public";
59pub const PG_CATALOG_SCHEMA_NAME: &str = "pg_catalog";
60pub const INFORMATION_SCHEMA_SCHEMA_NAME: &str = "information_schema";
61pub const RW_CATALOG_SCHEMA_NAME: &str = "rw_catalog";
62pub const RESERVED_PG_SCHEMA_PREFIX: &str = "pg_";
63pub const DEFAULT_SUPER_USER: &str = "root";
64pub const DEFAULT_SUPER_USER_ID: UserId = UserId::new(1);
65// This is for compatibility with customized utils for PostgreSQL.
66pub const DEFAULT_SUPER_USER_FOR_PG: &str = "postgres";
67pub const DEFAULT_SUPER_USER_FOR_PG_ID: u32 = 2;
68
69// This is the default superuser for admin, which is used only for cloud control plane.
70pub const DEFAULT_SUPER_USER_FOR_ADMIN: &str = "rwadmin";
71pub const DEFAULT_SUPER_USER_FOR_ADMIN_ID: UserId = UserId::new(3);
72
73pub const NON_RESERVED_USER_ID: UserId = UserId::new(11);
74
75pub const MAX_SYS_CATALOG_NUM: i32 = 5000;
76pub const SYS_CATALOG_START_ID: i32 = i32::MAX - MAX_SYS_CATALOG_NUM;
77
78pub use risingwave_pb::id::OBJECT_ID_PLACEHOLDER;
79
80pub const SYSTEM_SCHEMAS: [&str; 3] = [
81    PG_CATALOG_SCHEMA_NAME,
82    INFORMATION_SCHEMA_SCHEMA_NAME,
83    RW_CATALOG_SCHEMA_NAME,
84];
85pub fn is_system_schema(schema_name: &str) -> bool {
86    SYSTEM_SCHEMAS.contains(&schema_name)
87}
88
89pub fn is_reserved_admin_user(user_name: &str) -> bool {
90    user_name == DEFAULT_SUPER_USER_FOR_ADMIN
91}
92
93pub const RW_RESERVED_COLUMN_NAME_PREFIX: &str = "_rw_";
94
95/// When there is no primary key specified while creating source, will use
96/// the message key as primary key in `BYTEA` type with this name.
97/// Note: the field has version to track, please refer to [`default_key_column_name_version_mapping`]
98pub const DEFAULT_KEY_COLUMN_NAME: &str = "_rw_key";
99
100pub fn default_key_column_name_version_mapping(version: &ColumnDescVersion) -> &str {
101    match version {
102        ColumnDescVersion::Unspecified => DEFAULT_KEY_COLUMN_NAME,
103        _ => DEFAULT_KEY_COLUMN_NAME,
104    }
105}
106
107/// For kafka source, we attach a hidden column [`KAFKA_TIMESTAMP_COLUMN_NAME`] to it, so that we
108/// can limit the timestamp range when querying it directly with batch query. The column type is
109/// [`crate::types::DataType::Timestamptz`]. For more details, please refer to
110/// [this rfc](https://github.com/risingwavelabs/rfcs/pull/20).
111pub const KAFKA_TIMESTAMP_COLUMN_NAME: &str = "_rw_kafka_timestamp";
112
113/// RisingWave iceberg table engine will create the column `_risingwave_iceberg_row_id` in the iceberg table.
114///
115/// Iceberg V3 spec use `_row_id` as a reserved column name for row lineage, so if the table without primary key,
116/// we can't use `_row_id` directly for iceberg, so use `_risingwave_iceberg_row_id` instead.
117pub const RISINGWAVE_ICEBERG_ROW_ID: &str = "_risingwave_iceberg_row_id";
118
119pub const ROW_ID_COLUMN_NAME: &str = "_row_id";
120/// The column ID preserved for the row ID column.
121pub const ROW_ID_COLUMN_ID: ColumnId = ColumnId::new(0);
122
123/// The column ID offset for user-defined columns.
124///
125/// All IDs of user-defined columns must be greater or equal to this value.
126pub const USER_COLUMN_ID_OFFSET: i32 = ROW_ID_COLUMN_ID.next().get_id();
127
128pub const RW_TIMESTAMP_COLUMN_NAME: &str = "_rw_timestamp";
129pub const RW_TIMESTAMP_COLUMN_ID: ColumnId = ColumnId::new(-1);
130
131/// The column name for the projected row ID in `ProjectSet`.
132/// This is a hidden column used to track row indices when expanding set-returning functions.
133pub const PROJECTED_ROW_ID_COLUMN_NAME: &str = "_rw_projected_row_id";
134
135pub const ICEBERG_SEQUENCE_NUM_COLUMN_NAME: &str = "_iceberg_sequence_number";
136pub const ICEBERG_FILE_PATH_COLUMN_NAME: &str = "_iceberg_file_path";
137pub const ICEBERG_FILE_POS_COLUMN_NAME: &str = "_iceberg_file_pos";
138
139pub const CDC_OFFSET_COLUMN_NAME: &str = "_rw_offset";
140/// The number of columns output by the cdc source job
141/// see [`ColumnCatalog::debezium_cdc_source_cols()`] for details
142pub const CDC_SOURCE_COLUMN_NUM: u32 = 3;
143pub const CDC_TABLE_NAME_COLUMN_NAME: &str = "_rw_table_name";
144
145pub const ICEBERG_SOURCE_PREFIX: &str = "__iceberg_source_";
146pub const ICEBERG_SINK_PREFIX: &str = "__iceberg_sink_";
147
148/// Iceberg snapshot summary property containing the highest RisingWave epoch committed by the
149/// sink. Append-only Iceberg engine batch reads use it to exclude committed rows from the sink
150/// log store.
151pub const RISINGWAVE_ICEBERG_COMMIT_EPOCH: &str = "risingwave.commit.epoch";
152
153/// The local system catalog reader in the frontend node.
154pub trait SysCatalogReader: Sync + Send + 'static {
155    /// Reads the data of the system catalog table.
156    fn read_table(&self, table_id: TableId) -> BoxStream<'_, Result<DataChunk, BoxedError>>;
157}
158
159pub type SysCatalogReaderRef = Arc<dyn SysCatalogReader>;
160
161#[derive(Clone, Debug, PartialEq, Default, Copy)]
162pub struct TableOption {
163    pub retention_seconds: Option<u32>, // second
164}
165
166impl From<&risingwave_pb::hummock::TableOption> for TableOption {
167    fn from(table_option: &risingwave_pb::hummock::TableOption) -> Self {
168        Self {
169            retention_seconds: table_option.retention_seconds,
170        }
171    }
172}
173
174impl From<&TableOption> for risingwave_pb::hummock::TableOption {
175    fn from(table_option: &TableOption) -> Self {
176        Self {
177            retention_seconds: table_option.retention_seconds,
178        }
179    }
180}
181
182impl TableOption {
183    pub fn new(retention_seconds: Option<u32>) -> Self {
184        // now we only support ttl for TableOption
185        TableOption { retention_seconds }
186    }
187}
188
189#[derive(Default, Debug, Clone, Copy, PartialEq, Eq, Hash)]
190pub enum ConflictBehavior {
191    #[default]
192    NoCheck,
193    Overwrite,
194    IgnoreConflict,
195    DoUpdateIfNotNull,
196}
197
198#[macro_export]
199macro_rules! _checked_conflict_behaviors {
200    () => {
201        ConflictBehavior::Overwrite
202            | ConflictBehavior::IgnoreConflict
203            | ConflictBehavior::DoUpdateIfNotNull
204    };
205}
206pub use _checked_conflict_behaviors as checked_conflict_behaviors;
207
208impl ConflictBehavior {
209    pub fn from_protobuf(tb_conflict_behavior: &PbHandleConflictBehavior) -> Self {
210        match tb_conflict_behavior {
211            PbHandleConflictBehavior::Overwrite => ConflictBehavior::Overwrite,
212            PbHandleConflictBehavior::Ignore => ConflictBehavior::IgnoreConflict,
213            PbHandleConflictBehavior::DoUpdateIfNotNull => ConflictBehavior::DoUpdateIfNotNull,
214            // This is for backward compatibility, in the previous version
215            // `HandleConflictBehavior::Unspecified` represented `NoCheck`, so just treat it as `NoCheck`.
216            PbHandleConflictBehavior::NoCheck | PbHandleConflictBehavior::Unspecified => {
217                ConflictBehavior::NoCheck
218            }
219        }
220    }
221
222    pub fn to_protobuf(self) -> PbHandleConflictBehavior {
223        match self {
224            ConflictBehavior::NoCheck => PbHandleConflictBehavior::NoCheck,
225            ConflictBehavior::Overwrite => PbHandleConflictBehavior::Overwrite,
226            ConflictBehavior::IgnoreConflict => PbHandleConflictBehavior::Ignore,
227            ConflictBehavior::DoUpdateIfNotNull => PbHandleConflictBehavior::DoUpdateIfNotNull,
228        }
229    }
230
231    pub fn debug_to_string(self) -> String {
232        match self {
233            ConflictBehavior::NoCheck => "NoCheck".to_owned(),
234            ConflictBehavior::Overwrite => "Overwrite".to_owned(),
235            ConflictBehavior::IgnoreConflict => "IgnoreConflict".to_owned(),
236            ConflictBehavior::DoUpdateIfNotNull => "DoUpdateIfNotNull".to_owned(),
237        }
238    }
239}
240
241#[derive(Default, Debug, Clone, Copy, PartialEq, Eq, Hash)]
242pub enum Engine {
243    #[default]
244    Hummock,
245    Iceberg,
246}
247
248impl Engine {
249    pub fn from_protobuf(engine: &PbEngine) -> Self {
250        match engine {
251            PbEngine::Hummock | PbEngine::Unspecified => Engine::Hummock,
252            PbEngine::Iceberg => Engine::Iceberg,
253        }
254    }
255
256    pub fn to_protobuf(self) -> PbEngine {
257        match self {
258            Engine::Hummock => PbEngine::Hummock,
259            Engine::Iceberg => PbEngine::Iceberg,
260        }
261    }
262
263    pub fn debug_to_string(self) -> String {
264        match self {
265            Engine::Hummock => "Hummock".to_owned(),
266            Engine::Iceberg => "Iceberg".to_owned(),
267        }
268    }
269}
270
271#[derive(Clone, Copy, Debug, Default, Display, Hash, PartialOrd, PartialEq, Eq, Ord)]
272pub enum StreamJobStatus {
273    #[default]
274    Creating,
275    Created,
276}
277
278impl StreamJobStatus {
279    pub fn from_proto(stream_job_status: PbStreamJobStatus) -> Self {
280        match stream_job_status {
281            PbStreamJobStatus::Creating => StreamJobStatus::Creating,
282            PbStreamJobStatus::Created | PbStreamJobStatus::Unspecified => StreamJobStatus::Created,
283        }
284    }
285
286    pub fn to_proto(self) -> PbStreamJobStatus {
287        match self {
288            StreamJobStatus::Creating => PbStreamJobStatus::Creating,
289            StreamJobStatus::Created => PbStreamJobStatus::Created,
290        }
291    }
292}
293
294#[derive(Clone, Copy, Debug, Display, Hash, PartialOrd, PartialEq, Eq, Ord, Default)]
295pub enum CreateType {
296    #[default]
297    Foreground,
298    Background,
299}
300
301impl CreateType {
302    pub fn from_proto(pb_create_type: PbCreateType) -> Self {
303        match pb_create_type {
304            PbCreateType::Foreground | PbCreateType::Unspecified => CreateType::Foreground,
305            PbCreateType::Background => CreateType::Background,
306        }
307    }
308
309    pub fn to_proto(self) -> PbCreateType {
310        match self {
311            CreateType::Foreground => PbCreateType::Foreground,
312            CreateType::Background => PbCreateType::Background,
313        }
314    }
315}
316
317#[derive(Clone, Debug)]
318pub enum AlterDatabaseParam {
319    // Barrier related parameters, per database.
320    // None represents the default value, which means it follows `SystemParams`.
321    BarrierIntervalMs(Option<u32>),
322    CheckpointFrequency(Option<u64>),
323}
324
325macro_rules! for_all_fragment_type_flags {
326    () => {
327        for_all_fragment_type_flags! {
328            {
329                Source,
330                Mview,
331                Sink,
332                Now,
333                StreamScan,
334                BarrierRecv,
335                Values,
336                Dml,
337                CdcFilter,
338                Skipped1,
339                SourceScan,
340                SnapshotBackfillStreamScan,
341                FsFetch,
342                CrossDbSnapshotBackfillStreamScan,
343                StreamCdcScan,
344                VectorIndexWrite,
345                UpstreamSinkUnion,
346                LocalityProvider
347            },
348            {},
349            0
350        }
351    };
352    (
353        {},
354        {
355            $(
356                {$flag:ident, $index:expr}
357            ),*
358        },
359        $next_index:expr
360    ) => {
361        #[derive(Clone, Copy, Debug, Display, Hash, PartialOrd, PartialEq, Eq)]
362        #[repr(u32)]
363        pub enum FragmentTypeFlag {
364            $(
365                $flag = (1 << $index),
366            )*
367        }
368
369        pub const FRAGMENT_TYPE_FLAG_LIST: [FragmentTypeFlag; $next_index] = [
370            $(
371                FragmentTypeFlag::$flag,
372            )*
373        ];
374
375        impl TryFrom<u32> for FragmentTypeFlag {
376            type Error = String;
377
378            fn try_from(value: u32) -> Result<Self, Self::Error> {
379                match value {
380                    $(
381                        value if value == (FragmentTypeFlag::$flag as u32) => Ok(FragmentTypeFlag::$flag),
382                    )*
383                    _ => Err(format!("Invalid FragmentTypeFlag value: {}", value)),
384                }
385            }
386        }
387
388        impl FragmentTypeFlag {
389            pub fn as_str_name(&self) -> &'static str {
390                match self {
391                    $(
392                        FragmentTypeFlag::$flag => paste::paste!{stringify!( [< $flag:snake:upper >] )},
393                    )*
394                }
395            }
396        }
397    };
398    (
399        {$first:ident $(, $rest:ident)*},
400        {
401            $(
402                {$flag:ident, $index:expr}
403            ),*
404        },
405        $next_index:expr
406    ) => {
407        for_all_fragment_type_flags! {
408            {$($rest),*},
409            {
410                $({$flag, $index},)*
411                {$first, $next_index}
412            },
413            $next_index + 1
414        }
415    };
416}
417
418for_all_fragment_type_flags!();
419
420impl FragmentTypeFlag {
421    pub fn raw_flag(flags: impl IntoIterator<Item = FragmentTypeFlag>) -> u32 {
422        flags.into_iter().fold(0, |acc, flag| acc | (flag as u32))
423    }
424
425    /// Fragments that may be affected by `BACKFILL_RATE_LIMIT`.
426    pub fn backfill_rate_limit_fragments() -> impl Iterator<Item = FragmentTypeFlag> {
427        [
428            FragmentTypeFlag::SourceScan,
429            FragmentTypeFlag::StreamScan,
430            FragmentTypeFlag::LocalityProvider,
431        ]
432        .into_iter()
433    }
434
435    /// Fragments that may be affected by `SOURCE_RATE_LIMIT`.
436    /// Note: for `FsFetch`, old fragments don't have this flag set, so don't use this to check.
437    pub fn source_rate_limit_fragments() -> impl Iterator<Item = FragmentTypeFlag> {
438        [FragmentTypeFlag::Source, FragmentTypeFlag::FsFetch].into_iter()
439    }
440
441    /// Fragments that may be affected by `BACKFILL_RATE_LIMIT`.
442    pub fn sink_rate_limit_fragments() -> impl Iterator<Item = FragmentTypeFlag> {
443        [FragmentTypeFlag::Sink].into_iter()
444    }
445
446    /// Note: this doesn't include `FsFetch` created in old versions.
447    pub fn rate_limit_fragments() -> impl Iterator<Item = FragmentTypeFlag> {
448        Self::backfill_rate_limit_fragments()
449            .chain(Self::source_rate_limit_fragments())
450            .chain(Self::sink_rate_limit_fragments())
451            .chain(Self::dml_rate_limit_fragments())
452    }
453
454    pub fn dml_rate_limit_fragments() -> impl Iterator<Item = FragmentTypeFlag> {
455        [FragmentTypeFlag::Dml].into_iter()
456    }
457}
458
459#[derive(Clone, Copy, Debug, Hash, PartialOrd, PartialEq, Eq, Default)]
460pub struct FragmentTypeMask(u32);
461
462impl Binary for FragmentTypeMask {
463    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
464        write!(f, "{:b}", self.0)
465    }
466}
467
468impl From<i32> for FragmentTypeMask {
469    fn from(value: i32) -> Self {
470        Self(value as u32)
471    }
472}
473
474impl From<u32> for FragmentTypeMask {
475    fn from(value: u32) -> Self {
476        Self(value)
477    }
478}
479
480impl From<FragmentTypeMask> for u32 {
481    fn from(value: FragmentTypeMask) -> Self {
482        value.0
483    }
484}
485
486impl From<FragmentTypeMask> for i32 {
487    fn from(value: FragmentTypeMask) -> Self {
488        value.0 as _
489    }
490}
491
492impl FragmentTypeMask {
493    pub fn empty() -> Self {
494        FragmentTypeMask(0)
495    }
496
497    pub fn add(&mut self, flag: FragmentTypeFlag) {
498        self.0 |= flag as u32;
499    }
500
501    pub fn contains_any(&self, flags: impl IntoIterator<Item = FragmentTypeFlag>) -> bool {
502        let flag = FragmentTypeFlag::raw_flag(flags);
503        (self.0 & flag) != 0
504    }
505
506    pub fn contains(&self, flag: FragmentTypeFlag) -> bool {
507        self.contains_any([flag])
508    }
509}
510
511#[cfg(test)]
512mod tests {
513    use itertools::Itertools;
514    use risingwave_common::catalog::FRAGMENT_TYPE_FLAG_LIST;
515
516    use crate::catalog::FragmentTypeFlag;
517
518    #[test]
519    fn test_all_fragment_type_flag() {
520        expect_test::expect![[r#"
521            [
522                (
523                    Source,
524                    1,
525                    "SOURCE",
526                ),
527                (
528                    Mview,
529                    2,
530                    "MVIEW",
531                ),
532                (
533                    Sink,
534                    4,
535                    "SINK",
536                ),
537                (
538                    Now,
539                    8,
540                    "NOW",
541                ),
542                (
543                    StreamScan,
544                    16,
545                    "STREAM_SCAN",
546                ),
547                (
548                    BarrierRecv,
549                    32,
550                    "BARRIER_RECV",
551                ),
552                (
553                    Values,
554                    64,
555                    "VALUES",
556                ),
557                (
558                    Dml,
559                    128,
560                    "DML",
561                ),
562                (
563                    CdcFilter,
564                    256,
565                    "CDC_FILTER",
566                ),
567                (
568                    Skipped1,
569                    512,
570                    "SKIPPED1",
571                ),
572                (
573                    SourceScan,
574                    1024,
575                    "SOURCE_SCAN",
576                ),
577                (
578                    SnapshotBackfillStreamScan,
579                    2048,
580                    "SNAPSHOT_BACKFILL_STREAM_SCAN",
581                ),
582                (
583                    FsFetch,
584                    4096,
585                    "FS_FETCH",
586                ),
587                (
588                    CrossDbSnapshotBackfillStreamScan,
589                    8192,
590                    "CROSS_DB_SNAPSHOT_BACKFILL_STREAM_SCAN",
591                ),
592                (
593                    StreamCdcScan,
594                    16384,
595                    "STREAM_CDC_SCAN",
596                ),
597                (
598                    VectorIndexWrite,
599                    32768,
600                    "VECTOR_INDEX_WRITE",
601                ),
602                (
603                    UpstreamSinkUnion,
604                    65536,
605                    "UPSTREAM_SINK_UNION",
606                ),
607                (
608                    LocalityProvider,
609                    131072,
610                    "LOCALITY_PROVIDER",
611                ),
612            ]
613        "#]]
614        .assert_debug_eq(
615            &FRAGMENT_TYPE_FLAG_LIST
616                .into_iter()
617                .map(|flag| (flag, flag as u32, flag.as_str_name()))
618                .collect_vec(),
619        );
620        for flag in FRAGMENT_TYPE_FLAG_LIST {
621            assert_eq!(FragmentTypeFlag::try_from(flag as u32).unwrap(), flag);
622        }
623    }
624}