1use anyhow::{Context, anyhow};
16use fancy_regex::Regex;
17use itertools::Itertools;
18use risingwave_common::catalog::{CdcKeyComparison, ColumnCatalog};
19use risingwave_connector::WithOptionsSecResolved;
20use risingwave_connector::source::UPSTREAM_SOURCE_KEY;
21use risingwave_connector::source::cdc::external::{
22 DATABASE_NAME_KEY, ExternalTableConfig, ExternalTableImpl, SCHEMA_NAME_KEY, SchemaTableName,
23 TABLE_NAME_KEY,
24};
25use risingwave_connector::source::cdc::{
26 MYSQL_CDC_CONNECTOR, POSTGRES_CDC_CONNECTOR, SQL_SERVER_CDC_CONNECTOR,
27};
28use risingwave_sqlparser::ast::{ColumnDef, ColumnOption, SourceWatermark, TableConstraint};
29use thiserror_ext::AsReport;
30
31use crate::error::{ErrorCode, Result, RwError};
32use crate::handler::create_source::reject_variant_columns;
33use crate::handler::create_table::{bind_sql_columns, bind_sql_pk_names, bind_table_constraints};
34
35pub(crate) fn derive_with_options_for_cdc_table(
43 source_with_properties: &WithOptionsSecResolved,
44 external_table_name: String,
45) -> Result<(WithOptionsSecResolved, String)> {
46 let source_database_name: &str = source_with_properties
48 .get("database.name")
49 .ok_or_else(|| anyhow!("The source with properties does not contain 'database.name'"))?
50 .as_str();
51 let mut with_options = source_with_properties.clone();
52 if let Some(connector) = source_with_properties.get(UPSTREAM_SOURCE_KEY) {
53 match connector.as_str() {
54 MYSQL_CDC_CONNECTOR => {
55 let (db_name, table_name) = external_table_name.split_once('.').ok_or_else(|| {
58 anyhow!("The upstream table name must contain database name prefix, e.g. 'database.table'")
59 })?;
60 if !source_database_name
62 .split(',')
63 .map(|s| s.trim())
64 .any(|name| name == db_name)
65 {
66 return Err(anyhow!(
67 "The database name `{}` in the FROM clause is not included in the database name `{}` in source definition",
68 db_name,
69 source_database_name
70 ).into());
71 }
72 with_options.insert(DATABASE_NAME_KEY.into(), db_name.into());
73 with_options.insert(TABLE_NAME_KEY.into(), table_name.into());
74 return Ok((with_options, external_table_name));
76 }
77 POSTGRES_CDC_CONNECTOR => {
78 let (schema_name, table_name) =
79 parse_postgres_cdc_external_table_name(&external_table_name)?;
80
81 with_options.insert(SCHEMA_NAME_KEY.into(), schema_name);
83 with_options.insert(TABLE_NAME_KEY.into(), table_name);
84 return Ok((with_options, external_table_name));
86 }
87 SQL_SERVER_CDC_CONNECTOR => {
88 let parts: Vec<&str> = external_table_name.split('.').collect();
96 let (schema_name, table_name) = match parts.len() {
97 3 => {
98 let db_name = parts[0];
101 let schema_name = parts[1];
102 let table_name = parts[2];
103
104 if db_name != source_database_name {
105 return Err(anyhow!(
106 "The database name '{}' in FROM clause does not match the database name '{}' specified in source definition. \
107 You can either use 'schema.table' format (recommended) or ensure the database name matches.",
108 db_name,
109 source_database_name
110 ).into());
111 }
112 (schema_name, table_name)
113 }
114 2 => {
115 let schema_name = parts[0];
118 let table_name = parts[1];
119 (schema_name, table_name)
120 }
121 1 => {
122 return Err(anyhow!(
125 "Invalid table name format '{}'. For SQL Server CDC, you must specify the schema name. \
126 Use 'schema.table' format (e.g., 'dbo.{}') or 'database.schema.table' format (e.g., '{}.dbo.{}').",
127 external_table_name,
128 external_table_name,
129 source_database_name,
130 external_table_name
131 ).into());
132 }
133 _ => {
134 return Err(anyhow!(
136 "Invalid table name format '{}'. Expected 'schema.table' or 'database.schema.table'.",
137 external_table_name
138 ).into());
139 }
140 };
141
142 with_options.insert(SCHEMA_NAME_KEY.into(), schema_name.into());
144 with_options.insert(TABLE_NAME_KEY.into(), table_name.into());
145
146 let normalized_external_table_name = format!("{}.{}", schema_name, table_name);
149 return Ok((with_options, normalized_external_table_name));
150 }
151 _ => {
152 return Err(RwError::from(anyhow!(
153 "connector {} is not supported for cdc table",
154 connector
155 )));
156 }
157 };
158 }
159 unreachable!("All valid CDC connectors should have returned by now")
160}
161
162fn parse_postgres_cdc_external_table_name(external_table_name: &str) -> Result<(String, String)> {
167 let mut parts = vec![];
168 let mut current = String::new();
169 let mut chars = external_table_name.chars().peekable();
170 let mut in_quote = false;
171 let mut just_closed_quote = false;
172
173 while let Some(ch) = chars.next() {
174 if in_quote {
175 if ch == '"' {
176 if chars.peek() == Some(&'"') {
177 current.push('"');
178 chars.next();
179 } else {
180 in_quote = false;
181 just_closed_quote = true;
182 }
183 } else {
184 current.push(ch);
185 }
186 } else {
187 match ch {
188 '.' => {
189 if current.is_empty() {
190 return Err(anyhow!(
191 "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
192 external_table_name
193 )
194 .into());
195 }
196 parts.push(std::mem::take(&mut current));
197 just_closed_quote = false;
198 }
199 '"' if current.is_empty() => {
200 in_quote = true;
201 }
202 '"' => {
203 return Err(anyhow!(
204 "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
205 external_table_name
206 )
207 .into());
208 }
209 _ if just_closed_quote => {
210 return Err(anyhow!(
211 "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
212 external_table_name
213 )
214 .into());
215 }
216 _ => current.push(ch),
217 }
218 }
219 }
220
221 if in_quote || current.is_empty() {
222 return Err(anyhow!(
223 "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
224 external_table_name
225 )
226 .into());
227 }
228 parts.push(current);
229
230 if let [schema_name, table_name] = parts.as_slice() {
231 Ok((schema_name.clone(), table_name.clone()))
232 } else {
233 Err(
234 anyhow!("The upstream table name must contain schema name prefix, e.g. 'public.table'")
235 .into(),
236 )
237 }
238}
239
240pub(crate) fn reject_pk_filtered_by_debezium_column_filter(
253 pk_names: &[String],
254 cdc_with_options: &WithOptionsSecResolved,
255) -> Result<()> {
256 const EXCLUDE_KEY: &str = "debezium.column.exclude.list";
257 const INCLUDE_KEY: &str = "debezium.column.include.list";
258
259 let st = SchemaTableName::from_properties(cdc_with_options.as_plaintext());
260 reject_pk_filtered_by_debezium_column_filter_inner(
261 pk_names,
262 &st,
263 cdc_with_options.get(EXCLUDE_KEY).map(String::as_str),
264 cdc_with_options.get(INCLUDE_KEY).map(String::as_str),
265 )
266}
267
268fn reject_pk_filtered_by_debezium_column_filter_inner(
269 pk_names: &[String],
270 st: &SchemaTableName,
271 exclude_list: Option<&str>,
272 include_list: Option<&str>,
273) -> Result<()> {
274 const EXCLUDE_KEY: &str = "debezium.column.exclude.list";
275 const INCLUDE_KEY: &str = "debezium.column.include.list";
276
277 let pk_full_names = pk_names
278 .iter()
279 .map(|pk| (pk, format!("{}.{}.{}", st.schema_name, st.table_name, pk)))
280 .collect_vec();
281
282 if let Some(exclude_list) = exclude_list {
283 let patterns = compile_debezium_column_filter_patterns(EXCLUDE_KEY, exclude_list)?;
284 for (pk, pk_full_name) in &pk_full_names {
285 for (pattern, regex) in &patterns {
286 if regex.is_match(pk_full_name).map_err(|err| {
287 ErrorCode::InvalidInputSyntax(format!(
288 "failed to evaluate Debezium column filter pattern `{pattern}` in `{EXCLUDE_KEY}`: {}",
289 err.as_report()
290 ))
291 })? {
292 return Err(ErrorCode::InvalidInputSyntax(format!(
293 "primary key column `{pk}` is excluded by `{EXCLUDE_KEY}` pattern \
294 `{pattern}`. Excluding a PK column causes silent data corruption: \
295 Debezium keeps the PK in the message key but drops it from the payload, \
296 so RisingWave cannot match UPDATE/DELETE events against the original row."
297 ))
298 .into());
299 }
300 }
301 }
302 }
303
304 if let Some(include_list) = include_list {
305 let patterns = compile_debezium_column_filter_patterns(INCLUDE_KEY, include_list)?;
306 for (pk, pk_full_name) in &pk_full_names {
307 let mut included = false;
308 for (_, regex) in &patterns {
309 if regex.is_match(pk_full_name).map_err(|err| {
310 ErrorCode::InvalidInputSyntax(format!(
311 "failed to evaluate Debezium column filter pattern in `{INCLUDE_KEY}`: {}",
312 err.as_report()
313 ))
314 })? {
315 included = true;
316 break;
317 }
318 }
319 if !included {
320 return Err(ErrorCode::InvalidInputSyntax(format!(
321 "primary key column `{pk}` is not included by `{INCLUDE_KEY}`. Omitting a PK \
322 column causes silent data corruption: Debezium keeps the PK in the message key \
323 but drops it from the payload, so RisingWave cannot match UPDATE/DELETE events \
324 against the original row."
325 ))
326 .into());
327 }
328 }
329 }
330
331 Ok(())
332}
333
334fn compile_debezium_column_filter_patterns(
335 key: &str,
336 filter_list: &str,
337) -> Result<Vec<(String, Regex)>> {
338 filter_list
339 .split(',')
340 .map(str::trim)
341 .filter(|pattern| !pattern.is_empty())
342 .map(|pattern| {
343 let anchored_pattern = format!("(?i:^(?:{pattern})$)");
344 let regex = Regex::new(&anchored_pattern).map_err(|err| {
345 ErrorCode::InvalidInputSyntax(format!(
346 "invalid Debezium column filter pattern `{pattern}` in `{key}`: {}",
347 err.as_report()
348 ))
349 })?;
350 Ok((pattern.to_owned(), regex))
351 })
352 .collect()
353}
354
355pub(crate) fn not_null_check_for_cdc_table(
357 wildcard_idx: &Option<usize>,
358 column_defs: &Vec<ColumnDef>,
359) -> Result<()> {
360 if !wildcard_idx.is_some()
361 && column_defs.iter().any(|col| {
362 col.options
363 .iter()
364 .any(|opt| matches!(opt.option, ColumnOption::NotNull))
365 })
366 {
367 return Err(ErrorCode::NotSupported(
368 "CDC table with NOT NULL constraint is not supported".to_owned(),
369 "Please remove the NOT NULL constraint for columns".to_owned(),
370 )
371 .into());
372 }
373 Ok(())
374}
375
376pub(crate) fn sanity_check_for_table_on_cdc_source(
378 append_only: bool,
379 column_defs: &Vec<ColumnDef>,
380 wildcard_idx: &Option<usize>,
381 constraints: &Vec<TableConstraint>,
382 source_watermarks: &Vec<SourceWatermark>,
383) -> Result<()> {
384 if wildcard_idx.is_some() && !column_defs.is_empty() {
386 return Err(ErrorCode::NotSupported(
387 "wildcard(*) and column definitions cannot be used together".to_owned(),
388 "Remove the wildcard or column definitions".to_owned(),
389 )
390 .into());
391 }
392
393 if !wildcard_idx.is_some()
395 && !constraints.iter().any(|c| {
396 matches!(
397 c,
398 TableConstraint::Unique {
399 is_primary: true,
400 ..
401 }
402 )
403 })
404 && !column_defs.iter().any(|col| {
405 col.options
406 .iter()
407 .any(|opt| matches!(opt.option, ColumnOption::Unique { is_primary: true }))
408 })
409 {
410 return Err(ErrorCode::NotSupported(
411 "CDC table without primary key constraint is not supported".to_owned(),
412 "Please define a primary key".to_owned(),
413 )
414 .into());
415 }
416
417 if append_only {
418 return Err(ErrorCode::NotSupported(
419 "append only modifier on the table created from a CDC source".into(),
420 "Remove the APPEND ONLY clause".into(),
421 )
422 .into());
423 }
424
425 if !source_watermarks.is_empty()
426 && source_watermarks
427 .iter()
428 .any(|watermark| !watermark.with_ttl)
429 {
430 return Err(ErrorCode::NotSupported(
431 "non-TTL watermark defined on the table created from a CDC source".into(),
432 "Use `WATERMARK ... WITH TTL` instead.".into(),
433 )
434 .into());
435 }
436
437 Ok(())
438}
439
440pub(crate) async fn bind_cdc_table_schema_externally(
442 cdc_with_options: WithOptionsSecResolved,
443) -> Result<(Vec<ColumnCatalog>, Vec<String>, Vec<CdcKeyComparison>)> {
444 let (options, secret_refs) = cdc_with_options.into_parts();
445 let config = ExternalTableConfig::try_from_btreemap(options, secret_refs)
446 .context("failed to extract external table config")?;
447
448 let table = ExternalTableImpl::connect(config)
449 .await
450 .context("failed to auto derive table schema")?;
451
452 let pk_names = table.pk_names().clone();
453 let pk_comparisons = table.pk_column_comparisons(&pk_names)?;
454 Ok((
455 table
456 .column_descs()
457 .iter()
458 .cloned()
459 .map(|column_desc| ColumnCatalog {
460 column_desc,
461 is_hidden: false,
462 })
463 .collect(),
464 pk_names,
465 pk_comparisons,
466 ))
467}
468
469pub(crate) async fn bind_cdc_pk_comparisons_externally(
470 cdc_with_options: WithOptionsSecResolved,
471 pk_names: &[String],
472) -> Result<Vec<CdcKeyComparison>> {
473 let (options, secret_refs) = cdc_with_options.into_parts();
476 let config = ExternalTableConfig::try_from_btreemap(options, secret_refs)
477 .context("failed to extract external table config")?;
478 Ok(ExternalTableImpl::discover_pk_column_comparisons(&config, pk_names).await?)
479}
480
481pub(crate) fn bind_cdc_table_schema(
483 column_defs: &Vec<ColumnDef>,
484 constraints: &Vec<TableConstraint>,
485 is_for_replace_plan: bool,
486) -> Result<(Vec<ColumnCatalog>, Vec<String>)> {
487 let columns = bind_sql_columns(column_defs, is_for_replace_plan)?;
488 reject_variant_columns(&columns, "on a table created from a CDC source")?;
490
491 let pk_names = bind_sql_pk_names(column_defs, bind_table_constraints(constraints)?)?;
492 Ok((columns, pk_names))
493}
494
495#[cfg(test)]
496mod tests {
497 use super::*;
498
499 fn test_schema_table_name() -> SchemaTableName {
500 SchemaTableName {
501 schema_name: "public".to_owned(),
502 table_name: "orders".to_owned(),
503 }
504 }
505
506 fn pk_names() -> Vec<String> {
507 vec!["plan_id".to_owned(), "site_id".to_owned()]
508 }
509
510 #[test]
511 fn test_debezium_filter_rejects_literal_excluded_pk() {
512 let err = reject_pk_filtered_by_debezium_column_filter_inner(
513 &pk_names(),
514 &test_schema_table_name(),
515 Some("public.orders.site_id"),
516 None,
517 )
518 .unwrap_err();
519
520 assert!(err.to_report_string().contains("site_id"));
521 assert!(
522 err.to_report_string()
523 .contains("debezium.column.exclude.list")
524 );
525 }
526
527 #[test]
528 fn test_debezium_filter_rejects_include_list_missing_pk() {
529 let err = reject_pk_filtered_by_debezium_column_filter_inner(
530 &pk_names(),
531 &test_schema_table_name(),
532 None,
533 Some("public.orders.plan_id,public.orders.payload"),
534 )
535 .unwrap_err();
536
537 assert!(err.to_report_string().contains("site_id"));
538 assert!(
539 err.to_report_string()
540 .contains("debezium.column.include.list")
541 );
542 }
543
544 #[test]
545 fn test_debezium_filter_accepts_include_list_covering_all_pks() {
546 reject_pk_filtered_by_debezium_column_filter_inner(
547 &pk_names(),
548 &test_schema_table_name(),
549 None,
550 Some("public.orders.plan_id,public.orders.site_id,public.orders.payload"),
551 )
552 .unwrap();
553 }
554
555 #[test]
556 fn test_debezium_filter_matches_regex_patterns() {
557 reject_pk_filtered_by_debezium_column_filter_inner(
558 &pk_names(),
559 &test_schema_table_name(),
560 None,
561 Some(r"public[.]orders[.](plan_id|site_id),public[.]orders[.]payload"),
562 )
563 .unwrap();
564
565 let err = reject_pk_filtered_by_debezium_column_filter_inner(
566 &pk_names(),
567 &test_schema_table_name(),
568 Some(r".*[.]orders[.]site_id"),
569 None,
570 )
571 .unwrap_err();
572
573 assert!(err.to_report_string().contains("site_id"));
574 }
575
576 #[test]
577 fn test_debezium_filter_matches_patterns_case_insensitively() {
578 let err = reject_pk_filtered_by_debezium_column_filter_inner(
579 &pk_names(),
580 &test_schema_table_name(),
581 Some("PUBLIC.ORDERS.SITE_ID"),
582 None,
583 )
584 .unwrap_err();
585
586 assert!(err.to_report_string().contains("site_id"));
587
588 reject_pk_filtered_by_debezium_column_filter_inner(
589 &pk_names(),
590 &test_schema_table_name(),
591 None,
592 Some("Public.Orders.Plan_ID,Public.Orders.Site_ID"),
593 )
594 .unwrap();
595 }
596
597 #[test]
598 fn test_parse_postgres_cdc_external_table_name() {
599 for (input, expected) in [
600 ("public.Note", ("public", "Note")),
601 ("public.\"Note\"", ("public", "Note")),
602 (
603 "\"Mixed.Schema\".\"Note.Table\"",
604 ("Mixed.Schema", "Note.Table"),
605 ),
606 ("public.\"Note\"\"Archive\"", ("public", "Note\"Archive")),
607 ] {
608 assert_eq!(
609 parse_postgres_cdc_external_table_name(input).unwrap(),
610 (expected.0.to_owned(), expected.1.to_owned()),
611 "input: {input}"
612 );
613 }
614
615 for input in [
616 "Note",
617 "public.",
618 ".Note",
619 "public.\"Note",
620 "public.\"Note\"Archive",
621 "public.Note.Archive",
622 ] {
623 assert!(
624 parse_postgres_cdc_external_table_name(input).is_err(),
625 "input should be rejected: {input}"
626 );
627 }
628 }
629}