Skip to main content

laminar_connectors/lakehouse/
delta_config.rs

1//! Delta Lake sink config. Parsed from SQL `WITH (...)` via
2//! [`DeltaLakeSinkConfig::from_config`].
3#![allow(clippy::disallowed_types)] // cold path: lakehouse configuration
4
5use 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/// Configuration for the Delta Lake sink connector.
20///
21/// Parsed from SQL `WITH (...)` clause options or constructed programmatically.
22#[derive(Debug, Clone)]
23pub struct DeltaLakeSinkConfig {
24    /// Path to the Delta Lake table (local, `s3://`, `az://`, `gs://`).
25    pub table_path: String,
26
27    /// Columns to partition by (e.g., `["trade_date", "hour"]`).
28    pub partition_columns: Vec<String>,
29
30    /// Target Parquet file size in bytes (default: 128 MB).
31    pub target_file_size: usize,
32
33    /// Maximum number of records to buffer before flushing to Parquet.
34    pub max_buffer_records: usize,
35
36    /// Maximum time to buffer records before flushing.
37    pub max_buffer_duration: Duration,
38
39    /// Whether to enable schema evolution (auto-merge new columns).
40    pub schema_evolution: bool,
41
42    /// Write mode: Append, Overwrite, or Upsert (CDC merge).
43    pub write_mode: DeltaWriteMode,
44
45    /// Key columns for upsert/merge operations (required for Upsert mode).
46    pub merge_key_columns: Vec<String>,
47
48    /// Storage options (S3 credentials, Azure keys, etc.).
49    pub storage_options: HashMap<String, String>,
50
51    /// Delivery guarantee: `AtLeastOnce` or `ExactlyOnce`.
52    pub delivery_guarantee: DeliveryGuarantee,
53
54    /// Catalog type for table discovery.
55    pub catalog_type: DeltaCatalogType,
56
57    /// Catalog database name (required for Glue).
58    pub catalog_database: Option<String>,
59
60    /// Catalog name (required for Unity).
61    pub catalog_name: Option<String>,
62
63    /// Catalog schema name (required for Unity).
64    pub catalog_schema: Option<String>,
65
66    /// Storage location for auto-created Unity Catalog external tables.
67    /// When set and the `uc://` table doesn't exist, the sink creates it
68    /// via the Unity Catalog REST API at this storage location.
69    pub catalog_storage_location: Option<String>,
70
71    /// End-to-end timeout for one Delta write, including table reopen and all
72    /// optimistic commit retries (default: 30s).
73    pub write_timeout: Duration,
74
75    /// Parquet writer properties (compression, bloom filters, statistics, etc.).
76    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, // 128 MB
85            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    /// Creates a minimal config for testing.
105    #[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    /// Parses a sink config from a [`ConnectorConfig`] (SQL WITH clause).
114    ///
115    /// # Required keys
116    ///
117    /// - `table.path` - Path to Delta Lake table
118    ///
119    /// # Errors
120    ///
121    /// Returns `ConnectorError::MissingConfig` if required keys are absent,
122    /// or `ConnectorError::ConfigurationError` on invalid values.
123    #[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        // ── Catalog configuration ──
189        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        // Unity-specific: populate workspace_url and access_token into the enum variant.
206        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        // ── Parquet writer configuration ──
229        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        // Resolve storage credentials: explicit options + environment variable fallbacks.
275        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        // Map LogStore configuration keys to delta-rs storage options.
280        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    /// Formats the storage options for safe logging with secrets redacted.
294    #[must_use]
295    pub fn display_storage_options(&self) -> String {
296        SecretMasker::display_map(&self.storage_options)
297    }
298
299    /// Validates the configuration for consistency.
300    ///
301    /// # Errors
302    ///
303    /// Returns `ConnectorError::ConfigurationError` on invalid combinations.
304    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        // Eagerly validate that WriterProperties can be built so invalid
357        // codec/level combos are caught at config time, not first write.
358        #[cfg(feature = "delta-lake")]
359        {
360            self.parquet.to_writer_properties()?;
361        }
362        self.validate_catalog()?;
363
364        // Validate cloud storage credentials (skip when catalog resolves the path).
365        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    /// Validates catalog-specific requirements.
383    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/// Delta Lake write mode.
447#[derive(Debug, Clone, Copy, PartialEq, Eq)]
448pub enum DeltaWriteMode {
449    /// Append-only: all records are inserts. Most efficient for immutable streams.
450    Append,
451    /// Overwrite: replace partition contents. Used for batch-style recomputation.
452    Overwrite,
453    /// Upsert/Merge: CDC-style insert/update/delete via MERGE statement.
454    /// Requires `merge_key_columns` to be set. Integrates with Z-sets.
455    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/// Delta Lake catalog type for table discovery.
482///
483/// Catalogs enable referencing tables by logical names instead of raw paths.
484#[derive(Debug, Clone, PartialEq, Eq, Default)]
485pub enum DeltaCatalogType {
486    /// No catalog — table path is a direct file or cloud URI.
487    #[default]
488    None,
489    /// AWS Glue Data Catalog.
490    Glue,
491    /// Databricks Unity Catalog.
492    Unity {
493        /// Databricks workspace URL (e.g., `https://xxx.cloud.databricks.com`).
494        workspace_url: String,
495        /// Databricks access token.
496        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/// Configuration for Parquet writer properties (compression, dictionary
527/// encoding, statistics, bloom filters, row group sizing).
528#[derive(Debug, Clone)]
529pub struct ParquetWriteConfig {
530    /// Compression codec: `"zstd"`, `"snappy"`, `"lz4"`, `"gzip"`, or `"none"`.
531    pub compression: String,
532    /// Compression level (default: 1 — ZSTD L1 for hot writes).
533    pub compression_level: i32,
534    /// Whether to enable dictionary encoding (default: true).
535    pub dictionary_enabled: bool,
536    /// Statistics granularity: `"none"`, `"chunk"`, or `"page"` (default: `"page"`).
537    pub statistics: String,
538    /// Columns to build bloom filters for (default: empty).
539    pub bloom_filter_columns: Vec<String>,
540    /// Bloom filter false-positive probability (default: 0.01).
541    pub bloom_filter_fpp: f64,
542    /// Bloom filter expected number of distinct values (0 = parquet default).
543    pub bloom_filter_ndv: u64,
544    /// Maximum rows per row group (default: 1,000,000).
545    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    /// Builds `WriterProperties` for hot-path writes (uses `compression_level`).
566    ///
567    /// # Errors
568    ///
569    /// Returns `ConnectorError::ConfigurationError` on invalid codec/level.
570    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    /// Shared builder: maps string codec → `Compression`, sets dictionary,
577    /// statistics, bloom filters, and row group size.
578    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    // ── Config parsing tests ──
664
665    #[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    // ── Enum tests ──
808
809    #[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    // ── Cloud storage integration tests ──
874
875    #[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        // Missing credentials is a warning (IAM fallback), not a hard error.
900        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        // GCS missing credentials is warning-only (Application Default Credentials).
942        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    // ── Catalog tests ──
975
976    #[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    // ── Parquet config tests ──
1124
1125    #[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        // fpp = 0.0 should be rejected
1184        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        // fpp = 1.0 should be rejected
1190        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}