Skip to main content

rustium_config/
lib.rs

1//! Strict, versioned Rustium configuration.
2
3mod debezium;
4
5use std::{collections::BTreeMap, env, fs, path::Path, time::Duration};
6
7use regex::Regex;
8use rustium_core::{Error, Result, RetryPolicy};
9use serde::{Deserialize, Serialize};
10use sha2::{Digest, Sha256};
11use url::Url;
12
13const API_VERSION: &str = "rustium.io/v1alpha1";
14
15#[derive(Debug, Clone, Serialize, Deserialize)]
16#[serde(deny_unknown_fields)]
17pub struct Config {
18    pub api_version: String,
19    pub kind: String,
20    pub metadata: Metadata,
21    pub source: SourceConfig,
22    #[serde(default)]
23    pub snapshot: SnapshotConfig,
24    #[serde(default)]
25    pub format: FormatConfig,
26    pub sink: SinkConfig,
27    #[serde(default)]
28    pub state: StateConfig,
29    #[serde(default)]
30    pub runtime: RuntimeSettings,
31    #[serde(default)]
32    pub server: ServerConfig,
33    #[serde(default)]
34    pub observability: ObservabilityConfig,
35    #[serde(skip)]
36    pub compatibility_warnings: Vec<String>,
37}
38
39impl Config {
40    pub fn load(path: impl AsRef<Path>) -> Result<Self> {
41        let path = path.as_ref();
42        let raw = fs::read_to_string(path)?;
43        if path
44            .extension()
45            .is_some_and(|extension| extension == "properties")
46            || !raw
47                .lines()
48                .any(|line| line.trim_start().starts_with("api_version:"))
49        {
50            Self::from_debezium_properties(&raw)
51        } else {
52            Self::from_yaml(&raw)
53        }
54    }
55
56    pub fn from_yaml(raw: &str) -> Result<Self> {
57        let interpolated = interpolate_environment(raw)?;
58        let config: Self = serde_yaml::from_str(&interpolated)
59            .map_err(|error| Error::Configuration(error.to_string()))?;
60        config.validate()?;
61        Ok(config)
62    }
63
64    pub fn from_debezium_properties(raw: &str) -> Result<Self> {
65        debezium::parse(raw)
66    }
67
68    pub fn validate(&self) -> Result<()> {
69        if self.api_version != API_VERSION {
70            return Err(Error::Configuration(format!(
71                "unsupported api_version {:?}; expected {API_VERSION:?}",
72                self.api_version
73            )));
74        }
75        if self.kind != "Connector" {
76            return Err(Error::Configuration(
77                "kind must be exactly \"Connector\"".into(),
78            ));
79        }
80        validate_name(&self.metadata.name, "metadata.name")?;
81        self.source.validate()?;
82        self.sink.validate()?;
83        self.format.validate(&self.sink)?;
84        if self.snapshot.fetch_size == 0 {
85            return Err(Error::Configuration(
86                "snapshot.fetch_size must be greater than zero".into(),
87            ));
88        }
89        self.snapshot.validate()?;
90        if self.runtime.channel_capacity == 0 {
91            return Err(Error::Configuration(
92                "runtime.channel_capacity must be greater than zero".into(),
93            ));
94        }
95        if self.runtime.max_batch_size == 0
96            || self.runtime.max_batch_size > self.runtime.channel_capacity
97        {
98            return Err(Error::Configuration(
99                "runtime.max_batch_size must be between 1 and channel_capacity".into(),
100            ));
101        }
102        if self.runtime.errors_max_retries < -1 {
103            return Err(Error::Configuration(
104                "runtime.errors_max_retries must be -1, 0, or a positive integer".into(),
105            ));
106        }
107        if self.runtime.errors_retry_delay_initial.is_zero() {
108            return Err(Error::Configuration(
109                "runtime.errors_retry_delay_initial must be greater than zero".into(),
110            ));
111        }
112        if self.runtime.errors_retry_delay_max < self.runtime.errors_retry_delay_initial {
113            return Err(Error::Configuration(
114                "runtime.errors_retry_delay_max must be greater than or equal to errors_retry_delay_initial"
115                    .into(),
116            ));
117        }
118        Ok(())
119    }
120
121    #[must_use]
122    pub fn fingerprint(&self) -> String {
123        let semantic = serde_json::json!({
124            "api_version": self.api_version,
125            "name": self.metadata.name,
126            "source": self.source.semantic_config(),
127            "snapshot": self.snapshot.semantic_config(),
128            "format": self.format.semantic_config(),
129            "sink": self.sink.semantic_config(),
130        });
131        let bytes = serde_json::to_vec(&semantic).expect("configuration serialization cannot fail");
132        hex_digest(&bytes)
133    }
134}
135
136#[derive(Debug, Clone, Serialize, Deserialize)]
137#[serde(deny_unknown_fields)]
138pub struct Metadata {
139    pub name: String,
140    #[serde(default)]
141    pub labels: BTreeMap<String, String>,
142}
143
144#[derive(Debug, Clone, Serialize, Deserialize)]
145#[serde(tag = "type", rename_all = "snake_case")]
146#[allow(clippy::large_enum_variant)]
147pub enum SourceConfig {
148    Postgresql(Box<PostgresSourceConfig>),
149    Mysql(MySqlSourceConfig),
150    Sqlserver(SqlServerSourceConfig),
151    Oracle(OracleSourceConfig),
152    Mongodb(MongoDbSourceConfig),
153    Mariadb(DebeziumSourceConfig),
154    Db2(DebeziumSourceConfig),
155    Cassandra(DebeziumSourceConfig),
156    Vitess(DebeziumSourceConfig),
157    Spanner(DebeziumSourceConfig),
158    Informix(DebeziumSourceConfig),
159    Cockroachdb(DebeziumSourceConfig),
160    Yashandb(DebeziumSourceConfig),
161}
162
163impl SourceConfig {
164    fn validate(&self) -> Result<()> {
165        match self {
166            Self::Postgresql(config) => config.validate(),
167            Self::Mysql(config) => config.validate(),
168            Self::Sqlserver(config) => config.validate(),
169            Self::Oracle(config) => config.validate(),
170            Self::Mongodb(config) => config.validate(),
171            Self::Mariadb(config) => config.validate(DebeziumConnectorKind::Mariadb),
172            Self::Db2(config) => config.validate(DebeziumConnectorKind::Db2),
173            Self::Cassandra(config) => config.validate(DebeziumConnectorKind::Cassandra),
174            Self::Vitess(config) => config.validate(DebeziumConnectorKind::Vitess),
175            Self::Spanner(config) => config.validate(DebeziumConnectorKind::Spanner),
176            Self::Informix(config) => config.validate(DebeziumConnectorKind::Informix),
177            Self::Cockroachdb(config) => config.validate(DebeziumConnectorKind::CockroachDb),
178            Self::Yashandb(config) => config.validate(DebeziumConnectorKind::YashanDb),
179        }
180    }
181
182    #[must_use]
183    pub fn as_postgresql(&self) -> Option<&PostgresSourceConfig> {
184        match self {
185            Self::Postgresql(config) => Some(config),
186            _ => None,
187        }
188    }
189
190    #[must_use]
191    pub fn as_mysql(&self) -> Option<&MySqlSourceConfig> {
192        match self {
193            Self::Mysql(config) => Some(config),
194            _ => None,
195        }
196    }
197
198    #[must_use]
199    pub fn as_sqlserver(&self) -> Option<&SqlServerSourceConfig> {
200        match self {
201            Self::Sqlserver(config) => Some(config),
202            _ => None,
203        }
204    }
205
206    #[must_use]
207    pub fn as_oracle(&self) -> Option<&OracleSourceConfig> {
208        match self {
209            Self::Oracle(config) => Some(config),
210            _ => None,
211        }
212    }
213
214    #[must_use]
215    pub fn as_mongodb(&self) -> Option<&MongoDbSourceConfig> {
216        match self {
217            Self::Mongodb(config) => Some(config),
218            _ => None,
219        }
220    }
221
222    #[must_use]
223    pub fn as_debezium(&self) -> Option<(DebeziumConnectorKind, &DebeziumSourceConfig)> {
224        match self {
225            Self::Mariadb(config) => Some((DebeziumConnectorKind::Mariadb, config)),
226            Self::Db2(config) => Some((DebeziumConnectorKind::Db2, config)),
227            Self::Cassandra(config) => Some((DebeziumConnectorKind::Cassandra, config)),
228            Self::Vitess(config) => Some((DebeziumConnectorKind::Vitess, config)),
229            Self::Spanner(config) => Some((DebeziumConnectorKind::Spanner, config)),
230            Self::Informix(config) => Some((DebeziumConnectorKind::Informix, config)),
231            Self::Cockroachdb(config) => Some((DebeziumConnectorKind::CockroachDb, config)),
232            Self::Yashandb(config) => Some((DebeziumConnectorKind::YashanDb, config)),
233            _ => None,
234        }
235    }
236
237    fn semantic_config(&self) -> serde_json::Value {
238        match self {
239            Self::Postgresql(config) => {
240                let mut semantic = serde_json::json!({
241                    "type": "postgresql",
242                    "hostname": config.hostname,
243                    "port": config.port,
244                    "database": config.database,
245                    "publication": config.publication,
246                    "slot_name": config.slot_name,
247                    "tables": config.tables,
248                });
249                if config.publication_autocreate_mode != PublicationAutoCreateMode::Disabled {
250                    semantic
251                        .as_object_mut()
252                        .expect("source semantic is an object")
253                        .insert(
254                            "publication_autocreate_mode".into(),
255                            serde_json::json!(config.publication_autocreate_mode),
256                        );
257                }
258                if !config.replica_identity_autoset_values.is_empty() {
259                    semantic
260                        .as_object_mut()
261                        .expect("source semantic is an object")
262                        .insert(
263                            "replica_identity_autoset_values".into(),
264                            serde_json::json!(config.replica_identity_autoset_values),
265                        );
266                }
267                if config.publish_via_partition_root {
268                    semantic
269                        .as_object_mut()
270                        .expect("source semantic is an object")
271                        .insert("publish_via_partition_root".into(), true.into());
272                }
273                if config.slot_failover {
274                    semantic
275                        .as_object_mut()
276                        .expect("source semantic is an object")
277                        .insert("slot_failover".into(), true.into());
278                }
279                if config.offset_mismatch_strategy != PostgresOffsetMismatchStrategy::NoValidation {
280                    semantic
281                        .as_object_mut()
282                        .expect("source semantic is an object")
283                        .insert(
284                            "offset_mismatch_strategy".into(),
285                            serde_json::json!(config.offset_mismatch_strategy),
286                        );
287                }
288                if config.lsn_flush_mode != PostgresLsnFlushMode::Connector {
289                    semantic
290                        .as_object_mut()
291                        .expect("source semantic is an object")
292                        .insert(
293                            "lsn_flush_mode".into(),
294                            serde_json::json!(config.lsn_flush_mode),
295                        );
296                }
297                if !config.snapshot_isolation_mode.imports_snapshot() {
298                    semantic
299                        .as_object_mut()
300                        .expect("source semantic is an object")
301                        .insert(
302                            "snapshot_isolation_mode".into(),
303                            serde_json::json!(config.snapshot_isolation_mode),
304                        );
305                }
306                if !config.xmin_fetch_interval.is_zero() {
307                    semantic
308                        .as_object_mut()
309                        .expect("source semantic is an object")
310                        .insert(
311                            "xmin_fetch_interval".into(),
312                            serde_json::json!(config.xmin_fetch_interval),
313                        );
314                }
315                if !config.slot_stream_params.is_empty() {
316                    semantic
317                        .as_object_mut()
318                        .expect("source semantic is an object")
319                        .insert(
320                            "slot_stream_params".into(),
321                            serde_json::json!(config.slot_stream_params),
322                        );
323                }
324                if !config.database_initial_statements.is_empty() {
325                    semantic
326                        .as_object_mut()
327                        .expect("source semantic is an object")
328                        .insert(
329                            "database_initial_statements".into(),
330                            serde_json::json!(config.database_initial_statements),
331                        );
332                }
333                add_heartbeat_semantics(
334                    &mut semantic,
335                    config.heartbeat_interval,
336                    config.heartbeat_action_query.as_deref(),
337                    &config.heartbeat_topics_prefix,
338                    config.heartbeat_topic_name.as_deref(),
339                );
340                if config.signal_data_collection.is_some() || config.read_only {
341                    semantic
342                        .as_object_mut()
343                        .expect("source semantic is an object")
344                        .insert(
345                            "incremental_snapshot".into(),
346                            serde_json::json!({
347                                "signal_data_collection": config.signal_data_collection,
348                                "chunk_size": config.incremental_snapshot_chunk_size,
349                                "watermarking_strategy": config.incremental_snapshot_watermarking_strategy,
350                                "read_only": config.read_only,
351                            }),
352                        );
353                }
354                if config.signal_enabled_channels != default_signal_enabled_channels()
355                    || config
356                        .signal_enabled_channels
357                        .iter()
358                        .any(|channel| channel == "file")
359                {
360                    semantic
361                        .as_object_mut()
362                        .expect("source semantic is an object")
363                        .insert(
364                            "signals".into(),
365                            serde_json::json!({
366                                "enabled_channels": config.signal_enabled_channels,
367                                "file": config.signal_file,
368                                "poll_interval_ms": config.signal_poll_interval.as_millis(),
369                                "kafka_topic": config.signal_kafka_topic,
370                                "kafka_group_id": config.signal_kafka_group_id,
371                            }),
372                        );
373                }
374                semantic
375                    .as_object_mut()
376                    .expect("source semantic is an object")
377                    .insert(
378                        "hstore_handling_mode".into(),
379                        config.hstore_handling_mode.clone().into(),
380                    );
381                if config.interval_handling_mode != default_postgres_interval_handling_mode() {
382                    semantic
383                        .as_object_mut()
384                        .expect("source semantic is an object")
385                        .insert(
386                            "interval_handling_mode".into(),
387                            config.interval_handling_mode.clone().into(),
388                        );
389                }
390                if config.include_unknown_datatypes {
391                    semantic
392                        .as_object_mut()
393                        .expect("source semantic is an object")
394                        .insert("include_unknown_datatypes".into(), true.into());
395                }
396                if config.money_fraction_digits != default_postgres_money_fraction_digits() {
397                    semantic
398                        .as_object_mut()
399                        .expect("source semantic is an object")
400                        .insert(
401                            "money_fraction_digits".into(),
402                            config.money_fraction_digits.into(),
403                        );
404                }
405                if config.captures_logical_decoding_messages() {
406                    semantic
407                        .as_object_mut()
408                        .expect("source semantic is an object")
409                        .insert(
410                            "logical_decoding_messages".into(),
411                            serde_json::json!({
412                                "include": config.message_prefix_include_list,
413                                "exclude": config.message_prefix_exclude_list,
414                            }),
415                        );
416                }
417                if !config.column_transformations.is_empty() {
418                    semantic
419                        .as_object_mut()
420                        .expect("source semantic is an object")
421                        .insert(
422                            "column_transformations".into(),
423                            serde_json::Value::Array(column_transformation_semantics(
424                                &config.column_transformations,
425                            )),
426                        );
427                }
428                semantic
429            }
430            Self::Mysql(config) => {
431                let mut semantic = serde_json::json!({
432                    "type": "mysql",
433                    "hostname": config.hostname,
434                    "port": config.port,
435                    "databases": config.databases,
436                    "server_id": config.server_id,
437                    "tables": config.tables,
438                    "ssl_ca": config.ssl_ca,
439                    "ssl_cert": config.ssl_cert,
440                    "ssl_key": config.ssl_key,
441                    "schema_history_skip_unparseable_ddl": config.schema_history_skip_unparseable_ddl,
442                    "gtid_source_includes": config.gtid_source_includes,
443                    "gtid_source_excludes": config.gtid_source_excludes,
444                    "gtid_source_filter_dml_events": config.gtid_source_filter_dml_events,
445                });
446                if config.ssl_keystore.is_some() || config.ssl_truststore.is_some() {
447                    semantic
448                        .as_object_mut()
449                        .expect("source semantic is an object")
450                        .insert(
451                            "java_ssl_stores".into(),
452                            serde_json::json!({
453                                "keystore": config.ssl_keystore,
454                                "truststore": config.ssl_truststore,
455                            }),
456                        );
457                }
458                add_heartbeat_semantics(
459                    &mut semantic,
460                    config.heartbeat_interval,
461                    config.heartbeat_action_query.as_deref(),
462                    &config.heartbeat_topics_prefix,
463                    config.heartbeat_topic_name.as_deref(),
464                );
465                semantic
466                    .as_object_mut()
467                    .expect("source semantic is an object")
468                    .insert(
469                        "signals".into(),
470                        serde_json::json!({
471                            "data_collection": config.signal_data_collection,
472                            "enabled_channels": config.signal_enabled_channels,
473                            "file": config.signal_file,
474                            "poll_interval_ms": config.signal_poll_interval.as_millis(),
475                            "chunk_size": config.incremental_snapshot_chunk_size,
476                            "watermarking_strategy": config.incremental_snapshot_watermarking_strategy,
477                            "kafka_topic": config.signal_kafka_topic,
478                            "kafka_bootstrap_servers": config.signal_kafka_bootstrap_servers,
479                            "kafka_group_id": config.signal_kafka_group_id,
480                        }),
481                    );
482                if !config.column_transformations.is_empty() {
483                    semantic
484                        .as_object_mut()
485                        .expect("source semantic is an object")
486                        .insert(
487                            "column_transformations".into(),
488                            serde_json::Value::Array(column_transformation_semantics(
489                                &config.column_transformations,
490                            )),
491                        );
492                }
493                semantic
494            }
495            Self::Sqlserver(config) => {
496                let mut semantic = serde_json::json!({
497                "type": "sqlserver",
498                "hostname": config.hostname,
499                "port": config.port,
500                "databases": config.databases,
501                "tables": config.tables,
502                "encrypt": config.encrypt,
503                });
504                add_heartbeat_semantics(
505                    &mut semantic,
506                    config.heartbeat_interval,
507                    config.heartbeat_action_query.as_deref(),
508                    &config.heartbeat_topics_prefix,
509                    config.heartbeat_topic_name.as_deref(),
510                );
511                semantic
512                    .as_object_mut()
513                    .expect("source semantic is an object")
514                    .insert(
515                        "signals".into(),
516                        serde_json::json!({
517                            "data_collection": config.signal_data_collection,
518                            "enabled_channels": config.signal_enabled_channels,
519                            "file": config.signal_file,
520                            "poll_interval_ms": config.signal_poll_interval.as_millis(),
521                            "chunk_size": config.incremental_snapshot_chunk_size,
522                            "watermarking_strategy": config.incremental_snapshot_watermarking_strategy,
523                            "kafka_topic": config.signal_kafka_topic,
524                            "kafka_bootstrap_servers": config.signal_kafka_bootstrap_servers,
525                            "kafka_group_id": config.signal_kafka_group_id,
526                        }),
527                    );
528                if !config.column_transformations.is_empty() {
529                    semantic
530                        .as_object_mut()
531                        .expect("source semantic is an object")
532                        .insert(
533                            "column_transformations".into(),
534                            serde_json::Value::Array(column_transformation_semantics(
535                                &config.column_transformations,
536                            )),
537                        );
538                }
539                semantic
540            }
541            Self::Oracle(config) => {
542                let mut semantic = serde_json::json!({
543                    "type": "oracle",
544                    "hostname": config.hostname,
545                    "port": config.port,
546                    "database": config.database,
547                    "pdb_name": config.pdb_name,
548                    "schemas": config.schemas,
549                    "tables": config.tables,
550                    "log_mining_strategy": config.log_mining_strategy,
551                    "archive_log_only_mode": config.archive_log_only_mode,
552                });
553                add_heartbeat_semantics(
554                    &mut semantic,
555                    config.heartbeat_interval,
556                    config.heartbeat_action_query.as_deref(),
557                    &config.heartbeat_topics_prefix,
558                    config.heartbeat_topic_name.as_deref(),
559                );
560                semantic
561            }
562            Self::Mongodb(config) => {
563                let mut semantic = serde_json::json!({
564                    "type": "mongodb",
565                    "connection_string": config.connection_string,
566                    "databases": config.databases,
567                    "collections": config.collections,
568                    "full_document": config.full_document,
569                    "full_document_before_change": config.full_document_before_change,
570                });
571                add_heartbeat_semantics(
572                    &mut semantic,
573                    config.heartbeat_interval,
574                    None,
575                    &config.heartbeat_topics_prefix,
576                    config.heartbeat_topic_name.as_deref(),
577                );
578                semantic
579            }
580            Self::Mariadb(config) => config.semantic_config(DebeziumConnectorKind::Mariadb),
581            Self::Db2(config) => config.semantic_config(DebeziumConnectorKind::Db2),
582            Self::Cassandra(config) => config.semantic_config(DebeziumConnectorKind::Cassandra),
583            Self::Vitess(config) => config.semantic_config(DebeziumConnectorKind::Vitess),
584            Self::Spanner(config) => config.semantic_config(DebeziumConnectorKind::Spanner),
585            Self::Informix(config) => config.semantic_config(DebeziumConnectorKind::Informix),
586            Self::Cockroachdb(config) => config.semantic_config(DebeziumConnectorKind::CockroachDb),
587            Self::Yashandb(config) => config.semantic_config(DebeziumConnectorKind::YashanDb),
588        }
589    }
590}
591
592#[derive(Debug, Clone, Serialize, Deserialize)]
593#[serde(deny_unknown_fields)]
594pub struct PostgresSourceConfig {
595    #[serde(default = "default_hostname")]
596    pub hostname: String,
597    #[serde(default = "default_postgres_port")]
598    pub port: u16,
599    pub database: String,
600    pub username: String,
601    pub password: String,
602    pub publication: String,
603    #[serde(default)]
604    pub publication_autocreate_mode: PublicationAutoCreateMode,
605    #[serde(default)]
606    pub replica_identity_autoset_values: Vec<PostgresReplicaIdentityRule>,
607    #[serde(default)]
608    pub publish_via_partition_root: bool,
609    #[serde(default = "default_slot_name")]
610    pub slot_name: String,
611    #[serde(default)]
612    pub drop_slot_on_stop: bool,
613    #[serde(default)]
614    pub slot_failover: bool,
615    #[serde(default)]
616    pub slot_ownership: SlotOwnership,
617    #[serde(default)]
618    pub offset_mismatch_strategy: PostgresOffsetMismatchStrategy,
619    #[serde(default)]
620    pub lsn_flush_mode: PostgresLsnFlushMode,
621    #[serde(default = "default_postgres_lsn_flush_timeout")]
622    #[serde(with = "humantime_serde")]
623    pub lsn_flush_timeout: Duration,
624    #[serde(default)]
625    pub lsn_flush_timeout_action: PostgresLsnFlushTimeoutAction,
626    #[serde(default)]
627    pub slot_stream_params: BTreeMap<String, String>,
628    #[serde(default)]
629    pub database_initial_statements: Vec<String>,
630    #[serde(default)]
631    pub snapshot_locking_mode: PostgresSnapshotLockingMode,
632    #[serde(default = "default_postgres_snapshot_lock_timeout")]
633    #[serde(with = "humantime_serde")]
634    pub snapshot_lock_timeout: Duration,
635    #[serde(default)]
636    pub snapshot_isolation_mode: PostgresSnapshotIsolationMode,
637    #[serde(default)]
638    #[serde(with = "humantime_serde")]
639    pub xmin_fetch_interval: Duration,
640    #[serde(default)]
641    pub tables: TableSelection,
642    #[serde(default = "default_ssl_mode")]
643    pub ssl_mode: String,
644    #[serde(default)]
645    pub ssl_root_cert: Option<String>,
646    #[serde(default)]
647    pub ssl_cert: Option<String>,
648    #[serde(default)]
649    pub ssl_key: Option<String>,
650    #[serde(default)]
651    pub ssl_key_password: Option<String>,
652    #[serde(default = "default_connect_timeout")]
653    #[serde(with = "humantime_serde")]
654    pub connect_timeout: Duration,
655    #[serde(default = "default_postgres_status_update_interval")]
656    #[serde(with = "humantime_serde")]
657    pub status_update_interval: Duration,
658    #[serde(default = "default_true")]
659    pub tcp_keepalive: bool,
660    #[serde(default)]
661    #[serde(with = "humantime_serde")]
662    pub heartbeat_interval: Duration,
663    #[serde(default)]
664    pub heartbeat_action_query: Option<String>,
665    #[serde(default = "default_heartbeat_topics_prefix")]
666    pub heartbeat_topics_prefix: String,
667    #[serde(default)]
668    pub heartbeat_topic_name: Option<String>,
669    #[serde(default)]
670    pub signal_data_collection: Option<String>,
671    #[serde(default = "default_signal_enabled_channels")]
672    pub signal_enabled_channels: Vec<String>,
673    #[serde(default = "default_signal_file")]
674    pub signal_file: String,
675    #[serde(default = "default_signal_poll_interval")]
676    #[serde(with = "humantime_serde")]
677    pub signal_poll_interval: Duration,
678    #[serde(default)]
679    pub signal_kafka_topic: Option<String>,
680    #[serde(default)]
681    pub signal_kafka_bootstrap_servers: Vec<String>,
682    #[serde(default = "default_signal_kafka_group_id")]
683    pub signal_kafka_group_id: String,
684    #[serde(default = "default_signal_kafka_poll_timeout")]
685    #[serde(with = "humantime_serde")]
686    pub signal_kafka_poll_timeout: Duration,
687    #[serde(default)]
688    pub signal_kafka_consumer_properties: BTreeMap<String, String>,
689    #[serde(default = "default_incremental_snapshot_chunk_size")]
690    pub incremental_snapshot_chunk_size: usize,
691    #[serde(default = "default_incremental_snapshot_watermarking_strategy")]
692    pub incremental_snapshot_watermarking_strategy: String,
693    #[serde(default)]
694    pub read_only: bool,
695    #[serde(default = "default_hstore_handling_mode")]
696    pub hstore_handling_mode: String,
697    #[serde(default = "default_postgres_interval_handling_mode")]
698    pub interval_handling_mode: String,
699    #[serde(default)]
700    pub include_unknown_datatypes: bool,
701    #[serde(default = "default_postgres_money_fraction_digits")]
702    pub money_fraction_digits: i16,
703    #[serde(default)]
704    pub schema_refresh_mode: PostgresSchemaRefreshMode,
705    #[serde(default)]
706    pub logical_decoding_messages: bool,
707    #[serde(default)]
708    pub message_prefix_include_list: Vec<String>,
709    #[serde(default)]
710    pub message_prefix_exclude_list: Vec<String>,
711    #[serde(default)]
712    pub column_transformations: Vec<ColumnTransformRule>,
713}
714
715impl PostgresSourceConfig {
716    fn validate(&self) -> Result<()> {
717        for (value, field) in [
718            (&self.database, "source.database"),
719            (&self.username, "source.username"),
720            (&self.publication, "source.publication"),
721            (&self.slot_name, "source.slot_name"),
722        ] {
723            if value.trim().is_empty() {
724                return Err(Error::Configuration(format!("{field} must not be empty")));
725            }
726        }
727        validate_name(&self.slot_name, "source.slot_name")?;
728        validate_name(&self.publication, "source.publication")?;
729        if self.drop_slot_on_stop && self.slot_ownership == SlotOwnership::External {
730            return Err(Error::Configuration(
731                "source.drop_slot_on_stop can be enabled only for a managed replication slot"
732                    .into(),
733            ));
734        }
735        if self.slot_failover && self.slot_ownership == SlotOwnership::External {
736            return Err(Error::Configuration(
737                "source.slot_failover can be enabled only for a managed replication slot".into(),
738            ));
739        }
740        if self.slot_ownership == SlotOwnership::External
741            && self.offset_mismatch_strategy.advances_slot()
742        {
743            return Err(Error::Configuration(
744                "source.offset_mismatch_strategy=trust_offset or trust_greater_lsn requires slot_ownership=managed because it can advance the replication slot"
745                    .into(),
746            ));
747        }
748        if self.connect_timeout.is_zero() {
749            return Err(Error::Configuration(
750                "source.connect_timeout must be greater than zero".into(),
751            ));
752        }
753        if self.status_update_interval.is_zero() {
754            return Err(Error::Configuration(
755                "source.status_update_interval must be greater than zero".into(),
756            ));
757        }
758        if self.lsn_flush_timeout.is_zero() {
759            return Err(Error::Configuration(
760                "source.lsn_flush_timeout must be greater than zero".into(),
761            ));
762        }
763        if !matches!(
764            self.ssl_mode.as_str(),
765            "disable" | "allow" | "prefer" | "require" | "verify-ca" | "verify-full"
766        ) {
767            return Err(Error::Configuration(
768                "source.ssl_mode must be disable, allow, prefer, require, verify-ca, or verify-full"
769                    .into(),
770            ));
771        }
772        for (value, field) in [
773            (self.ssl_root_cert.as_deref(), "source.ssl_root_cert"),
774            (self.ssl_cert.as_deref(), "source.ssl_cert"),
775            (self.ssl_key.as_deref(), "source.ssl_key"),
776            (self.ssl_key_password.as_deref(), "source.ssl_key_password"),
777        ] {
778            if value.is_some_and(|value| value.trim().is_empty()) {
779                return Err(Error::Configuration(format!(
780                    "{field} must not be blank when configured"
781                )));
782            }
783        }
784        if self.ssl_cert.is_some() || self.ssl_key.is_some() || self.ssl_key_password.is_some() {
785            return Err(Error::Configuration(
786                "source.ssl_cert, source.ssl_key, and source.ssl_key_password are not supported by the current pg_walstream rustls transport; use server certificate verification with source.ssl_root_cert or omit client-certificate authentication"
787                    .into(),
788            ));
789        }
790        for (name, value) in &self.slot_stream_params {
791            if name != "origin" {
792                return Err(Error::Configuration(format!(
793                    "source.slot_stream_params currently supports only the pgoutput origin parameter; found {name:?}"
794                )));
795            }
796            if !matches!(value.as_str(), "any" | "none") {
797                return Err(Error::Configuration(format!(
798                    "source.slot_stream_params.origin must be any or none; found {value:?}"
799                )));
800            }
801        }
802        if self
803            .database_initial_statements
804            .iter()
805            .any(|statement| statement.trim().is_empty())
806        {
807            return Err(Error::Configuration(
808                "source.database_initial_statements must not contain blank statements".into(),
809            ));
810        }
811        if self.snapshot_lock_timeout.as_millis() > i32::MAX as u128 {
812            return Err(Error::Configuration(format!(
813                "source.snapshot_lock_timeout must not exceed {}ms",
814                i32::MAX
815            )));
816        }
817        for pattern in self.tables.include.iter().chain(self.tables.exclude.iter()) {
818            if Regex::new(pattern).is_err() {
819                return Err(Error::Configuration(format!(
820                    "table selector {pattern:?} is not a valid regular expression"
821                )));
822            }
823        }
824        for rule in &self.replica_identity_autoset_values {
825            if rule.table.trim().is_empty() || Regex::new(&rule.table).is_err() {
826                return Err(Error::Configuration(format!(
827                    "PostgreSQL replica identity table selector {:?} is not a valid regular expression",
828                    rule.table
829                )));
830            }
831            match (&rule.identity, rule.index.as_deref()) {
832                (PostgresReplicaIdentity::Index, Some(index)) if !index.trim().is_empty() => {}
833                (PostgresReplicaIdentity::Index, _) => {
834                    return Err(Error::Configuration(
835                        "PostgreSQL replica identity index mode requires a non-empty index name"
836                            .into(),
837                    ));
838                }
839                (_, Some(_)) => {
840                    return Err(Error::Configuration(
841                        "PostgreSQL replica identity index is valid only with identity=index"
842                            .into(),
843                    ));
844                }
845                (_, None) => {}
846            }
847        }
848        validate_heartbeat(
849            &self.heartbeat_topics_prefix,
850            self.heartbeat_topic_name.as_deref(),
851            self.heartbeat_action_query.as_deref(),
852        )?;
853        if self.incremental_snapshot_chunk_size == 0 {
854            return Err(Error::Configuration(
855                "source.incremental_snapshot_chunk_size must be greater than zero".into(),
856            ));
857        }
858        if self.incremental_snapshot_watermarking_strategy != "insert_insert" {
859            return Err(Error::Configuration(
860                "source.incremental_snapshot_watermarking_strategy currently supports only insert_insert"
861                    .into(),
862            ));
863        }
864        if !matches!(self.hstore_handling_mode.as_str(), "json" | "map") {
865            return Err(Error::Configuration(
866                "source.hstore_handling_mode must be json or map".into(),
867            ));
868        }
869        if !matches!(
870            self.interval_handling_mode.as_str(),
871            "postgres" | "numeric" | "string"
872        ) {
873            return Err(Error::Configuration(
874                "source.interval_handling_mode must be postgres, numeric, or string".into(),
875            ));
876        }
877        if !self.message_prefix_include_list.is_empty()
878            && !self.message_prefix_exclude_list.is_empty()
879        {
880            return Err(Error::Configuration(
881                "source.message_prefix_include_list and source.message_prefix_exclude_list are mutually exclusive"
882                    .into(),
883            ));
884        }
885        for pattern in self
886            .message_prefix_include_list
887            .iter()
888            .chain(&self.message_prefix_exclude_list)
889        {
890            if Regex::new(pattern).is_err() {
891                return Err(Error::Configuration(format!(
892                    "PostgreSQL logical decoding message prefix selector {pattern:?} is not a valid regular expression"
893                )));
894            }
895        }
896        for rule in &self.column_transformations {
897            rule.validate()?;
898        }
899        if self.signal_poll_interval.is_zero() {
900            return Err(Error::Configuration(
901                "source.signal_poll_interval must be greater than zero".into(),
902            ));
903        }
904        if self.signal_enabled_channels.iter().any(|channel| {
905            channel.trim().is_empty()
906                || !matches!(channel.as_str(), "source" | "file" | "in-process" | "kafka")
907        }) {
908            return Err(Error::Configuration(
909                "source.signal_enabled_channels currently supports source, file, in-process, and kafka"
910                    .into(),
911            ));
912        }
913        if self
914            .signal_enabled_channels
915            .iter()
916            .any(|channel| channel == "file")
917            && self.signal_file.trim().is_empty()
918        {
919            return Err(Error::Configuration(
920                "source.signal_file must not be empty when the file signal channel is enabled"
921                    .into(),
922            ));
923        }
924        if self
925            .signal_enabled_channels
926            .iter()
927            .any(|channel| channel == "kafka")
928        {
929            if self.signal_kafka_bootstrap_servers.is_empty()
930                || self
931                    .signal_kafka_bootstrap_servers
932                    .iter()
933                    .any(|server| server.trim().is_empty())
934            {
935                return Err(Error::Configuration(
936                    "source.signal_kafka_bootstrap_servers must contain at least one server when the kafka signal channel is enabled"
937                        .into(),
938                ));
939            }
940            if self.signal_kafka_group_id.trim().is_empty() {
941                return Err(Error::Configuration(
942                    "source.signal_kafka_group_id must not be empty when the kafka signal channel is enabled"
943                        .into(),
944                ));
945            }
946            if self
947                .signal_kafka_topic
948                .as_deref()
949                .is_some_and(|topic| topic.trim().is_empty())
950            {
951                return Err(Error::Configuration(
952                    "source.signal_kafka_topic must not be empty when configured".into(),
953                ));
954            }
955            for property in ["enable.auto.commit", "enable.auto.offset.store"] {
956                if self
957                    .signal_kafka_consumer_properties
958                    .get(property)
959                    .is_some_and(|value| !value.eq_ignore_ascii_case("false"))
960                {
961                    return Err(Error::Configuration(format!(
962                        "source.signal_kafka_consumer_properties.{property} must be false so signal offsets follow Rustium checkpoints"
963                    )));
964                }
965            }
966        }
967        if let Some(collection) = &self.signal_data_collection {
968            let mut parts = collection.split('.');
969            let schema = parts.next().unwrap_or_default();
970            let table = parts.next().unwrap_or_default();
971            if schema.is_empty() || table.is_empty() || parts.next().is_some() {
972                return Err(Error::Configuration(
973                    "source.signal_data_collection must be a schema-qualified PostgreSQL table"
974                        .into(),
975                ));
976            }
977            validate_name(schema, "source.signal_data_collection schema")?;
978            validate_name(table, "source.signal_data_collection table")?;
979        }
980        Ok(())
981    }
982
983    #[must_use]
984    pub fn captures_logical_decoding_messages(&self) -> bool {
985        self.logical_decoding_messages
986            || !self.message_prefix_include_list.is_empty()
987            || !self.message_prefix_exclude_list.is_empty()
988    }
989
990    #[must_use]
991    pub fn includes_message_prefix(&self, prefix: &str) -> bool {
992        if !self.captures_logical_decoding_messages() {
993            return false;
994        }
995        let included = self.message_prefix_include_list.is_empty()
996            || self
997                .message_prefix_include_list
998                .iter()
999                .any(|pattern| regex_matches(pattern, prefix));
1000        included
1001            && !self
1002                .message_prefix_exclude_list
1003                .iter()
1004                .any(|pattern| regex_matches(pattern, prefix))
1005    }
1006
1007    pub fn connection_url(&self, replication: bool) -> Result<String> {
1008        let mut url = Url::parse("postgresql://localhost")
1009            .map_err(|error| Error::Configuration(error.to_string()))?;
1010        url.set_host(Some(&self.hostname))
1011            .map_err(|_| Error::Configuration("invalid source.hostname".into()))?;
1012        url.set_port(Some(self.port))
1013            .map_err(|_| Error::Configuration("invalid source.port".into()))?;
1014        url.set_username(&self.username)
1015            .map_err(|_| Error::Configuration("invalid source.username".into()))?;
1016        url.set_password(Some(&self.password))
1017            .map_err(|_| Error::Configuration("invalid source.password".into()))?;
1018        url.set_path(&self.database);
1019        {
1020            let mut query = url.query_pairs_mut();
1021            query.append_pair("sslmode", &self.ssl_mode);
1022            if let Some(root_cert) = &self.ssl_root_cert {
1023                query.append_pair("sslrootcert", root_cert);
1024            }
1025            query.append_pair(
1026                "connect_timeout",
1027                &self.connect_timeout.as_secs().max(1).to_string(),
1028            );
1029            query.append_pair("keepalives", if self.tcp_keepalive { "1" } else { "0" });
1030            if replication {
1031                query.append_pair("replication", "database");
1032            }
1033        }
1034        Ok(url.into())
1035    }
1036}
1037
1038#[derive(Debug, Clone, Serialize, Deserialize)]
1039#[serde(deny_unknown_fields)]
1040pub struct MySqlSourceConfig {
1041    #[serde(default = "default_hostname")]
1042    pub hostname: String,
1043    #[serde(default = "default_mysql_port")]
1044    pub port: u16,
1045    pub username: String,
1046    pub password: String,
1047    #[serde(default)]
1048    pub databases: Vec<String>,
1049    #[serde(default = "default_mysql_server_id")]
1050    pub server_id: u32,
1051    #[serde(default)]
1052    pub tables: TableSelection,
1053    #[serde(default = "default_mysql_ssl_mode")]
1054    pub ssl_mode: String,
1055    #[serde(default = "default_mysql_connection_time_zone")]
1056    pub connection_time_zone: String,
1057    #[serde(default)]
1058    pub ssl_ca: Option<String>,
1059    #[serde(default)]
1060    pub ssl_cert: Option<String>,
1061    #[serde(default)]
1062    pub ssl_key: Option<String>,
1063    #[serde(default)]
1064    pub ssl_keystore: Option<String>,
1065    #[serde(default)]
1066    pub ssl_keystore_password: Option<String>,
1067    #[serde(default)]
1068    pub ssl_truststore: Option<String>,
1069    #[serde(default)]
1070    pub ssl_truststore_password: Option<String>,
1071    #[serde(default = "default_connect_timeout")]
1072    #[serde(with = "humantime_serde")]
1073    pub connect_timeout: Duration,
1074    #[serde(default = "default_true")]
1075    pub connect_keep_alive: bool,
1076    #[serde(default = "default_mysql_keep_alive_interval")]
1077    #[serde(with = "humantime_serde")]
1078    pub connect_keep_alive_interval: Duration,
1079    #[serde(default = "default_mysql_reconnect_max_attempts")]
1080    pub reconnect_max_attempts: u32,
1081    #[serde(default)]
1082    pub schema_history_skip_unparseable_ddl: bool,
1083    #[serde(default)]
1084    pub gtid_source_includes: Vec<String>,
1085    #[serde(default)]
1086    pub gtid_source_excludes: Vec<String>,
1087    #[serde(default = "default_true")]
1088    pub gtid_source_filter_dml_events: bool,
1089    #[serde(default)]
1090    #[serde(with = "humantime_serde")]
1091    pub heartbeat_interval: Duration,
1092    #[serde(default)]
1093    pub heartbeat_action_query: Option<String>,
1094    #[serde(default = "default_heartbeat_topics_prefix")]
1095    pub heartbeat_topics_prefix: String,
1096    #[serde(default)]
1097    pub heartbeat_topic_name: Option<String>,
1098    #[serde(default)]
1099    pub signal_data_collection: Option<String>,
1100    #[serde(default = "default_signal_enabled_channels")]
1101    pub signal_enabled_channels: Vec<String>,
1102    #[serde(default = "default_signal_file")]
1103    pub signal_file: String,
1104    #[serde(default = "default_signal_poll_interval")]
1105    #[serde(with = "humantime_serde")]
1106    pub signal_poll_interval: Duration,
1107    #[serde(default = "default_incremental_snapshot_chunk_size")]
1108    pub incremental_snapshot_chunk_size: usize,
1109    #[serde(default = "default_incremental_snapshot_watermarking_strategy")]
1110    pub incremental_snapshot_watermarking_strategy: String,
1111    #[serde(default)]
1112    pub signal_kafka_topic: Option<String>,
1113    #[serde(default)]
1114    pub signal_kafka_bootstrap_servers: Vec<String>,
1115    #[serde(default = "default_signal_kafka_group_id")]
1116    pub signal_kafka_group_id: String,
1117    #[serde(default = "default_signal_kafka_poll_timeout")]
1118    #[serde(with = "humantime_serde")]
1119    pub signal_kafka_poll_timeout: Duration,
1120    #[serde(default)]
1121    pub signal_kafka_consumer_properties: BTreeMap<String, String>,
1122    #[serde(default)]
1123    pub column_transformations: Vec<ColumnTransformRule>,
1124}
1125
1126impl MySqlSourceConfig {
1127    fn validate(&self) -> Result<()> {
1128        if self.username.trim().is_empty() {
1129            return Err(Error::Configuration(
1130                "source.username must not be empty".into(),
1131            ));
1132        }
1133        if self.server_id == 0 {
1134            return Err(Error::Configuration(
1135                "source.server_id must be greater than zero".into(),
1136            ));
1137        }
1138        if self.connect_keep_alive_interval.is_zero() {
1139            return Err(Error::Configuration(
1140                "source.connect_keep_alive_interval must be greater than zero".into(),
1141            ));
1142        }
1143        if self.connect_keep_alive && self.reconnect_max_attempts == 0 {
1144            return Err(Error::Configuration(
1145                "source.reconnect_max_attempts must be greater than zero when connect_keep_alive is enabled"
1146                    .into(),
1147            ));
1148        }
1149        self.session_time_zone()?;
1150        validate_heartbeat(
1151            &self.heartbeat_topics_prefix,
1152            self.heartbeat_topic_name.as_deref(),
1153            self.heartbeat_action_query.as_deref(),
1154        )?;
1155        if self.incremental_snapshot_chunk_size == 0 {
1156            return Err(Error::Configuration(
1157                "source.incremental_snapshot_chunk_size must be greater than zero".into(),
1158            ));
1159        }
1160        if self.incremental_snapshot_watermarking_strategy != "insert_insert" {
1161            return Err(Error::Configuration(
1162                "source.incremental_snapshot_watermarking_strategy currently supports only insert_insert".into(),
1163            ));
1164        }
1165        if self.signal_poll_interval.is_zero() {
1166            return Err(Error::Configuration(
1167                "source.signal_poll_interval must be greater than zero".into(),
1168            ));
1169        }
1170        if self.signal_enabled_channels.iter().any(|channel| {
1171            channel.trim().is_empty()
1172                || !matches!(channel.as_str(), "source" | "file" | "in-process" | "kafka")
1173        }) {
1174            return Err(Error::Configuration(
1175                "source.signal_enabled_channels currently supports source, file, in-process, and kafka for MySQL".into(),
1176            ));
1177        }
1178        if self
1179            .signal_enabled_channels
1180            .iter()
1181            .any(|channel| channel == "file")
1182            && self.signal_file.trim().is_empty()
1183        {
1184            return Err(Error::Configuration(
1185                "source.signal_file must not be empty when the file signal channel is enabled"
1186                    .into(),
1187            ));
1188        }
1189        if let Some(collection) = &self.signal_data_collection {
1190            let mut parts = collection.split('.');
1191            let database = parts.next().unwrap_or_default();
1192            let table = parts.next().unwrap_or_default();
1193            if database.is_empty() || table.is_empty() || parts.next().is_some() {
1194                return Err(Error::Configuration(
1195                    "source.signal_data_collection must be a database-qualified MySQL table".into(),
1196                ));
1197            }
1198            validate_name(database, "source.signal_data_collection database")?;
1199            validate_name(table, "source.signal_data_collection table")?;
1200        }
1201        if self
1202            .signal_enabled_channels
1203            .iter()
1204            .any(|channel| channel == "kafka")
1205        {
1206            if self.signal_kafka_bootstrap_servers.is_empty()
1207                || self
1208                    .signal_kafka_bootstrap_servers
1209                    .iter()
1210                    .any(|server| server.trim().is_empty())
1211            {
1212                return Err(Error::Configuration(
1213                    "source.signal_kafka_bootstrap_servers must contain at least one server when the kafka signal channel is enabled".into(),
1214                ));
1215            }
1216            if self.signal_kafka_group_id.trim().is_empty() {
1217                return Err(Error::Configuration(
1218                    "source.signal_kafka_group_id must not be empty when the kafka signal channel is enabled".into(),
1219                ));
1220            }
1221            if self
1222                .signal_kafka_topic
1223                .as_deref()
1224                .is_some_and(|topic| topic.trim().is_empty())
1225            {
1226                return Err(Error::Configuration(
1227                    "source.signal_kafka_topic must not be empty when configured".into(),
1228                ));
1229            }
1230            for property in ["enable.auto.commit", "enable.auto.offset.store"] {
1231                if self
1232                    .signal_kafka_consumer_properties
1233                    .get(property)
1234                    .is_some_and(|value| !value.eq_ignore_ascii_case("false"))
1235                {
1236                    return Err(Error::Configuration(format!(
1237                        "source.signal_kafka_consumer_properties.{property} must be false so signal offsets follow Rustium checkpoints"
1238                    )));
1239                }
1240            }
1241        }
1242        for rule in &self.column_transformations {
1243            rule.validate()?;
1244        }
1245        if !matches!(
1246            self.ssl_mode.as_str(),
1247            "disabled" | "preferred" | "required" | "verify_ca" | "verify_identity"
1248        ) {
1249            return Err(Error::Configuration(
1250                "source.ssl_mode must be one of disabled, preferred, required, verify_ca, or verify_identity"
1251                .into(),
1252            ));
1253        }
1254        if self.ssl_cert.is_some() != self.ssl_key.is_some() {
1255            return Err(Error::Configuration(
1256                "source.ssl_cert and source.ssl_key must be configured together".into(),
1257            ));
1258        }
1259        if self.ssl_keystore.is_some() && (self.ssl_cert.is_some() || self.ssl_key.is_some()) {
1260            return Err(Error::Configuration(
1261                "source.ssl_keystore cannot be combined with source.ssl_cert/source.ssl_key".into(),
1262            ));
1263        }
1264        if self.ssl_truststore.is_some() && self.ssl_ca.is_some() {
1265            return Err(Error::Configuration(
1266                "source.ssl_truststore cannot be combined with source.ssl_ca".into(),
1267            ));
1268        }
1269        if self.ssl_keystore.is_none() && self.ssl_keystore_password.is_some() {
1270            return Err(Error::Configuration(
1271                "source.ssl_keystore_password requires source.ssl_keystore".into(),
1272            ));
1273        }
1274        if self.ssl_truststore.is_none() && self.ssl_truststore_password.is_some() {
1275            return Err(Error::Configuration(
1276                "source.ssl_truststore_password requires source.ssl_truststore".into(),
1277            ));
1278        }
1279        if self.ssl_mode == "disabled"
1280            && (self.ssl_ca.is_some()
1281                || self.ssl_cert.is_some()
1282                || self.ssl_key.is_some()
1283                || self.ssl_keystore.is_some()
1284                || self.ssl_truststore.is_some())
1285        {
1286            return Err(Error::Configuration(
1287                "source.ssl_ca/ssl_cert/ssl_key/ssl_keystore/ssl_truststore require an enabled MySQL TLS mode".into(),
1288            ));
1289        }
1290        if self
1291            .databases
1292            .iter()
1293            .any(|database| database.trim().is_empty())
1294        {
1295            return Err(Error::Configuration(
1296                "source.databases must not contain empty names".into(),
1297            ));
1298        }
1299        if !self.gtid_source_includes.is_empty() && !self.gtid_source_excludes.is_empty() {
1300            return Err(Error::Configuration(
1301                "source.gtid_source_includes and source.gtid_source_excludes cannot both be configured"
1302                    .into(),
1303            ));
1304        }
1305        for pattern in self
1306            .gtid_source_includes
1307            .iter()
1308            .chain(self.gtid_source_excludes.iter())
1309        {
1310            if pattern.trim().is_empty() || Regex::new(pattern).is_err() {
1311                return Err(Error::Configuration(format!(
1312                    "GTID source filter {pattern:?} is not a valid regular expression"
1313                )));
1314            }
1315        }
1316        validate_table_patterns(&self.tables)
1317    }
1318
1319    pub fn session_time_zone(&self) -> Result<&'static str> {
1320        let value = self.connection_time_zone.trim();
1321        if value == "+00:00"
1322            || value.eq_ignore_ascii_case("UTC")
1323            || value.eq_ignore_ascii_case("Z")
1324            || value.eq_ignore_ascii_case("Etc/UTC")
1325        {
1326            return Ok("+00:00");
1327        }
1328
1329        Err(Error::Configuration(format!(
1330            "source.connection_time_zone currently supports only UTC, Z, Etc/UTC, or +00:00; found {:?}",
1331            self.connection_time_zone
1332        )))
1333    }
1334
1335    pub fn connection_url(&self) -> Result<String> {
1336        let mut url = Url::parse("mysql://localhost")
1337            .map_err(|error| Error::Configuration(error.to_string()))?;
1338        url.set_host(Some(&self.hostname))
1339            .map_err(|_| Error::Configuration("invalid source.hostname".into()))?;
1340        url.set_port(Some(self.port))
1341            .map_err(|_| Error::Configuration("invalid source.port".into()))?;
1342        url.set_username(&self.username)
1343            .map_err(|_| Error::Configuration("invalid source.username".into()))?;
1344        url.set_password(Some(&self.password))
1345            .map_err(|_| Error::Configuration("invalid source.password".into()))?;
1346        Ok(url.into())
1347    }
1348}
1349
1350#[derive(Debug, Clone, Serialize, Deserialize)]
1351#[serde(deny_unknown_fields)]
1352pub struct SqlServerSourceConfig {
1353    #[serde(default = "default_hostname")]
1354    pub hostname: String,
1355    #[serde(default = "default_sqlserver_port")]
1356    pub port: u16,
1357    pub username: String,
1358    pub password: String,
1359    pub databases: Vec<String>,
1360    #[serde(default)]
1361    pub tables: TableSelection,
1362    #[serde(default = "default_connect_timeout")]
1363    #[serde(with = "humantime_serde")]
1364    pub connect_timeout: Duration,
1365    #[serde(default = "default_true")]
1366    pub encrypt: bool,
1367    #[serde(default)]
1368    pub trust_server_certificate: bool,
1369    #[serde(default = "default_poll_interval")]
1370    #[serde(with = "humantime_serde")]
1371    pub poll_interval: Duration,
1372    #[serde(default = "default_streaming_fetch_size")]
1373    pub streaming_fetch_size: usize,
1374    #[serde(default = "default_sqlserver_snapshot_isolation")]
1375    pub snapshot_isolation_mode: String,
1376    #[serde(default)]
1377    #[serde(with = "humantime_serde")]
1378    pub heartbeat_interval: Duration,
1379    #[serde(default)]
1380    pub heartbeat_action_query: Option<String>,
1381    #[serde(default = "default_heartbeat_topics_prefix")]
1382    pub heartbeat_topics_prefix: String,
1383    #[serde(default)]
1384    pub heartbeat_topic_name: Option<String>,
1385    #[serde(default)]
1386    pub signal_data_collection: Option<String>,
1387    #[serde(default = "default_signal_enabled_channels")]
1388    pub signal_enabled_channels: Vec<String>,
1389    #[serde(default = "default_signal_file")]
1390    pub signal_file: String,
1391    #[serde(default = "default_signal_poll_interval")]
1392    #[serde(with = "humantime_serde")]
1393    pub signal_poll_interval: Duration,
1394    #[serde(default = "default_incremental_snapshot_chunk_size")]
1395    pub incremental_snapshot_chunk_size: usize,
1396    #[serde(default = "default_incremental_snapshot_watermarking_strategy")]
1397    pub incremental_snapshot_watermarking_strategy: String,
1398    #[serde(default)]
1399    pub signal_kafka_topic: Option<String>,
1400    #[serde(default)]
1401    pub signal_kafka_bootstrap_servers: Vec<String>,
1402    #[serde(default = "default_signal_kafka_group_id")]
1403    pub signal_kafka_group_id: String,
1404    #[serde(default = "default_signal_kafka_poll_timeout")]
1405    #[serde(with = "humantime_serde")]
1406    pub signal_kafka_poll_timeout: Duration,
1407    #[serde(default)]
1408    pub signal_kafka_consumer_properties: BTreeMap<String, String>,
1409    #[serde(default)]
1410    pub column_transformations: Vec<ColumnTransformRule>,
1411}
1412
1413impl SqlServerSourceConfig {
1414    fn validate(&self) -> Result<()> {
1415        if self.username.trim().is_empty() || self.databases.len() != 1 {
1416            return Err(Error::Configuration(
1417                "SQL Server source currently requires exactly one database per connector".into(),
1418            ));
1419        }
1420        if self.streaming_fetch_size == 0 {
1421            return Err(Error::Configuration(
1422                "source.streaming_fetch_size must be greater than zero".into(),
1423            ));
1424        }
1425        validate_heartbeat(
1426            &self.heartbeat_topics_prefix,
1427            self.heartbeat_topic_name.as_deref(),
1428            self.heartbeat_action_query.as_deref(),
1429        )?;
1430        if self.incremental_snapshot_chunk_size == 0 {
1431            return Err(Error::Configuration(
1432                "source.incremental_snapshot_chunk_size must be greater than zero".into(),
1433            ));
1434        }
1435        if self.incremental_snapshot_watermarking_strategy != "insert_insert" {
1436            return Err(Error::Configuration(
1437                "source.incremental_snapshot_watermarking_strategy currently supports only insert_insert"
1438                    .into(),
1439            ));
1440        }
1441        if self.signal_poll_interval.is_zero() {
1442            return Err(Error::Configuration(
1443                "source.signal_poll_interval must be greater than zero".into(),
1444            ));
1445        }
1446        if self.signal_enabled_channels.iter().any(|channel| {
1447            channel.trim().is_empty()
1448                || !matches!(channel.as_str(), "source" | "file" | "in-process" | "kafka")
1449        }) {
1450            return Err(Error::Configuration(
1451                "source.signal_enabled_channels currently supports source, file, in-process, and kafka for SQL Server"
1452                    .into(),
1453            ));
1454        }
1455        if self
1456            .signal_enabled_channels
1457            .iter()
1458            .any(|channel| channel == "file")
1459            && self.signal_file.trim().is_empty()
1460        {
1461            return Err(Error::Configuration(
1462                "source.signal_file must not be empty when the file signal channel is enabled"
1463                    .into(),
1464            ));
1465        }
1466        if let Some(collection) = &self.signal_data_collection {
1467            validate_sqlserver_signal_collection(collection, &self.databases[0])?;
1468        }
1469        if self
1470            .signal_enabled_channels
1471            .iter()
1472            .any(|channel| channel == "kafka")
1473        {
1474            if self.signal_kafka_bootstrap_servers.is_empty()
1475                || self
1476                    .signal_kafka_bootstrap_servers
1477                    .iter()
1478                    .any(|server| server.trim().is_empty())
1479            {
1480                return Err(Error::Configuration(
1481                    "source.signal_kafka_bootstrap_servers must contain at least one server when the kafka signal channel is enabled"
1482                        .into(),
1483                ));
1484            }
1485            if self.signal_kafka_group_id.trim().is_empty() {
1486                return Err(Error::Configuration(
1487                    "source.signal_kafka_group_id must not be empty when the kafka signal channel is enabled"
1488                        .into(),
1489                ));
1490            }
1491            if self
1492                .signal_kafka_topic
1493                .as_deref()
1494                .is_some_and(|topic| topic.trim().is_empty())
1495            {
1496                return Err(Error::Configuration(
1497                    "source.signal_kafka_topic must not be empty when configured".into(),
1498                ));
1499            }
1500            for property in ["enable.auto.commit", "enable.auto.offset.store"] {
1501                if self
1502                    .signal_kafka_consumer_properties
1503                    .get(property)
1504                    .is_some_and(|value| !value.eq_ignore_ascii_case("false"))
1505                {
1506                    return Err(Error::Configuration(format!(
1507                        "source.signal_kafka_consumer_properties.{property} must be false so signal offsets follow Rustium checkpoints"
1508                    )));
1509                }
1510            }
1511        }
1512        for rule in &self.column_transformations {
1513            rule.validate()?;
1514        }
1515        if !matches!(
1516            self.snapshot_isolation_mode.as_str(),
1517            "exclusive" | "snapshot" | "repeatable_read" | "read_committed" | "read_uncommitted"
1518        ) {
1519            return Err(Error::Configuration(
1520                "source.snapshot_isolation_mode is unsupported".into(),
1521            ));
1522        }
1523        validate_table_patterns(&self.tables)
1524    }
1525}
1526
1527#[derive(Debug, Clone, Serialize, Deserialize)]
1528#[serde(deny_unknown_fields)]
1529pub struct OracleSourceConfig {
1530    #[serde(default = "default_hostname")]
1531    pub hostname: String,
1532    #[serde(default = "default_oracle_port")]
1533    pub port: u16,
1534    pub username: String,
1535    pub password: String,
1536    pub database: String,
1537    #[serde(default)]
1538    pub pdb_name: Option<String>,
1539    #[serde(default)]
1540    pub schemas: Vec<String>,
1541    #[serde(default)]
1542    pub tables: TableSelection,
1543    #[serde(default = "default_connect_timeout")]
1544    #[serde(with = "humantime_serde")]
1545    pub connect_timeout: Duration,
1546    #[serde(default = "default_poll_interval")]
1547    #[serde(with = "humantime_serde")]
1548    pub poll_interval: Duration,
1549    #[serde(default = "default_streaming_fetch_size")]
1550    pub batch_size: usize,
1551    #[serde(default = "default_oracle_log_mining_strategy")]
1552    pub log_mining_strategy: String,
1553    #[serde(default)]
1554    pub archive_log_only_mode: bool,
1555    #[serde(default)]
1556    #[serde(with = "humantime_serde")]
1557    pub heartbeat_interval: Duration,
1558    #[serde(default)]
1559    pub heartbeat_action_query: Option<String>,
1560    #[serde(default = "default_heartbeat_topics_prefix")]
1561    pub heartbeat_topics_prefix: String,
1562    #[serde(default)]
1563    pub heartbeat_topic_name: Option<String>,
1564}
1565
1566impl OracleSourceConfig {
1567    fn validate(&self) -> Result<()> {
1568        for (value, field) in [
1569            (&self.username, "source.username"),
1570            (&self.database, "source.database"),
1571        ] {
1572            if value.trim().is_empty() {
1573                return Err(Error::Configuration(format!("{field} must not be empty")));
1574            }
1575        }
1576        if self
1577            .pdb_name
1578            .as_deref()
1579            .is_some_and(|value| value.trim().is_empty())
1580        {
1581            return Err(Error::Configuration(
1582                "source.pdb_name must not be blank when configured".into(),
1583            ));
1584        }
1585        if self.connect_timeout.is_zero() || self.poll_interval.is_zero() {
1586            return Err(Error::Configuration(
1587                "Oracle connect_timeout and poll_interval must be greater than zero".into(),
1588            ));
1589        }
1590        if self.batch_size == 0 {
1591            return Err(Error::Configuration(
1592                "source.batch_size must be greater than zero".into(),
1593            ));
1594        }
1595        if self.log_mining_strategy != "online_catalog" {
1596            return Err(Error::Configuration(
1597                "source.log_mining_strategy currently supports only online_catalog".into(),
1598            ));
1599        }
1600        if self.schemas.iter().any(|schema| schema.trim().is_empty()) {
1601            return Err(Error::Configuration(
1602                "source.schemas must not contain blank schema names".into(),
1603            ));
1604        }
1605        validate_heartbeat(
1606            &self.heartbeat_topics_prefix,
1607            self.heartbeat_topic_name.as_deref(),
1608            self.heartbeat_action_query.as_deref(),
1609        )?;
1610        validate_table_patterns(&self.tables)
1611    }
1612}
1613
1614#[derive(Debug, Clone, Serialize, Deserialize)]
1615#[serde(deny_unknown_fields)]
1616pub struct MongoDbSourceConfig {
1617    pub connection_string: String,
1618    #[serde(default)]
1619    pub databases: Vec<String>,
1620    #[serde(default)]
1621    pub collections: TableSelection,
1622    #[serde(default = "default_connect_timeout")]
1623    #[serde(with = "humantime_serde")]
1624    pub connect_timeout: Duration,
1625    #[serde(default = "default_poll_interval")]
1626    #[serde(with = "humantime_serde")]
1627    pub poll_interval: Duration,
1628    #[serde(default = "default_streaming_fetch_size")]
1629    pub batch_size: usize,
1630    #[serde(default = "default_mongodb_full_document")]
1631    pub full_document: String,
1632    #[serde(default = "default_mongodb_full_document_before_change")]
1633    pub full_document_before_change: String,
1634    #[serde(default)]
1635    #[serde(with = "humantime_serde")]
1636    pub heartbeat_interval: Duration,
1637    #[serde(default = "default_heartbeat_topics_prefix")]
1638    pub heartbeat_topics_prefix: String,
1639    #[serde(default)]
1640    pub heartbeat_topic_name: Option<String>,
1641}
1642
1643impl MongoDbSourceConfig {
1644    fn validate(&self) -> Result<()> {
1645        let uri = Url::parse(&self.connection_string).map_err(|error| {
1646            Error::Configuration(format!("invalid MongoDB connection string: {error}"))
1647        })?;
1648        if !matches!(uri.scheme(), "mongodb" | "mongodb+srv") {
1649            return Err(Error::Configuration(
1650                "source.connection_string must use mongodb:// or mongodb+srv://".into(),
1651            ));
1652        }
1653        if self.connect_timeout.is_zero() || self.poll_interval.is_zero() {
1654            return Err(Error::Configuration(
1655                "MongoDB connect_timeout and poll_interval must be greater than zero".into(),
1656            ));
1657        }
1658        if self.batch_size == 0 {
1659            return Err(Error::Configuration(
1660                "source.batch_size must be greater than zero".into(),
1661            ));
1662        }
1663        if self
1664            .databases
1665            .iter()
1666            .any(|database| database.trim().is_empty())
1667        {
1668            return Err(Error::Configuration(
1669                "source.databases must not contain blank database names".into(),
1670            ));
1671        }
1672        if !matches!(self.full_document.as_str(), "default" | "update_lookup") {
1673            return Err(Error::Configuration(
1674                "source.full_document must be default or update_lookup".into(),
1675            ));
1676        }
1677        if !matches!(
1678            self.full_document_before_change.as_str(),
1679            "off" | "when_available" | "required"
1680        ) {
1681            return Err(Error::Configuration(
1682                "source.full_document_before_change must be off, when_available, or required"
1683                    .into(),
1684            ));
1685        }
1686        validate_table_patterns(&self.collections)
1687    }
1688}
1689
1690#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1691#[serde(rename_all = "snake_case")]
1692pub enum DebeziumConnectorKind {
1693    Mariadb,
1694    Db2,
1695    Cassandra,
1696    Vitess,
1697    Spanner,
1698    Informix,
1699    CockroachDb,
1700    YashanDb,
1701}
1702
1703impl DebeziumConnectorKind {
1704    #[must_use]
1705    pub const fn source_type(self) -> &'static str {
1706        match self {
1707            Self::Mariadb => "mariadb",
1708            Self::Db2 => "db2",
1709            Self::Cassandra => "cassandra",
1710            Self::Vitess => "vitess",
1711            Self::Spanner => "spanner",
1712            Self::Informix => "informix",
1713            Self::CockroachDb => "cockroachdb",
1714            Self::YashanDb => "yashandb",
1715        }
1716    }
1717
1718    #[must_use]
1719    pub const fn connector_class(self) -> Option<&'static str> {
1720        match self {
1721            Self::Mariadb => Some("io.debezium.connector.mariadb.MariaDbConnector"),
1722            Self::Db2 => Some("io.debezium.connector.db2.Db2Connector"),
1723            // Cassandra is a standalone per-node process rather than a Kafka Connect source.
1724            Self::Cassandra => None,
1725            Self::Vitess => Some("io.debezium.connector.vitess.VitessConnector"),
1726            Self::Spanner => Some("io.debezium.connector.spanner.SpannerConnector"),
1727            Self::Informix => Some("io.debezium.connector.informix.InformixConnector"),
1728            Self::CockroachDb => Some("io.debezium.connector.cockroachdb.CockroachDBConnector"),
1729            Self::YashanDb => Some("io.debezium.connector.yashandb.YashanDbConnector"),
1730        }
1731    }
1732}
1733
1734#[derive(Debug, Clone, Serialize, Deserialize)]
1735#[serde(deny_unknown_fields)]
1736pub struct DebeziumSourceConfig {
1737    #[serde(default)]
1738    pub bridge: DebeziumBridgeConfig,
1739    #[serde(default)]
1740    pub command: Option<String>,
1741    #[serde(default)]
1742    pub command_args: Vec<String>,
1743    #[serde(default)]
1744    pub command_environment: BTreeMap<String, String>,
1745    #[serde(default)]
1746    pub properties: BTreeMap<String, String>,
1747    #[serde(default = "default_debezium_offset_file")]
1748    pub offset_file: String,
1749    #[serde(default = "default_debezium_schema_history_file")]
1750    pub schema_history_file: String,
1751    #[serde(default = "default_heartbeat_topics_prefix")]
1752    pub heartbeat_topics_prefix: String,
1753    #[serde(default)]
1754    pub heartbeat_topic_name: Option<String>,
1755}
1756
1757impl DebeziumSourceConfig {
1758    fn validate(&self, kind: DebeziumConnectorKind) -> Result<()> {
1759        if self
1760            .command
1761            .as_deref()
1762            .is_some_and(|command| command.trim().is_empty())
1763        {
1764            return Err(Error::Configuration(
1765                "source.command must not be blank when configured".into(),
1766            ));
1767        }
1768        if self.command.is_some()
1769            && (self.offset_file.trim().is_empty() || self.schema_history_file.trim().is_empty())
1770        {
1771            return Err(Error::Configuration(
1772                "managed Debezium offset_file and schema_history_file must not be blank".into(),
1773            ));
1774        }
1775        if matches!(self.bridge, DebeziumBridgeConfig::Kafka { .. }) && self.command.is_some() {
1776            return Err(Error::Configuration(
1777                "managed Debezium command execution currently requires bridge.type=http; run Kafka-producing connectors externally"
1778                    .into(),
1779            ));
1780        }
1781        if self.command.is_some()
1782            && matches!(
1783                self.bridge,
1784                DebeziumBridgeConfig::Http {
1785                    authentication_token: Some(_),
1786                    ..
1787                }
1788            )
1789        {
1790            return Err(Error::Configuration(
1791                "managed Debezium HTTP mode does not support authentication_token; keep the default loopback listener or run Debezium behind an authenticated proxy"
1792                    .into(),
1793            ));
1794        }
1795        if self.command.is_some()
1796            && self
1797                .properties
1798                .get("topic.prefix")
1799                .is_none_or(|value| value.trim().is_empty())
1800        {
1801            return Err(Error::Configuration(
1802                "managed Debezium execution requires source.properties.topic.prefix".into(),
1803            ));
1804        }
1805        if kind == DebeziumConnectorKind::Cassandra && self.command.is_some() {
1806            return Err(Error::Configuration(
1807                "the Debezium Cassandra connector is a standalone per-node process; run it externally and use bridge.type=kafka"
1808                    .into(),
1809            ));
1810        }
1811        if let Some(configured) = self.properties.get("connector.class")
1812            && kind
1813                .connector_class()
1814                .is_some_and(|expected| configured != expected)
1815        {
1816            return Err(Error::Configuration(format!(
1817                "source.properties connector.class {configured:?} does not match source.type={} ({:?})",
1818                kind.source_type(),
1819                kind.connector_class().unwrap_or_default()
1820            )));
1821        }
1822        if self.properties.keys().any(|key| {
1823            key.starts_with("debezium.sink.")
1824                || key.starts_with("rustium.")
1825                || key.starts_with("debezium.format.")
1826        }) {
1827            return Err(Error::Configuration(
1828                "source.properties accepts Debezium connector properties only; bridge sink, format, and rustium.* properties are owned by Rustium"
1829                    .into(),
1830            ));
1831        }
1832        self.bridge.validate()?;
1833        validate_heartbeat(
1834            &self.heartbeat_topics_prefix,
1835            self.heartbeat_topic_name.as_deref(),
1836            None,
1837        )
1838    }
1839
1840    fn semantic_config(&self, kind: DebeziumConnectorKind) -> serde_json::Value {
1841        serde_json::json!({
1842            "type": kind.source_type(),
1843            "bridge": self.bridge.semantic_config(),
1844            "command": self.command,
1845            "command_args": self.command_args,
1846            "command_environment_keys": self.command_environment.keys().collect::<Vec<_>>(),
1847            "properties": redacted_properties(&self.properties),
1848            "offset_file": self.offset_file,
1849            "schema_history_file": self.schema_history_file,
1850            "heartbeat_topics_prefix": self.heartbeat_topics_prefix,
1851            "heartbeat_topic_name": self.heartbeat_topic_name,
1852        })
1853    }
1854}
1855
1856#[derive(Debug, Clone, Serialize, Deserialize)]
1857#[serde(tag = "type", rename_all = "snake_case", deny_unknown_fields)]
1858pub enum DebeziumBridgeConfig {
1859    Http {
1860        #[serde(default = "default_debezium_bridge_listen")]
1861        listen: String,
1862        #[serde(default = "default_debezium_bridge_path")]
1863        path: String,
1864        #[serde(default)]
1865        authentication_token: Option<String>,
1866        #[serde(default = "default_debezium_bridge_request_timeout")]
1867        #[serde(with = "humantime_serde")]
1868        request_timeout: Duration,
1869        #[serde(default = "default_debezium_bridge_max_body_size")]
1870        max_body_size: usize,
1871    },
1872    Kafka {
1873        bootstrap_servers: Vec<String>,
1874        topics: Vec<String>,
1875        #[serde(default = "default_debezium_bridge_group_id")]
1876        group_id: String,
1877        #[serde(default)]
1878        consumer_properties: BTreeMap<String, String>,
1879        #[serde(default = "default_debezium_bridge_request_timeout")]
1880        #[serde(with = "humantime_serde")]
1881        poll_timeout: Duration,
1882    },
1883}
1884
1885impl Default for DebeziumBridgeConfig {
1886    fn default() -> Self {
1887        Self::Http {
1888            listen: default_debezium_bridge_listen(),
1889            path: default_debezium_bridge_path(),
1890            authentication_token: None,
1891            request_timeout: default_debezium_bridge_request_timeout(),
1892            max_body_size: default_debezium_bridge_max_body_size(),
1893        }
1894    }
1895}
1896
1897impl DebeziumBridgeConfig {
1898    fn validate(&self) -> Result<()> {
1899        match self {
1900            Self::Http {
1901                listen,
1902                path,
1903                authentication_token,
1904                request_timeout,
1905                max_body_size,
1906            } => {
1907                listen.parse::<std::net::SocketAddr>().map_err(|error| {
1908                    Error::Configuration(format!(
1909                        "source.bridge.listen must be a socket address: {error}"
1910                    ))
1911                })?;
1912                if !path.starts_with('/') || path.contains(['?', '#']) {
1913                    return Err(Error::Configuration(
1914                        "source.bridge.path must be an absolute HTTP path without a query or fragment"
1915                            .into(),
1916                    ));
1917                }
1918                if authentication_token
1919                    .as_deref()
1920                    .is_some_and(|token| token.trim().is_empty())
1921                {
1922                    return Err(Error::Configuration(
1923                        "source.bridge.authentication_token must not be blank".into(),
1924                    ));
1925                }
1926                if request_timeout.is_zero() || *max_body_size == 0 {
1927                    return Err(Error::Configuration(
1928                        "source.bridge request_timeout and max_body_size must be greater than zero"
1929                            .into(),
1930                    ));
1931                }
1932            }
1933            Self::Kafka {
1934                bootstrap_servers,
1935                topics,
1936                group_id,
1937                consumer_properties,
1938                poll_timeout,
1939            } => {
1940                if bootstrap_servers.is_empty()
1941                    || bootstrap_servers
1942                        .iter()
1943                        .any(|server| server.trim().is_empty())
1944                    || topics.is_empty()
1945                    || topics.iter().any(|topic| topic.trim().is_empty())
1946                    || group_id.trim().is_empty()
1947                    || poll_timeout.is_zero()
1948                {
1949                    return Err(Error::Configuration(
1950                        "source.bridge Kafka bootstrap_servers, topics, group_id, and poll_timeout must be non-empty"
1951                            .into(),
1952                    ));
1953                }
1954                for protected in ["group.id", "enable.auto.commit", "enable.auto.offset.store"] {
1955                    if consumer_properties.contains_key(protected) {
1956                        return Err(Error::Configuration(format!(
1957                            "source.bridge.consumer_properties.{protected} is owned by Rustium"
1958                        )));
1959                    }
1960                }
1961            }
1962        }
1963        Ok(())
1964    }
1965
1966    fn semantic_config(&self) -> serde_json::Value {
1967        match self {
1968            Self::Http {
1969                listen,
1970                path,
1971                request_timeout,
1972                max_body_size,
1973                ..
1974            } => serde_json::json!({
1975                "type": "http",
1976                "listen": listen,
1977                "path": path,
1978                "request_timeout_ms": request_timeout.as_millis(),
1979                "max_body_size": max_body_size,
1980            }),
1981            Self::Kafka {
1982                bootstrap_servers,
1983                topics,
1984                group_id,
1985                consumer_properties,
1986                poll_timeout,
1987            } => {
1988                serde_json::json!({
1989                    "type": "kafka",
1990                    "bootstrap_servers": bootstrap_servers,
1991                    "topics": topics,
1992                    "group_id": group_id,
1993                    "consumer_properties": redacted_properties(consumer_properties),
1994                    "poll_timeout_ms": poll_timeout.as_millis(),
1995                })
1996            }
1997        }
1998    }
1999}
2000
2001fn redacted_properties(
2002    properties: &BTreeMap<String, String>,
2003) -> serde_json::Map<String, serde_json::Value> {
2004    properties
2005        .iter()
2006        .map(|(key, value)| {
2007            let value = if sensitive_property(key) {
2008                serde_json::Value::Null
2009            } else {
2010                value.clone().into()
2011            };
2012            (key.clone(), value)
2013        })
2014        .collect()
2015}
2016
2017fn sensitive_property(key: &str) -> bool {
2018    let key = key.to_ascii_lowercase();
2019    ["password", "secret", "token", "credential", "private.key"]
2020        .iter()
2021        .any(|marker| key.contains(marker))
2022}
2023
2024fn validate_sqlserver_signal_collection(collection: &str, database: &str) -> Result<()> {
2025    let parts = collection.split('.').collect::<Vec<_>>();
2026    let (schema, table) = match parts.as_slice() {
2027        [schema, table] => (*schema, *table),
2028        [configured_database, schema, table] if configured_database == &database => {
2029            (*schema, *table)
2030        }
2031        [configured_database, _, _] => {
2032            return Err(Error::Configuration(format!(
2033                "source.signal_data_collection database {configured_database:?} does not match {database:?}"
2034            )));
2035        }
2036        _ => {
2037            return Err(Error::Configuration(
2038                "source.signal_data_collection must be schema.table or database.schema.table for SQL Server"
2039                    .into(),
2040            ));
2041        }
2042    };
2043    validate_name(schema, "source.signal_data_collection schema")?;
2044    validate_name(table, "source.signal_data_collection table")
2045}
2046
2047#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2048#[serde(rename_all = "snake_case")]
2049pub enum SlotOwnership {
2050    #[default]
2051    Managed,
2052    External,
2053}
2054
2055#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2056#[serde(rename_all = "snake_case")]
2057pub enum PostgresOffsetMismatchStrategy {
2058    #[default]
2059    NoValidation,
2060    TrustOffset,
2061    TrustSlot,
2062    TrustGreaterLsn,
2063}
2064
2065#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2066#[serde(rename_all = "snake_case")]
2067pub enum PostgresLsnFlushMode {
2068    #[default]
2069    Connector,
2070    Manual,
2071    ConnectorAndDriver,
2072}
2073
2074#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2075#[serde(rename_all = "snake_case")]
2076pub enum PostgresLsnFlushTimeoutAction {
2077    #[default]
2078    Fail,
2079    Warn,
2080    Ignore,
2081}
2082
2083#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2084#[serde(rename_all = "snake_case")]
2085pub enum PostgresSchemaRefreshMode {
2086    #[default]
2087    ColumnsDiff,
2088    ColumnsDiffExcludeUnchangedToast,
2089}
2090
2091#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2092#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
2093pub enum ColumnTransformRule {
2094    Truncate {
2095        length: u32,
2096        columns: Vec<String>,
2097    },
2098    Mask {
2099        length: u32,
2100        columns: Vec<String>,
2101    },
2102    Hash {
2103        algorithm: ColumnHashAlgorithm,
2104        salt: String,
2105        #[serde(default)]
2106        version: ColumnHashVersion,
2107        columns: Vec<String>,
2108    },
2109}
2110
2111impl ColumnTransformRule {
2112    #[must_use]
2113    pub const fn priority(&self) -> u8 {
2114        match self {
2115            Self::Truncate { .. } => 0,
2116            Self::Mask { .. } => 1,
2117            Self::Hash {
2118                version: ColumnHashVersion::V1,
2119                ..
2120            } => 2,
2121            Self::Hash {
2122                version: ColumnHashVersion::V2,
2123                ..
2124            } => 3,
2125        }
2126    }
2127
2128    #[must_use]
2129    pub fn columns(&self) -> &[String] {
2130        match self {
2131            Self::Truncate { columns, .. }
2132            | Self::Mask { columns, .. }
2133            | Self::Hash { columns, .. } => columns,
2134        }
2135    }
2136
2137    fn validate(&self) -> Result<()> {
2138        let length = match self {
2139            Self::Truncate { length, .. } | Self::Mask { length, .. } => Some(*length),
2140            Self::Hash { salt, .. } => {
2141                if salt.is_empty() {
2142                    return Err(Error::Configuration(
2143                        "source.column_transformations hash salt must not be empty".into(),
2144                    ));
2145                }
2146                None
2147            }
2148        };
2149        if length.is_some_and(|length| length > i32::MAX as u32) {
2150            return Err(Error::Configuration(format!(
2151                "source.column_transformations length must not exceed {}",
2152                i32::MAX
2153            )));
2154        }
2155        if self.columns().is_empty() {
2156            return Err(Error::Configuration(
2157                "source.column_transformations columns must not be empty".into(),
2158            ));
2159        }
2160        for selector in self.columns() {
2161            if selector.trim().is_empty() || Regex::new(&format!("^(?:{selector})$")).is_err() {
2162                return Err(Error::Configuration(format!(
2163                    "column transformation selector {selector:?} is not a valid regular expression"
2164                )));
2165            }
2166        }
2167        Ok(())
2168    }
2169}
2170
2171#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
2172#[serde(rename_all = "snake_case")]
2173pub enum ColumnHashAlgorithm {
2174    Md2,
2175    Md5,
2176    Sha1,
2177    Sha224,
2178    Sha256,
2179    Sha384,
2180    Sha512,
2181    Sha512_224,
2182    Sha512_256,
2183    Sha3_224,
2184    Sha3_256,
2185    Sha3_384,
2186    Sha3_512,
2187}
2188
2189#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2190#[serde(rename_all = "snake_case")]
2191pub enum ColumnHashVersion {
2192    #[default]
2193    V1,
2194    V2,
2195}
2196
2197#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2198#[serde(rename_all = "snake_case")]
2199pub enum PostgresSnapshotLockingMode {
2200    #[default]
2201    None,
2202    Shared,
2203}
2204
2205#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2206#[serde(rename_all = "snake_case")]
2207pub enum PostgresSnapshotIsolationMode {
2208    #[default]
2209    Serializable,
2210    RepeatableRead,
2211    ReadCommitted,
2212    ReadUncommitted,
2213}
2214
2215impl PostgresSnapshotIsolationMode {
2216    #[must_use]
2217    pub const fn imports_snapshot(self) -> bool {
2218        matches!(self, Self::Serializable | Self::RepeatableRead)
2219    }
2220}
2221
2222impl PostgresOffsetMismatchStrategy {
2223    #[must_use]
2224    pub const fn advances_slot(self) -> bool {
2225        matches!(self, Self::TrustOffset | Self::TrustGreaterLsn)
2226    }
2227}
2228
2229#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2230#[serde(rename_all = "snake_case")]
2231pub enum PublicationAutoCreateMode {
2232    #[default]
2233    Disabled,
2234    AllTables,
2235    Filtered,
2236    NoTables,
2237}
2238
2239#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2240#[serde(deny_unknown_fields)]
2241pub struct PostgresReplicaIdentityRule {
2242    pub table: String,
2243    pub identity: PostgresReplicaIdentity,
2244    #[serde(default)]
2245    pub index: Option<String>,
2246}
2247
2248#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
2249#[serde(rename_all = "snake_case")]
2250pub enum PostgresReplicaIdentity {
2251    Default,
2252    Full,
2253    Nothing,
2254    Index,
2255}
2256
2257#[derive(Debug, Clone, Default, Serialize, Deserialize)]
2258#[serde(deny_unknown_fields)]
2259pub struct TableSelection {
2260    #[serde(default)]
2261    pub include: Vec<String>,
2262    #[serde(default)]
2263    pub exclude: Vec<String>,
2264}
2265
2266impl TableSelection {
2267    #[must_use]
2268    pub fn includes(&self, schema: &str, table: &str) -> bool {
2269        let name = format!("{schema}.{table}");
2270        let included = self.include.is_empty()
2271            || self
2272                .include
2273                .iter()
2274                .any(|pattern| regex_matches(pattern, &name));
2275        included
2276            && !self
2277                .exclude
2278                .iter()
2279                .any(|pattern| regex_matches(pattern, &name))
2280    }
2281}
2282
2283fn validate_table_patterns(tables: &TableSelection) -> Result<()> {
2284    for pattern in tables.include.iter().chain(tables.exclude.iter()) {
2285        if Regex::new(pattern).is_err() {
2286            return Err(Error::Configuration(format!(
2287                "table selector {pattern:?} is not a valid regular expression"
2288            )));
2289        }
2290    }
2291    Ok(())
2292}
2293
2294fn regex_matches(pattern: &str, value: &str) -> bool {
2295    Regex::new(&format!("^(?:{pattern})$")).is_ok_and(|pattern| pattern.is_match(value))
2296}
2297
2298#[derive(Debug, Clone, Serialize, Deserialize)]
2299#[serde(deny_unknown_fields)]
2300pub struct SnapshotConfig {
2301    #[serde(default)]
2302    pub mode: SnapshotMode,
2303    #[serde(default = "default_snapshot_fetch_size")]
2304    pub fetch_size: usize,
2305    #[serde(default)]
2306    pub include_collections: Vec<String>,
2307}
2308
2309impl Default for SnapshotConfig {
2310    fn default() -> Self {
2311        Self {
2312            mode: SnapshotMode::Initial,
2313            fetch_size: default_snapshot_fetch_size(),
2314            include_collections: Vec::new(),
2315        }
2316    }
2317}
2318
2319impl SnapshotConfig {
2320    fn validate(&self) -> Result<()> {
2321        for pattern in &self.include_collections {
2322            if Regex::new(pattern).is_err() {
2323                return Err(Error::Configuration(format!(
2324                    "snapshot include selector {pattern:?} is not a valid regular expression"
2325                )));
2326            }
2327        }
2328        Ok(())
2329    }
2330
2331    #[must_use]
2332    pub fn includes_collection(&self, collection: &str) -> bool {
2333        self.include_collections.is_empty()
2334            || self
2335                .include_collections
2336                .iter()
2337                .any(|pattern| regex_matches(pattern, collection))
2338    }
2339
2340    fn semantic_config(&self) -> serde_json::Value {
2341        let mut semantic = serde_json::json!({
2342            "mode": self.mode,
2343            "fetch_size": self.fetch_size,
2344        });
2345        if !self.include_collections.is_empty() {
2346            semantic
2347                .as_object_mut()
2348                .expect("snapshot semantic is an object")
2349                .insert(
2350                    "include_collections".into(),
2351                    serde_json::json!(self.include_collections),
2352                );
2353        }
2354        semantic
2355    }
2356}
2357
2358#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2359#[serde(rename_all = "snake_case")]
2360pub enum SnapshotMode {
2361    #[default]
2362    Initial,
2363    Never,
2364    WhenNeeded,
2365}
2366
2367#[derive(Debug, Clone, Serialize, Deserialize)]
2368#[serde(deny_unknown_fields)]
2369pub struct FormatConfig {
2370    #[serde(rename = "type", default)]
2371    pub kind: FormatType,
2372    #[serde(default = "default_unavailable_value")]
2373    pub unavailable_value: String,
2374    #[serde(default = "default_true")]
2375    pub tombstones_on_delete: bool,
2376    #[serde(default, skip_serializing_if = "Option::is_none")]
2377    pub schema_registry: Option<SchemaRegistryConfig>,
2378}
2379
2380impl Default for FormatConfig {
2381    fn default() -> Self {
2382        Self {
2383            kind: FormatType::DebeziumJson,
2384            unavailable_value: default_unavailable_value(),
2385            tombstones_on_delete: true,
2386            schema_registry: None,
2387        }
2388    }
2389}
2390
2391impl FormatConfig {
2392    fn validate(&self, sink: &SinkConfig) -> Result<()> {
2393        match (self.kind, &self.schema_registry) {
2394            (
2395                FormatType::DebeziumJsonSchema
2396                | FormatType::DebeziumAvro
2397                | FormatType::DebeziumProtobuf,
2398                Some(registry),
2399            ) => {
2400                if !matches!(sink, SinkConfig::Kafka { .. }) {
2401                    return Err(Error::Configuration(format!(
2402                        "format.type={} requires sink.type=kafka",
2403                        self.kind.as_str()
2404                    )));
2405                }
2406                registry.validate()
2407            }
2408            (
2409                FormatType::DebeziumJsonSchema
2410                | FormatType::DebeziumAvro
2411                | FormatType::DebeziumProtobuf,
2412                None,
2413            ) => Err(Error::Configuration(format!(
2414                "format.type={} requires format.schema_registry",
2415                self.kind.as_str()
2416            ))),
2417            (_, Some(_)) => Err(Error::Configuration(
2418                "format.schema_registry is only valid with a schema-registry format".into(),
2419            )),
2420            (_, None) => Ok(()),
2421        }
2422    }
2423
2424    fn semantic_config(&self) -> serde_json::Value {
2425        serde_json::json!({
2426            "type": self.kind,
2427            "unavailable_value": self.unavailable_value,
2428            "tombstones_on_delete": self.tombstones_on_delete,
2429            "schema_registry": self.schema_registry.as_ref().map(SchemaRegistryConfig::semantic_config),
2430        })
2431    }
2432}
2433
2434#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2435#[serde(rename_all = "snake_case")]
2436pub enum FormatType {
2437    RustiumJson,
2438    #[default]
2439    DebeziumJson,
2440    DebeziumJsonSchema,
2441    DebeziumAvro,
2442    DebeziumProtobuf,
2443}
2444
2445impl FormatType {
2446    const fn as_str(self) -> &'static str {
2447        match self {
2448            Self::RustiumJson => "rustium_json",
2449            Self::DebeziumJson => "debezium_json",
2450            Self::DebeziumJsonSchema => "debezium_json_schema",
2451            Self::DebeziumAvro => "debezium_avro",
2452            Self::DebeziumProtobuf => "debezium_protobuf",
2453        }
2454    }
2455}
2456
2457#[derive(Debug, Clone, Serialize, Deserialize)]
2458#[serde(deny_unknown_fields)]
2459pub struct SchemaRegistryConfig {
2460    pub urls: Vec<String>,
2461    #[serde(default)]
2462    pub username: Option<String>,
2463    #[serde(default)]
2464    pub password: Option<String>,
2465    #[serde(default = "default_schema_registry_timeout")]
2466    #[serde(with = "humantime_serde")]
2467    pub request_timeout: Duration,
2468    #[serde(default = "default_schema_registry_cache_capacity")]
2469    pub cache_capacity: usize,
2470}
2471
2472impl SchemaRegistryConfig {
2473    fn validate(&self) -> Result<()> {
2474        if self.urls.is_empty() {
2475            return Err(Error::Configuration(
2476                "format.schema_registry.urls must contain at least one URL".into(),
2477            ));
2478        }
2479        for raw in &self.urls {
2480            let url = Url::parse(raw).map_err(|error| {
2481                Error::Configuration(format!(
2482                    "format.schema_registry URL {raw:?} is invalid: {error}"
2483                ))
2484            })?;
2485            if !matches!(url.scheme(), "http" | "https") || url.host_str().is_none() {
2486                return Err(Error::Configuration(format!(
2487                    "format.schema_registry URL {raw:?} must be an absolute HTTP(S) URL"
2488                )));
2489            }
2490        }
2491        if self.password.is_some() && self.username.is_none() {
2492            return Err(Error::Configuration(
2493                "format.schema_registry.password requires username".into(),
2494            ));
2495        }
2496        if self.request_timeout.is_zero() {
2497            return Err(Error::Configuration(
2498                "format.schema_registry.request_timeout must be greater than zero".into(),
2499            ));
2500        }
2501        if self.cache_capacity == 0 {
2502            return Err(Error::Configuration(
2503                "format.schema_registry.cache_capacity must be greater than zero".into(),
2504            ));
2505        }
2506        Ok(())
2507    }
2508
2509    fn semantic_config(&self) -> serde_json::Value {
2510        serde_json::json!({
2511            "urls": self.urls,
2512            "username": self.username,
2513            "request_timeout_ms": self.request_timeout.as_millis(),
2514            "cache_capacity": self.cache_capacity,
2515        })
2516    }
2517}
2518
2519#[derive(Debug, Clone, Serialize, Deserialize)]
2520#[serde(tag = "type", rename_all = "snake_case", deny_unknown_fields)]
2521pub enum SinkConfig {
2522    Stdout {
2523        #[serde(default = "default_topic_prefix")]
2524        topic_prefix: String,
2525    },
2526    Kafka {
2527        bootstrap_servers: Vec<String>,
2528        topic_prefix: String,
2529        #[serde(default = "default_kafka_acks")]
2530        acks: String,
2531        #[serde(default = "default_kafka_compression")]
2532        compression: String,
2533        #[serde(default = "default_delivery_timeout")]
2534        #[serde(with = "humantime_serde")]
2535        delivery_timeout: Duration,
2536        #[serde(default)]
2537        properties: BTreeMap<String, String>,
2538    },
2539}
2540
2541impl SinkConfig {
2542    fn validate(&self) -> Result<()> {
2543        match self {
2544            Self::Stdout { topic_prefix } => validate_name(topic_prefix, "sink.topic_prefix"),
2545            Self::Kafka {
2546                bootstrap_servers,
2547                topic_prefix,
2548                acks,
2549                delivery_timeout,
2550                properties,
2551                ..
2552            } => {
2553                if bootstrap_servers.is_empty()
2554                    || bootstrap_servers
2555                        .iter()
2556                        .any(|server| server.trim().is_empty())
2557                {
2558                    return Err(Error::Configuration(
2559                        "sink.bootstrap_servers must contain at least one server".into(),
2560                    ));
2561                }
2562                if !matches!(acks.as_str(), "all" | "-1") {
2563                    return Err(Error::Configuration(
2564                        "sink.acks must be all or -1 so Kafka replicates a batch before checkpointing"
2565                            .into(),
2566                    ));
2567                }
2568                if delivery_timeout.is_zero() || delivery_timeout.as_millis() > i32::MAX as u128 {
2569                    return Err(Error::Configuration(format!(
2570                        "sink.delivery_timeout must be between 1 and {} milliseconds",
2571                        i32::MAX
2572                    )));
2573                }
2574                const RESERVED: &[&str] = &[
2575                    "acks",
2576                    "bootstrap.servers",
2577                    "compression.codec",
2578                    "compression.type",
2579                    "delivery.timeout.ms",
2580                    "enable.idempotence",
2581                    "message.timeout.ms",
2582                    "metadata.broker.list",
2583                    "request.required.acks",
2584                ];
2585                if let Some(property) = properties
2586                    .keys()
2587                    .find(|property| RESERVED.contains(&property.as_str()))
2588                {
2589                    return Err(Error::Configuration(format!(
2590                        "sink.properties key {property:?} is managed by Rustium and cannot be overridden"
2591                    )));
2592                }
2593                validate_name(topic_prefix, "sink.topic_prefix")
2594            }
2595        }
2596    }
2597
2598    #[must_use]
2599    pub fn topic_prefix(&self) -> &str {
2600        match self {
2601            Self::Stdout { topic_prefix } | Self::Kafka { topic_prefix, .. } => topic_prefix,
2602        }
2603    }
2604
2605    fn semantic_config(&self) -> serde_json::Value {
2606        match self {
2607            Self::Stdout { topic_prefix } => {
2608                serde_json::json!({"type": "stdout", "topic_prefix": topic_prefix})
2609            }
2610            Self::Kafka {
2611                bootstrap_servers,
2612                topic_prefix,
2613                ..
2614            } => serde_json::json!({
2615                "type": "kafka",
2616                "bootstrap_servers": bootstrap_servers,
2617                "topic_prefix": topic_prefix,
2618            }),
2619        }
2620    }
2621}
2622
2623#[derive(Debug, Clone, Serialize, Deserialize)]
2624#[serde(deny_unknown_fields)]
2625pub struct StateConfig {
2626    #[serde(rename = "type", default)]
2627    pub kind: StateType,
2628    #[serde(default = "default_state_path")]
2629    pub path: String,
2630}
2631
2632impl Default for StateConfig {
2633    fn default() -> Self {
2634        Self {
2635            kind: StateType::Sqlite,
2636            path: default_state_path(),
2637        }
2638    }
2639}
2640
2641#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2642#[serde(rename_all = "snake_case")]
2643pub enum StateType {
2644    #[default]
2645    Sqlite,
2646}
2647
2648#[derive(Debug, Clone, Serialize, Deserialize)]
2649#[serde(deny_unknown_fields)]
2650pub struct RuntimeSettings {
2651    #[serde(default = "default_channel_capacity")]
2652    pub channel_capacity: usize,
2653    #[serde(default = "default_batch_size")]
2654    pub max_batch_size: usize,
2655    #[serde(default = "default_flush_interval")]
2656    #[serde(with = "humantime_serde")]
2657    pub flush_interval: Duration,
2658    #[serde(default = "default_shutdown_timeout")]
2659    #[serde(with = "humantime_serde")]
2660    pub shutdown_timeout: Duration,
2661    #[serde(default = "default_errors_max_retries")]
2662    pub errors_max_retries: i32,
2663    #[serde(default = "default_errors_retry_delay_initial")]
2664    #[serde(with = "humantime_serde")]
2665    pub errors_retry_delay_initial: Duration,
2666    #[serde(default = "default_errors_retry_delay_max")]
2667    #[serde(with = "humantime_serde")]
2668    pub errors_retry_delay_max: Duration,
2669}
2670
2671impl Default for RuntimeSettings {
2672    fn default() -> Self {
2673        Self {
2674            channel_capacity: default_channel_capacity(),
2675            max_batch_size: default_batch_size(),
2676            flush_interval: default_flush_interval(),
2677            shutdown_timeout: default_shutdown_timeout(),
2678            errors_max_retries: default_errors_max_retries(),
2679            errors_retry_delay_initial: default_errors_retry_delay_initial(),
2680            errors_retry_delay_max: default_errors_retry_delay_max(),
2681        }
2682    }
2683}
2684
2685impl RuntimeSettings {
2686    #[must_use]
2687    pub const fn retry_policy(&self) -> RetryPolicy {
2688        RetryPolicy {
2689            max_retries: self.errors_max_retries,
2690            initial_delay: self.errors_retry_delay_initial,
2691            max_delay: self.errors_retry_delay_max,
2692        }
2693    }
2694}
2695
2696#[derive(Debug, Clone, Serialize, Deserialize)]
2697#[serde(deny_unknown_fields)]
2698pub struct ServerConfig {
2699    #[serde(default = "default_bind")]
2700    pub bind: String,
2701    #[serde(default)]
2702    pub enable_mutations: bool,
2703}
2704
2705impl Default for ServerConfig {
2706    fn default() -> Self {
2707        Self {
2708            bind: default_bind(),
2709            enable_mutations: false,
2710        }
2711    }
2712}
2713
2714#[derive(Debug, Clone, Serialize, Deserialize)]
2715#[serde(deny_unknown_fields)]
2716pub struct ObservabilityConfig {
2717    #[serde(default)]
2718    pub log_format: LogFormat,
2719    #[serde(default = "default_log_level")]
2720    pub log_level: String,
2721    #[serde(default = "default_true")]
2722    pub metrics: bool,
2723}
2724
2725impl Default for ObservabilityConfig {
2726    fn default() -> Self {
2727        Self {
2728            log_format: LogFormat::Json,
2729            log_level: default_log_level(),
2730            metrics: true,
2731        }
2732    }
2733}
2734
2735#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2736#[serde(rename_all = "snake_case")]
2737pub enum LogFormat {
2738    #[default]
2739    Json,
2740    Pretty,
2741}
2742
2743fn interpolate_environment(input: &str) -> Result<String> {
2744    let pattern = Regex::new(r"\$\{([A-Za-z_][A-Za-z0-9_]*)(?::-([^}]*))?\}")
2745        .expect("environment interpolation regex is valid");
2746    let mut missing = Vec::new();
2747    let output = pattern.replace_all(input, |captures: &regex::Captures<'_>| {
2748        let name = &captures[1];
2749        match env::var(name) {
2750            Ok(value) => value,
2751            Err(_) => captures.get(2).map_or_else(
2752                || {
2753                    missing.push(name.to_string());
2754                    String::new()
2755                },
2756                |value| value.as_str().to_string(),
2757            ),
2758        }
2759    });
2760    if missing.is_empty() {
2761        Ok(output.into_owned())
2762    } else {
2763        Err(Error::Configuration(format!(
2764            "missing environment variables: {}",
2765            missing.join(", ")
2766        )))
2767    }
2768}
2769
2770fn validate_name(value: &str, field: &str) -> Result<()> {
2771    let valid = Regex::new(r"^[A-Za-z_][A-Za-z0-9_.-]{0,254}$").expect("name regex is valid");
2772    if valid.is_match(value) {
2773        Ok(())
2774    } else {
2775        Err(Error::Configuration(format!(
2776            "{field} contains unsupported characters"
2777        )))
2778    }
2779}
2780
2781fn validate_heartbeat(
2782    topics_prefix: &str,
2783    topic_name: Option<&str>,
2784    action_query: Option<&str>,
2785) -> Result<()> {
2786    if topics_prefix.trim().is_empty() {
2787        return Err(Error::Configuration(
2788            "source.heartbeat_topics_prefix must not be empty".into(),
2789        ));
2790    }
2791    if topic_name.is_some_and(|name| name.trim().is_empty()) {
2792        return Err(Error::Configuration(
2793            "source.heartbeat_topic_name must not be empty when set".into(),
2794        ));
2795    }
2796    if action_query.is_some_and(|query| query.trim().is_empty()) {
2797        return Err(Error::Configuration(
2798            "source.heartbeat_action_query must not be empty when set".into(),
2799        ));
2800    }
2801    Ok(())
2802}
2803
2804fn add_heartbeat_semantics(
2805    semantic: &mut serde_json::Value,
2806    interval: Duration,
2807    action_query: Option<&str>,
2808    topics_prefix: &str,
2809    topic_name: Option<&str>,
2810) {
2811    if interval.is_zero()
2812        && action_query.is_none()
2813        && topics_prefix == "__debezium-heartbeat"
2814        && topic_name.is_none()
2815    {
2816        return;
2817    }
2818    semantic
2819        .as_object_mut()
2820        .expect("source semantic is an object")
2821        .insert(
2822            "heartbeat".into(),
2823            serde_json::json!({
2824                "interval": interval,
2825                "action_query": action_query,
2826                "topics_prefix": topics_prefix,
2827                "topic_name": topic_name,
2828            }),
2829        );
2830}
2831
2832fn column_transformation_semantics(rules: &[ColumnTransformRule]) -> Vec<serde_json::Value> {
2833    let mut ordered = rules.iter().enumerate().collect::<Vec<_>>();
2834    ordered.sort_by_key(|(index, rule)| (rule.priority(), *index));
2835    ordered
2836        .into_iter()
2837        .map(|(_, rule)| match rule {
2838            ColumnTransformRule::Truncate { length, columns } => serde_json::json!({
2839                "kind": "truncate",
2840                "length": length,
2841                "columns": columns,
2842            }),
2843            ColumnTransformRule::Mask { length, columns } => serde_json::json!({
2844                "kind": "mask",
2845                "length": length,
2846                "columns": columns,
2847            }),
2848            ColumnTransformRule::Hash {
2849                algorithm,
2850                salt,
2851                version,
2852                columns,
2853            } => serde_json::json!({
2854                "kind": "hash",
2855                "algorithm": algorithm,
2856                "salt_sha256": hex_digest(salt.as_bytes()),
2857                "version": version,
2858                "columns": columns,
2859            }),
2860        })
2861        .collect()
2862}
2863
2864fn hex_digest(input: &[u8]) -> String {
2865    let mut hasher = Sha256::new();
2866    hasher.update(input);
2867    hasher
2868        .finalize()
2869        .iter()
2870        .map(|byte| format!("{byte:02x}"))
2871        .collect()
2872}
2873
2874fn default_hostname() -> String {
2875    "localhost".into()
2876}
2877const fn default_postgres_port() -> u16 {
2878    5432
2879}
2880const fn default_mysql_port() -> u16 {
2881    3306
2882}
2883const fn default_sqlserver_port() -> u16 {
2884    1433
2885}
2886const fn default_oracle_port() -> u16 {
2887    1521
2888}
2889const fn default_mysql_server_id() -> u32 {
2890    5_401
2891}
2892fn default_mysql_ssl_mode() -> String {
2893    "preferred".into()
2894}
2895fn default_mysql_connection_time_zone() -> String {
2896    "UTC".into()
2897}
2898fn default_mysql_keep_alive_interval() -> Duration {
2899    Duration::from_secs(60)
2900}
2901const fn default_mysql_reconnect_max_attempts() -> u32 {
2902    10
2903}
2904fn default_poll_interval() -> Duration {
2905    Duration::from_millis(500)
2906}
2907const fn default_streaming_fetch_size() -> usize {
2908    10_240
2909}
2910fn default_sqlserver_snapshot_isolation() -> String {
2911    "repeatable_read".into()
2912}
2913fn default_oracle_log_mining_strategy() -> String {
2914    "online_catalog".into()
2915}
2916fn default_mongodb_full_document() -> String {
2917    "update_lookup".into()
2918}
2919fn default_mongodb_full_document_before_change() -> String {
2920    "when_available".into()
2921}
2922fn default_slot_name() -> String {
2923    "rustium".into()
2924}
2925fn default_ssl_mode() -> String {
2926    "prefer".into()
2927}
2928const fn default_connect_timeout() -> Duration {
2929    Duration::from_secs(30)
2930}
2931const fn default_postgres_status_update_interval() -> Duration {
2932    Duration::from_secs(10)
2933}
2934const fn default_postgres_lsn_flush_timeout() -> Duration {
2935    Duration::from_secs(30)
2936}
2937const fn default_postgres_snapshot_lock_timeout() -> Duration {
2938    Duration::from_secs(10)
2939}
2940const fn default_snapshot_fetch_size() -> usize {
2941    10_000
2942}
2943const fn default_incremental_snapshot_chunk_size() -> usize {
2944    1_024
2945}
2946fn default_incremental_snapshot_watermarking_strategy() -> String {
2947    "insert_insert".into()
2948}
2949fn default_hstore_handling_mode() -> String {
2950    "json".into()
2951}
2952fn default_postgres_interval_handling_mode() -> String {
2953    "postgres".into()
2954}
2955const fn default_postgres_money_fraction_digits() -> i16 {
2956    2
2957}
2958fn default_signal_enabled_channels() -> Vec<String> {
2959    vec!["source".into()]
2960}
2961fn default_signal_file() -> String {
2962    "file-signals.txt".into()
2963}
2964const fn default_signal_poll_interval() -> Duration {
2965    Duration::from_secs(5)
2966}
2967fn default_signal_kafka_group_id() -> String {
2968    "kafka-signal".into()
2969}
2970const fn default_signal_kafka_poll_timeout() -> Duration {
2971    Duration::from_millis(100)
2972}
2973fn default_unavailable_value() -> String {
2974    "__rustium_unavailable_value".into()
2975}
2976fn default_heartbeat_topics_prefix() -> String {
2977    "__debezium-heartbeat".into()
2978}
2979fn default_topic_prefix() -> String {
2980    "rustium".into()
2981}
2982fn default_kafka_acks() -> String {
2983    "all".into()
2984}
2985fn default_kafka_compression() -> String {
2986    "lz4".into()
2987}
2988const fn default_delivery_timeout() -> Duration {
2989    Duration::from_secs(30)
2990}
2991
2992const fn default_schema_registry_timeout() -> Duration {
2993    Duration::from_secs(10)
2994}
2995
2996const fn default_schema_registry_cache_capacity() -> usize {
2997    1_000
2998}
2999fn default_state_path() -> String {
3000    "rustium.db".into()
3001}
3002const fn default_channel_capacity() -> usize {
3003    2_048
3004}
3005const fn default_batch_size() -> usize {
3006    512
3007}
3008const fn default_flush_interval() -> Duration {
3009    Duration::from_millis(100)
3010}
3011const fn default_shutdown_timeout() -> Duration {
3012    Duration::from_secs(30)
3013}
3014const fn default_errors_max_retries() -> i32 {
3015    10
3016}
3017const fn default_errors_retry_delay_initial() -> Duration {
3018    Duration::from_millis(300)
3019}
3020const fn default_errors_retry_delay_max() -> Duration {
3021    Duration::from_secs(10)
3022}
3023fn default_bind() -> String {
3024    "127.0.0.1:8080".into()
3025}
3026fn default_debezium_bridge_listen() -> String {
3027    "127.0.0.1:18080".into()
3028}
3029fn default_debezium_bridge_path() -> String {
3030    "/events".into()
3031}
3032const fn default_debezium_bridge_request_timeout() -> Duration {
3033    Duration::from_secs(60)
3034}
3035const fn default_debezium_bridge_max_body_size() -> usize {
3036    16 * 1024 * 1024
3037}
3038fn default_debezium_bridge_group_id() -> String {
3039    "rustium-debezium-bridge".into()
3040}
3041fn default_debezium_offset_file() -> String {
3042    "rustium-debezium.offsets".into()
3043}
3044fn default_debezium_schema_history_file() -> String {
3045    "rustium-debezium-schema-history.dat".into()
3046}
3047fn default_log_level() -> String {
3048    "info".into()
3049}
3050const fn default_true() -> bool {
3051    true
3052}
3053
3054#[cfg(test)]
3055mod tests {
3056    use super::*;
3057
3058    const CONFIG: &str = r#"
3059api_version: rustium.io/v1alpha1
3060kind: Connector
3061metadata:
3062  name: orders-cdc
3063source:
3064  type: postgresql
3065  database: app
3066  username: rustium
3067  password: secret
3068  publication: rustium_pub
3069  slot_name: rustium_orders
3070  tables:
3071    include: [public.orders]
3072sink:
3073  type: stdout
3074  topic_prefix: app
3075"#;
3076
3077    #[test]
3078    fn parses_minimal_config() {
3079        let config = Config::from_yaml(CONFIG).unwrap();
3080        assert_eq!(config.metadata.name, "orders-cdc");
3081        assert_eq!(config.runtime.max_batch_size, 512);
3082        assert_eq!(config.runtime.errors_max_retries, 10);
3083        assert_eq!(
3084            config.runtime.errors_retry_delay_initial,
3085            Duration::from_millis(300)
3086        );
3087        assert_eq!(
3088            config.runtime.errors_retry_delay_max,
3089            Duration::from_secs(10)
3090        );
3091        assert_eq!(config.runtime.retry_policy(), RetryPolicy::default());
3092        assert!(config.format.tombstones_on_delete);
3093        let postgres = config.source.as_postgresql().unwrap();
3094        assert_eq!(
3095            postgres.status_update_interval,
3096            default_postgres_status_update_interval()
3097        );
3098        assert!(postgres.tcp_keepalive);
3099        assert_eq!(
3100            postgres.offset_mismatch_strategy,
3101            PostgresOffsetMismatchStrategy::NoValidation
3102        );
3103        assert!(!postgres.drop_slot_on_stop);
3104        assert_eq!(postgres.lsn_flush_mode, PostgresLsnFlushMode::Connector);
3105        assert_eq!(
3106            postgres.lsn_flush_timeout,
3107            default_postgres_lsn_flush_timeout()
3108        );
3109        assert_eq!(
3110            postgres.lsn_flush_timeout_action,
3111            PostgresLsnFlushTimeoutAction::Fail
3112        );
3113        assert!(postgres.slot_stream_params.is_empty());
3114        assert!(postgres.database_initial_statements.is_empty());
3115        assert_eq!(
3116            postgres.snapshot_locking_mode,
3117            PostgresSnapshotLockingMode::None
3118        );
3119        assert_eq!(
3120            postgres.snapshot_lock_timeout,
3121            default_postgres_snapshot_lock_timeout()
3122        );
3123        assert_eq!(
3124            postgres.snapshot_isolation_mode,
3125            PostgresSnapshotIsolationMode::Serializable
3126        );
3127        assert!(postgres.xmin_fetch_interval.is_zero());
3128        assert!(postgres.ssl_root_cert.is_none());
3129        assert!(!postgres.include_unknown_datatypes);
3130        assert_eq!(
3131            postgres.money_fraction_digits,
3132            default_postgres_money_fraction_digits()
3133        );
3134        assert_eq!(
3135            postgres.schema_refresh_mode,
3136            PostgresSchemaRefreshMode::ColumnsDiff
3137        );
3138        let connection_url = Url::parse(&postgres.connection_url(true).unwrap()).unwrap();
3139        assert!(
3140            connection_url
3141                .query_pairs()
3142                .any(|(key, value)| key == "replication" && value == "database")
3143        );
3144        assert!(
3145            connection_url
3146                .query_pairs()
3147                .any(|(key, value)| key == "keepalives" && value == "1")
3148        );
3149        assert_eq!(
3150            config.source.as_postgresql().unwrap().hstore_handling_mode,
3151            "json"
3152        );
3153        assert_eq!(
3154            config
3155                .source
3156                .as_postgresql()
3157                .unwrap()
3158                .interval_handling_mode,
3159            "postgres"
3160        );
3161        assert!(
3162            !config
3163                .source
3164                .as_postgresql()
3165                .unwrap()
3166                .captures_logical_decoding_messages()
3167        );
3168        assert_eq!(
3169            config
3170                .source
3171                .as_postgresql()
3172                .unwrap()
3173                .publication_autocreate_mode,
3174            PublicationAutoCreateMode::Disabled
3175        );
3176        assert!(
3177            config.source.semantic_config()["publication_autocreate_mode"].is_null(),
3178            "the native default must preserve the pre-autocreate fingerprint shape"
3179        );
3180        assert!(
3181            config.source.semantic_config()["replica_identity_autoset_values"].is_null(),
3182            "empty replica identity rules must preserve the old fingerprint shape"
3183        );
3184        assert!(
3185            config.source.semantic_config()["publish_via_partition_root"].is_null(),
3186            "disabled partition-root publication must preserve the old fingerprint shape"
3187        );
3188        assert!(
3189            config.source.semantic_config()["slot_failover"].is_null(),
3190            "disabled failover slots must preserve the old fingerprint shape"
3191        );
3192        assert!(
3193            config.source.semantic_config()["offset_mismatch_strategy"].is_null(),
3194            "no_validation must preserve the old fingerprint shape"
3195        );
3196        assert!(
3197            config.source.semantic_config()["lsn_flush_mode"].is_null(),
3198            "connector LSN flushing must preserve the old fingerprint shape"
3199        );
3200        assert!(
3201            config.source.semantic_config()["snapshot_isolation_mode"].is_null(),
3202            "the default PostgreSQL snapshot isolation mode must preserve the old fingerprint shape"
3203        );
3204        assert!(
3205            config.source.semantic_config()["slot_stream_params"].is_null(),
3206            "empty slot stream parameters must preserve the old fingerprint shape"
3207        );
3208        assert!(
3209            config.source.semantic_config()["database_initial_statements"].is_null(),
3210            "empty database initial statements must preserve the old fingerprint shape"
3211        );
3212        assert!(
3213            config.source.semantic_config()["interval_handling_mode"].is_null(),
3214            "the native PostgreSQL interval mode must preserve the old fingerprint shape"
3215        );
3216        assert!(
3217            config.source.semantic_config()["include_unknown_datatypes"].is_null(),
3218            "omitting unknown PostgreSQL datatypes must preserve the old fingerprint shape"
3219        );
3220        assert!(
3221            config.source.semantic_config()["money_fraction_digits"].is_null(),
3222            "the default PostgreSQL money scale must preserve the old fingerprint shape"
3223        );
3224        assert!(
3225            config.source.semantic_config()["schema_refresh_mode"].is_null(),
3226            "the pgoutput schema refresh compatibility mode must not change fingerprints"
3227        );
3228        assert!(
3229            config.source.semantic_config()["logical_decoding_messages"].is_null(),
3230            "disabled native logical decoding messages must preserve the old fingerprint shape"
3231        );
3232    }
3233
3234    #[test]
3235    fn parses_native_postgresql_logical_decoding_message_filters() {
3236        let configured = Config::from_yaml(&CONFIG.replace(
3237            "  slot_name: rustium_orders\n",
3238            "  slot_name: rustium_orders\n  logical_decoding_messages: true\n  message_prefix_include_list: [orders\\..*, audit]\n",
3239        ))
3240        .unwrap();
3241        let source = configured.source.as_postgresql().unwrap();
3242        assert!(source.includes_message_prefix("orders.created"));
3243        assert!(source.includes_message_prefix("audit"));
3244        assert!(!source.includes_message_prefix("orders"));
3245        assert_ne!(
3246            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3247            configured.fingerprint()
3248        );
3249
3250        let conflicting = Config::from_yaml(&CONFIG.replace(
3251            "  slot_name: rustium_orders\n",
3252            "  slot_name: rustium_orders\n  message_prefix_include_list: [orders]\n  message_prefix_exclude_list: [audit]\n",
3253        ))
3254        .unwrap_err();
3255        assert!(conflicting.to_string().contains("mutually exclusive"));
3256
3257        let invalid = Config::from_yaml(&CONFIG.replace(
3258            "  slot_name: rustium_orders\n",
3259            "  slot_name: rustium_orders\n  message_prefix_include_list: ['[']\n",
3260        ))
3261        .unwrap_err();
3262        assert!(
3263            invalid
3264                .to_string()
3265                .contains("not a valid regular expression")
3266        );
3267    }
3268
3269    #[test]
3270    fn validates_native_postgresql_replica_identity_rules() {
3271        let configured = Config::from_yaml(&CONFIG.replace(
3272            "  slot_name: rustium_orders\n",
3273            "  slot_name: rustium_orders\n  replica_identity_autoset_values:\n    - table: public\\.orders\n      identity: full\n    - table: public\\.customers\n      identity: index\n      index: customers_replica_key\n",
3274        ))
3275        .unwrap();
3276        let rules = &configured
3277            .source
3278            .as_postgresql()
3279            .unwrap()
3280            .replica_identity_autoset_values;
3281        assert_eq!(rules.len(), 2);
3282        assert_eq!(rules[0].identity, PostgresReplicaIdentity::Full);
3283        assert_eq!(rules[1].index.as_deref(), Some("customers_replica_key"));
3284        assert_ne!(
3285            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3286            configured.fingerprint()
3287        );
3288
3289        let missing_index = Config::from_yaml(&CONFIG.replace(
3290            "  slot_name: rustium_orders\n",
3291            "  slot_name: rustium_orders\n  replica_identity_autoset_values:\n    - table: public\\.orders\n      identity: index\n",
3292        ))
3293        .unwrap_err();
3294        assert!(
3295            missing_index
3296                .to_string()
3297                .contains("requires a non-empty index")
3298        );
3299
3300        let unexpected_index = Config::from_yaml(&CONFIG.replace(
3301            "  slot_name: rustium_orders\n",
3302            "  slot_name: rustium_orders\n  replica_identity_autoset_values:\n    - table: public\\.orders\n      identity: full\n      index: orders_key\n",
3303        ))
3304        .unwrap_err();
3305        assert!(unexpected_index.to_string().contains("valid only"));
3306    }
3307
3308    #[test]
3309    fn parses_native_postgresql_partition_root_publication() {
3310        let configured = Config::from_yaml(&CONFIG.replace(
3311            "  publication: rustium_pub\n",
3312            "  publication: rustium_pub\n  publish_via_partition_root: true\n",
3313        ))
3314        .unwrap();
3315        assert!(
3316            configured
3317                .source
3318                .as_postgresql()
3319                .unwrap()
3320                .publish_via_partition_root
3321        );
3322        assert_ne!(
3323            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3324            configured.fingerprint()
3325        );
3326    }
3327
3328    #[test]
3329    fn parses_native_postgresql_failover_slot() {
3330        let configured = Config::from_yaml(&CONFIG.replace(
3331            "  slot_name: rustium_orders\n",
3332            "  slot_name: rustium_orders\n  slot_failover: true\n",
3333        ))
3334        .unwrap();
3335        assert!(configured.source.as_postgresql().unwrap().slot_failover);
3336        assert_ne!(
3337            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3338            configured.fingerprint()
3339        );
3340
3341        let external = Config::from_yaml(&CONFIG.replace(
3342            "  slot_name: rustium_orders\n",
3343            "  slot_name: rustium_orders\n  slot_failover: true\n  slot_ownership: external\n",
3344        ))
3345        .unwrap_err();
3346        assert!(external.to_string().contains("only for a managed"));
3347    }
3348
3349    #[test]
3350    fn parses_native_postgresql_drop_slot_on_stop_without_changing_fingerprint() {
3351        let configured = Config::from_yaml(&CONFIG.replace(
3352            "  slot_name: rustium_orders\n",
3353            "  slot_name: rustium_orders\n  drop_slot_on_stop: true\n",
3354        ))
3355        .unwrap();
3356        assert!(configured.source.as_postgresql().unwrap().drop_slot_on_stop);
3357        assert_eq!(
3358            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3359            configured.fingerprint()
3360        );
3361
3362        let external = Config::from_yaml(&CONFIG.replace(
3363            "  slot_name: rustium_orders\n",
3364            "  slot_name: rustium_orders\n  drop_slot_on_stop: true\n  slot_ownership: external\n",
3365        ))
3366        .unwrap_err();
3367        assert!(external.to_string().contains("only for a managed"));
3368    }
3369
3370    #[test]
3371    fn parses_native_postgresql_offset_mismatch_strategy() {
3372        let configured = Config::from_yaml(&CONFIG.replace(
3373            "  slot_name: rustium_orders\n",
3374            "  slot_name: rustium_orders\n  offset_mismatch_strategy: trust_slot\n",
3375        ))
3376        .unwrap();
3377        assert_eq!(
3378            configured
3379                .source
3380                .as_postgresql()
3381                .unwrap()
3382                .offset_mismatch_strategy,
3383            PostgresOffsetMismatchStrategy::TrustSlot
3384        );
3385        assert_eq!(
3386            configured.source.semantic_config()["offset_mismatch_strategy"],
3387            "trust_slot"
3388        );
3389        assert_ne!(
3390            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3391            configured.fingerprint()
3392        );
3393
3394        let external = Config::from_yaml(&CONFIG.replace(
3395            "  slot_name: rustium_orders\n",
3396            "  slot_name: rustium_orders\n  slot_ownership: external\n  offset_mismatch_strategy: trust_offset\n",
3397        ))
3398        .unwrap_err();
3399        assert!(
3400            external
3401                .to_string()
3402                .contains("requires slot_ownership=managed")
3403        );
3404    }
3405
3406    #[test]
3407    fn parses_native_postgresql_lsn_flush_mode() {
3408        let configured = Config::from_yaml(&CONFIG.replace(
3409            "  slot_name: rustium_orders\n",
3410            "  slot_name: rustium_orders\n  lsn_flush_mode: connector_and_driver\n",
3411        ))
3412        .unwrap();
3413        assert_eq!(
3414            configured.source.as_postgresql().unwrap().lsn_flush_mode,
3415            PostgresLsnFlushMode::ConnectorAndDriver
3416        );
3417        assert_eq!(
3418            configured.source.semantic_config()["lsn_flush_mode"],
3419            "connector_and_driver"
3420        );
3421        assert_ne!(
3422            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3423            configured.fingerprint()
3424        );
3425    }
3426
3427    #[test]
3428    fn validates_native_postgresql_lsn_flush_timeout() {
3429        let baseline = Config::from_yaml(CONFIG).unwrap();
3430        let configured = Config::from_yaml(&CONFIG.replace(
3431            "  slot_name: rustium_orders\n",
3432            "  slot_name: rustium_orders\n  lsn_flush_timeout: 125ms\n  lsn_flush_timeout_action: warn\n",
3433        ))
3434        .unwrap();
3435        let source = configured.source.as_postgresql().unwrap();
3436        assert_eq!(source.lsn_flush_timeout, Duration::from_millis(125));
3437        assert_eq!(
3438            source.lsn_flush_timeout_action,
3439            PostgresLsnFlushTimeoutAction::Warn
3440        );
3441        assert_eq!(baseline.fingerprint(), configured.fingerprint());
3442
3443        let error = Config::from_yaml(&CONFIG.replace(
3444            "  slot_name: rustium_orders\n",
3445            "  slot_name: rustium_orders\n  lsn_flush_timeout: 0ms\n",
3446        ))
3447        .unwrap_err();
3448        assert!(error.to_string().contains("source.lsn_flush_timeout"));
3449    }
3450
3451    #[test]
3452    fn parses_native_postgresql_slot_stream_origin() {
3453        let configured = Config::from_yaml(&CONFIG.replace(
3454            "  slot_name: rustium_orders\n",
3455            "  slot_name: rustium_orders\n  slot_stream_params:\n    origin: none\n",
3456        ))
3457        .unwrap();
3458        assert_eq!(
3459            configured
3460                .source
3461                .as_postgresql()
3462                .unwrap()
3463                .slot_stream_params
3464                .get("origin")
3465                .map(String::as_str),
3466            Some("none")
3467        );
3468        assert_eq!(
3469            configured.source.semantic_config()["slot_stream_params"]["origin"],
3470            "none"
3471        );
3472        assert_ne!(
3473            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3474            configured.fingerprint()
3475        );
3476
3477        for params in ["    origin: local\n", "    add-tables: public.orders\n"] {
3478            let invalid = Config::from_yaml(&CONFIG.replace(
3479                "  slot_name: rustium_orders\n",
3480                &format!("  slot_name: rustium_orders\n  slot_stream_params:\n{params}"),
3481            ))
3482            .unwrap_err();
3483            assert!(invalid.to_string().contains("slot_stream_params"));
3484        }
3485    }
3486
3487    #[test]
3488    fn parses_native_postgresql_database_initial_statements() {
3489        let configured = Config::from_yaml(&CONFIG.replace(
3490            "  slot_name: rustium_orders\n",
3491            "  slot_name: rustium_orders\n  database_initial_statements:\n    - SET application_name = 'rustium-cdc'\n    - SET statement_timeout = '5s'\n",
3492        ))
3493        .unwrap();
3494        assert_eq!(
3495            configured
3496                .source
3497                .as_postgresql()
3498                .unwrap()
3499                .database_initial_statements,
3500            [
3501                "SET application_name = 'rustium-cdc'",
3502                "SET statement_timeout = '5s'"
3503            ]
3504        );
3505        assert_eq!(
3506            configured.source.semantic_config()["database_initial_statements"][0],
3507            "SET application_name = 'rustium-cdc'"
3508        );
3509        assert_ne!(
3510            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3511            configured.fingerprint()
3512        );
3513
3514        let blank = Config::from_yaml(&CONFIG.replace(
3515            "  slot_name: rustium_orders\n",
3516            "  slot_name: rustium_orders\n  database_initial_statements: [' ']\n",
3517        ))
3518        .unwrap_err();
3519        assert!(blank.to_string().contains("database_initial_statements"));
3520    }
3521
3522    #[test]
3523    fn parses_native_postgresql_snapshot_locking_without_changing_fingerprint() {
3524        let baseline = Config::from_yaml(CONFIG).unwrap();
3525        let configured = Config::from_yaml(&CONFIG.replace(
3526            "  slot_name: rustium_orders\n",
3527            "  slot_name: rustium_orders\n  snapshot_locking_mode: shared\n  snapshot_lock_timeout: 250ms\n",
3528        ))
3529        .unwrap();
3530        let source = configured.source.as_postgresql().unwrap();
3531        assert_eq!(
3532            source.snapshot_locking_mode,
3533            PostgresSnapshotLockingMode::Shared
3534        );
3535        assert_eq!(source.snapshot_lock_timeout, Duration::from_millis(250));
3536        assert_eq!(baseline.fingerprint(), configured.fingerprint());
3537
3538        let too_large = Config::from_yaml(&CONFIG.replace(
3539            "  slot_name: rustium_orders\n",
3540            "  slot_name: rustium_orders\n  snapshot_lock_timeout: 2147483648ms\n",
3541        ))
3542        .unwrap_err();
3543        assert!(too_large.to_string().contains("must not exceed"));
3544    }
3545
3546    #[test]
3547    fn parses_native_postgresql_snapshot_isolation_modes_into_fingerprint() {
3548        let baseline = Config::from_yaml(CONFIG).unwrap();
3549        for mode in [
3550            "serializable",
3551            "repeatable_read",
3552            "read_committed",
3553            "read_uncommitted",
3554        ] {
3555            let configured = Config::from_yaml(&CONFIG.replace(
3556                "  slot_name: rustium_orders\n",
3557                &format!("  slot_name: rustium_orders\n  snapshot_isolation_mode: {mode}\n"),
3558            ))
3559            .unwrap();
3560            let expected = match mode {
3561                "serializable" => PostgresSnapshotIsolationMode::Serializable,
3562                "repeatable_read" => PostgresSnapshotIsolationMode::RepeatableRead,
3563                "read_committed" => PostgresSnapshotIsolationMode::ReadCommitted,
3564                "read_uncommitted" => PostgresSnapshotIsolationMode::ReadUncommitted,
3565                _ => unreachable!(),
3566            };
3567            assert_eq!(
3568                configured
3569                    .source
3570                    .as_postgresql()
3571                    .unwrap()
3572                    .snapshot_isolation_mode,
3573                expected
3574            );
3575            if matches!(mode, "serializable" | "repeatable_read") {
3576                assert_eq!(baseline.fingerprint(), configured.fingerprint());
3577            } else {
3578                assert_ne!(baseline.fingerprint(), configured.fingerprint());
3579            }
3580        }
3581
3582        let invalid = Config::from_yaml(&CONFIG.replace(
3583            "  slot_name: rustium_orders\n",
3584            "  slot_name: rustium_orders\n  snapshot_isolation_mode: dirty_read\n",
3585        ))
3586        .unwrap_err();
3587        assert!(invalid.to_string().contains("dirty_read"));
3588    }
3589
3590    #[test]
3591    fn parses_native_postgresql_column_transformations_into_fingerprint() {
3592        let baseline = Config::from_yaml(CONFIG).unwrap();
3593        let configured = Config::from_yaml(&CONFIG.replace(
3594            "  slot_name: rustium_orders\n",
3595            "  slot_name: rustium_orders\n  column_transformations:\n    - kind: truncate\n      length: 3\n      columns: ['public\\.orders\\.customer_name']\n    - kind: hash\n      algorithm: sha256\n      salt: private-salt\n      version: v2\n      columns: ['public\\.orders\\.email']\n",
3596        ))
3597        .unwrap();
3598        let source = configured.source.as_postgresql().unwrap();
3599        assert_eq!(source.column_transformations.len(), 2);
3600        assert!(matches!(
3601            &source.column_transformations[0],
3602            ColumnTransformRule::Truncate { length: 3, .. }
3603        ));
3604        assert_ne!(baseline.fingerprint(), configured.fingerprint());
3605        let semantic = configured.source.semantic_config().to_string();
3606        assert!(!semantic.contains("private-salt"));
3607
3608        let invalid = Config::from_yaml(&CONFIG.replace(
3609            "  slot_name: rustium_orders\n",
3610            "  slot_name: rustium_orders\n  column_transformations:\n    - kind: hash\n      algorithm: sha256\n      salt: ''\n      columns: ['public.orders.email']\n",
3611        ))
3612        .unwrap_err();
3613        assert!(invalid.to_string().contains("salt"));
3614    }
3615
3616    #[test]
3617    fn parses_native_postgresql_xmin_fetch_interval_into_fingerprint() {
3618        let baseline = Config::from_yaml(CONFIG).unwrap();
3619        let configured = Config::from_yaml(&CONFIG.replace(
3620            "  slot_name: rustium_orders\n",
3621            "  slot_name: rustium_orders\n  xmin_fetch_interval: 25ms\n",
3622        ))
3623        .unwrap();
3624        assert_eq!(
3625            configured
3626                .source
3627                .as_postgresql()
3628                .unwrap()
3629                .xmin_fetch_interval,
3630            Duration::from_millis(25)
3631        );
3632        assert_ne!(baseline.fingerprint(), configured.fingerprint());
3633
3634        let explicit_default = Config::from_yaml(&CONFIG.replace(
3635            "  slot_name: rustium_orders\n",
3636            "  slot_name: rustium_orders\n  xmin_fetch_interval: 0ms\n",
3637        ))
3638        .unwrap();
3639        assert_eq!(baseline.fingerprint(), explicit_default.fingerprint());
3640    }
3641
3642    #[test]
3643    fn validates_native_postgresql_tls_material() {
3644        let configured = Config::from_yaml(&CONFIG.replace(
3645            "  slot_name: rustium_orders\n",
3646            "  slot_name: rustium_orders\n  ssl_mode: verify-full\n  ssl_root_cert: /run/secrets/postgres-ca.pem\n",
3647        ))
3648        .unwrap();
3649        let source = configured.source.as_postgresql().unwrap();
3650        assert_eq!(source.ssl_mode, "verify-full");
3651        assert_eq!(
3652            source.ssl_root_cert.as_deref(),
3653            Some("/run/secrets/postgres-ca.pem")
3654        );
3655        let connection_url = Url::parse(&source.connection_url(true).unwrap()).unwrap();
3656        assert!(connection_url.query_pairs().any(|(key, value)| {
3657            key == "sslrootcert" && value == "/run/secrets/postgres-ca.pem"
3658        }));
3659        assert_eq!(
3660            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3661            configured.fingerprint(),
3662            "PostgreSQL TLS transport settings must not change event semantics"
3663        );
3664
3665        let invalid_mode = Config::from_yaml(&CONFIG.replace(
3666            "  slot_name: rustium_orders\n",
3667            "  slot_name: rustium_orders\n  ssl_mode: verify_identity\n",
3668        ))
3669        .unwrap_err();
3670        assert!(invalid_mode.to_string().contains("source.ssl_mode"));
3671
3672        let client_key = Config::from_yaml(&CONFIG.replace(
3673            "  slot_name: rustium_orders\n",
3674            "  slot_name: rustium_orders\n  ssl_mode: require\n  ssl_cert: /run/secrets/client.pem\n  ssl_key: /run/secrets/client.key\n",
3675        ))
3676        .unwrap_err();
3677        assert!(
3678            client_key
3679                .to_string()
3680                .contains("pg_walstream rustls transport")
3681        );
3682    }
3683
3684    #[test]
3685    fn parses_native_postgresql_interval_handling_mode() {
3686        let configured = Config::from_yaml(&CONFIG.replace(
3687            "  slot_name: rustium_orders\n",
3688            "  slot_name: rustium_orders\n  interval_handling_mode: numeric\n",
3689        ))
3690        .unwrap();
3691        assert_eq!(
3692            configured
3693                .source
3694                .as_postgresql()
3695                .unwrap()
3696                .interval_handling_mode,
3697            "numeric"
3698        );
3699        assert_ne!(
3700            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3701            configured.fingerprint()
3702        );
3703
3704        let invalid = Config::from_yaml(&CONFIG.replace(
3705            "  slot_name: rustium_orders\n",
3706            "  slot_name: rustium_orders\n  interval_handling_mode: invalid\n",
3707        ))
3708        .unwrap_err();
3709        assert!(invalid.to_string().contains("postgres, numeric, or string"));
3710    }
3711
3712    #[test]
3713    fn parses_native_postgresql_unknown_datatype_mode() {
3714        let configured = Config::from_yaml(&CONFIG.replace(
3715            "  slot_name: rustium_orders\n",
3716            "  slot_name: rustium_orders\n  include_unknown_datatypes: true\n",
3717        ))
3718        .unwrap();
3719        assert!(
3720            configured
3721                .source
3722                .as_postgresql()
3723                .unwrap()
3724                .include_unknown_datatypes
3725        );
3726        assert_eq!(
3727            configured.source.semantic_config()["include_unknown_datatypes"],
3728            true
3729        );
3730        assert_ne!(
3731            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3732            configured.fingerprint()
3733        );
3734    }
3735
3736    #[test]
3737    fn parses_native_postgresql_money_fraction_digits() {
3738        let configured = Config::from_yaml(&CONFIG.replace(
3739            "  slot_name: rustium_orders\n",
3740            "  slot_name: rustium_orders\n  money_fraction_digits: 4\n",
3741        ))
3742        .unwrap();
3743        assert_eq!(
3744            configured
3745                .source
3746                .as_postgresql()
3747                .unwrap()
3748                .money_fraction_digits,
3749            4
3750        );
3751        assert_eq!(
3752            configured.source.semantic_config()["money_fraction_digits"],
3753            4
3754        );
3755        assert_ne!(
3756            Config::from_yaml(CONFIG).unwrap().fingerprint(),
3757            configured.fingerprint()
3758        );
3759
3760        let negative = Config::from_yaml(&CONFIG.replace(
3761            "  slot_name: rustium_orders\n",
3762            "  slot_name: rustium_orders\n  money_fraction_digits: -1\n",
3763        ))
3764        .unwrap();
3765        assert_eq!(
3766            negative
3767                .source
3768                .as_postgresql()
3769                .unwrap()
3770                .money_fraction_digits,
3771            -1
3772        );
3773        assert!(
3774            Config::from_yaml(&CONFIG.replace(
3775                "  slot_name: rustium_orders\n",
3776                "  slot_name: rustium_orders\n  money_fraction_digits: 32768\n",
3777            ))
3778            .is_err()
3779        );
3780    }
3781
3782    #[test]
3783    fn parses_native_postgresql_schema_refresh_modes_without_changing_fingerprints() {
3784        let baseline = Config::from_yaml(CONFIG).unwrap();
3785        let configured = Config::from_yaml(&CONFIG.replace(
3786            "  slot_name: rustium_orders\n",
3787            "  slot_name: rustium_orders\n  schema_refresh_mode: columns_diff_exclude_unchanged_toast\n",
3788        ))
3789        .unwrap();
3790        assert_eq!(
3791            configured
3792                .source
3793                .as_postgresql()
3794                .unwrap()
3795                .schema_refresh_mode,
3796            PostgresSchemaRefreshMode::ColumnsDiffExcludeUnchangedToast
3797        );
3798        assert_eq!(baseline.fingerprint(), configured.fingerprint());
3799        assert!(
3800            configured.source.semantic_config()["schema_refresh_mode"].is_null(),
3801            "schema.refresh.mode is operational compatibility under pgoutput"
3802        );
3803
3804        assert!(
3805            Config::from_yaml(&CONFIG.replace(
3806                "  slot_name: rustium_orders\n",
3807                "  slot_name: rustium_orders\n  schema_refresh_mode: invalid\n",
3808            ))
3809            .is_err()
3810        );
3811    }
3812
3813    #[test]
3814    fn validates_native_schema_registry_formats() {
3815        let config = Config::from_yaml(
3816            r#"
3817api_version: rustium.io/v1alpha1
3818kind: Connector
3819metadata:
3820  name: orders-cdc
3821source:
3822  type: postgresql
3823  database: app
3824  username: rustium
3825  password: secret
3826  publication: rustium_pub
3827  slot_name: rustium_orders
3828format:
3829  type: debezium_json_schema
3830  schema_registry:
3831    urls: [https://registry-1:8081, https://registry-2:8081]
3832    username: registry-user
3833    password: registry-secret
3834    request_timeout: 5s
3835sink:
3836  type: kafka
3837  bootstrap_servers: [kafka:9092]
3838  topic_prefix: app
3839"#,
3840        )
3841        .unwrap();
3842        assert_eq!(config.format.kind, FormatType::DebeziumJsonSchema);
3843        assert_eq!(
3844            config
3845                .format
3846                .schema_registry
3847                .as_ref()
3848                .unwrap()
3849                .request_timeout,
3850            Duration::from_secs(5)
3851        );
3852
3853        let error = Config::from_yaml(
3854            &CONFIG.replace(
3855                "sink:\n",
3856                "format:\n  type: debezium_json_schema\n  schema_registry:\n    urls: [http://registry:8081]\nsink:\n",
3857            ),
3858        )
3859        .unwrap_err();
3860        assert!(error.to_string().contains("requires sink.type=kafka"));
3861
3862        let avro = Config::from_yaml(
3863            r#"
3864api_version: rustium.io/v1alpha1
3865kind: Connector
3866metadata:
3867  name: orders-avro
3868source:
3869  type: mysql
3870  databases: [app]
3871  username: rustium
3872  password: secret
3873format:
3874  type: debezium_avro
3875  schema_registry:
3876    urls: [http://registry:8081]
3877    cache_capacity: 32
3878sink:
3879  type: kafka
3880  bootstrap_servers: [kafka:9092]
3881  topic_prefix: app
3882"#,
3883        )
3884        .unwrap();
3885        assert_eq!(avro.format.kind, FormatType::DebeziumAvro);
3886        assert_eq!(
3887            avro.format.schema_registry.as_ref().unwrap().cache_capacity,
3888            32
3889        );
3890
3891        let protobuf = Config::from_yaml(
3892            r#"
3893api_version: rustium.io/v1alpha1
3894kind: Connector
3895metadata:
3896  name: orders-protobuf
3897source:
3898  type: mysql
3899  databases: [app]
3900  username: rustium
3901  password: secret
3902format:
3903  type: debezium_protobuf
3904  schema_registry:
3905    urls: [http://registry:8081]
3906sink:
3907  type: kafka
3908  bootstrap_servers: [kafka:9092]
3909  topic_prefix: app
3910"#,
3911        )
3912        .unwrap();
3913        assert_eq!(protobuf.format.kind, FormatType::DebeziumProtobuf);
3914    }
3915
3916    #[test]
3917    fn rejects_unknown_fields() {
3918        let error = Config::from_yaml(&format!("{CONFIG}\nunknown: true\n")).unwrap_err();
3919        assert!(error.to_string().contains("unknown field"));
3920    }
3921
3922    #[test]
3923    fn rejects_zero_snapshot_fetch_size() {
3924        let error =
3925            Config::from_yaml(&CONFIG.replace("sink:\n", "snapshot:\n  fetch_size: 0\nsink:\n"))
3926                .unwrap_err();
3927        assert!(error.to_string().contains("snapshot.fetch_size"));
3928    }
3929
3930    #[test]
3931    fn validates_anchored_snapshot_collection_filters() {
3932        let config = Config::from_yaml(&CONFIG.replace(
3933            "sink:\n",
3934            "snapshot:\n  include_collections: [public\\.orders]\nsink:\n",
3935        ))
3936        .unwrap();
3937        assert!(config.snapshot.includes_collection("public.orders"));
3938        assert!(!config.snapshot.includes_collection("archive.public.orders"));
3939        assert!(!config.snapshot.includes_collection("public.orders_history"));
3940
3941        let error = Config::from_yaml(&CONFIG.replace(
3942            "sink:\n",
3943            "snapshot:\n  include_collections: ['[invalid']\nsink:\n",
3944        ))
3945        .unwrap_err();
3946        assert!(error.to_string().contains("snapshot include selector"));
3947    }
3948
3949    #[test]
3950    fn preserves_snapshot_fingerprint_shape_without_a_filter() {
3951        assert_eq!(
3952            SnapshotConfig::default().semantic_config(),
3953            serde_json::json!({
3954                "mode": "initial",
3955                "fetch_size": default_snapshot_fetch_size(),
3956            })
3957        );
3958
3959        let filtered = SnapshotConfig {
3960            include_collections: vec![r"public\.orders".into()],
3961            ..SnapshotConfig::default()
3962        };
3963        assert_eq!(
3964            filtered.semantic_config()["include_collections"],
3965            serde_json::json!([r"public\.orders"])
3966        );
3967    }
3968
3969    #[test]
3970    fn rejects_zero_postgresql_connect_timeout() {
3971        let error = Config::from_yaml(&CONFIG.replace(
3972            "  database: app\n",
3973            "  database: app\n  connect_timeout: 0s\n",
3974        ))
3975        .unwrap_err();
3976        assert!(error.to_string().contains("source.connect_timeout"));
3977    }
3978
3979    #[test]
3980    fn validates_postgresql_connection_tuning() {
3981        let configured = Config::from_yaml(&CONFIG.replace(
3982            "  database: app\n",
3983            "  database: app\n  status_update_interval: 125ms\n  tcp_keepalive: false\n",
3984        ))
3985        .unwrap();
3986        let source = configured.source.as_postgresql().unwrap();
3987        assert_eq!(source.status_update_interval, Duration::from_millis(125));
3988        assert!(!source.tcp_keepalive);
3989        let connection_url = Url::parse(&source.connection_url(false).unwrap()).unwrap();
3990        assert!(
3991            connection_url
3992                .query_pairs()
3993                .any(|(key, value)| key == "keepalives" && value == "0")
3994        );
3995
3996        let error = Config::from_yaml(&CONFIG.replace(
3997            "  database: app\n",
3998            "  database: app\n  status_update_interval: 0ms\n",
3999        ))
4000        .unwrap_err();
4001        assert!(error.to_string().contains("source.status_update_interval"));
4002    }
4003
4004    #[test]
4005    fn validates_native_retry_settings() {
4006        let configured = Config::from_yaml(&format!(
4007            "{CONFIG}\nruntime:\n  errors_max_retries: -1\n  errors_retry_delay_initial: 25ms\n  errors_retry_delay_max: 1s\n"
4008        ))
4009        .unwrap();
4010        assert_eq!(configured.runtime.errors_max_retries, -1);
4011        assert_eq!(
4012            configured.runtime.errors_retry_delay_initial,
4013            Duration::from_millis(25)
4014        );
4015        assert_eq!(
4016            configured.runtime.retry_policy(),
4017            RetryPolicy {
4018                max_retries: -1,
4019                initial_delay: Duration::from_millis(25),
4020                max_delay: Duration::from_secs(1),
4021            }
4022        );
4023
4024        for (settings, expected) in [
4025            ("errors_max_retries: -2", "errors_max_retries"),
4026            (
4027                "errors_retry_delay_initial: 0ms",
4028                "errors_retry_delay_initial",
4029            ),
4030            (
4031                "errors_retry_delay_initial: 2s\n  errors_retry_delay_max: 1s",
4032                "errors_retry_delay_max",
4033            ),
4034        ] {
4035            let error = Config::from_yaml(&format!("{CONFIG}\nruntime:\n  {settings}\n"))
4036                .expect_err("invalid retry settings must fail");
4037            assert!(error.to_string().contains(expected));
4038        }
4039    }
4040
4041    #[test]
4042    fn fingerprint_ignores_password() {
4043        let first = Config::from_yaml(CONFIG).unwrap();
4044        let second = Config::from_yaml(&CONFIG.replace("secret", "rotated")).unwrap();
4045        assert_eq!(first.fingerprint(), second.fingerprint());
4046    }
4047
4048    #[test]
4049    fn parses_native_tombstone_override() {
4050        let config = Config::from_yaml(
4051            &CONFIG.replace("sink:\n", "format:\n  tombstones_on_delete: false\nsink:\n"),
4052        )
4053        .unwrap();
4054        assert!(!config.format.tombstones_on_delete);
4055    }
4056
4057    #[test]
4058    fn rejects_non_durable_kafka_sink_settings() {
4059        let kafka = CONFIG.replace(
4060            "sink:\n  type: stdout\n  topic_prefix: app",
4061            "sink:\n  type: kafka\n  bootstrap_servers: [127.0.0.1:9092]\n  topic_prefix: app\n  acks: all\n  delivery_timeout: 3s",
4062        );
4063        Config::from_yaml(&kafka).unwrap();
4064
4065        let acknowledgements = Config::from_yaml(&kafka.replace("acks: all", "acks: '1'"))
4066            .expect_err("non-durable Kafka acknowledgements must fail");
4067        assert!(acknowledgements.to_string().contains("sink.acks"));
4068
4069        let timeout =
4070            Config::from_yaml(&kafka.replace("delivery_timeout: 3s", "delivery_timeout: 0s"))
4071                .expect_err("zero Kafka delivery timeout must fail");
4072        assert!(timeout.to_string().contains("sink.delivery_timeout"));
4073
4074        let override_property = Config::from_yaml(&format!(
4075            "{kafka}\n  properties:\n    enable.idempotence: 'false'\n"
4076        ))
4077        .expect_err("managed Kafka producer property must fail");
4078        assert!(
4079            override_property
4080                .to_string()
4081                .contains("cannot be overridden")
4082        );
4083    }
4084
4085    #[test]
4086    fn parses_native_mysql_heartbeat_settings() {
4087        let raw = r#"
4088api_version: rustium.io/v1alpha1
4089kind: Connector
4090metadata:
4091  name: inventory-mysql
4092source:
4093  type: mysql
4094  username: rustium
4095  password: secret
4096  databases: [inventory]
4097  connection_time_zone: Etc/UTC
4098  heartbeat_interval: 5s
4099  heartbeat_topics_prefix: __heartbeat
4100  heartbeat_topic_name: shared-heartbeat
4101sink:
4102  type: stdout
4103  topic_prefix: inventory
4104"#;
4105        let config = Config::from_yaml(raw).unwrap();
4106        let source = config.source.as_mysql().unwrap();
4107        assert_eq!(source.connection_time_zone, "Etc/UTC");
4108        assert_eq!(source.session_time_zone().unwrap(), "+00:00");
4109        assert_eq!(source.heartbeat_interval, Duration::from_secs(5));
4110        assert_eq!(source.heartbeat_topics_prefix, "__heartbeat");
4111        assert_eq!(
4112            source.heartbeat_topic_name.as_deref(),
4113            Some("shared-heartbeat")
4114        );
4115        let implicit_utc =
4116            Config::from_yaml(&raw.replace("  connection_time_zone: Etc/UTC\n", "")).unwrap();
4117        assert_eq!(config.fingerprint(), implicit_utc.fingerprint());
4118    }
4119
4120    #[test]
4121    fn parses_native_mysql_column_transformations_into_fingerprint() {
4122        let raw = r#"
4123api_version: rustium.io/v1alpha1
4124kind: Connector
4125metadata:
4126  name: inventory-mysql
4127source:
4128  type: mysql
4129  username: rustium
4130  password: secret
4131  databases: [inventory]
4132sink:
4133  type: stdout
4134  topic_prefix: inventory
4135"#;
4136        let baseline = Config::from_yaml(raw).unwrap();
4137        let configured = Config::from_yaml(&raw.replace(
4138            "sink:\n",
4139            "  column_transformations:\n    - kind: truncate\n      length: 3\n      columns: ['inventory\\.orders\\.customer']\n    - kind: hash\n      algorithm: sha256\n      salt: mysql-private-salt\n      version: v2\n      columns: ['inventory\\.orders\\.email']\nsink:\n",
4140        ))
4141        .unwrap();
4142        let source = configured.source.as_mysql().unwrap();
4143        assert_eq!(source.column_transformations.len(), 2);
4144        assert_ne!(baseline.fingerprint(), configured.fingerprint());
4145        let semantic = configured.source.semantic_config().to_string();
4146        assert!(!semantic.contains("mysql-private-salt"));
4147    }
4148
4149    #[test]
4150    fn parses_native_sqlserver_column_transformations_into_fingerprint() {
4151        let raw = r#"
4152api_version: rustium.io/v1alpha1
4153kind: Connector
4154metadata:
4155  name: inventory-sqlserver
4156source:
4157  type: sqlserver
4158  username: rustium
4159  password: secret
4160  databases: [inventory]
4161sink:
4162  type: stdout
4163  topic_prefix: inventory
4164"#;
4165        let baseline = Config::from_yaml(raw).unwrap();
4166        let configured = Config::from_yaml(&raw.replace(
4167            "sink:\n",
4168            "  column_transformations:\n    - kind: truncate\n      length: 7\n      columns: ['inventory\\.dbo\\.customers\\.summary']\n    - kind: hash\n      algorithm: sha256\n      salt: sqlserver-private-salt\n      version: v2\n      columns: ['dbo\\.customers\\.email']\nsink:\n",
4169        ))
4170        .unwrap();
4171        let source = configured.source.as_sqlserver().unwrap();
4172        assert_eq!(source.column_transformations.len(), 2);
4173        assert_ne!(baseline.fingerprint(), configured.fingerprint());
4174        let semantic = configured.source.semantic_config().to_string();
4175        assert!(!semantic.contains("sqlserver-private-salt"));
4176        assert!(semantic.contains(&hex_digest(b"sqlserver-private-salt")));
4177    }
4178
4179    #[test]
4180    fn parses_native_mysql_java_tls_stores_without_fingerprinting_passwords() {
4181        let raw = r#"
4182api_version: rustium.io/v1alpha1
4183kind: Connector
4184metadata:
4185  name: inventory-mysql
4186source:
4187  type: mysql
4188  username: rustium
4189  password: secret
4190  databases: [inventory]
4191  ssl_mode: verify_identity
4192  ssl_keystore: /run/secrets/client.p12
4193  ssl_keystore_password: client-secret
4194  ssl_truststore: /run/secrets/trust.jks
4195  ssl_truststore_password: trust-secret
4196sink:
4197  type: stdout
4198  topic_prefix: inventory
4199"#;
4200        let config = Config::from_yaml(raw).unwrap();
4201        let source = config.source.as_mysql().unwrap();
4202        assert_eq!(
4203            source.ssl_keystore.as_deref(),
4204            Some("/run/secrets/client.p12")
4205        );
4206        assert_eq!(
4207            source.ssl_truststore.as_deref(),
4208            Some("/run/secrets/trust.jks")
4209        );
4210
4211        let rotated = Config::from_yaml(
4212            &raw.replace("client-secret", "rotated-client")
4213                .replace("trust-secret", "rotated-trust"),
4214        )
4215        .unwrap();
4216        assert_eq!(config.fingerprint(), rotated.fingerprint());
4217
4218        let different_store =
4219            Config::from_yaml(&raw.replace("client.p12", "replacement.p12")).unwrap();
4220        assert_ne!(config.fingerprint(), different_store.fingerprint());
4221    }
4222
4223    #[test]
4224    fn parses_native_postgresql_heartbeat_settings() {
4225        let default = Config::from_yaml(CONFIG).unwrap();
4226        let config = Config::from_yaml(
4227            &CONFIG.replace(
4228                "  password: secret\n",
4229                "  password: secret\n  heartbeat_interval: 5s\n  heartbeat_action_query: UPDATE public.heartbeat SET touched_at = now()\n  heartbeat_topics_prefix: __heartbeat\n  heartbeat_topic_name: shared-heartbeat\n  read_only: true\n",
4230            ),
4231        )
4232        .unwrap();
4233        let source = config.source.as_postgresql().unwrap();
4234        assert_eq!(source.heartbeat_interval, Duration::from_secs(5));
4235        assert_eq!(
4236            source.heartbeat_action_query.as_deref(),
4237            Some("UPDATE public.heartbeat SET touched_at = now()")
4238        );
4239        assert_eq!(source.heartbeat_topics_prefix, "__heartbeat");
4240        assert_eq!(
4241            source.heartbeat_topic_name.as_deref(),
4242            Some("shared-heartbeat")
4243        );
4244        assert!(source.read_only);
4245        assert!(default.source.semantic_config().get("heartbeat").is_none());
4246        assert_ne!(default.fingerprint(), config.fingerprint());
4247    }
4248
4249    #[test]
4250    fn parses_native_debezium_bridge_and_excludes_secrets_from_fingerprint() {
4251        let raw = r#"
4252api_version: rustium.io/v1alpha1
4253kind: Connector
4254metadata:
4255  name: db2-bridge
4256source:
4257  type: db2
4258  bridge:
4259    type: http
4260    listen: 127.0.0.1:18080
4261    path: /events
4262  properties:
4263    connector.class: io.debezium.connector.db2.Db2Connector
4264    topic.prefix: inventory
4265    database.hostname: db2
4266    database.password: secret-one
4267sink:
4268  type: stdout
4269  topic_prefix: inventory
4270"#;
4271        let config = Config::from_yaml(raw).unwrap();
4272        let (kind, source) = config.source.as_debezium().unwrap();
4273        assert_eq!(kind, DebeziumConnectorKind::Db2);
4274        assert_eq!(source.properties["database.hostname"], "db2");
4275        let rotated = Config::from_yaml(&raw.replace("secret-one", "secret-two")).unwrap();
4276        assert_eq!(config.fingerprint(), rotated.fingerprint());
4277
4278        let kafka = raw.replace(
4279            "    type: http\n    listen: 127.0.0.1:18080\n    path: /events",
4280            "    type: kafka\n    bootstrap_servers: [kafka:9092]\n    topics: [inventory.DB2INST1.CUSTOMERS]\n    group_id: rustium-db2\n    consumer_properties:\n      security.protocol: SASL_SSL\n      sasl.username: rustium\n      sasl.password: kafka-secret-one",
4281        );
4282        let kafka_config = Config::from_yaml(&kafka).unwrap();
4283        let kafka_rotated =
4284            Config::from_yaml(&kafka.replace("kafka-secret-one", "kafka-secret-two")).unwrap();
4285        assert_eq!(kafka_config.fingerprint(), kafka_rotated.fingerprint());
4286    }
4287}