1#![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 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
105pub 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 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 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 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 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
235pub 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 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 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 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 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 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 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 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 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 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 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 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 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 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
412pub 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 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 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
458pub 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
463pub 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
468pub 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
473pub fn check_source_allow_alter_on_fly_fields(
476 connector_name: &str,
477 fields: &[String],
478) -> crate::error::ConnectorResult<()> {
479 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 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
527pub fn check_sink_allow_alter_on_fly_fields(
530 sink_name: &str,
531 fields: &[String],
532) -> crate::error::ConnectorResult<()> {
533 let allowed_fields = if sink_name == JdbcSink::SINK_NAME {
539 CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS.get(JdbcSink::SINK_NAME)
540 } else {
541 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}