1use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
16use std::sync::{Arc, LazyLock};
17
18use anyhow::Context;
19use either::Either;
20use iceberg::arrow::type_to_arrow_type;
21use iceberg::spec::Transform;
22use itertools::Itertools;
23use maplit::{convert_args, hashmap, hashset};
24use pgwire::pg_response::{PgResponse, StatementType};
25use risingwave_common::array::arrow::IcebergArrowConvert;
26use risingwave_common::array::arrow::arrow_schema_iceberg::DataType as ArrowDataType;
27use risingwave_common::bail;
28use risingwave_common::catalog::{
29 ColumnCatalog, ICEBERG_SINK_PREFIX, ObjectId, RISINGWAVE_ICEBERG_ROW_ID, ROW_ID_COLUMN_NAME,
30 Schema,
31};
32use risingwave_common::license::Feature;
33use risingwave_common::secret::LocalSecretManager;
34use risingwave_common::system_param::reader::SystemParamsRead;
35use risingwave_common::types::{DataType, Timestamptz};
36use risingwave_common::util::epoch::Epoch;
37use risingwave_connector::sink::catalog::{SinkCatalog, SinkFormatDesc};
38use risingwave_connector::sink::file_sink::s3::SnowflakeSink;
39use risingwave_connector::sink::iceberg::{ICEBERG_SINK, IcebergConfig};
40use risingwave_connector::sink::kafka::KAFKA_SINK;
41use risingwave_connector::sink::snowflake_redshift::redshift::RedshiftSink;
42use risingwave_connector::sink::snowflake_redshift::snowflake::SnowflakeV2Sink;
43use risingwave_connector::sink::{
44 CONNECTOR_TYPE_KEY, SINK_SNAPSHOT_OPTION, SINK_TYPE_OPTION, SINK_USER_FORCE_APPEND_ONLY_OPTION,
45 SINK_USER_IGNORE_DELETE_OPTION, Sink, enforce_secret_sink,
46};
47use risingwave_connector::{
48 AUTO_SCHEMA_CHANGE_KEY, SINK_CREATE_TABLE_IF_NOT_EXISTS_KEY, SINK_INTERMEDIATE_TABLE_NAME,
49 SINK_TARGET_TABLE_NAME, WithPropertiesExt,
50};
51use risingwave_pb::catalog::connection_params::PbConnectionType;
52use risingwave_pb::telemetry::TelemetryDatabaseObject;
53use risingwave_sqlparser::ast::{
54 CreateSink, CreateSinkStatement, EmitMode, Encode, ExplainOptions, Format, FormatEncodeOptions,
55 ObjectName, Query, Statement,
56};
57use risingwave_sqlparser::parser::Parser;
58
59use super::RwPgResponse;
60use super::create_mv::get_column_names;
61use super::create_source::UPSTREAM_SOURCE_KEY;
62use super::util::gen_query_from_table_name;
63use crate::binder::{Binder, Relation};
64use crate::catalog::root_catalog::SchemaPath;
65use crate::catalog::table_catalog::TableType;
66use crate::error::{ErrorCode, Result, RwError};
67use crate::expr::{ExprImpl, InputRef, rewrite_now_to_proctime};
68use crate::handler::HandlerArgs;
69use crate::handler::alter_table_column::fetch_table_catalog_for_alter;
70use crate::handler::create_mv::{
71 extract_streaming_job_resource_options, parse_column_names, resolve_streaming_job_resource_type,
72};
73use crate::handler::util::{
74 LongRunningNotificationAction, check_connector_match_connection_type,
75 ensure_connection_type_allowed, ensure_local_fs_connector_allowed,
76 execute_with_long_running_notification, get_table_catalog_by_table_name,
77 reject_internal_table_dependencies,
78};
79use crate::optimizer::backfill_order_strategy::plan_backfill_order;
80use crate::optimizer::plan_node::{
81 IcebergPartitionInfo, LogicalSource, PartitionComputeInfo, StreamPlanRef as PlanRef,
82 StreamProject, ensure_sync_log_store_fragment_root, generic,
83};
84use crate::optimizer::{OptimizerContext, RelationCollectorVisitor};
85use crate::scheduler::streaming_manager::CreatingStreamingJobInfo;
86use crate::session::SessionImpl;
87use crate::session::current::notice_to_user;
88use crate::stream_fragmenter::{GraphJobType, build_graph_with_strategy};
89use crate::utils::{resolve_connection_ref_and_secret_ref, resolve_privatelink_in_with_option};
90use crate::{Explain, Planner, TableCatalog, WithOptions, WithOptionsSecResolved};
91
92static SINK_ALLOWED_CONNECTION_CONNECTOR: LazyLock<HashSet<PbConnectionType>> =
93 LazyLock::new(|| {
94 hashset! {
95 PbConnectionType::Unspecified,
96 PbConnectionType::Kafka,
97 PbConnectionType::Iceberg,
98 PbConnectionType::Elasticsearch,
99 }
100 });
101
102static SINK_ALLOWED_CONNECTION_SCHEMA_REGISTRY: LazyLock<HashSet<PbConnectionType>> =
103 LazyLock::new(|| {
104 hashset! {
105 PbConnectionType::Unspecified,
106 PbConnectionType::SchemaRegistry,
107 }
108 });
109
110const SINK_SINCE_TIMESTAMP_OPTION: &str = "since_timestamp";
111
112pub struct SinkPlanContext {
114 pub query: Box<Query>,
115 pub sink_plan: PlanRef,
116 pub sink_catalog: SinkCatalog,
117 pub target_table_catalog: Option<Arc<TableCatalog>>,
118 pub dependencies: HashSet<ObjectId>,
119 pub since_timestamp_epoch: Option<u64>,
120}
121
122pub async fn gen_sink_plan(
123 handler_args: HandlerArgs,
124 stmt: CreateSinkStatement,
125 explain_options: Option<ExplainOptions>,
126 is_iceberg_engine_internal: bool,
127) -> Result<SinkPlanContext> {
128 let session = handler_args.session.clone();
129 let session = session.as_ref();
130 let user_specified_columns = !stmt.columns.is_empty();
131 let db_name = &session.database();
132 let (sink_schema_name, sink_table_name) =
133 Binder::resolve_schema_qualified_name(db_name, &stmt.sink_name)?;
134
135 let mut with_options = handler_args.with_options.clone();
136 extract_streaming_job_resource_options(&mut with_options);
139
140 if session
141 .env()
142 .system_params_manager()
143 .get_params()
144 .load()
145 .enforce_secret()
146 && Feature::SecretManagement.check_available().is_ok()
147 {
148 enforce_secret_sink(&with_options)?;
149 }
150
151 resolve_privatelink_in_with_option(&mut with_options)?;
152 let (mut resolved_with_options, connection_type, connector_conn_ref) =
153 resolve_connection_ref_and_secret_ref(
154 with_options,
155 session,
156 Some(TelemetryDatabaseObject::Sink),
157 )?;
158
159 let since_timestamp_epoch = resolved_with_options
160 .remove(SINK_SINCE_TIMESTAMP_OPTION)
161 .map(|value| {
162 let timestamp = value.parse::<Timestamptz>().map_err(|err| {
163 ErrorCode::InvalidInputSyntax(format!(
164 "invalid value {value:?} of '{SINK_SINCE_TIMESTAMP_OPTION}' option: {err}; \
165 expected a timestamptz string with an explicit time zone, \
166 for example '2024-01-01 00:00:00Z'"
167 ))
168 })?;
169 let timestamp_millis = u64::try_from(timestamp.timestamp_millis()).unwrap_or(0);
170 Ok::<_, RwError>(Epoch::from_unix_millis_or_earliest(timestamp_millis).0)
171 })
172 .transpose()?;
173 if since_timestamp_epoch.is_some() {
174 Feature::SinkSinceTimestamp.check_available()?;
175 }
176
177 ensure_connection_type_allowed(connection_type, &SINK_ALLOWED_CONNECTION_CONNECTOR)?;
178
179 if !matches!(connection_type, PbConnectionType::Unspecified) {
181 let Some(connector) = resolved_with_options.get_connector() else {
182 return Err(RwError::from(ErrorCode::ProtocolError(format!(
183 "missing field '{}' in WITH clause",
184 CONNECTOR_TYPE_KEY
185 ))));
186 };
187 check_connector_match_connection_type(connector.as_str(), &connection_type)?;
188 }
189
190 let partition_info = get_partition_compute_info(&resolved_with_options).await?;
191
192 let context = if let Some(explain_options) = explain_options {
193 OptimizerContext::new(handler_args.clone(), explain_options)
194 } else {
195 OptimizerContext::from_handler_args(handler_args.clone())
196 };
197 let is_auto_schema_change = resolved_with_options
198 .get(AUTO_SCHEMA_CHANGE_KEY)
199 .map(|value| {
200 value.parse::<bool>().map_err(|_| {
201 ErrorCode::InvalidInputSyntax(format!(
202 "invalid value {} of '{}' option, expect",
203 value, AUTO_SCHEMA_CHANGE_KEY
204 ))
205 })
206 })
207 .transpose()?
208 .unwrap_or(false);
209
210 if is_auto_schema_change && !is_iceberg_engine_internal {
211 Feature::SinkAutoSchemaChange.check_available()?;
212 }
213
214 let sink_into_table_name = stmt.into_table_name.as_ref().map(|name| name.real_value());
215 if sink_into_table_name.is_some() {
216 let prev = resolved_with_options.insert(CONNECTOR_TYPE_KEY.to_owned(), "table".to_owned());
217
218 if prev.is_some() {
219 return Err(RwError::from(ErrorCode::BindError(
220 "In the case of sinking into table, the 'connector' parameter should not be provided.".to_owned(),
221 )));
222 }
223 }
224 let connector = resolved_with_options
225 .get(CONNECTOR_TYPE_KEY)
226 .cloned()
227 .ok_or_else(|| ErrorCode::BindError(format!("missing field '{CONNECTOR_TYPE_KEY}'")))?;
228 ensure_local_fs_connector_allowed(session, &connector)?;
229
230 let sink_from_table_name;
232 let direct_sink_from_name: Option<(ObjectName, bool)>;
235 let mut query = match stmt.sink_from {
236 CreateSink::From(from_name) => {
237 sink_from_table_name = from_name.0.last().unwrap().real_value();
238 direct_sink_from_name = Some((from_name.clone(), is_auto_schema_change));
239 if is_auto_schema_change && sink_into_table_name.is_some() {
240 return Err(RwError::from(ErrorCode::InvalidInputSyntax(
241 "auto schema change not supported for sink-into-table".to_owned(),
242 )));
243 }
244 if resolved_with_options
245 .value_eq_ignore_case(SINK_CREATE_TABLE_IF_NOT_EXISTS_KEY, "true")
246 && connector == RedshiftSink::SINK_NAME
247 || connector == SnowflakeV2Sink::SINK_NAME
248 {
249 if let Some(table_name) = resolved_with_options.get(SINK_TARGET_TABLE_NAME) {
250 if resolved_with_options
252 .get(SINK_INTERMEDIATE_TABLE_NAME)
253 .is_none()
254 {
255 let intermediate_table_name = format!(
257 "rw_{}_{}_{}",
258 sink_table_name,
259 table_name,
260 uuid::Uuid::new_v4()
261 );
262 resolved_with_options.insert(
263 SINK_INTERMEDIATE_TABLE_NAME.to_owned(),
264 intermediate_table_name,
265 );
266 }
267 } else {
268 return Err(RwError::from(ErrorCode::BindError(
269 "'table.name' option must be specified.".to_owned(),
270 )));
271 }
272 }
273 Box::new(gen_query_from_table_name(from_name))
274 }
275 CreateSink::AsQuery(query) => {
276 if is_auto_schema_change {
277 return Err(RwError::from(ErrorCode::InvalidInputSyntax(
278 "auto schema change not supported for CREATE SINK AS QUERY".to_owned(),
279 )));
280 }
281 sink_from_table_name = sink_table_name.clone();
282 direct_sink_from_name = None;
283 query
284 }
285 };
286
287 if is_iceberg_engine_internal && let Some((from_name, _)) = &direct_sink_from_name {
288 let (table, _) = get_table_catalog_by_table_name(session, from_name)?;
289 let pk_names = table.pk_column_names();
290 if pk_names.len() == 1 && pk_names[0].eq(ROW_ID_COLUMN_NAME) {
291 let [stmt]: [_; 1] = Parser::parse_sql(&format!(
292 "select {} as {}, * from {}",
293 ROW_ID_COLUMN_NAME, RISINGWAVE_ICEBERG_ROW_ID, from_name
294 ))
295 .context("unable to parse query")?
296 .try_into()
297 .unwrap();
298 let Statement::Query(parsed_query) = stmt else {
299 panic!("unexpected statement: {:?}", stmt);
300 };
301 query = parsed_query;
302 }
303 }
304
305 let (sink_database_id, sink_schema_id) =
306 session.get_database_and_schema_id_for_create(sink_schema_name.clone())?;
307
308 if since_timestamp_epoch.is_some() {
309 if sink_into_table_name.is_some() {
310 return Err(ErrorCode::BindError(format!(
311 "`{SINK_SINCE_TIMESTAMP_OPTION}` does not support `CREATE SINK INTO TABLE`"
312 ))
313 .into());
314 }
315 if is_iceberg_engine_internal {
316 return Err(ErrorCode::BindError(format!(
317 "`{SINK_SINCE_TIMESTAMP_OPTION}` does not support iceberg engine internal sinks"
318 ))
319 .into());
320 }
321 if let Some((from_name, _)) = &direct_sink_from_name {
322 let (table, _) = get_table_catalog_by_table_name(session, from_name)?;
323 if table.database_id != sink_database_id {
324 return Err(ErrorCode::NotSupported(
325 format!(
326 "`{SINK_SINCE_TIMESTAMP_OPTION}` does not support cross-database sinks"
327 ),
328 "Please create the sink in the same database as the upstream table.".to_owned(),
329 )
330 .into());
331 }
332 }
333 }
334
335 let (
336 dependent_relations,
337 dependent_udfs,
338 dependent_secrets,
339 bound,
340 auto_refresh_schema_from_table,
341 ) = {
342 let mut binder = Binder::new_for_stream(session);
343 let auto_refresh_schema_from_table = if let Some((from_name, true)) = &direct_sink_from_name
344 {
345 let from_relation = binder.bind_relation_by_name(from_name, None, None, true)?;
346 if let Relation::BaseTable(table) = from_relation {
347 if table.table_catalog.table_type != TableType::Table {
348 return Err(ErrorCode::InvalidInputSyntax(format!(
349 "auto schema change is supported only on TABLE, but got {:?}",
350 table.table_catalog.table_type
351 ))
352 .into());
353 }
354 if table.table_catalog.database_id != sink_database_id {
355 return Err(ErrorCode::InvalidInputSyntax(
356 "auto schema change sinks do not support cross-database tables".to_owned(),
357 )
358 .into());
359 }
360 for col in &table.table_catalog.columns {
361 if !col.is_hidden() && (col.is_generated() || col.is_rw_sys_column()) {
362 return Err(ErrorCode::InvalidInputSyntax(format!(
363 "auto schema change is not supported for tables with visible generated columns or visible system columns, but found column {}",
364 col.name()
365 ))
366 .into());
367 }
368 }
369 Some(table.table_catalog)
370 } else {
371 return Err(RwError::from(ErrorCode::NotSupported(
372 "auto schema change only supported for TABLE".to_owned(),
373 "try recreating the sink from table".to_owned(),
374 )));
375 }
376 } else {
377 None
378 };
379
380 let bound = binder.bind_query(&query)?;
381
382 (
383 binder.included_relations().clone(),
384 binder.included_udfs().clone(),
385 binder.included_secrets().clone(),
386 bound,
387 auto_refresh_schema_from_table,
388 )
389 };
390
391 reject_internal_table_dependencies(session, &dependent_relations, "CREATE SINK")?;
392
393 let col_names = if sink_into_table_name.is_some() {
394 parse_column_names(&stmt.columns)
395 } else {
396 get_column_names(&bound, stmt.columns)?
398 };
399
400 let emit_on_window_close = stmt.emit_mode == Some(EmitMode::OnWindowClose);
401 if emit_on_window_close {
402 context.warn_to_user("EMIT ON WINDOW CLOSE is currently an experimental feature. Please use it with caution.");
403 }
404
405 let format_desc = match stmt.sink_schema {
406 Some(f) => {
408 validate_compatibility(&connector, &f)?;
409 Some(bind_sink_format_desc(session,f)?)
410 }
411 None => match resolved_with_options.get(SINK_TYPE_OPTION) {
412 Some(t) => SinkFormatDesc::from_legacy_type(&connector, t)?.map(|mut f| {
414 session.notice_to_user("Consider using the newer syntax `FORMAT ... ENCODE ...` instead of `type = '...'`.");
415 if let Some(v) = resolved_with_options.get(SINK_USER_FORCE_APPEND_ONLY_OPTION) {
416 f.options.insert(SINK_USER_FORCE_APPEND_ONLY_OPTION.into(), v.into());
417 }
418 if let Some(v) = resolved_with_options.get(SINK_USER_IGNORE_DELETE_OPTION) {
419 f.options.insert(SINK_USER_IGNORE_DELETE_OPTION.into(), v.into());
420 }
421 f
422 }),
423 None => None,
425 },
426 };
427
428 let definition = context.normalized_sql().to_owned();
429 let mut plan_root = if is_iceberg_engine_internal {
430 Planner::new_for_iceberg_table_engine_sink(context.into()).plan_query(bound)?
431 } else {
432 Planner::new_for_stream(context.into()).plan_query(bound)?
433 };
434 if let Some(col_names) = &col_names {
435 plan_root.set_out_names(col_names.clone())?;
436 };
437
438 let without_snapshot = matches!(
439 resolved_with_options.remove(SINK_SNAPSHOT_OPTION),
440 Some(flag) if flag.eq_ignore_ascii_case("false")
441 );
442
443 if since_timestamp_epoch.is_some() && !without_snapshot {
444 return Err(ErrorCode::BindError(format!(
445 "`{SINK_SINCE_TIMESTAMP_OPTION}` requires `snapshot = false`"
446 ))
447 .into());
448 }
449
450 let target_table_catalog = stmt
451 .into_table_name
452 .as_ref()
453 .map(|table_name| fetch_table_catalog_for_alter(session, table_name).map(|t| t.0))
454 .transpose()?;
455
456 if let Some(target_table_catalog) = &target_table_catalog {
457 if let Some(col_names) = col_names {
458 let target_table_columns = target_table_catalog
459 .columns()
460 .iter()
461 .map(|c| c.name())
462 .collect::<BTreeSet<_>>();
463 for c in col_names {
464 if !target_table_columns.contains(c.as_str()) {
465 return Err(RwError::from(ErrorCode::BindError(format!(
466 "Column {} not found in table {}",
467 c,
468 target_table_catalog.name()
469 ))));
470 }
471 }
472 }
473 if target_table_catalog
474 .columns()
475 .iter()
476 .any(|col| !col.nullable())
477 {
478 notice_to_user(format!(
479 "The target table `{}` contains NOT NULL columns. Rows written by the sink that violate those constraints will be ignored silently.",
480 target_table_catalog.name(),
481 ));
482 }
483 }
484
485 let sink_plan = plan_root.gen_sink_plan(
486 sink_table_name,
487 definition,
488 resolved_with_options,
489 emit_on_window_close,
490 db_name.to_owned(),
491 sink_from_table_name,
492 format_desc,
493 without_snapshot,
494 since_timestamp_epoch.is_some(),
495 is_iceberg_engine_internal,
496 target_table_catalog.clone(),
497 partition_info,
498 user_specified_columns,
499 auto_refresh_schema_from_table,
500 )?;
501
502 let sink_desc = sink_plan.sink_desc().clone();
503
504 let mut sink_plan: PlanRef = sink_plan.into_stream_plan()?;
505 sink_plan = ensure_sync_log_store_fragment_root(sink_plan);
506
507 let ctx = sink_plan.ctx();
508 let explain_trace = ctx.is_explain_trace();
509 if explain_trace {
510 ctx.trace("Create Sink:");
511 ctx.trace(sink_plan.explain_to_string());
512 }
513 tracing::trace!("sink_plan: {:?}", sink_plan.explain_to_string());
514
515 let dependencies =
518 RelationCollectorVisitor::collect_with(dependent_relations, sink_plan.clone())
519 .into_iter()
520 .chain(dependent_udfs.iter().copied().map_into())
521 .chain(
522 dependent_secrets
523 .iter()
524 .copied()
525 .map(|id| id.as_object_id()),
526 )
527 .collect();
528
529 let sink_catalog = sink_desc.into_catalog(
530 sink_schema_id,
531 sink_database_id,
532 session.user_id(),
533 connector_conn_ref,
534 );
535
536 if let Some(table_catalog) = &target_table_catalog {
537 for column in sink_catalog.full_columns() {
538 if !column.can_dml() {
539 unreachable!(
540 "cannot derive generated columns or the `_rw_timestamp` system column in a sink catalog, but found one"
541 );
542 }
543 }
544
545 let table_columns_without_rw_timestamp = table_catalog.columns_without_rw_timestamp();
546 let exprs = derive_default_column_project_for_sink(
547 &sink_catalog,
548 sink_plan.schema(),
549 &table_columns_without_rw_timestamp,
550 user_specified_columns,
551 )?;
552
553 let logical_project = generic::Project::new(exprs, sink_plan);
554
555 sink_plan = StreamProject::new(logical_project).into();
556
557 let exprs = LogicalSource::derive_output_exprs_from_generated_columns(
558 &table_columns_without_rw_timestamp,
559 )?;
560
561 if let Some(exprs) = exprs {
562 let logical_project = generic::Project::new(exprs, sink_plan);
563 sink_plan = StreamProject::new(logical_project).into();
564 }
565 };
566
567 Ok(SinkPlanContext {
568 query,
569 sink_plan,
570 sink_catalog,
571 target_table_catalog,
572 dependencies,
573 since_timestamp_epoch,
574 })
575}
576
577pub async fn get_partition_compute_info(
582 with_options: &WithOptionsSecResolved,
583) -> Result<Option<PartitionComputeInfo>> {
584 let (options, secret_refs) = with_options.clone().into_parts();
585 let Some(connector) = options.get(UPSTREAM_SOURCE_KEY).cloned() else {
586 return Ok(None);
587 };
588 let properties = LocalSecretManager::global().fill_secrets(options, secret_refs)?;
589 match connector.as_str() {
590 ICEBERG_SINK => {
591 let iceberg_config = IcebergConfig::from_btreemap(properties)?;
592 get_partition_compute_info_for_iceberg(&iceberg_config).await
593 }
594 _ => Ok(None),
595 }
596}
597
598async fn get_partition_compute_info_for_iceberg(
599 _iceberg_config: &IcebergConfig,
600) -> Result<Option<PartitionComputeInfo>> {
601 if _iceberg_config.create_table_if_not_exists {
603 return Ok(None);
604 }
605 let table = _iceberg_config.load_table().await?;
606 let partition_spec = table.metadata().default_partition_spec();
607 if partition_spec.is_unpartitioned() {
608 return Ok(None);
609 }
610
611 let has_sparse_partition = partition_spec.fields().iter().any(|f| match f.transform {
616 Transform::Identity | Transform::Truncate(_) | Transform::Bucket(_) => true,
618 Transform::Year
620 | Transform::Month
621 | Transform::Day
622 | Transform::Hour
623 | Transform::Void
624 | Transform::Unknown => false,
625 });
626 if !has_sparse_partition {
627 return Ok(None);
628 }
629
630 let arrow_type = type_to_arrow_type(&iceberg::spec::Type::Struct(
631 table.metadata().default_partition_type().clone(),
632 ))
633 .map_err(|_| {
634 RwError::from(ErrorCode::SinkError(
635 "Failed to convert the Iceberg partition type to an Arrow type".into(),
636 ))
637 })?;
638 let ArrowDataType::Struct(struct_fields) = arrow_type else {
639 return Err(RwError::from(ErrorCode::SinkError(
640 "The Iceberg partition type must be a struct type".into(),
641 )));
642 };
643
644 let schema = table.metadata().current_schema();
645 let partition_fields = partition_spec
646 .fields()
647 .iter()
648 .map(|f| {
649 let source_f =
650 schema
651 .field_by_id(f.source_id)
652 .ok_or(RwError::from(ErrorCode::SinkError(
653 "Failed to look up the Iceberg partition field".into(),
654 )))?;
655 Ok((source_f.name.clone(), f.transform))
656 })
657 .collect::<Result<Vec<_>>>()?;
658
659 Ok(Some(PartitionComputeInfo::Iceberg(IcebergPartitionInfo {
660 partition_type: IcebergArrowConvert.struct_from_fields(&struct_fields)?,
661 partition_fields,
662 })))
663}
664
665pub async fn handle_create_sink(
666 mut handle_args: HandlerArgs,
667 stmt: CreateSinkStatement,
668 is_iceberg_engine_internal: bool,
669) -> Result<RwPgResponse> {
670 let session = handle_args.session.clone();
671
672 session.check_cluster_limits().await?;
673
674 let mode = if stmt.or_replace {
675 prepare_replace_sink(&mut handle_args, &stmt)?
676 } else {
677 let if_not_exists = stmt.if_not_exists;
678 if let Either::Right(resp) = session.check_relation_name_duplicated(
679 stmt.sink_name.clone(),
680 StatementType::CREATE_SINK,
681 if_not_exists,
682 )? {
683 return Ok(resp);
684 }
685
686 if stmt.sink_name.base_name().starts_with(ICEBERG_SINK_PREFIX) {
687 return Err(RwError::from(ErrorCode::InvalidInputSyntax(format!(
688 "Sink name cannot start with reserved prefix '{}'",
689 ICEBERG_SINK_PREFIX
690 ))));
691 }
692
693 SinkCreateMode::Create { if_not_exists }
694 };
695
696 create_sink_or_replace(handle_args, stmt, is_iceberg_engine_internal, mode).await
697}
698
699enum SinkCreateMode {
700 Create { if_not_exists: bool },
701 Replace { original_sink: Arc<SinkCatalog> },
702}
703
704impl SinkCreateMode {
705 fn statement_name(&self) -> &'static str {
706 match self {
707 SinkCreateMode::Create { .. } => "CREATE SINK",
708 SinkCreateMode::Replace { .. } => "REPLACE SINK",
709 }
710 }
711}
712
713async fn create_sink_or_replace(
714 mut handle_args: HandlerArgs,
715 stmt: CreateSinkStatement,
716 is_iceberg_engine_internal: bool,
717 mode: SinkCreateMode,
718) -> Result<RwPgResponse> {
719 let session = handle_args.session.clone();
720
721 let resource_type =
722 resolve_streaming_job_resource_type(session.as_ref(), &mut handle_args.with_options)?;
723
724 let (sink, graph, dependencies, since_timestamp_epoch) = {
725 let backfill_order_strategy = handle_args.with_options.backfill_order_strategy();
726 let SinkPlanContext {
727 query,
728 sink_plan: plan,
729 sink_catalog: mut sink,
730 target_table_catalog,
731 dependencies,
732 since_timestamp_epoch,
733 } = gen_sink_plan(handle_args, stmt, None, is_iceberg_engine_internal).await?;
734
735 let has_order_by = !query.order_by.is_empty();
736 if has_order_by {
737 plan.ctx().warn_to_user(
738 r#"The ORDER BY clause in the CREATE SINK statement has no effect at all."#
739 .to_owned(),
740 );
741 }
742
743 match &mode {
744 SinkCreateMode::Create { .. } => {
745 if let Some(table_catalog) = &target_table_catalog {
746 sink.original_target_columns = table_catalog.columns_without_rw_timestamp();
747 }
748 }
749 SinkCreateMode::Replace { original_sink } => {
750 if target_table_catalog.is_some() {
751 return Err(ErrorCode::NotSupported(
752 "REPLACE SINK INTO TABLE is not supported yet".to_owned(),
753 "replace ordinary sinks first".to_owned(),
754 )
755 .into());
756 }
757
758 sink.schema_id = original_sink.schema_id;
759 sink.database_id = original_sink.database_id;
760 sink.name = original_sink.name.clone();
761 sink.owner = original_sink.owner;
762 }
763 }
764
765 let backfill_order =
766 plan_backfill_order(session.as_ref(), backfill_order_strategy, plan.clone())?;
767 let graph =
768 build_graph_with_strategy(plan, Some(GraphJobType::Sink), Some(backfill_order))?;
769
770 (sink, graph, dependencies, since_timestamp_epoch)
771 };
772
773 let statement_name = mode.statement_name();
774 let catalog_writer = session.catalog_writer()?;
775 match mode {
776 SinkCreateMode::Create { if_not_exists } => {
777 let _job_guard = session.env().creating_streaming_job_tracker().guard(
778 CreatingStreamingJobInfo::new(
779 session.session_id(),
780 sink.database_id,
781 sink.schema_id,
782 sink.name.clone(),
783 ),
784 );
785
786 execute_with_long_running_notification(
787 catalog_writer.create_sink(
788 sink.to_proto(),
789 graph,
790 dependencies,
791 resource_type,
792 if_not_exists,
793 since_timestamp_epoch,
794 ),
795 &session,
796 statement_name,
797 LongRunningNotificationAction::MonitorBackfillJob,
798 )
799 .await?;
800 }
801 SinkCreateMode::Replace { original_sink } => {
802 let original_sink_id = original_sink.id;
803 execute_with_long_running_notification(
804 catalog_writer.replace_sink(
805 original_sink_id,
806 sink.to_proto(),
807 graph,
808 dependencies,
809 resource_type,
810 ),
811 &session,
812 statement_name,
813 LongRunningNotificationAction::DiagnoseBarrierLatency,
814 )
815 .await?;
816
817 tracing::info!(
818 old_sink_id = %original_sink_id,
819 sink_name = %sink.name,
820 "replace sink plan submitted"
821 );
822 }
823 }
824
825 Ok(PgResponse::empty_result(StatementType::CREATE_SINK))
826}
827
828fn sink_replace_requires_exactly_once_state(sink: &SinkCatalog) -> bool {
829 match sink.properties.get("is_exactly_once") {
830 Some(value) => value.eq_ignore_ascii_case("true"),
831 None => sink
832 .properties
833 .get(CONNECTOR_TYPE_KEY)
834 .is_some_and(|connector| connector.eq_ignore_ascii_case(ICEBERG_SINK)),
835 }
836}
837
838fn prepare_replace_sink(
839 handle_args: &mut HandlerArgs,
840 stmt: &CreateSinkStatement,
841) -> Result<SinkCreateMode> {
842 let session = handle_args.session.clone();
843 if stmt.if_not_exists {
844 return Err(ErrorCode::InvalidInputSyntax(
845 "REPLACE SINK does not support IF NOT EXISTS".to_owned(),
846 )
847 .into());
848 }
849 if !matches!(&stmt.sink_from, CreateSink::From(_)) {
850 return Err(ErrorCode::NotSupported(
851 "REPLACE SINK currently only supports REPLACE SINK ... FROM table_or_mv".to_owned(),
852 "use REPLACE SINK name FROM existing_relation ...".to_owned(),
853 )
854 .into());
855 }
856 if stmt.into_table_name.is_some() {
857 return Err(ErrorCode::NotSupported(
858 "REPLACE SINK INTO TABLE is not supported yet".to_owned(),
859 "replace ordinary sinks first".to_owned(),
860 )
861 .into());
862 }
863 if handle_args
864 .with_options
865 .get(AUTO_SCHEMA_CHANGE_KEY)
866 .is_some_and(|value| value.eq_ignore_ascii_case("true"))
867 {
868 return Err(ErrorCode::NotSupported(
869 "REPLACE SINK with auto schema change is not supported yet".to_owned(),
870 "disable auto schema change for this replacement".to_owned(),
871 )
872 .into());
873 }
874 if handle_args
875 .with_options
876 .contains_key(SINK_SINCE_TIMESTAMP_OPTION)
877 {
878 return Err(ErrorCode::NotSupported(
879 "REPLACE SINK with since_timestamp is not supported yet".to_owned(),
880 "create a new sink with since_timestamp instead".to_owned(),
881 )
882 .into());
883 }
884 match handle_args.with_options.get(SINK_SNAPSHOT_OPTION) {
885 Some(value) if !value.eq_ignore_ascii_case("false") => {
886 return Err(ErrorCode::InvalidInputSyntax(
887 "REPLACE SINK must not enable snapshot backfill".to_owned(),
888 )
889 .into());
890 }
891 Some(_) => {}
892 None => {
893 handle_args
894 .with_options
895 .insert(SINK_SNAPSHOT_OPTION.to_owned(), "false".to_owned());
896 }
897 }
898
899 let db_name = session.database();
900 let (sink_schema_name, sink_table_name) =
901 Binder::resolve_schema_qualified_name(&db_name, &stmt.sink_name)?;
902 let original_sink = {
903 let search_path = session.config().search_path();
904 let user_name = session.user_name();
905 let schema_path = SchemaPath::new(sink_schema_name.as_deref(), &search_path, &user_name);
906 let reader = session.env().catalog_reader().read_guard();
907 let (sink, schema_name) =
908 reader.get_created_sink_by_name(&db_name, schema_path, &sink_table_name)?;
909 session.check_privilege_for_drop_alter(schema_name, &**sink)?;
910 if sink.target_table.is_some() {
911 return Err(ErrorCode::NotSupported(
912 "REPLACE SINK INTO TABLE is not supported yet".to_owned(),
913 "replace ordinary sinks first".to_owned(),
914 )
915 .into());
916 }
917 if sink.auto_refresh_schema_from_table.is_some() {
918 return Err(ErrorCode::NotSupported(
919 "REPLACE SINK with auto schema change is not supported yet".to_owned(),
920 "drop and recreate this auto schema change sink".to_owned(),
921 )
922 .into());
923 }
924 if sink_replace_requires_exactly_once_state(sink) {
925 return Err(ErrorCode::NotSupported(
926 "REPLACE SINK does not support exactly-once sinks yet".to_owned(),
927 "set is_exactly_once=false or recreate the sink manually".to_owned(),
928 )
929 .into());
930 }
931 sink.clone()
932 };
933
934 Ok(SinkCreateMode::Replace { original_sink })
935}
936
937pub fn fetch_incoming_sinks(
938 session: &Arc<SessionImpl>,
939 table: &TableCatalog,
940) -> Result<Vec<Arc<SinkCatalog>>> {
941 let reader = session.env().catalog_reader().read_guard();
942 let schema = reader.get_schema_by_id(table.database_id, table.schema_id)?;
943 let Some(incoming_sinks) = schema.table_incoming_sinks(table.id) else {
944 return Ok(vec![]);
945 };
946 let mut sinks = vec![];
947 for sink_id in incoming_sinks {
948 sinks.push(
949 schema
950 .get_sink_by_id(*sink_id)
951 .expect("should exist")
952 .clone(),
953 );
954 }
955 Ok(sinks)
956}
957
958fn derive_sink_to_table_expr(
959 sink_schema: &Schema,
960 idx: usize,
961 target_type: &DataType,
962) -> Result<ExprImpl> {
963 let input_type = &sink_schema.fields()[idx].data_type;
964
965 if !target_type.equals_datatype(input_type) {
966 bail!(
967 "column type mismatch: {:?} vs {:?}, column name: {:?}",
968 target_type,
969 input_type,
970 sink_schema.fields()[idx].name
971 );
972 } else {
973 Ok(ExprImpl::InputRef(Box::new(InputRef::new(
974 idx,
975 input_type.clone(),
976 ))))
977 }
978}
979
980pub(crate) fn derive_default_column_project_for_sink(
981 sink: &SinkCatalog,
982 sink_schema: &Schema,
983 columns: &[ColumnCatalog],
984 user_specified_columns: bool,
985) -> Result<Vec<ExprImpl>> {
986 assert_eq!(sink.full_schema().len(), sink_schema.len());
987
988 let default_column_exprs = TableCatalog::default_column_exprs(columns);
989
990 let mut exprs = vec![];
991
992 let sink_visible_col_idxes = sink
993 .full_columns()
994 .iter()
995 .positions(|c| !c.is_hidden())
996 .collect_vec();
997 let sink_visible_col_idxes_by_name = sink
998 .full_columns()
999 .iter()
1000 .enumerate()
1001 .filter(|(_, c)| !c.is_hidden())
1002 .map(|(i, c)| (c.name(), i))
1003 .collect::<BTreeMap<_, _>>();
1004
1005 for (idx, column) in columns.iter().enumerate() {
1006 if !column.can_dml() {
1007 continue;
1008 }
1009
1010 let default_col_expr =
1011 || -> ExprImpl { rewrite_now_to_proctime(default_column_exprs[idx].clone()) };
1012
1013 let sink_col_expr = |sink_col_idx: usize| -> Result<ExprImpl> {
1014 derive_sink_to_table_expr(sink_schema, sink_col_idx, column.data_type())
1015 };
1016
1017 if user_specified_columns {
1021 if let Some(idx) = sink_visible_col_idxes_by_name.get(column.name()) {
1022 exprs.push(sink_col_expr(*idx)?);
1023 } else {
1024 exprs.push(default_col_expr());
1025 }
1026 } else {
1027 if idx < sink_visible_col_idxes.len() {
1028 exprs.push(sink_col_expr(sink_visible_col_idxes[idx])?);
1029 } else {
1030 exprs.push(default_col_expr());
1031 };
1032 }
1033 }
1034 Ok(exprs)
1035}
1036
1037fn bind_sink_format_desc(
1041 session: &SessionImpl,
1042 value: FormatEncodeOptions,
1043) -> Result<SinkFormatDesc> {
1044 use risingwave_connector::sink::catalog::{SinkEncode, SinkFormat};
1045 use risingwave_connector::sink::encoder::TimestamptzHandlingMode;
1046 use risingwave_sqlparser::ast::{Encode as E, Format as F};
1047
1048 let format = match value.format {
1049 F::Plain => SinkFormat::AppendOnly,
1050 F::Upsert => SinkFormat::Upsert,
1051 F::Debezium => SinkFormat::Debezium,
1052 f @ (F::Native | F::DebeziumMongo | F::Maxwell | F::Canal | F::None) => {
1053 return Err(ErrorCode::BindError(format!("sink format unsupported: {f}")).into());
1054 }
1055 };
1056 let encode = match value.row_encode {
1057 E::Json => SinkEncode::Json,
1058 E::Protobuf => SinkEncode::Protobuf,
1059 E::Avro => SinkEncode::Avro,
1060 E::Template => SinkEncode::Template,
1061 E::Parquet => SinkEncode::Parquet,
1062 E::Bytes => SinkEncode::Bytes,
1063 e @ (E::Native | E::Csv | E::None | E::Text) => {
1064 return Err(ErrorCode::BindError(format!("sink encode unsupported: {e}")).into());
1065 }
1066 };
1067
1068 let mut key_encode = None;
1069 if let Some(encode) = value.key_encode {
1070 match encode {
1071 E::Text => key_encode = Some(SinkEncode::Text),
1072 E::Bytes => key_encode = Some(SinkEncode::Bytes),
1073 _ => {
1074 return Err(ErrorCode::BindError(format!(
1075 "sink key encode unsupported: {encode}, only TEXT and BYTES supported"
1076 ))
1077 .into());
1078 }
1079 }
1080 }
1081
1082 let (props, connection_type_flag, schema_registry_conn_ref) =
1083 resolve_connection_ref_and_secret_ref(
1084 WithOptions::try_from(value.row_options.as_slice())?,
1085 session,
1086 Some(TelemetryDatabaseObject::Sink),
1087 )?;
1088 ensure_connection_type_allowed(
1089 connection_type_flag,
1090 &SINK_ALLOWED_CONNECTION_SCHEMA_REGISTRY,
1091 )?;
1092 let (mut options, secret_refs) = props.into_parts();
1093
1094 options
1095 .entry(TimestamptzHandlingMode::OPTION_KEY.to_owned())
1096 .or_insert(TimestamptzHandlingMode::FRONTEND_DEFAULT.to_owned());
1097
1098 Ok(SinkFormatDesc {
1099 format,
1100 encode,
1101 options,
1102 secret_refs,
1103 key_encode,
1104 connection_id: schema_registry_conn_ref,
1105 })
1106}
1107
1108static CONNECTORS_COMPATIBLE_FORMATS: LazyLock<HashMap<String, HashMap<Format, Vec<Encode>>>> =
1109 LazyLock::new(|| {
1110 use risingwave_connector::sink::Sink as _;
1111 use risingwave_connector::sink::file_sink::azblob::AzblobSink;
1112 use risingwave_connector::sink::file_sink::fs::FsSink;
1113 use risingwave_connector::sink::file_sink::gcs::GcsSink;
1114 use risingwave_connector::sink::file_sink::opendal_sink::FileSink;
1115 use risingwave_connector::sink::file_sink::s3::S3Sink;
1116 use risingwave_connector::sink::file_sink::webhdfs::WebhdfsSink;
1117 use risingwave_connector::sink::google_pubsub::GooglePubSubSink;
1118 use risingwave_connector::sink::kafka::KafkaSink;
1119 use risingwave_connector::sink::kinesis::KinesisSink;
1120 use risingwave_connector::sink::mqtt::MqttSink;
1121 use risingwave_connector::sink::pulsar::PulsarSink;
1122 use risingwave_connector::sink::redis::RedisSink;
1123
1124 convert_args!(hashmap!(
1125 GooglePubSubSink::SINK_NAME => hashmap!(
1126 Format::Plain => vec![Encode::Json],
1127 ),
1128 KafkaSink::SINK_NAME => hashmap!(
1129 Format::Plain => vec![Encode::Json, Encode::Avro, Encode::Protobuf, Encode::Bytes],
1130 Format::Upsert => vec![Encode::Json, Encode::Avro, Encode::Protobuf],
1131 Format::Debezium => vec![Encode::Json],
1132 ),
1133 FileSink::<S3Sink>::SINK_NAME => hashmap!(
1134 Format::Plain => vec![Encode::Parquet, Encode::Json],
1135 ),
1136 FileSink::<SnowflakeSink>::SINK_NAME => hashmap!(
1137 Format::Plain => vec![Encode::Parquet, Encode::Json],
1138 ),
1139 FileSink::<GcsSink>::SINK_NAME => hashmap!(
1140 Format::Plain => vec![Encode::Parquet, Encode::Json],
1141 ),
1142 FileSink::<AzblobSink>::SINK_NAME => hashmap!(
1143 Format::Plain => vec![Encode::Parquet, Encode::Json],
1144 ),
1145 FileSink::<WebhdfsSink>::SINK_NAME => hashmap!(
1146 Format::Plain => vec![Encode::Parquet, Encode::Json],
1147 ),
1148 FileSink::<FsSink>::SINK_NAME => hashmap!(
1149 Format::Plain => vec![Encode::Parquet, Encode::Json],
1150 ),
1151 KinesisSink::SINK_NAME => hashmap!(
1152 Format::Plain => vec![Encode::Json],
1153 Format::Upsert => vec![Encode::Json],
1154 Format::Debezium => vec![Encode::Json],
1155 ),
1156 MqttSink::SINK_NAME => hashmap!(
1157 Format::Plain => vec![Encode::Json, Encode::Protobuf],
1158 ),
1159 PulsarSink::SINK_NAME => hashmap!(
1160 Format::Plain => vec![Encode::Json],
1161 Format::Upsert => vec![Encode::Json],
1162 Format::Debezium => vec![Encode::Json],
1163 ),
1164 RedisSink::SINK_NAME => hashmap!(
1165 Format::Plain => vec![Encode::Json, Encode::Template],
1166 Format::Upsert => vec![Encode::Json, Encode::Template],
1167 ),
1168 ))
1169 });
1170
1171pub fn validate_compatibility(connector: &str, format_desc: &FormatEncodeOptions) -> Result<()> {
1172 let compatible_formats = CONNECTORS_COMPATIBLE_FORMATS
1173 .get(connector)
1174 .ok_or_else(|| {
1175 ErrorCode::BindError(format!(
1176 "connector {} is not supported by FORMAT ... ENCODE ... syntax",
1177 connector
1178 ))
1179 })?;
1180 let compatible_encodes = compatible_formats.get(&format_desc.format).ok_or_else(|| {
1181 ErrorCode::BindError(format!(
1182 "connector {} does not support format {:?}",
1183 connector, format_desc.format
1184 ))
1185 })?;
1186 if !compatible_encodes.contains(&format_desc.row_encode) {
1187 return Err(ErrorCode::BindError(format!(
1188 "connector {} does not support format {:?} with encode {:?}",
1189 connector, format_desc.format, format_desc.row_encode
1190 ))
1191 .into());
1192 }
1193
1194 if let Some(encode) = &format_desc.key_encode
1196 && connector != KAFKA_SINK
1197 && matches!(encode, Encode::Bytes)
1198 {
1199 return Err(ErrorCode::BindError(format!(
1200 "key encode bytes only works with kafka connector, but found {}",
1201 connector
1202 ))
1203 .into());
1204 }
1205
1206 Ok(())
1207}
1208
1209#[cfg(test)]
1210pub mod tests {
1211 use risingwave_common::catalog::{CreateType, DEFAULT_DATABASE_NAME, DEFAULT_SCHEMA_NAME};
1212 use risingwave_common::config::FrontendConfig;
1213
1214 use crate::catalog::root_catalog::SchemaPath;
1215 use crate::test_utils::{LocalFrontend, PROTO_FILE_DATA, create_proto_file};
1216
1217 #[tokio::test]
1218 async fn test_create_sink_handler() {
1219 let proto_file = create_proto_file(PROTO_FILE_DATA);
1220 let sql = format!(
1221 r#"CREATE SOURCE t1
1222 WITH (connector = 'kafka', kafka.topic = 'abc', kafka.brokers = 'localhost:1001')
1223 FORMAT PLAIN ENCODE PROTOBUF (message = '.test.TestRecord', schema.location = 'file://{}')"#,
1224 proto_file.path().to_str().unwrap()
1225 );
1226 let frontend = LocalFrontend::new(Default::default()).await;
1227 frontend.run_sql(sql).await.unwrap();
1228
1229 let sql = "create materialized view mv1 as select t1.country from t1;";
1230 frontend.run_sql(sql).await.unwrap();
1231
1232 let sql = r#"CREATE SINK snk1 FROM mv1
1233 WITH (connector = 'jdbc', mysql.endpoint = '127.0.0.1:3306', mysql.table =
1234 '<table_name>', mysql.database = '<database_name>', mysql.user = '<user_name>',
1235 mysql.password = '<password>', type = 'append-only', force_append_only = 'true');"#.to_owned();
1236 frontend.run_sql(sql).await.unwrap();
1237
1238 let session = frontend.session_ref();
1239 let catalog_reader = session.env().catalog_reader().read_guard();
1240 let schema_path = SchemaPath::Name(DEFAULT_SCHEMA_NAME);
1241
1242 let (source, _) = catalog_reader
1244 .get_source_by_name(DEFAULT_DATABASE_NAME, schema_path, "t1")
1245 .unwrap();
1246 assert_eq!(source.name, "t1");
1247
1248 let (table, schema_name) = catalog_reader
1250 .get_created_table_by_name(DEFAULT_DATABASE_NAME, schema_path, "mv1")
1251 .unwrap();
1252 assert_eq!(table.name(), "mv1");
1253 let schema_name = schema_name.to_owned();
1254
1255 let (sink, _) = catalog_reader
1257 .get_created_sink_by_name(
1258 DEFAULT_DATABASE_NAME,
1259 SchemaPath::Name(&schema_name),
1260 "snk1",
1261 )
1262 .unwrap();
1263 assert_eq!(sink.name, "snk1");
1264 drop(catalog_reader);
1265
1266 let sql = r#"REPLACE SINK snk1 FROM mv1
1267 WITH (connector = 'jdbc', mysql.endpoint = '127.0.0.1:3306', mysql.table =
1268 '<table_name>', mysql.database = '<database_name>', mysql.user = '<user_name>',
1269 mysql.password = '<password>', type = 'append-only', force_append_only = 'true');"#.to_owned();
1270 frontend.run_sql(sql).await.unwrap();
1271
1272 let catalog_reader = session.env().catalog_reader().read_guard();
1273 let (sink, _) = catalog_reader
1274 .get_created_sink_by_name(
1275 DEFAULT_DATABASE_NAME,
1276 SchemaPath::Name(&schema_name),
1277 "snk1",
1278 )
1279 .unwrap();
1280 assert_eq!(sink.name, "snk1");
1281 assert_eq!(sink.create_type, CreateType::Foreground);
1284 }
1285
1286 #[tokio::test]
1287 async fn test_create_fs_sink_requires_frontend_config() {
1288 let frontend = LocalFrontend::with_frontend_config(
1289 Default::default(),
1290 FrontendConfig {
1291 unsafe_enable_local_fs_connector: false,
1292 ..Default::default()
1293 },
1294 )
1295 .await;
1296 frontend.run_sql("CREATE TABLE t(v int);").await.unwrap();
1297 frontend
1298 .run_sql("CREATE MATERIALIZED VIEW mv AS SELECT * FROM t;")
1299 .await
1300 .unwrap();
1301
1302 let err = frontend
1303 .run_sql(
1304 r#"CREATE SINK local_sink FROM mv
1305 WITH (
1306 connector = 'fs',
1307 fs.path = '/tmp/rw-local-sink',
1308 type = 'append-only',
1309 force_append_only = 'true'
1310 ) FORMAT PLAIN ENCODE JSON (force_append_only = 'true');"#
1311 .to_owned(),
1312 )
1313 .await
1314 .unwrap_err();
1315
1316 assert!(
1317 err.to_string()
1318 .contains("frontend.unsafe_enable_local_fs_connector = true"),
1319 "{err:?}"
1320 );
1321 }
1322}