1#![allow(clippy::disallowed_types)] use std::collections::HashMap;
6use std::fmt;
7use std::str::FromStr;
8use std::time::Duration;
9
10use crate::connector::DeliveryGuarantee;
11
12use crate::config::ConnectorConfig;
13use crate::error::ConnectorError;
14use crate::storage::{
15 CloudConfigValidator, ResolvedStorageOptions, SecretMasker, StorageCredentialResolver,
16 StorageProvider,
17};
18
19#[derive(Debug, Clone)]
23pub struct DeltaLakeSinkConfig {
24 pub table_path: String,
26
27 pub partition_columns: Vec<String>,
29
30 pub target_file_size: usize,
32
33 pub max_buffer_records: usize,
35
36 pub max_buffer_duration: Duration,
38
39 pub schema_evolution: bool,
41
42 pub write_mode: DeltaWriteMode,
44
45 pub merge_key_columns: Vec<String>,
47
48 pub storage_options: HashMap<String, String>,
50
51 pub delivery_guarantee: DeliveryGuarantee,
53
54 pub catalog_type: DeltaCatalogType,
56
57 pub catalog_database: Option<String>,
59
60 pub catalog_name: Option<String>,
62
63 pub catalog_schema: Option<String>,
65
66 pub catalog_storage_location: Option<String>,
70
71 pub write_timeout: Duration,
74
75 pub parquet: ParquetWriteConfig,
77}
78
79impl Default for DeltaLakeSinkConfig {
80 fn default() -> Self {
81 Self {
82 table_path: String::new(),
83 partition_columns: Vec::new(),
84 target_file_size: 128 * 1024 * 1024, max_buffer_records: 100_000,
86 max_buffer_duration: Duration::from_secs(60),
87 schema_evolution: false,
88 write_mode: DeltaWriteMode::Append,
89 merge_key_columns: Vec::new(),
90 storage_options: HashMap::new(),
91 delivery_guarantee: DeliveryGuarantee::AtLeastOnce,
92 catalog_type: DeltaCatalogType::None,
93 catalog_database: None,
94 catalog_name: None,
95 catalog_schema: None,
96 catalog_storage_location: None,
97 write_timeout: Duration::from_secs(30),
98 parquet: ParquetWriteConfig::default(),
99 }
100 }
101}
102
103impl DeltaLakeSinkConfig {
104 #[must_use]
106 pub fn new(table_path: &str) -> Self {
107 Self {
108 table_path: table_path.to_string(),
109 ..Default::default()
110 }
111 }
112
113 #[allow(clippy::too_many_lines)]
124 pub fn from_config(config: &ConnectorConfig) -> Result<Self, ConnectorError> {
125 let mut cfg = Self {
126 table_path: config.require("table.path")?.to_string(),
127 ..Self::default()
128 };
129
130 if let Some(v) = config.get("partition.columns") {
131 cfg.partition_columns = v
132 .split(',')
133 .map(|c| c.trim().to_string())
134 .filter(|c| !c.is_empty())
135 .collect();
136 }
137 if let Some(v) = config.get("target.file.size") {
138 cfg.target_file_size = v.parse().map_err(|_| {
139 ConnectorError::ConfigurationError(format!("invalid target.file.size: '{v}'"))
140 })?;
141 }
142 if let Some(v) = config.get("max.buffer.records") {
143 cfg.max_buffer_records = v.parse().map_err(|_| {
144 ConnectorError::ConfigurationError(format!("invalid max.buffer.records: '{v}'"))
145 })?;
146 }
147 if let Some(v) = config.get("max.buffer.duration.ms") {
148 let ms: u64 = v.parse().map_err(|_| {
149 ConnectorError::ConfigurationError(format!("invalid max.buffer.duration.ms: '{v}'"))
150 })?;
151 cfg.max_buffer_duration = Duration::from_millis(ms);
152 }
153 if let Some(v) = config.get("schema.evolution") {
154 cfg.schema_evolution = v.eq_ignore_ascii_case("true");
155 }
156 if let Some(v) = config.get("write.mode") {
157 cfg.write_mode = v.parse().map_err(|_| {
158 ConnectorError::ConfigurationError(format!(
159 "invalid write.mode: '{v}' (expected 'append', 'overwrite', or 'upsert')"
160 ))
161 })?;
162 }
163 if let Some(v) = config.get("merge.key.columns") {
164 cfg.merge_key_columns = v
165 .split(',')
166 .map(|c| c.trim().to_string())
167 .filter(|c| !c.is_empty())
168 .collect();
169 }
170 if let Some(v) = config.get("delivery.guarantee") {
171 cfg.delivery_guarantee = v.parse().map_err(|_| {
172 ConnectorError::ConfigurationError(format!(
173 "invalid delivery.guarantee: '{v}' \
174 (expected 'exactly-once' or 'at-least-once')"
175 ))
176 })?;
177 }
178 if cfg.delivery_guarantee == DeliveryGuarantee::ExactlyOnce
179 && cfg.write_mode != DeltaWriteMode::Append
180 {
181 return Err(ConnectorError::ConfigurationError(
182 "Delta exactly-once is supported only for coordinated append mode; \
183 upsert/overwrite do not expose a certified distributed committable"
184 .into(),
185 ));
186 }
187
188 if let Some(v) = config.get("catalog.type") {
190 cfg.catalog_type = v.parse().map_err(|_| {
191 ConnectorError::ConfigurationError(format!(
192 "invalid catalog.type: '{v}' (expected 'none', 'glue', or 'unity')"
193 ))
194 })?;
195 }
196 if let Some(v) = config.get("catalog.database") {
197 cfg.catalog_database = Some(v.to_string());
198 }
199 if let Some(v) = config.get("catalog.name") {
200 cfg.catalog_name = Some(v.to_string());
201 }
202 if let Some(v) = config.get("catalog.schema") {
203 cfg.catalog_schema = Some(v.to_string());
204 }
205 if let DeltaCatalogType::Unity {
207 ref mut workspace_url,
208 ref mut access_token,
209 } = cfg.catalog_type
210 {
211 if let Some(v) = config.get("catalog.workspace_url") {
212 *workspace_url = v.to_string();
213 }
214 if let Some(v) = config.get("catalog.access_token") {
215 *access_token = v.to_string();
216 }
217 }
218 if let Some(v) = config.get("catalog.storage.location") {
219 cfg.catalog_storage_location = Some(v.to_string());
220 }
221 if let Some(v) = config.get("write.timeout.ms") {
222 let ms: u64 = v.parse().map_err(|_| {
223 ConnectorError::ConfigurationError(format!("invalid write.timeout.ms: '{v}'"))
224 })?;
225 cfg.write_timeout = Duration::from_millis(ms);
226 }
227
228 if let Some(v) = config.get("parquet.compression") {
230 cfg.parquet.compression = v.to_string();
231 }
232 if let Some(v) = config.get("parquet.compression.level") {
233 cfg.parquet.compression_level = v.parse().map_err(|_| {
234 ConnectorError::ConfigurationError(format!(
235 "invalid parquet.compression.level: '{v}'"
236 ))
237 })?;
238 }
239 if let Some(v) = config.get("parquet.dictionary.enabled") {
240 cfg.parquet.dictionary_enabled = v.eq_ignore_ascii_case("true");
241 }
242 if let Some(v) = config.get("parquet.statistics") {
243 cfg.parquet.statistics = v.to_string();
244 }
245 if let Some(v) = config.get("parquet.bloom.filter.columns") {
246 cfg.parquet.bloom_filter_columns = v
247 .split(',')
248 .map(|c| c.trim().to_string())
249 .filter(|c| !c.is_empty())
250 .collect();
251 }
252 if let Some(v) = config.get("parquet.bloom.filter.fpp") {
253 cfg.parquet.bloom_filter_fpp = v.parse().map_err(|_| {
254 ConnectorError::ConfigurationError(format!(
255 "invalid parquet.bloom.filter.fpp: '{v}'"
256 ))
257 })?;
258 }
259 if let Some(v) = config.get("parquet.bloom.filter.ndv") {
260 cfg.parquet.bloom_filter_ndv = v.parse().map_err(|_| {
261 ConnectorError::ConfigurationError(format!(
262 "invalid parquet.bloom.filter.ndv: '{v}'"
263 ))
264 })?;
265 }
266 if let Some(v) = config.get("parquet.max.row.group.size") {
267 cfg.parquet.max_row_group_size = v.parse().map_err(|_| {
268 ConnectorError::ConfigurationError(format!(
269 "invalid parquet.max.row.group.size: '{v}'"
270 ))
271 })?;
272 }
273
274 let explicit_storage = config.properties_with_prefix("storage.");
276 let resolved = StorageCredentialResolver::resolve(&cfg.table_path, &explicit_storage);
277 cfg.storage_options = resolved.options;
278
279 if let Some(v) = config.get("storage.s3_locking_provider") {
281 cfg.storage_options
282 .insert("AWS_S3_LOCKING_PROVIDER".to_string(), v.to_string());
283 }
284 if let Some(v) = config.get("storage.dynamodb_table_name") {
285 cfg.storage_options
286 .insert("DELTA_DYNAMO_TABLE_NAME".to_string(), v.to_string());
287 }
288
289 cfg.validate()?;
290 Ok(cfg)
291 }
292
293 #[must_use]
295 pub fn display_storage_options(&self) -> String {
296 SecretMasker::display_map(&self.storage_options)
297 }
298
299 pub fn validate(&self) -> Result<(), ConnectorError> {
305 if self.table_path.is_empty() {
306 return Err(ConnectorError::missing_config("table.path"));
307 }
308 if self.write_mode == DeltaWriteMode::Upsert && self.merge_key_columns.is_empty() {
309 return Err(ConnectorError::ConfigurationError(
310 "upsert mode requires 'merge.key.columns' to be set".into(),
311 ));
312 }
313 if self.max_buffer_records == 0 {
314 return Err(ConnectorError::ConfigurationError(
315 "max.buffer.records must be > 0".into(),
316 ));
317 }
318 if self.target_file_size == 0 {
319 return Err(ConnectorError::ConfigurationError(
320 "target.file.size must be > 0".into(),
321 ));
322 }
323 if self.write_timeout < Duration::from_secs(5) {
324 return Err(ConnectorError::ConfigurationError(
325 "write.timeout.ms must be >= 5000 (5 seconds)".into(),
326 ));
327 }
328
329 match self.parquet.compression.to_lowercase().as_str() {
330 "zstd" | "snappy" | "lz4" | "gzip" | "none" | "uncompressed" => {}
331 other => {
332 return Err(ConnectorError::ConfigurationError(format!(
333 "unknown parquet.compression: '{other}' \
334 (expected 'zstd', 'snappy', 'lz4', 'gzip', or 'none')"
335 )));
336 }
337 }
338 match self.parquet.statistics.to_lowercase().as_str() {
339 "none" | "chunk" | "page" => {}
340 other => {
341 return Err(ConnectorError::ConfigurationError(format!(
342 "unknown parquet.statistics: '{other}' (expected 'none', 'chunk', or 'page')"
343 )));
344 }
345 }
346 if self.parquet.bloom_filter_fpp <= 0.0 || self.parquet.bloom_filter_fpp >= 1.0 {
347 return Err(ConnectorError::ConfigurationError(
348 "parquet.bloom.filter.fpp must be in (0.0, 1.0)".into(),
349 ));
350 }
351 if self.parquet.max_row_group_size == 0 {
352 return Err(ConnectorError::ConfigurationError(
353 "parquet.max.row.group.size must be > 0".into(),
354 ));
355 }
356 #[cfg(feature = "delta-lake")]
359 {
360 self.parquet.to_writer_properties()?;
361 }
362 self.validate_catalog()?;
363
364 if self.catalog_type == DeltaCatalogType::None {
366 let resolved = ResolvedStorageOptions {
367 provider: StorageProvider::detect(&self.table_path),
368 options: self.storage_options.clone(),
369 env_resolved_keys: Vec::new(),
370 };
371 let cloud_result = CloudConfigValidator::validate(&resolved);
372 if !cloud_result.is_valid() {
373 return Err(ConnectorError::ConfigurationError(
374 cloud_result.error_message(),
375 ));
376 }
377 }
378
379 Ok(())
380 }
381
382 fn validate_catalog(&self) -> Result<(), ConnectorError> {
384 match &self.catalog_type {
385 DeltaCatalogType::None => {}
386 DeltaCatalogType::Glue => {
387 #[cfg(not(feature = "delta-lake-glue"))]
388 return Err(ConnectorError::ConfigurationError(
389 "Glue catalog requires the 'delta-lake-glue' feature. \
390 Build with: cargo build --features delta-lake-glue"
391 .into(),
392 ));
393 #[cfg(feature = "delta-lake-glue")]
394 if self.catalog_database.is_none() {
395 return Err(ConnectorError::ConfigurationError(
396 "Glue catalog requires 'catalog.database' to be set".into(),
397 ));
398 }
399 }
400 DeltaCatalogType::Unity {
401 workspace_url,
402 access_token,
403 } => {
404 #[cfg(not(feature = "delta-lake-unity"))]
405 {
406 let _ = (workspace_url, access_token);
407 return Err(ConnectorError::ConfigurationError(
408 "Unity catalog requires the 'delta-lake-unity' feature. \
409 Build with: cargo build --features delta-lake-unity"
410 .into(),
411 ));
412 }
413 #[cfg(feature = "delta-lake-unity")]
414 {
415 if workspace_url.is_empty() {
416 return Err(ConnectorError::ConfigurationError(
417 "Unity catalog requires 'catalog.workspace_url' to be set".into(),
418 ));
419 }
420 if access_token.is_empty() {
421 return Err(ConnectorError::ConfigurationError(
422 "Unity catalog requires 'catalog.access_token' to be set".into(),
423 ));
424 }
425 if self.catalog_storage_location.is_some() {
426 if self.catalog_name.is_none() {
427 return Err(ConnectorError::ConfigurationError(
428 "Unity catalog auto-create requires 'catalog.name' to be set"
429 .into(),
430 ));
431 }
432 if self.catalog_schema.is_none() {
433 return Err(ConnectorError::ConfigurationError(
434 "Unity catalog auto-create requires 'catalog.schema' to be set"
435 .into(),
436 ));
437 }
438 }
439 }
440 }
441 }
442 Ok(())
443 }
444}
445
446#[derive(Debug, Clone, Copy, PartialEq, Eq)]
448pub enum DeltaWriteMode {
449 Append,
451 Overwrite,
453 Upsert,
456}
457
458impl FromStr for DeltaWriteMode {
459 type Err = String;
460
461 fn from_str(s: &str) -> Result<Self, Self::Err> {
462 match s.to_lowercase().as_str() {
463 "append" => Ok(Self::Append),
464 "overwrite" => Ok(Self::Overwrite),
465 "upsert" | "merge" => Ok(Self::Upsert),
466 other => Err(format!("unknown write mode: '{other}'")),
467 }
468 }
469}
470
471impl fmt::Display for DeltaWriteMode {
472 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
473 match self {
474 Self::Append => write!(f, "append"),
475 Self::Overwrite => write!(f, "overwrite"),
476 Self::Upsert => write!(f, "upsert"),
477 }
478 }
479}
480
481#[derive(Debug, Clone, PartialEq, Eq, Default)]
485pub enum DeltaCatalogType {
486 #[default]
488 None,
489 Glue,
491 Unity {
493 workspace_url: String,
495 access_token: String,
497 },
498}
499
500impl FromStr for DeltaCatalogType {
501 type Err = String;
502
503 fn from_str(s: &str) -> Result<Self, Self::Err> {
504 match s.to_lowercase().as_str() {
505 "none" | "" => Ok(Self::None),
506 "glue" => Ok(Self::Glue),
507 "unity" => Ok(Self::Unity {
508 workspace_url: String::new(),
509 access_token: String::new(),
510 }),
511 other => Err(format!("unknown catalog type: '{other}'")),
512 }
513 }
514}
515
516impl fmt::Display for DeltaCatalogType {
517 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
518 match self {
519 Self::None => write!(f, "none"),
520 Self::Glue => write!(f, "glue"),
521 Self::Unity { .. } => write!(f, "unity"),
522 }
523 }
524}
525
526#[derive(Debug, Clone)]
529pub struct ParquetWriteConfig {
530 pub compression: String,
532 pub compression_level: i32,
534 pub dictionary_enabled: bool,
536 pub statistics: String,
538 pub bloom_filter_columns: Vec<String>,
540 pub bloom_filter_fpp: f64,
542 pub bloom_filter_ndv: u64,
544 pub max_row_group_size: usize,
546}
547
548impl Default for ParquetWriteConfig {
549 fn default() -> Self {
550 Self {
551 compression: "zstd".to_string(),
552 compression_level: 1,
553 dictionary_enabled: true,
554 statistics: "page".to_string(),
555 bloom_filter_columns: Vec::new(),
556 bloom_filter_fpp: 0.01,
557 bloom_filter_ndv: 0,
558 max_row_group_size: 1_000_000,
559 }
560 }
561}
562
563#[cfg(feature = "delta-lake")]
564impl ParquetWriteConfig {
565 pub fn to_writer_properties(
571 &self,
572 ) -> Result<deltalake::parquet::file::properties::WriterProperties, ConnectorError> {
573 self.build_properties(self.compression_level)
574 }
575
576 fn build_properties(
579 &self,
580 level: i32,
581 ) -> Result<deltalake::parquet::file::properties::WriterProperties, ConnectorError> {
582 use deltalake::parquet::basic::{Compression, GzipLevel, ZstdLevel};
583 use deltalake::parquet::file::properties::{EnabledStatistics, WriterProperties};
584 use deltalake::parquet::schema::types::ColumnPath;
585
586 let compression = match self.compression.to_lowercase().as_str() {
587 "zstd" => {
588 let zstd_level = ZstdLevel::try_new(level).map_err(|e| {
589 ConnectorError::ConfigurationError(format!("invalid ZSTD level {level}: {e}"))
590 })?;
591 Compression::ZSTD(zstd_level)
592 }
593 "snappy" => Compression::SNAPPY,
594 "lz4" => Compression::LZ4_RAW,
595 "gzip" => {
596 let level_u32: u32 = level.try_into().map_err(|_| {
597 ConnectorError::ConfigurationError(format!(
598 "invalid GZIP level {level}: must be non-negative"
599 ))
600 })?;
601 let gzip_level = GzipLevel::try_new(level_u32).map_err(|e| {
602 ConnectorError::ConfigurationError(format!("invalid GZIP level {level}: {e}"))
603 })?;
604 Compression::GZIP(gzip_level)
605 }
606 "none" | "uncompressed" => Compression::UNCOMPRESSED,
607 other => {
608 return Err(ConnectorError::ConfigurationError(format!(
609 "unknown parquet.compression: '{other}' \
610 (expected 'zstd', 'snappy', 'lz4', 'gzip', or 'none')"
611 )));
612 }
613 };
614
615 let statistics = match self.statistics.to_lowercase().as_str() {
616 "none" => EnabledStatistics::None,
617 "chunk" => EnabledStatistics::Chunk,
618 "page" => EnabledStatistics::Page,
619 other => {
620 return Err(ConnectorError::ConfigurationError(format!(
621 "unknown parquet.statistics: '{other}' (expected 'none', 'chunk', or 'page')"
622 )));
623 }
624 };
625
626 let mut builder = WriterProperties::builder()
627 .set_compression(compression)
628 .set_dictionary_enabled(self.dictionary_enabled)
629 .set_statistics_enabled(statistics)
630 .set_max_row_group_size(self.max_row_group_size);
631
632 for col_name in &self.bloom_filter_columns {
633 let col_path = ColumnPath::from(col_name.as_str());
634 builder = builder
635 .set_column_bloom_filter_enabled(col_path.clone(), true)
636 .set_column_bloom_filter_fpp(col_path.clone(), self.bloom_filter_fpp);
637 if self.bloom_filter_ndv > 0 {
638 builder = builder.set_column_bloom_filter_ndv(col_path, self.bloom_filter_ndv);
639 }
640 }
641
642 Ok(builder.build())
643 }
644}
645
646#[cfg(test)]
647#[allow(clippy::field_reassign_with_default)]
648mod tests {
649 use super::*;
650
651 fn make_config(pairs: &[(&str, &str)]) -> ConnectorConfig {
652 let mut config = ConnectorConfig::new("delta-lake");
653 for (k, v) in pairs {
654 config.set(*k, *v);
655 }
656 config
657 }
658
659 fn required_pairs() -> Vec<(&'static str, &'static str)> {
660 vec![("table.path", "/data/warehouse/trades")]
661 }
662
663 #[test]
666 fn test_parse_required_fields() {
667 let config = make_config(&required_pairs());
668 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
669 assert_eq!(cfg.table_path, "/data/warehouse/trades");
670 assert_eq!(cfg.write_mode, DeltaWriteMode::Append);
671 assert_eq!(cfg.delivery_guarantee, DeliveryGuarantee::AtLeastOnce);
672 assert!(cfg.partition_columns.is_empty());
673 assert!(cfg.merge_key_columns.is_empty());
674 assert_eq!(cfg.target_file_size, 128 * 1024 * 1024);
675 assert_eq!(cfg.max_buffer_records, 100_000);
676 assert!(!cfg.schema_evolution);
677 }
678
679 #[test]
680 fn test_missing_table_path() {
681 let config = ConnectorConfig::new("delta-lake");
682 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
683 }
684
685 #[test]
686 fn test_parse_all_optional_fields() {
687 let mut pairs = required_pairs();
688 pairs.extend_from_slice(&[
689 ("partition.columns", "trade_date, hour"),
690 ("target.file.size", "67108864"),
691 ("max.buffer.records", "50000"),
692 ("max.buffer.duration.ms", "30000"),
693 ("schema.evolution", "true"),
694 ("write.mode", "upsert"),
695 ("merge.key.columns", "customer_id, order_id"),
696 ("delivery.guarantee", "at-least-once"),
697 ("storage.aws_access_key_id", "AKID123"),
698 ("storage.aws_region", "us-east-1"),
699 ]);
700 let config = make_config(&pairs);
701 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
702
703 assert_eq!(cfg.partition_columns, vec!["trade_date", "hour"]);
704 assert_eq!(cfg.target_file_size, 67_108_864);
705 assert_eq!(cfg.max_buffer_records, 50_000);
706 assert_eq!(cfg.max_buffer_duration, Duration::from_secs(30));
707 assert!(cfg.schema_evolution);
708 assert_eq!(cfg.write_mode, DeltaWriteMode::Upsert);
709 assert_eq!(cfg.merge_key_columns, vec!["customer_id", "order_id"]);
710 assert_eq!(cfg.delivery_guarantee, DeliveryGuarantee::AtLeastOnce);
711 assert_eq!(
712 cfg.storage_options.get("aws_access_key_id"),
713 Some(&"AKID123".to_string())
714 );
715 assert_eq!(
716 cfg.storage_options.get("aws_region"),
717 Some(&"us-east-1".to_string())
718 );
719 }
720
721 #[test]
722 fn test_upsert_requires_merge_key() {
723 let mut pairs = required_pairs();
724 pairs.push(("write.mode", "upsert"));
725 let config = make_config(&pairs);
726 let result = DeltaLakeSinkConfig::from_config(&config);
727 assert!(result.is_err());
728 let err = result.unwrap_err().to_string();
729 assert!(err.contains("merge.key.columns"), "error: {err}");
730 }
731
732 #[test]
733 fn test_empty_table_path_rejected() {
734 let mut cfg = DeltaLakeSinkConfig::default();
735 cfg.table_path = String::new();
736 assert!(cfg.validate().is_err());
737 }
738
739 #[test]
740 fn test_zero_max_buffer_records_rejected() {
741 let mut pairs = required_pairs();
742 pairs.push(("max.buffer.records", "0"));
743 let config = make_config(&pairs);
744 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
745 }
746
747 #[test]
748 fn test_zero_target_file_size_rejected() {
749 let mut pairs = required_pairs();
750 pairs.push(("target.file.size", "0"));
751 let config = make_config(&pairs);
752 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
753 }
754
755 #[test]
756 fn test_invalid_target_file_size() {
757 let mut pairs = required_pairs();
758 pairs.push(("target.file.size", "abc"));
759 let config = make_config(&pairs);
760 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
761 }
762
763 #[test]
764 fn test_invalid_write_mode() {
765 let mut pairs = required_pairs();
766 pairs.push(("write.mode", "unknown"));
767 let config = make_config(&pairs);
768 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
769 }
770
771 #[test]
772 fn test_storage_options_prefix_stripping() {
773 let mut pairs = required_pairs();
774 pairs.push(("storage.aws_access_key_id", "AKID"));
775 pairs.push(("storage.aws_secret_access_key", "SECRET"));
776 pairs.push(("table.path", "/data/test"));
777 let config = make_config(&pairs);
778 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
779
780 assert_eq!(cfg.storage_options.len(), 2);
781 assert!(cfg.storage_options.contains_key("aws_access_key_id"));
782 assert!(cfg.storage_options.contains_key("aws_secret_access_key"));
783 assert!(!cfg
784 .storage_options
785 .contains_key("storage.aws_access_key_id"));
786 }
787
788 #[test]
789 fn test_defaults() {
790 let cfg = DeltaLakeSinkConfig::default();
791 assert!(cfg.table_path.is_empty());
792 assert_eq!(cfg.target_file_size, 128 * 1024 * 1024);
793 assert_eq!(cfg.max_buffer_records, 100_000);
794 assert_eq!(cfg.max_buffer_duration, Duration::from_secs(60));
795 assert!(!cfg.schema_evolution);
796 assert_eq!(cfg.write_mode, DeltaWriteMode::Append);
797 assert_eq!(cfg.delivery_guarantee, DeliveryGuarantee::AtLeastOnce);
798 }
799
800 #[test]
801 fn test_new_helper() {
802 let cfg = DeltaLakeSinkConfig::new("/tmp/test_table");
803 assert_eq!(cfg.table_path, "/tmp/test_table");
804 assert_eq!(cfg.write_mode, DeltaWriteMode::Append);
805 }
806
807 #[test]
810 fn test_write_mode_parse() {
811 assert_eq!(
812 "append".parse::<DeltaWriteMode>().unwrap(),
813 DeltaWriteMode::Append
814 );
815 assert_eq!(
816 "overwrite".parse::<DeltaWriteMode>().unwrap(),
817 DeltaWriteMode::Overwrite
818 );
819 assert_eq!(
820 "upsert".parse::<DeltaWriteMode>().unwrap(),
821 DeltaWriteMode::Upsert
822 );
823 assert_eq!(
824 "merge".parse::<DeltaWriteMode>().unwrap(),
825 DeltaWriteMode::Upsert
826 );
827 assert!("unknown".parse::<DeltaWriteMode>().is_err());
828 }
829
830 #[test]
831 fn test_write_mode_display() {
832 assert_eq!(DeltaWriteMode::Append.to_string(), "append");
833 assert_eq!(DeltaWriteMode::Overwrite.to_string(), "overwrite");
834 assert_eq!(DeltaWriteMode::Upsert.to_string(), "upsert");
835 }
836
837 #[test]
838 fn test_delivery_guarantee_parse() {
839 assert_eq!(
840 "at-least-once".parse::<DeliveryGuarantee>().unwrap(),
841 DeliveryGuarantee::AtLeastOnce
842 );
843 assert_eq!(
844 "at_least_once".parse::<DeliveryGuarantee>().unwrap(),
845 DeliveryGuarantee::AtLeastOnce
846 );
847 assert_eq!(
848 "exactly-once".parse::<DeliveryGuarantee>().unwrap(),
849 DeliveryGuarantee::ExactlyOnce
850 );
851 assert_eq!(
852 "exactly_once".parse::<DeliveryGuarantee>().unwrap(),
853 DeliveryGuarantee::ExactlyOnce
854 );
855 assert!("unknown".parse::<DeliveryGuarantee>().is_err());
856 }
857
858 #[test]
859 fn test_delivery_guarantee_display() {
860 assert_eq!(DeliveryGuarantee::AtLeastOnce.to_string(), "at-least-once");
861 assert_eq!(DeliveryGuarantee::ExactlyOnce.to_string(), "exactly-once");
862 }
863
864 #[test]
865 fn test_partition_columns_empty_filter() {
866 let mut pairs = required_pairs();
867 pairs.push(("partition.columns", "a,,b, ,c"));
868 let config = make_config(&pairs);
869 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
870 assert_eq!(cfg.partition_columns, vec!["a", "b", "c"]);
871 }
872
873 #[test]
876 fn test_s3_path_requires_region() {
877 let config = make_config(&[("table.path", "s3://my-bucket/trades")]);
878 let result = DeltaLakeSinkConfig::from_config(&config);
879 assert!(result.is_err());
880 let err = result.unwrap_err().to_string();
881 assert!(err.contains("aws_region"), "error: {err}");
882 }
883
884 #[test]
885 fn test_s3_path_with_region_and_credentials() {
886 let config = make_config(&[
887 ("table.path", "s3://my-bucket/trades"),
888 ("storage.aws_region", "us-east-1"),
889 ("storage.aws_access_key_id", "AKID123"),
890 ("storage.aws_secret_access_key", "SECRET"),
891 ]);
892 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
893 assert_eq!(cfg.storage_options["aws_region"], "us-east-1");
894 assert_eq!(cfg.storage_options["aws_access_key_id"], "AKID123");
895 }
896
897 #[test]
898 fn test_s3_path_with_region_only_warns_no_error() {
899 let config = make_config(&[
901 ("table.path", "s3://my-bucket/trades"),
902 ("storage.aws_region", "us-east-1"),
903 ]);
904 assert!(DeltaLakeSinkConfig::from_config(&config).is_ok());
905 }
906
907 #[test]
908 fn test_s3_path_access_key_without_secret_errors() {
909 let config = make_config(&[
910 ("table.path", "s3://my-bucket/trades"),
911 ("storage.aws_region", "us-east-1"),
912 ("storage.aws_access_key_id", "AKID123"),
913 ]);
914 let result = DeltaLakeSinkConfig::from_config(&config);
915 assert!(result.is_err());
916 let err = result.unwrap_err().to_string();
917 assert!(err.contains("aws_secret_access_key"), "error: {err}");
918 }
919
920 #[test]
921 fn test_azure_path_requires_account_name() {
922 let config = make_config(&[("table.path", "az://my-container/trades")]);
923 let result = DeltaLakeSinkConfig::from_config(&config);
924 assert!(result.is_err());
925 let err = result.unwrap_err().to_string();
926 assert!(err.contains("azure_storage_account_name"), "error: {err}");
927 }
928
929 #[test]
930 fn test_azure_path_with_account_name_and_key() {
931 let config = make_config(&[
932 ("table.path", "az://my-container/trades"),
933 ("storage.azure_storage_account_name", "myaccount"),
934 ("storage.azure_storage_account_key", "base64key=="),
935 ]);
936 assert!(DeltaLakeSinkConfig::from_config(&config).is_ok());
937 }
938
939 #[test]
940 fn test_gcs_path_always_valid() {
941 let config = make_config(&[("table.path", "gs://my-bucket/trades")]);
943 assert!(DeltaLakeSinkConfig::from_config(&config).is_ok());
944 }
945
946 #[test]
947 fn test_local_path_no_cloud_validation() {
948 let config = make_config(&[("table.path", "/data/warehouse/trades")]);
949 assert!(DeltaLakeSinkConfig::from_config(&config).is_ok());
950 }
951
952 #[test]
953 fn test_display_storage_options_redacts_secrets() {
954 let mut cfg = DeltaLakeSinkConfig::new("s3://bucket/path");
955 cfg.storage_options
956 .insert("aws_region".to_string(), "us-east-1".to_string());
957 cfg.storage_options.insert(
958 "aws_secret_access_key".to_string(),
959 "TOP_SECRET".to_string(),
960 );
961
962 let display = cfg.display_storage_options();
963 assert!(display.contains("aws_region=us-east-1"));
964 assert!(display.contains("aws_secret_access_key=***"));
965 assert!(!display.contains("TOP_SECRET"));
966 }
967
968 #[test]
969 fn test_display_storage_options_empty() {
970 let cfg = DeltaLakeSinkConfig::new("/local/path");
971 assert!(cfg.display_storage_options().is_empty());
972 }
973
974 #[test]
977 fn test_catalog_type_parse() {
978 assert_eq!(
979 "none".parse::<DeltaCatalogType>().unwrap(),
980 DeltaCatalogType::None
981 );
982 assert_eq!(
983 "glue".parse::<DeltaCatalogType>().unwrap(),
984 DeltaCatalogType::Glue
985 );
986 assert!(matches!(
987 "unity".parse::<DeltaCatalogType>().unwrap(),
988 DeltaCatalogType::Unity { .. }
989 ));
990 assert!("unknown".parse::<DeltaCatalogType>().is_err());
991 }
992
993 #[test]
994 fn test_catalog_type_display() {
995 assert_eq!(DeltaCatalogType::None.to_string(), "none");
996 assert_eq!(DeltaCatalogType::Glue.to_string(), "glue");
997 assert_eq!(
998 DeltaCatalogType::Unity {
999 workspace_url: "url".into(),
1000 access_token: "tok".into()
1001 }
1002 .to_string(),
1003 "unity"
1004 );
1005 }
1006
1007 #[test]
1008 fn test_catalog_none_default() {
1009 let config = make_config(&required_pairs());
1010 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1011 assert_eq!(cfg.catalog_type, DeltaCatalogType::None);
1012 assert!(cfg.catalog_database.is_none());
1013 assert!(cfg.catalog_name.is_none());
1014 assert!(cfg.catalog_schema.is_none());
1015 assert!(cfg.catalog_storage_location.is_none());
1016 }
1017
1018 #[cfg(feature = "delta-lake-glue")]
1019 #[test]
1020 fn test_catalog_glue_valid() {
1021 let mut pairs = required_pairs();
1022 pairs.extend_from_slice(&[
1023 ("catalog.type", "glue"),
1024 ("catalog.database", "my_database"),
1025 ]);
1026 let config = make_config(&pairs);
1027 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1028 assert_eq!(cfg.catalog_type, DeltaCatalogType::Glue);
1029 assert_eq!(cfg.catalog_database.as_deref(), Some("my_database"));
1030 }
1031
1032 #[cfg(feature = "delta-lake-glue")]
1033 #[test]
1034 fn test_catalog_glue_missing_database() {
1035 let mut pairs = required_pairs();
1036 pairs.push(("catalog.type", "glue"));
1037 let config = make_config(&pairs);
1038 let result = DeltaLakeSinkConfig::from_config(&config);
1039 assert!(result.is_err());
1040 let err = result.unwrap_err().to_string();
1041 assert!(err.contains("catalog.database"), "error: {err}");
1042 }
1043
1044 #[cfg(feature = "delta-lake-unity")]
1045 #[test]
1046 fn test_catalog_unity_valid() {
1047 let mut pairs = required_pairs();
1048 pairs.extend_from_slice(&[
1049 ("catalog.type", "unity"),
1050 ("catalog.workspace_url", "https://my.databricks.com"),
1051 ("catalog.access_token", "dapi123"),
1052 ("catalog.name", "main"),
1053 ("catalog.schema", "default"),
1054 ]);
1055 let config = make_config(&pairs);
1056 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1057 assert!(matches!(
1058 cfg.catalog_type,
1059 DeltaCatalogType::Unity {
1060 ref workspace_url,
1061 ref access_token
1062 }
1063 if workspace_url == "https://my.databricks.com"
1064 && access_token == "dapi123"
1065 ));
1066 assert_eq!(cfg.catalog_name.as_deref(), Some("main"));
1067 assert_eq!(cfg.catalog_schema.as_deref(), Some("default"));
1068 }
1069
1070 #[cfg(feature = "delta-lake-unity")]
1071 #[test]
1072 fn test_catalog_unity_missing_workspace_url() {
1073 let mut pairs = required_pairs();
1074 pairs.extend_from_slice(&[
1075 ("catalog.type", "unity"),
1076 ("catalog.access_token", "dapi123"),
1077 ("catalog.name", "main"),
1078 ("catalog.schema", "default"),
1079 ]);
1080 let config = make_config(&pairs);
1081 let result = DeltaLakeSinkConfig::from_config(&config);
1082 assert!(result.is_err());
1083 let err = result.unwrap_err().to_string();
1084 assert!(err.contains("workspace_url"), "error: {err}");
1085 }
1086
1087 #[cfg(feature = "delta-lake-unity")]
1088 #[test]
1089 fn test_catalog_unity_missing_access_token() {
1090 let mut pairs = required_pairs();
1091 pairs.extend_from_slice(&[
1092 ("catalog.type", "unity"),
1093 ("catalog.workspace_url", "https://my.databricks.com"),
1094 ("catalog.name", "main"),
1095 ("catalog.schema", "default"),
1096 ]);
1097 let config = make_config(&pairs);
1098 let result = DeltaLakeSinkConfig::from_config(&config);
1099 assert!(result.is_err());
1100 let err = result.unwrap_err().to_string();
1101 assert!(err.contains("access_token"), "error: {err}");
1102 }
1103
1104 #[test]
1105 fn test_catalog_storage_location_default_none() {
1106 let config = make_config(&required_pairs());
1107 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1108 assert!(cfg.catalog_storage_location.is_none());
1109 }
1110
1111 #[test]
1112 fn test_catalog_storage_location_parsed() {
1113 let mut pairs = required_pairs();
1114 pairs.push(("catalog.storage.location", "s3://bucket/warehouse/table"));
1115 let config = make_config(&pairs);
1116 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1117 assert_eq!(
1118 cfg.catalog_storage_location.as_deref(),
1119 Some("s3://bucket/warehouse/table")
1120 );
1121 }
1122
1123 #[test]
1126 fn test_parquet_config_defaults() {
1127 let cfg = ParquetWriteConfig::default();
1128 assert_eq!(cfg.compression, "zstd");
1129 assert_eq!(cfg.compression_level, 1);
1130 assert!(cfg.dictionary_enabled);
1131 assert_eq!(cfg.statistics, "page");
1132 assert!(cfg.bloom_filter_columns.is_empty());
1133 assert!((cfg.bloom_filter_fpp - 0.01).abs() < f64::EPSILON);
1134 assert_eq!(cfg.bloom_filter_ndv, 0);
1135 assert_eq!(cfg.max_row_group_size, 1_000_000);
1136 }
1137
1138 #[test]
1139 fn test_parquet_compression_parsing() {
1140 for codec in &["zstd", "snappy", "lz4", "gzip", "none"] {
1141 let mut pairs = required_pairs();
1142 pairs.push(("parquet.compression", codec));
1143 let config = make_config(&pairs);
1144 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1145 assert_eq!(cfg.parquet.compression, *codec);
1146 }
1147 }
1148
1149 #[test]
1150 fn test_parquet_compression_level_parsing() {
1151 let mut pairs = required_pairs();
1152 pairs.push(("parquet.compression.level", "5"));
1153 let config = make_config(&pairs);
1154 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1155 assert_eq!(cfg.parquet.compression_level, 5);
1156 }
1157
1158 #[test]
1159 fn test_parquet_compression_level_invalid() {
1160 let mut pairs = required_pairs();
1161 pairs.push(("parquet.compression.level", "abc"));
1162 let config = make_config(&pairs);
1163 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
1164 }
1165
1166 #[test]
1167 fn test_parquet_bloom_filter_columns_parsing() {
1168 let mut pairs = required_pairs();
1169 pairs.push((
1170 "parquet.bloom.filter.columns",
1171 " user_id , event_type , ts ",
1172 ));
1173 let config = make_config(&pairs);
1174 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1175 assert_eq!(
1176 cfg.parquet.bloom_filter_columns,
1177 vec!["user_id", "event_type", "ts"]
1178 );
1179 }
1180
1181 #[test]
1182 fn test_parquet_bloom_filter_fpp_validation() {
1183 let mut pairs = required_pairs();
1185 pairs.push(("parquet.bloom.filter.fpp", "0.0"));
1186 let config = make_config(&pairs);
1187 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
1188
1189 let mut pairs = required_pairs();
1191 pairs.push(("parquet.bloom.filter.fpp", "1.0"));
1192 let config = make_config(&pairs);
1193 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
1194 }
1195
1196 #[test]
1197 fn test_parquet_max_row_group_size_zero_rejected() {
1198 let mut pairs = required_pairs();
1199 pairs.push(("parquet.max.row.group.size", "0"));
1200 let config = make_config(&pairs);
1201 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
1202 }
1203
1204 #[test]
1205 fn test_parquet_statistics_parsing() {
1206 for stat in &["none", "chunk", "page"] {
1207 let mut pairs = required_pairs();
1208 pairs.push(("parquet.statistics", stat));
1209 let config = make_config(&pairs);
1210 let cfg = DeltaLakeSinkConfig::from_config(&config).unwrap();
1211 assert_eq!(cfg.parquet.statistics, *stat);
1212 }
1213 }
1214
1215 #[test]
1216 fn test_parquet_invalid_statistics_rejected() {
1217 let mut pairs = required_pairs();
1218 pairs.push(("parquet.statistics", "full"));
1219 let config = make_config(&pairs);
1220 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
1221 }
1222
1223 #[test]
1224 fn test_parquet_invalid_compression_rejected() {
1225 let mut pairs = required_pairs();
1226 pairs.push(("parquet.compression", "brotli"));
1227 let config = make_config(&pairs);
1228 assert!(DeltaLakeSinkConfig::from_config(&config).is_err());
1229 }
1230
1231 #[cfg(feature = "delta-lake")]
1232 #[test]
1233 fn test_writer_properties_default_zstd() {
1234 let cfg = ParquetWriteConfig::default();
1235 assert!(cfg.to_writer_properties().is_ok());
1236 }
1237
1238 #[cfg(feature = "delta-lake")]
1239 #[test]
1240 fn test_writer_properties_invalid_codec() {
1241 let mut cfg = ParquetWriteConfig::default();
1242 cfg.compression = "brotli".to_string();
1243 assert!(cfg.to_writer_properties().is_err());
1244 }
1245
1246 #[cfg(feature = "delta-lake")]
1247 #[test]
1248 fn test_writer_properties_with_bloom_filters() {
1249 let mut cfg = ParquetWriteConfig::default();
1250 cfg.bloom_filter_columns = vec!["user_id".to_string(), "event_type".to_string()];
1251 assert!(cfg.to_writer_properties().is_ok());
1252 }
1253}