Skip to main content

risingwave_connector/
allow_alter_on_fly_fields.rs

1// Copyright 2025 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
15// THIS FILE IS AUTO_GENERATED. DO NOT EDIT
16// UPDATE WITH: ./risedev generate-with-options
17// This file is rewritten by `tests::test_allow_alter_on_fly_fields_rust_up_to_date` with
18// `UPDATE_EXPECT=1`.
19// To update content, change source/sink/connection WITH options definitions (for example,
20// `#[with_option(allow_alter_on_fly)]` on struct fields), then run `./risedev generate-with-options`.
21// `./risedev generate-with-options` runs three UPDATE_EXPECT tests:
22// 1) refresh `with_options_{source,sink,connection}.yaml`;
23// 2) regenerate this file from those YAML files;
24// 3) regenerate the Iceberg Engine option classifier.
25
26#![rustfmt::skip]
27
28use std::collections::{HashMap, HashSet};
29use std::sync::LazyLock;
30use crate::error::ConnectorError;
31use crate::sink::remote::JdbcSink;
32use crate::sink::Sink;
33
34macro_rules! use_source_properties {
35    ({ $({ $variant_name:ident, $prop_name:ty, $split:ty }),* }) => {
36        $(
37            #[allow(unused_imports)]
38            pub(super) use $prop_name;
39        )*
40    };
41}
42
43mod source_properties {
44    use crate::for_all_sources;
45    use crate::source::base::SourceProperties;
46
47    for_all_sources!(use_source_properties);
48
49    /// Implements a function that maps a source name string to the Rust type name of the corresponding property type.
50    /// Usage: `impl_source_name_to_prop_type_name!();` will generate:
51    /// ```ignore
52    /// pub fn source_name_to_prop_type_name(source_name: &str) -> Option<&'static str>
53    /// ```
54    macro_rules! impl_source_name_to_prop_type_name_inner {
55        ({ $({$variant:ident, $prop_name:ty, $split:ty}),* }) => {
56            pub fn source_name_to_prop_type_name(source_name: &str) -> Option<&'static str> {
57                match source_name {
58                    $(
59                        <$prop_name>::SOURCE_NAME => Some(std::any::type_name::<$prop_name>()),
60                    )*
61                    _ => None,
62                }
63            }
64        };
65    }
66
67    macro_rules! impl_source_name_to_prop_type_name {
68        () => {
69            $crate::for_all_sources! { impl_source_name_to_prop_type_name_inner }
70        };
71    }
72
73    impl_source_name_to_prop_type_name!();
74}
75
76mod sink_properties {
77    use crate::use_all_sink_configs;
78    use crate::sink::Sink;
79    use crate::sink::file_sink::fs::FsSink;
80
81    use_all_sink_configs!();
82
83    macro_rules! impl_sink_name_to_config_type_name_inner {
84        ({ $({ $variant_name:ident, $sink_type:ty, $config_type:ty }),* }) => {
85            pub fn sink_name_to_config_type_name(sink_name: &str) -> Option<&'static str> {
86                match sink_name {
87                $(
88                    <$sink_type>::SINK_NAME => Some(std::any::type_name::<$config_type>()),
89                )*
90                    _ => None,
91                }
92            }
93        };
94    }
95
96    macro_rules! impl_sink_name_to_config_type_name {
97        () => {
98            $crate::for_all_sinks! { impl_sink_name_to_config_type_name_inner }
99        };
100    }
101
102    impl_sink_name_to_config_type_name!();
103}
104
105/// Map of source connector names to their `allow_alter_on_fly` field names
106pub static SOURCE_ALLOW_ALTER_ON_FLY_FIELDS: LazyLock<HashMap<String, HashSet<String>>> = LazyLock::new(|| {
107    use source_properties::*;
108    let mut map = HashMap::new();
109    // CDC Properties - added for schema.change.failure.policy
110    map.try_insert(
111        std::any::type_name::<MysqlCdcProperties>().to_owned(),
112        [
113            "cdc.source.wait.streaming.start.timeout".to_owned(),
114            "debezium.max.queue.size".to_owned(),
115            "debezium.queue.memory.ratio".to_owned(),
116            "hostname".to_owned(),
117            "port".to_owned(),
118            "password".to_owned(),
119        ].into_iter().collect(),
120    ).unwrap();
121    map.try_insert(
122        std::any::type_name::<PostgresCdcProperties>().to_owned(),
123        [
124            "cdc.source.wait.streaming.start.timeout".to_owned(),
125            "debezium.max.queue.size".to_owned(),
126            "debezium.queue.memory.ratio".to_owned(),
127            "debezium.heartbeat.interval.ms".to_owned(),
128            "hostname".to_owned(),
129            "port".to_owned(),
130            "password".to_owned(),
131        ].into_iter().collect(),
132    ).unwrap();
133    map.try_insert(
134        std::any::type_name::<SqlServerCdcProperties>().to_owned(),
135        [
136            "cdc.source.wait.streaming.start.timeout".to_owned(),
137            "debezium.max.queue.size".to_owned(),
138            "debezium.queue.memory.ratio".to_owned(),
139            "hostname".to_owned(),
140            "port".to_owned(),
141            "password".to_owned(),
142        ].into_iter().collect(),
143    ).unwrap();
144    map.try_insert(
145        std::any::type_name::<OracleCdcProperties>().to_owned(),
146        [
147            "cdc.source.wait.streaming.start.timeout".to_owned(),
148            "debezium.max.queue.size".to_owned(),
149            "debezium.queue.memory.ratio".to_owned(),
150            "debezium.heartbeat.interval.ms".to_owned(),
151            "heartbeat.table.name".to_owned(),
152            "hostname".to_owned(),
153            "port".to_owned(),
154            "password".to_owned(),
155        ].into_iter().collect(),
156    ).unwrap();
157
158    map.try_insert(
159        std::any::type_name::<MongodbCdcProperties>().to_owned(),
160        [
161            "cdc.source.wait.streaming.start.timeout".to_owned(),
162            "debezium.max.queue.size".to_owned(),
163            "debezium.queue.memory.ratio".to_owned(),
164        ].into_iter().collect(),
165    ).unwrap();
166    // KafkaProperties
167    map.try_insert(
168        std::any::type_name::<KafkaProperties>().to_owned(),
169        [
170            "group.id.prefix".to_owned(),
171            "properties.sync.call.timeout".to_owned(),
172            "properties.security.protocol".to_owned(),
173            "properties.ssl.endpoint.identification.algorithm".to_owned(),
174            "properties.ssl.ca.location".to_owned(),
175            "properties.ssl.ca.pem".to_owned(),
176            "properties.ssl.certificate.location".to_owned(),
177            "properties.ssl.certificate.pem".to_owned(),
178            "properties.ssl.key.location".to_owned(),
179            "properties.ssl.key.pem".to_owned(),
180            "properties.ssl.key.password".to_owned(),
181            "properties.sasl.mechanism".to_owned(),
182            "properties.sasl.username".to_owned(),
183            "properties.sasl.password".to_owned(),
184            "properties.sasl.kerberos.service.name".to_owned(),
185            "properties.sasl.kerberos.keytab".to_owned(),
186            "properties.sasl.kerberos.principal".to_owned(),
187            "properties.sasl.kerberos.kinit.cmd".to_owned(),
188            "properties.sasl.kerberos.min.time.before.relogin".to_owned(),
189            "properties.sasl.oauthbearer.config".to_owned(),
190            "properties.sasl.oauthbearer.method".to_owned(),
191            "properties.sasl.oauthbearer.client.id".to_owned(),
192            "properties.sasl.oauthbearer.client.secret".to_owned(),
193            "properties.sasl.oauthbearer.token.endpoint.url".to_owned(),
194            "properties.sasl.oauthbearer.scope".to_owned(),
195            "properties.sasl.oauthbearer.extensions".to_owned(),
196            "properties.message.max.bytes".to_owned(),
197            "properties.receive.message.max.bytes".to_owned(),
198            "properties.statistics.interval.ms".to_owned(),
199            "properties.client.id".to_owned(),
200            "properties.enable.ssl.certificate.verification".to_owned(),
201            "properties.reconnect.backoff.ms".to_owned(),
202            "properties.reconnect.backoff.max.ms".to_owned(),
203            "properties.socket.connection.setup.timeout.ms".to_owned(),
204            "properties.retry.backoff.ms".to_owned(),
205            "properties.retry.backoff.max.ms".to_owned(),
206            "properties.queued.min.messages".to_owned(),
207            "properties.queued.max.messages.kbytes".to_owned(),
208            "properties.fetch.wait.max.ms".to_owned(),
209            "properties.fetch.queue.backoff.ms".to_owned(),
210            "properties.fetch.max.bytes".to_owned(),
211            "properties.enable.auto.commit".to_owned(),
212            "properties.auto.commit.interval.ms".to_owned(),
213        ].into_iter().collect(),
214    ).unwrap();
215    // PubsubProperties
216    map.try_insert(
217        std::any::type_name::<PubsubProperties>().to_owned(),
218        [
219            "pubsub.ack_deadline_seconds".to_owned(),
220            "pubsub.max_outstanding_messages".to_owned(),
221            "pubsub.max_outstanding_bytes".to_owned(),
222        ].into_iter().collect(),
223    ).unwrap();
224    // PulsarProperties
225    map.try_insert(
226        std::any::type_name::<PulsarProperties>().to_owned(),
227        [
228            "pulsar.operation.retry.max.retries".to_owned(),
229            "pulsar.operation.retry.delay".to_owned(),
230        ].into_iter().collect(),
231    ).unwrap();
232    map
233});
234
235/// Map of sink connector names to their `allow_alter_on_fly` field names
236pub static SINK_ALLOW_ALTER_ON_FLY_FIELDS: LazyLock<HashMap<String, HashSet<String>>> = LazyLock::new(|| {
237    use sink_properties::*;
238    let mut map = HashMap::new();
239    // ClickHouseConfig
240    map.try_insert(
241        std::any::type_name::<ClickHouseConfig>().to_owned(),
242        [
243            "commit_checkpoint_interval".to_owned(),
244        ].into_iter().collect(),
245    ).unwrap();
246    // DeltaLakeConfig
247    map.try_insert(
248        std::any::type_name::<DeltaLakeConfig>().to_owned(),
249        [
250            "commit_checkpoint_interval".to_owned(),
251        ].into_iter().collect(),
252    ).unwrap();
253    // DorisConfig
254    map.try_insert(
255        std::any::type_name::<DorisConfig>().to_owned(),
256        [
257            "doris.stream_load.http.timeout.ms".to_owned(),
258        ].into_iter().collect(),
259    ).unwrap();
260    // ElasticSearchConfig
261    map.try_insert(
262        std::any::type_name::<ElasticSearchConfig>().to_owned(),
263        [
264            "batch_num_messages".to_owned(),
265            "batch_size_kb".to_owned(),
266            "concurrent_requests".to_owned(),
267        ].into_iter().collect(),
268    ).unwrap();
269    // IcebergConfig
270    map.try_insert(
271        std::any::type_name::<IcebergConfig>().to_owned(),
272        [
273            "commit_checkpoint_interval".to_owned(),
274            "enable_compaction".to_owned(),
275            "compaction_interval_sec".to_owned(),
276            "enable_snapshot_expiration".to_owned(),
277            "snapshot_expiration_max_age_millis".to_owned(),
278            "snapshot_expiration_retain_last".to_owned(),
279            "snapshot_expiration_clear_expired_files".to_owned(),
280            "snapshot_expiration_clear_expired_meta_data".to_owned(),
281            "enable_manifest_rewrite".to_owned(),
282            "manifest_rewrite_target_size_bytes".to_owned(),
283            "manifest_rewrite_min_count_to_merge".to_owned(),
284            "compaction.max_snapshots_num".to_owned(),
285            "compaction.small_files_threshold_mb".to_owned(),
286            "compaction.delete_files_count_threshold".to_owned(),
287            "compaction.trigger_snapshot_count".to_owned(),
288            "compaction.target_file_size_mb".to_owned(),
289            "compaction.type".to_owned(),
290            "compaction.write_parquet_compression".to_owned(),
291            "compaction.write_parquet_max_row_group_rows".to_owned(),
292            "compaction.write_parquet_max_row_group_bytes".to_owned(),
293        ].into_iter().collect(),
294    ).unwrap();
295    // KafkaConfig
296    map.try_insert(
297        std::any::type_name::<KafkaConfig>().to_owned(),
298        [
299            "properties.sync.call.timeout".to_owned(),
300            "properties.security.protocol".to_owned(),
301            "properties.ssl.endpoint.identification.algorithm".to_owned(),
302            "properties.ssl.ca.location".to_owned(),
303            "properties.ssl.ca.pem".to_owned(),
304            "properties.ssl.certificate.location".to_owned(),
305            "properties.ssl.certificate.pem".to_owned(),
306            "properties.ssl.key.location".to_owned(),
307            "properties.ssl.key.pem".to_owned(),
308            "properties.ssl.key.password".to_owned(),
309            "properties.sasl.mechanism".to_owned(),
310            "properties.sasl.username".to_owned(),
311            "properties.sasl.password".to_owned(),
312            "properties.sasl.kerberos.service.name".to_owned(),
313            "properties.sasl.kerberos.keytab".to_owned(),
314            "properties.sasl.kerberos.principal".to_owned(),
315            "properties.sasl.kerberos.kinit.cmd".to_owned(),
316            "properties.sasl.kerberos.min.time.before.relogin".to_owned(),
317            "properties.sasl.oauthbearer.config".to_owned(),
318            "properties.sasl.oauthbearer.method".to_owned(),
319            "properties.sasl.oauthbearer.client.id".to_owned(),
320            "properties.sasl.oauthbearer.client.secret".to_owned(),
321            "properties.sasl.oauthbearer.token.endpoint.url".to_owned(),
322            "properties.sasl.oauthbearer.scope".to_owned(),
323            "properties.sasl.oauthbearer.extensions".to_owned(),
324            "properties.message.max.bytes".to_owned(),
325            "properties.receive.message.max.bytes".to_owned(),
326            "properties.statistics.interval.ms".to_owned(),
327            "properties.client.id".to_owned(),
328            "properties.enable.ssl.certificate.verification".to_owned(),
329            "properties.reconnect.backoff.ms".to_owned(),
330            "properties.reconnect.backoff.max.ms".to_owned(),
331            "properties.socket.connection.setup.timeout.ms".to_owned(),
332            "properties.retry.backoff.ms".to_owned(),
333            "properties.retry.backoff.max.ms".to_owned(),
334            "properties.allow.auto.create.topics".to_owned(),
335            "properties.queue.buffering.max.messages".to_owned(),
336            "properties.queue.buffering.max.kbytes".to_owned(),
337            "properties.queue.buffering.max.ms".to_owned(),
338            "properties.enable.idempotence".to_owned(),
339            "properties.message.send.max.retries".to_owned(),
340            "properties.batch.num.messages".to_owned(),
341            "properties.batch.size".to_owned(),
342            "properties.message.timeout.ms".to_owned(),
343            "properties.max.in.flight.requests.per.connection".to_owned(),
344            "properties.request.required.acks".to_owned(),
345        ].into_iter().collect(),
346    ).unwrap();
347    // LanceDbConfig
348    map.try_insert(
349        std::any::type_name::<LanceDbConfig>().to_owned(),
350        [
351            "commit_checkpoint_interval".to_owned(),
352        ].into_iter().collect(),
353    ).unwrap();
354    // OpenSearchConfig
355    map.try_insert(
356        std::any::type_name::<OpenSearchConfig>().to_owned(),
357        [
358            "batch_num_messages".to_owned(),
359            "batch_size_kb".to_owned(),
360            "concurrent_requests".to_owned(),
361        ].into_iter().collect(),
362    ).unwrap();
363    // PulsarConfig
364    map.try_insert(
365        std::any::type_name::<PulsarConfig>().to_owned(),
366        [
367            "properties.routing.mode".to_owned(),
368            "properties.routing_mode".to_owned(),
369            "routing_mode".to_owned(),
370            "pulsar.routing_mode".to_owned(),
371            "pulsar.routing.mode".to_owned(),
372            "pulsar.properties.routing.mode".to_owned(),
373            "pulsar.properties.routing_mode".to_owned(),
374        ].into_iter().collect(),
375    ).unwrap();
376    // SnowflakeV2Config
377    map.try_insert(
378        std::any::type_name::<SnowflakeV2Config>().to_owned(),
379        [
380            "commit_checkpoint_interval".to_owned(),
381        ].into_iter().collect(),
382    ).unwrap();
383    // StarrocksConfig
384    map.try_insert(
385        std::any::type_name::<StarrocksConfig>().to_owned(),
386        [
387            "starrocks.stream_load.http.timeout.ms".to_owned(),
388            "commit_checkpoint_interval".to_owned(),
389            "starrocks.max_batch_size_bytes".to_owned(),
390        ].into_iter().collect(),
391    ).unwrap();
392    // TurbopufferConfig
393    map.try_insert(
394        std::any::type_name::<TurbopufferConfig>().to_owned(),
395        [
396            "write_batch_size".to_owned(),
397            "max_linger_second".to_owned(),
398        ].into_iter().collect(),
399    ).unwrap();
400    // Jdbc
401    map.try_insert(
402        JdbcSink::SINK_NAME.to_owned(),
403        [
404            "jdbc.url".to_owned(),
405            "user".to_owned(),
406            "password".to_owned(),
407        ].into_iter().collect(),
408    ).unwrap();
409    map
410});
411
412/// Map of connection names to their `allow_alter_on_fly` field names
413pub static CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS: LazyLock<HashMap<String, HashSet<String>>> = LazyLock::new(|| {
414    use crate::connector_common::*;
415    let mut map = HashMap::new();
416    // KafkaConnection
417    map.try_insert(
418        std::any::type_name::<KafkaConnection>().to_owned(),
419        [
420            "properties.security.protocol".to_owned(),
421            "properties.ssl.endpoint.identification.algorithm".to_owned(),
422            "properties.ssl.ca.location".to_owned(),
423            "properties.ssl.ca.pem".to_owned(),
424            "properties.ssl.certificate.location".to_owned(),
425            "properties.ssl.certificate.pem".to_owned(),
426            "properties.ssl.key.location".to_owned(),
427            "properties.ssl.key.pem".to_owned(),
428            "properties.ssl.key.password".to_owned(),
429            "properties.sasl.mechanism".to_owned(),
430            "properties.sasl.username".to_owned(),
431            "properties.sasl.password".to_owned(),
432            "properties.sasl.kerberos.service.name".to_owned(),
433            "properties.sasl.kerberos.keytab".to_owned(),
434            "properties.sasl.kerberos.principal".to_owned(),
435            "properties.sasl.kerberos.kinit.cmd".to_owned(),
436            "properties.sasl.kerberos.min.time.before.relogin".to_owned(),
437            "properties.sasl.oauthbearer.config".to_owned(),
438            "properties.sasl.oauthbearer.method".to_owned(),
439            "properties.sasl.oauthbearer.client.id".to_owned(),
440            "properties.sasl.oauthbearer.client.secret".to_owned(),
441            "properties.sasl.oauthbearer.token.endpoint.url".to_owned(),
442            "properties.sasl.oauthbearer.scope".to_owned(),
443            "properties.sasl.oauthbearer.extensions".to_owned(),
444        ].into_iter().collect(),
445    ).unwrap();
446    // Jdbc
447    map.try_insert(
448        JdbcSink::SINK_NAME.to_owned(),
449        [
450            "jdbc.url".to_owned(),
451            "user".to_owned(),
452            "password".to_owned(),
453        ].into_iter().collect(),
454    ).unwrap();
455    map
456});
457
458/// Get all source connector names that have `allow_alter_on_fly` fields
459pub fn get_source_connectors_with_allow_alter_on_fly_fields() -> Vec<&'static str> {
460    SOURCE_ALLOW_ALTER_ON_FLY_FIELDS.keys().map(|s| s.as_str()).collect()
461}
462
463/// Get all sink connector names that have `allow_alter_on_fly` fields
464pub fn get_sink_connectors_with_allow_alter_on_fly_fields() -> Vec<&'static str> {
465    SINK_ALLOW_ALTER_ON_FLY_FIELDS.keys().map(|s| s.as_str()).collect()
466}
467
468/// Get all connection names that have `allow_alter_on_fly` fields
469pub fn get_connection_names_with_allow_alter_on_fly_fields() -> Vec<&'static str> {
470    CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS.keys().map(|s| s.as_str()).collect()
471}
472
473/// Checks if all given fields are allowed to be altered on the fly for the specified source connector.
474/// Returns Ok(()) if all fields are allowed, otherwise returns a `ConnectorError`.
475pub fn check_source_allow_alter_on_fly_fields(
476    connector_name: &str,
477    fields: &[String],
478) -> crate::error::ConnectorResult<()> {
479    // Convert connector name to the type name key
480    let Some(type_name) = source_properties::source_name_to_prop_type_name(connector_name) else {
481        return Err(ConnectorError::from(anyhow::anyhow!(
482            "Unknown source connector: {connector_name}"
483        )));
484    };
485    let Some(allowed_fields) = SOURCE_ALLOW_ALTER_ON_FLY_FIELDS.get(type_name) else {
486    return Err(ConnectorError::from(anyhow::anyhow!(
487        "No allow_alter_on_fly fields registered for connector: {connector_name}"
488    )));
489    };
490    for field in fields {
491        if !allowed_fields.contains(field) {
492            return Err(ConnectorError::from(anyhow::anyhow!(
493                "Field '{field}' is not allowed to be altered on the fly for connector: {connector_name}"
494            )));
495        }
496    }
497    Ok(())
498}
499
500pub fn check_connection_allow_alter_on_fly_fields(
501    connection_name: &str,
502    fields: &[String],
503) -> crate::error::ConnectorResult<()> {
504    use crate::source::connection_name_to_prop_type_name;
505
506    // Convert connection name to the type name key
507    let Some(type_name) = connection_name_to_prop_type_name(connection_name) else {
508        return Err(ConnectorError::from(anyhow::anyhow!(
509            "Unknown connection: {connection_name}"
510        )));
511    };
512    let Some(allowed_fields) = CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS.get(type_name) else {
513        return Err(ConnectorError::from(anyhow::anyhow!(
514            "No allow_alter_on_fly fields registered for connection: {connection_name}"
515        )));
516    };
517    for field in fields {
518        if !allowed_fields.contains(field) {
519            return Err(ConnectorError::from(anyhow::anyhow!(
520                "Field '{field}' is not allowed to be altered on the fly for connection: {connection_name}"
521            )));
522        }
523    }
524    Ok(())
525}
526
527/// Checks if all given fields are allowed to be altered on the fly for the specified sink connector.
528/// Returns Ok(()) if all fields are allowed, otherwise returns a `ConnectorError`.
529pub fn check_sink_allow_alter_on_fly_fields(
530    sink_name: &str,
531    fields: &[String],
532) -> crate::error::ConnectorResult<()> {
533    // TODO(#24846): JDBC sink currently uses `()` as sink config type in `for_all_sinks!`,
534    // so it cannot have an isolated key in `SINK_ALLOW_ALTER_ON_FLY_FIELDS`.
535    // Reuse the JDBC entry in `CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS` for now.
536    // TODO(#24846): remove this special case after JDBC sink has a dedicated config type
537    // and allow-alter fields are generated directly into `SINK_ALLOW_ALTER_ON_FLY_FIELDS`.
538    let allowed_fields = if sink_name == JdbcSink::SINK_NAME {
539        CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS.get(JdbcSink::SINK_NAME)
540    } else {
541        // Convert sink name to the type name key
542        let Some(type_name) = sink_properties::sink_name_to_config_type_name(sink_name) else {
543            return Err(ConnectorError::from(anyhow::anyhow!(
544                "Unknown sink connector: {sink_name}"
545            )));
546        };
547        SINK_ALLOW_ALTER_ON_FLY_FIELDS.get(type_name)
548    };
549    let Some(allowed_fields) = allowed_fields else {
550        return Err(ConnectorError::from(anyhow::anyhow!(
551            "No allow_alter_on_fly fields registered for sink: {sink_name}"
552        )));
553    };
554    for field in fields {
555        if !allowed_fields.contains(field) {
556            return Err(ConnectorError::from(anyhow::anyhow!(
557                "Field '{field}' is not allowed to be altered on the fly for sink: {sink_name}"
558            )));
559        }
560    }
561    Ok(())
562}