1mod 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 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: ®ex::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}