risingwave_frontend/catalog/
source_catalog.rs1use risingwave_common::catalog::{ColumnCatalog, ICEBERG_SOURCE_PREFIX, SourceVersionId};
16use risingwave_common::util::epoch::Epoch;
17use risingwave_connector::{WithOptionsSecResolved, WithPropertiesExt};
18use risingwave_pb::catalog::{PbSource, StreamSourceInfo, WatermarkDesc};
19use risingwave_pb::plan_common::SourceRefreshMode;
20use risingwave_sqlparser::ast;
21use risingwave_sqlparser::parser::Parser;
22use thiserror_ext::AsReport as _;
23
24use super::purify::try_purify_table_source_create_sql_ast;
25use super::{ColumnId, ConnectionId, DatabaseId, OwnedByUserCatalog, SchemaId, SourceId};
26use crate::catalog::TableId;
27use crate::error::Result;
28use crate::session::current::notice_to_user;
29use crate::user::UserId;
30
31#[derive(Clone, Debug, PartialEq, Eq, Hash)]
34pub struct SourceCatalog {
35 pub id: SourceId,
36 pub name: String,
37 pub schema_id: SchemaId,
38 pub database_id: DatabaseId,
39 pub columns: Vec<ColumnCatalog>,
40 pub pk_col_ids: Vec<ColumnId>,
41 pub append_only: bool,
42 pub owner: UserId,
43 pub info: StreamSourceInfo,
44 pub row_id_index: Option<usize>,
45 pub with_properties: WithOptionsSecResolved,
46 pub watermark_descs: Vec<WatermarkDesc>,
47 pub associated_table_id: Option<TableId>,
48 pub definition: String,
49 pub connection_id: Option<ConnectionId>,
50 pub created_at_epoch: Option<Epoch>,
51 pub initialized_at_epoch: Option<Epoch>,
52 pub version: SourceVersionId,
53 pub created_at_cluster_version: Option<String>,
54 pub initialized_at_cluster_version: Option<String>,
55 pub rate_limit: Option<u32>,
56 pub refresh_mode: Option<SourceRefreshMode>,
57}
58
59impl SourceCatalog {
60 pub fn create_sql(&self) -> String {
62 self.definition.clone()
63 }
64
65 pub fn create_sql_ast(&self) -> Result<ast::Statement> {
69 Ok(Parser::parse_exactly_one(&self.definition)?)
70 }
71
72 pub fn to_prost(&self) -> PbSource {
73 let (with_properties, secret_refs) = self.with_properties.clone().into_parts();
74 PbSource {
75 id: self.id,
76 schema_id: self.schema_id,
77 database_id: self.database_id,
78 name: self.name.clone(),
79 row_id_index: self.row_id_index.map(|idx| idx as _),
80 columns: self.columns.iter().map(|c| c.to_protobuf()).collect(),
81 pk_column_ids: self.pk_col_ids.iter().map(Into::into).collect(),
82 with_properties,
83 owner: self.owner,
84 info: Some(self.info.clone()),
85 watermark_descs: self.watermark_descs.clone(),
86 definition: self.definition.clone(),
87 connection_id: self.connection_id,
88 initialized_at_epoch: self.initialized_at_epoch.map(|x| x.0),
89 created_at_epoch: self.created_at_epoch.map(|x| x.0),
90 optional_associated_table_id: self.associated_table_id.map(Into::into),
91 version: self.version,
92 created_at_cluster_version: self.created_at_cluster_version.clone(),
93 initialized_at_cluster_version: self.initialized_at_cluster_version.clone(),
94 secret_refs,
95 rate_limit: self.rate_limit,
96 refresh_mode: self.refresh_mode,
97 }
98 }
99
100 pub fn version(&self) -> SourceVersionId {
102 self.version
103 }
104
105 pub fn connector_name(&self) -> String {
106 self.with_properties
107 .get_connector()
108 .expect("connector name is missing")
109 }
110
111 pub fn is_iceberg_connector(&self) -> bool {
112 self.with_properties.is_iceberg_connector()
113 }
114
115 pub fn is_cdc_table_source(&self) -> bool {
117 self.info.external_table.is_some()
118 }
119
120 pub fn iceberg_table_name(&self) -> Option<String> {
122 if self.name.starts_with(ICEBERG_SOURCE_PREFIX) {
123 Some(self.name[ICEBERG_SOURCE_PREFIX.len()..].to_string())
124 } else {
125 None
126 }
127 }
128
129 pub fn create_sql_purified(&self) -> String {
131 self.create_sql_ast_purified()
132 .and_then(|stmt| stmt.try_to_string().map_err(Into::into))
133 .unwrap_or_else(|_| self.create_sql())
134 }
135
136 pub fn create_sql_ast_purified(&self) -> Result<ast::Statement> {
140 if self.with_properties.is_cdc_connector() {
141 return self.create_sql_ast();
144 }
145
146 match try_purify_table_source_create_sql_ast(
147 self.create_sql_ast()?,
148 &self.columns,
149 self.row_id_index,
150 &self.pk_col_ids,
151 ) {
152 Ok(stmt) => return Ok(stmt),
153 Err(e) => notice_to_user(format!(
154 "error occurred while purifying definition for source \"{}\", \
155 results may be inaccurate: {}",
156 self.name,
157 e.as_report()
158 )),
159 }
160
161 self.create_sql_ast()
162 }
163
164 pub fn fill_purified_create_sql(&mut self) {
170 self.definition = self.create_sql_purified();
171 }
172}
173
174impl From<&PbSource> for SourceCatalog {
175 fn from(prost: &PbSource) -> Self {
176 let id = prost.id;
177 let name = prost.name.clone();
178 let database_id = prost.database_id;
179 let schema_id = prost.schema_id;
180 let prost_columns = prost.columns.clone();
181 let pk_col_ids = prost
182 .pk_column_ids
183 .clone()
184 .into_iter()
185 .map(Into::into)
186 .collect();
187 let connector_props_with_secrets =
188 WithOptionsSecResolved::new(prost.with_properties.clone(), prost.secret_refs.clone());
189 let columns = prost_columns.into_iter().map(ColumnCatalog::from).collect();
190 let row_id_index = prost.row_id_index.map(|idx| idx as _);
191
192 let append_only = row_id_index.is_some();
193 let owner = prost.owner;
194 let watermark_descs = prost.get_watermark_descs().clone();
195
196 let associated_table_id = prost.optional_associated_table_id.map(Into::into);
197 let version = prost.version;
198
199 let connection_id = prost.connection_id;
200 let rate_limit = prost.rate_limit;
201
202 Self {
203 id,
204 name,
205 schema_id,
206 database_id,
207 columns,
208 pk_col_ids,
209 append_only,
210 owner,
211 info: prost.info.clone().unwrap(),
212 row_id_index,
213 with_properties: connector_props_with_secrets,
214 watermark_descs,
215 associated_table_id,
216 definition: prost.definition.clone(),
217 connection_id,
218 created_at_epoch: prost.created_at_epoch.map(Epoch::from),
219 initialized_at_epoch: prost.initialized_at_epoch.map(Epoch::from),
220 version,
221 created_at_cluster_version: prost.created_at_cluster_version.clone(),
222 initialized_at_cluster_version: prost.initialized_at_cluster_version.clone(),
223 rate_limit,
224 refresh_mode: prost.refresh_mode,
225 }
226 }
227}
228
229impl OwnedByUserCatalog for SourceCatalog {
230 fn owner(&self) -> UserId {
231 self.owner
232 }
233}