Skip to main content

laminar_connectors/lakehouse/
mod.rs

1//! Lakehouse connectors (Delta Lake, Apache Iceberg).
2
3// Versioned envelope for Delta coordinated-commit descriptors.
4#[cfg(feature = "delta-lake")]
5mod commit_descriptor;
6
7// Delta Lake modules
8pub mod delta;
9pub mod delta_config;
10#[cfg(feature = "delta-lake")]
11pub mod delta_io;
12#[cfg(feature = "delta-lake")]
13pub mod delta_lookup;
14pub mod delta_metrics;
15#[cfg(feature = "delta-lake")]
16pub mod delta_reference;
17pub mod delta_source;
18pub mod delta_source_config;
19#[cfg(feature = "delta-lake")]
20pub mod delta_table_provider;
21#[cfg(feature = "delta-lake-unity")]
22pub(crate) mod unity_catalog;
23
24// Apache Iceberg modules
25pub mod iceberg;
26pub mod iceberg_config;
27#[cfg(feature = "iceberg")]
28pub mod iceberg_incremental;
29#[cfg(feature = "iceberg")]
30pub mod iceberg_io;
31#[cfg(feature = "iceberg")]
32pub mod iceberg_lookup;
33#[cfg(feature = "iceberg")]
34pub mod iceberg_reference;
35pub mod iceberg_source;
36
37// Common metrics
38pub mod metrics;
39#[cfg(any(test, feature = "delta-lake", feature = "iceberg"))]
40mod snapshot_schema;
41
42// Re-export Delta Lake types at module level.
43pub use delta::DeltaLakeSink;
44pub use delta_config::{DeltaCatalogType, DeltaLakeSinkConfig, DeltaWriteMode};
45#[cfg(feature = "delta-lake")]
46pub use delta_lookup::{DeltaLookupSource, DeltaLookupSourceConfig};
47pub use delta_metrics::DeltaLakeSinkMetrics;
48#[cfg(feature = "delta-lake")]
49pub use delta_reference::DeltaReferenceTableSource;
50pub use delta_source::DeltaSource;
51pub use delta_source_config::{DeltaReadMode, DeltaSourceConfig, SchemaEvolutionAction};
52pub use metrics::LakehouseSinkMetrics;
53
54// Re-export Iceberg types at module level.
55pub use iceberg::IcebergSink;
56pub use iceberg_config::{
57    IcebergCatalogConfig, IcebergCatalogType, IcebergSinkConfig, IcebergSourceConfig,
58};
59#[cfg(feature = "iceberg")]
60pub use iceberg_lookup::{IcebergLookupSource, IcebergLookupSourceConfig};
61#[cfg(feature = "iceberg")]
62pub use iceberg_reference::IcebergReferenceTableSource;
63pub use iceberg_source::IcebergSource;
64
65use std::sync::Arc;
66
67use crate::config::{ConfigKeySpec, ConnectorInfo};
68use crate::registry::ConnectorRegistry;
69
70/// Registers the Delta Lake sink connector with the given registry.
71///
72/// # Errors
73///
74/// Returns the registry error when a name is already registered or the registry is frozen.
75pub fn register_delta_lake_sink(
76    registry: &ConnectorRegistry,
77) -> Result<(), crate::error::ConnectorError> {
78    let info = ConnectorInfo {
79        name: "delta-lake".to_string(),
80        display_name: "Delta Lake Sink".to_string(),
81        version: env!("CARGO_PKG_VERSION").to_string(),
82        is_source: false,
83        is_sink: true,
84        config_keys: delta_lake_config_keys(),
85    };
86
87    registry.register_sink(
88        "delta-lake",
89        info,
90        Arc::new(|config, registry: Option<&Arc<prometheus::Registry>>| {
91            Ok(Box::new(DeltaLakeSink::new(
92                DeltaLakeSinkConfig::from_config(config)?,
93                registry.map(Arc::as_ref),
94            )))
95        }),
96    )
97}
98
99/// Registers the Delta Lake source connector with the given registry.
100///
101/// # Errors
102///
103/// Returns the registry error when a name is already registered or the registry is frozen.
104pub fn register_delta_lake_source(
105    registry: &ConnectorRegistry,
106) -> Result<(), crate::error::ConnectorError> {
107    let info = ConnectorInfo {
108        name: "delta-lake".to_string(),
109        display_name: "Delta Lake Source".to_string(),
110        version: env!("CARGO_PKG_VERSION").to_string(),
111        is_source: true,
112        is_sink: false,
113        config_keys: delta_lake_source_config_keys(),
114    };
115
116    registry.register_source(
117        "delta-lake",
118        info.clone(),
119        Arc::new(|registry: Option<&Arc<prometheus::Registry>>| {
120            Ok(Box::new(DeltaSource::new(
121                DeltaSourceConfig::default(),
122                registry.map(Arc::as_ref),
123            )))
124        }),
125    )?;
126
127    // Register finite startup snapshots for replicated reference tables.
128    #[cfg(feature = "delta-lake")]
129    registry.register_table_source(
130        "delta-lake",
131        info.clone(),
132        Arc::new(|config, declared_schema| {
133            Ok(Box::new(DeltaReferenceTableSource::from_connector_config(
134                config,
135                declared_schema,
136            )?))
137        }),
138    )?;
139
140    // Register lookup source factory for on-demand/partial cache mode.
141    #[cfg(feature = "delta-lake")]
142    registry.register_lookup_source("delta-lake", info, Arc::new(DeltaLookupFactory))?;
143
144    Ok(())
145}
146
147#[cfg(feature = "delta-lake")]
148struct DeltaLookupFactory;
149
150#[cfg(feature = "delta-lake")]
151#[async_trait::async_trait]
152impl crate::registry::LookupSourceFactory for DeltaLookupFactory {
153    async fn build(
154        &self,
155        config: crate::config::ConnectorConfig,
156        _declared_schema: Option<arrow_schema::SchemaRef>,
157    ) -> Result<Arc<dyn laminar_core::lookup::source::LookupSourceDyn>, crate::error::ConnectorError>
158    {
159        use crate::lakehouse::delta_source_config::DeltaSourceConfig;
160        let pk_columns: Vec<String> = config
161            .get("_primary_key_columns")
162            .unwrap_or("")
163            .split(',')
164            .map(|s| s.trim().to_string())
165            .filter(|s| !s.is_empty())
166            .collect();
167
168        if pk_columns.is_empty() {
169            return Err(crate::error::ConnectorError::ConfigurationError(
170                "delta-lake lookup source requires primary key columns".into(),
171            ));
172        }
173
174        let src_config = DeltaSourceConfig::from_config(&config)?;
175
176        let (resolved_path, resolved_opts) = crate::lakehouse::delta_io::resolve_catalog_options(
177            &src_config.catalog_type,
178            src_config.catalog_database.as_deref(),
179            src_config.catalog_name.as_deref(),
180            src_config.catalog_schema.as_deref(),
181            &src_config.table_path,
182            &src_config.storage_options,
183        )
184        .await?;
185
186        let lookup_config = DeltaLookupSourceConfig {
187            table_path: resolved_path,
188            storage_options: resolved_opts,
189            primary_key_columns: pk_columns,
190            table_name: "delta_lookup".to_string(),
191        };
192
193        // `From<LookupError>` preserves transient/non-transient class.
194        let source = DeltaLookupSource::open(lookup_config).await?;
195
196        Ok(Arc::new(source) as Arc<dyn laminar_core::lookup::source::LookupSourceDyn>)
197    }
198}
199
200/// Registers the Iceberg sink connector with the given registry.
201///
202/// # Errors
203///
204/// Returns the registry error when a name is already registered or the registry is frozen.
205pub fn register_iceberg_sink(
206    registry: &ConnectorRegistry,
207) -> Result<(), crate::error::ConnectorError> {
208    let info = ConnectorInfo {
209        name: "iceberg".to_string(),
210        display_name: "Apache Iceberg Sink".to_string(),
211        version: env!("CARGO_PKG_VERSION").to_string(),
212        is_source: false,
213        is_sink: true,
214        config_keys: iceberg_sink_config_keys(),
215    };
216
217    registry.register_sink(
218        "iceberg",
219        info,
220        Arc::new(|config, registry: Option<&Arc<prometheus::Registry>>| {
221            Ok(Box::new(IcebergSink::new(
222                IcebergSinkConfig::from_config(config)?,
223                registry.map(Arc::as_ref),
224            )))
225        }),
226    )
227}
228
229/// Registers the Iceberg source connector with the given registry.
230///
231/// # Errors
232///
233/// Returns the registry error when a name is already registered or the registry is frozen.
234pub fn register_iceberg_source(
235    registry: &ConnectorRegistry,
236) -> Result<(), crate::error::ConnectorError> {
237    let info = ConnectorInfo {
238        name: "iceberg".to_string(),
239        display_name: "Apache Iceberg Source".to_string(),
240        version: env!("CARGO_PKG_VERSION").to_string(),
241        is_source: true,
242        is_sink: false,
243        config_keys: iceberg_source_config_keys(),
244    };
245
246    registry.register_source(
247        "iceberg",
248        info.clone(),
249        Arc::new(|registry: Option<&Arc<prometheus::Registry>>| {
250            Ok(Box::new(IcebergSource::new(
251                IcebergSourceConfig {
252                    catalog: IcebergCatalogConfig {
253                        catalog_type: IcebergCatalogType::Rest,
254                        catalog_uri: "http://localhost:8181".to_string(),
255                        warehouse: "s3://default/wh".to_string(),
256                        storage_type: None,
257                        namespace: "default".to_string(),
258                        table_name: "default".to_string(),
259                        properties: std::collections::HashMap::new(),
260                    },
261                    poll_interval: std::time::Duration::from_secs(60),
262                    snapshot_id: None,
263                    select_columns: Vec::new(),
264                },
265                registry.map(Arc::as_ref),
266            )))
267        }),
268    )?;
269
270    // Register finite startup snapshots for replicated reference tables.
271    #[cfg(feature = "iceberg")]
272    registry.register_table_source(
273        "iceberg",
274        info.clone(),
275        Arc::new(|config, declared_schema| {
276            Ok(Box::new(
277                IcebergReferenceTableSource::from_connector_config(config, declared_schema)?,
278            ))
279        }),
280    )?;
281
282    // Register lookup source factory for on-demand/partial cache mode.
283    #[cfg(feature = "iceberg")]
284    registry.register_lookup_source("iceberg", info, Arc::new(IcebergLookupFactory))?;
285
286    Ok(())
287}
288
289#[cfg(feature = "iceberg")]
290struct IcebergLookupFactory;
291
292#[cfg(feature = "iceberg")]
293#[async_trait::async_trait]
294impl crate::registry::LookupSourceFactory for IcebergLookupFactory {
295    async fn build(
296        &self,
297        config: crate::config::ConnectorConfig,
298        _declared_schema: Option<arrow_schema::SchemaRef>,
299    ) -> Result<Arc<dyn laminar_core::lookup::source::LookupSourceDyn>, crate::error::ConnectorError>
300    {
301        use crate::lakehouse::iceberg_config::IcebergCatalogConfig;
302        let pk_columns: Vec<String> = config
303            .get("_primary_key_columns")
304            .unwrap_or("")
305            .split(',')
306            .map(|s| s.trim().to_string())
307            .filter(|s| !s.is_empty())
308            .collect();
309
310        if pk_columns.is_empty() {
311            return Err(crate::error::ConnectorError::ConfigurationError(
312                "iceberg lookup source requires primary key columns".into(),
313            ));
314        }
315
316        let catalog = IcebergCatalogConfig::from_config(&config)?;
317        let lookup_config = IcebergLookupSourceConfig {
318            catalog,
319            primary_key_columns: pk_columns,
320        };
321
322        let source = IcebergLookupSource::open(lookup_config).await?;
323        Ok(Arc::new(source) as Arc<dyn laminar_core::lookup::source::LookupSourceDyn>)
324    }
325}
326
327/// Registers all lakehouse sink connectors (Delta Lake, Iceberg).
328///
329/// # Errors
330///
331/// Returns the first registry error.
332pub fn register_lakehouse_sinks(
333    registry: &ConnectorRegistry,
334) -> Result<(), crate::error::ConnectorError> {
335    register_delta_lake_sink(registry)?;
336    register_iceberg_sink(registry)
337}
338
339/// Registers all lakehouse source connectors.
340///
341/// # Errors
342///
343/// Returns the first registry error.
344pub fn register_lakehouse_sources(
345    registry: &ConnectorRegistry,
346) -> Result<(), crate::error::ConnectorError> {
347    register_delta_lake_source(registry)?;
348    register_iceberg_source(registry)
349}
350
351#[allow(clippy::too_many_lines)]
352fn delta_lake_config_keys() -> Vec<ConfigKeySpec> {
353    vec![
354        ConfigKeySpec::required(
355            "table.path",
356            "Path to Delta Lake table (local, s3://, az://, gs://)",
357        ),
358        ConfigKeySpec::optional(
359            "partition.columns",
360            "Comma-separated partition column names",
361            "",
362        ),
363        ConfigKeySpec::optional(
364            "target.file.size",
365            "Target Parquet file size in bytes",
366            "134217728",
367        ),
368        ConfigKeySpec::optional(
369            "max.buffer.records",
370            "Maximum records to buffer before flushing",
371            "100000",
372        ),
373        ConfigKeySpec::optional(
374            "max.buffer.duration.ms",
375            "Maximum time to buffer before flushing (ms)",
376            "60000",
377        ),
378        ConfigKeySpec::optional(
379            "schema.evolution",
380            "Enable automatic schema evolution (additive columns)",
381            "false",
382        ),
383        ConfigKeySpec::optional(
384            "write.mode",
385            "Write mode: append, overwrite, upsert",
386            "append",
387        ),
388        ConfigKeySpec::optional(
389            "merge.key.columns",
390            "Key columns for upsert MERGE (required for upsert mode)",
391            "",
392        ),
393        // ── Catalog configuration ──
394        ConfigKeySpec::optional("catalog.type", "Catalog type: none, glue, unity", "none"),
395        ConfigKeySpec::optional(
396            "catalog.database",
397            "Catalog database name (required for Glue)",
398            "",
399        ),
400        ConfigKeySpec::optional("catalog.name", "Catalog name (required for Unity)", ""),
401        ConfigKeySpec::optional(
402            "catalog.schema",
403            "Catalog schema name (required for Unity)",
404            "",
405        ),
406        ConfigKeySpec::optional(
407            "catalog.workspace_url",
408            "Databricks workspace URL (required for Unity)",
409            "",
410        ),
411        ConfigKeySpec::optional(
412            "catalog.access_token",
413            "Databricks access token (required for Unity)",
414            "",
415        ),
416        ConfigKeySpec::optional(
417            "catalog.storage.location",
418            "Storage location for auto-created UC external tables (e.g. s3://bucket/path)",
419            "",
420        ),
421        // ── LogStore configuration ──
422        ConfigKeySpec::optional(
423            "storage.s3_locking_provider",
424            "S3 locking provider: 'dynamodb' for DynamoDB-backed log store",
425            "",
426        ),
427        ConfigKeySpec::optional(
428            "storage.dynamodb_table_name",
429            "DynamoDB table name for S3 locking (default: delta_log)",
430            "",
431        ),
432        // ── Cloud storage credentials (resolved via StorageCredentialResolver) ──
433        ConfigKeySpec::optional(
434            "storage.aws_access_key_id",
435            "AWS access key ID (falls back to AWS_ACCESS_KEY_ID env var)",
436            "",
437        ),
438        ConfigKeySpec::optional(
439            "storage.aws_secret_access_key",
440            "AWS secret access key (falls back to AWS_SECRET_ACCESS_KEY env var)",
441            "",
442        ),
443        ConfigKeySpec::optional(
444            "storage.aws_region",
445            "AWS region for S3 paths (falls back to AWS_REGION env var)",
446            "",
447        ),
448        ConfigKeySpec::optional(
449            "storage.aws_session_token",
450            "AWS session token for temporary credentials (falls back to AWS_SESSION_TOKEN)",
451            "",
452        ),
453        ConfigKeySpec::optional(
454            "storage.aws_endpoint",
455            "Custom S3 endpoint (MinIO, LocalStack; falls back to AWS_ENDPOINT_URL)",
456            "",
457        ),
458        ConfigKeySpec::optional(
459            "storage.aws_profile",
460            "AWS profile name (falls back to AWS_PROFILE env var)",
461            "",
462        ),
463        ConfigKeySpec::optional(
464            "storage.azure_storage_account_name",
465            "Azure storage account name (falls back to AZURE_STORAGE_ACCOUNT_NAME)",
466            "",
467        ),
468        ConfigKeySpec::optional(
469            "storage.azure_storage_account_key",
470            "Azure storage account key (falls back to AZURE_STORAGE_ACCOUNT_KEY)",
471            "",
472        ),
473        ConfigKeySpec::optional(
474            "storage.azure_storage_sas_token",
475            "Azure SAS token (falls back to AZURE_STORAGE_SAS_TOKEN)",
476            "",
477        ),
478        ConfigKeySpec::optional(
479            "storage.azure_storage_client_id",
480            "Azure client ID for service principal auth (falls back to AZURE_CLIENT_ID)",
481            "",
482        ),
483        ConfigKeySpec::optional(
484            "storage.google_service_account_path",
485            "Path to GCS service account JSON (falls back to GOOGLE_APPLICATION_CREDENTIALS)",
486            "",
487        ),
488        ConfigKeySpec::optional(
489            "storage.google_service_account_key",
490            "Inline GCS service account JSON (falls back to GOOGLE_SERVICE_ACCOUNT_KEY)",
491            "",
492        ),
493    ]
494}
495
496fn delta_lake_source_config_keys() -> Vec<ConfigKeySpec> {
497    vec![
498        ConfigKeySpec::required(
499            "table.path",
500            "Path to Delta Lake table (local, s3://, az://, gs://)",
501        ),
502        ConfigKeySpec::optional(
503            "starting.version",
504            "Starting version to read from (default: latest)",
505            "",
506        ),
507        ConfigKeySpec::optional(
508            "poll.interval.ms",
509            "How often to poll for new versions (ms)",
510            "1000",
511        ),
512        ConfigKeySpec::optional(
513            "read.mode",
514            "Read mode: 'incremental' (changes only) or 'snapshot' (full re-read)",
515            "incremental",
516        ),
517        ConfigKeySpec::optional(
518            "partition.filter",
519            "SQL predicate for partition filter pushdown (e.g. \"date = '2024-01-01'\")",
520            "",
521        ),
522        ConfigKeySpec::optional(
523            "schema.evolution.action",
524            "Action on schema change: 'warn' or 'error'",
525            "warn",
526        ),
527        ConfigKeySpec::optional(
528            "cdf.enabled",
529            "Use Change Data Feed for incremental reads (requires CDF on table)",
530            "false",
531        ),
532        // ── Catalog configuration ──
533        ConfigKeySpec::optional("catalog.type", "Catalog type: none, glue, unity", "none"),
534        ConfigKeySpec::optional(
535            "catalog.database",
536            "Catalog database name (required for Glue)",
537            "",
538        ),
539        ConfigKeySpec::optional("catalog.name", "Catalog name (required for Unity)", ""),
540        ConfigKeySpec::optional(
541            "catalog.schema",
542            "Catalog schema name (required for Unity)",
543            "",
544        ),
545        ConfigKeySpec::optional(
546            "catalog.workspace_url",
547            "Databricks workspace URL (required for Unity)",
548            "",
549        ),
550        ConfigKeySpec::optional(
551            "catalog.access_token",
552            "Databricks access token (required for Unity)",
553            "",
554        ),
555        // ── LogStore configuration ──
556        ConfigKeySpec::optional(
557            "storage.s3_locking_provider",
558            "S3 locking provider: 'dynamodb' for DynamoDB-backed log store",
559            "",
560        ),
561        ConfigKeySpec::optional(
562            "storage.dynamodb_table_name",
563            "DynamoDB table name for S3 locking (default: delta_log)",
564            "",
565        ),
566        // ── Cloud storage credentials ──
567        ConfigKeySpec::optional("storage.aws_access_key_id", "AWS access key ID", ""),
568        ConfigKeySpec::optional("storage.aws_secret_access_key", "AWS secret access key", ""),
569        ConfigKeySpec::optional("storage.aws_region", "AWS region for S3 paths", ""),
570        ConfigKeySpec::optional(
571            "storage.azure_storage_account_name",
572            "Azure storage account name",
573            "",
574        ),
575        ConfigKeySpec::optional(
576            "storage.azure_storage_account_key",
577            "Azure storage account key",
578            "",
579        ),
580        ConfigKeySpec::optional(
581            "storage.google_service_account_path",
582            "Path to GCS service account JSON",
583            "",
584        ),
585    ]
586}
587
588fn iceberg_sink_config_keys() -> Vec<ConfigKeySpec> {
589    vec![
590        ConfigKeySpec::required(
591            "catalog.uri",
592            "REST catalog URI (e.g., http://polaris:8181)",
593        ),
594        ConfigKeySpec::required(
595            "warehouse",
596            "Warehouse name (REST catalog) or URL (Hadoop catalog, e.g. s3://bucket/wh)",
597        ),
598        ConfigKeySpec::required("namespace", "Iceberg namespace (e.g., prod)"),
599        ConfigKeySpec::required("table.name", "Table name within the namespace"),
600        ConfigKeySpec::optional("catalog.type", "Catalog type: rest", "rest"),
601        ConfigKeySpec::optional(
602            "storage.type",
603            "Storage backend (s3 | s3a | fs). Required when warehouse is a name, not a URL",
604            "",
605        ),
606        ConfigKeySpec::optional(
607            "compression",
608            "Parquet compression: zstd, snappy, none",
609            "zstd",
610        ),
611        ConfigKeySpec::optional("auto.create", "Auto-create table if not exists", "false"),
612    ]
613}
614
615fn iceberg_source_config_keys() -> Vec<ConfigKeySpec> {
616    vec![
617        ConfigKeySpec::required(
618            "catalog.uri",
619            "REST catalog URI (e.g., http://polaris:8181)",
620        ),
621        ConfigKeySpec::required(
622            "warehouse",
623            "Warehouse location (e.g., s3://bucket/warehouse)",
624        ),
625        ConfigKeySpec::required("namespace", "Iceberg namespace (e.g., prod)"),
626        ConfigKeySpec::required("table.name", "Table name within the namespace"),
627        ConfigKeySpec::optional("catalog.type", "Catalog type: rest", "rest"),
628        ConfigKeySpec::optional(
629            "poll.interval.ms",
630            "How often to poll for new snapshots (ms)",
631            "60000",
632        ),
633        ConfigKeySpec::optional("snapshot.id", "Pin to a specific snapshot ID", ""),
634        ConfigKeySpec::optional(
635            "select.columns",
636            "Comma-separated column names to select (empty = all)",
637            "",
638        ),
639    ]
640}
641
642#[cfg(test)]
643mod tests {
644    use super::*;
645
646    #[test]
647    fn test_register_delta_lake_sink() {
648        let registry = ConnectorRegistry::new();
649        register_delta_lake_sink(&registry).unwrap();
650
651        let info = registry.sink_info("delta-lake");
652        assert!(info.is_some());
653        let info = info.unwrap();
654        assert_eq!(info.name, "delta-lake");
655        assert!(info.is_sink);
656        assert!(!info.is_source);
657        assert!(!info.config_keys.is_empty());
658    }
659
660    #[test]
661    fn test_config_keys_required() {
662        let keys = delta_lake_config_keys();
663        let required: Vec<&str> = keys
664            .iter()
665            .filter(|k| k.required)
666            .map(|k| k.key.as_str())
667            .collect();
668        assert!(required.contains(&"table.path"));
669        assert_eq!(required.len(), 1);
670    }
671
672    #[test]
673    fn test_config_keys_include_cloud_storage() {
674        let keys = delta_lake_config_keys();
675        let key_names: Vec<&str> = keys.iter().map(|k| k.key.as_str()).collect();
676        assert!(key_names.contains(&"storage.aws_access_key_id"));
677        assert!(key_names.contains(&"storage.aws_secret_access_key"));
678        assert!(key_names.contains(&"storage.aws_region"));
679        assert!(key_names.contains(&"storage.azure_storage_account_name"));
680        assert!(key_names.contains(&"storage.azure_storage_account_key"));
681        assert!(key_names.contains(&"storage.google_service_account_path"));
682    }
683
684    #[test]
685    fn test_config_keys_optional_present() {
686        let keys = delta_lake_config_keys();
687        let optional: Vec<&str> = keys
688            .iter()
689            .filter(|k| !k.required)
690            .map(|k| k.key.as_str())
691            .collect();
692        assert!(optional.contains(&"partition.columns"));
693        assert!(optional.contains(&"target.file.size"));
694        assert!(optional.contains(&"write.mode"));
695        assert!(!optional.contains(&"delivery.guarantee"));
696        assert!(optional.contains(&"merge.key.columns"));
697        assert!(optional.contains(&"schema.evolution"));
698        assert!(!optional.contains(&"checkpoint.interval"));
699        assert!(!optional.contains(&"max.commit.retries"));
700        assert!(!optional.iter().any(|key| key.starts_with("compaction.")));
701        assert!(!optional.contains(&"vacuum.retention.hours"));
702        assert!(!optional.contains(&"writer.id"));
703        // Catalog keys
704        assert!(optional.contains(&"catalog.type"));
705        assert!(optional.contains(&"catalog.database"));
706        assert!(optional.contains(&"catalog.name"));
707        assert!(optional.contains(&"catalog.schema"));
708        assert!(optional.contains(&"catalog.workspace_url"));
709        assert!(optional.contains(&"catalog.access_token"));
710        assert!(optional.contains(&"catalog.storage.location"));
711    }
712
713    #[test]
714    fn test_factory_creates_sink() {
715        let registry = ConnectorRegistry::new();
716        register_delta_lake_sink(&registry).unwrap();
717
718        let mut config = crate::config::ConnectorConfig::new("delta-lake");
719        config.set("table.path", "/tmp/laminardb-factory-test");
720        let sink = registry.create_sink(&config, None);
721        assert!(sink.is_ok());
722    }
723
724    // ── Delta Lake source registration tests ──
725
726    #[test]
727    fn test_register_delta_lake_source() {
728        let registry = ConnectorRegistry::new();
729        register_delta_lake_source(&registry).unwrap();
730
731        let info = registry.source_info("delta-lake");
732        assert!(info.is_some());
733        let info = info.unwrap();
734        assert_eq!(info.name, "delta-lake");
735        assert!(info.is_source);
736        assert!(!info.is_sink);
737        assert!(!info.config_keys.is_empty());
738    }
739
740    #[test]
741    fn test_source_config_keys() {
742        let keys = delta_lake_source_config_keys();
743        let required: Vec<&str> = keys
744            .iter()
745            .filter(|k| k.required)
746            .map(|k| k.key.as_str())
747            .collect();
748        assert!(required.contains(&"table.path"));
749        assert_eq!(required.len(), 1);
750
751        let optional: Vec<&str> = keys
752            .iter()
753            .filter(|k| !k.required)
754            .map(|k| k.key.as_str())
755            .collect();
756        assert!(optional.contains(&"starting.version"));
757        assert!(optional.contains(&"poll.interval.ms"));
758        // Catalog keys
759        assert!(optional.contains(&"catalog.type"));
760        assert!(optional.contains(&"catalog.database"));
761    }
762
763    #[test]
764    fn test_factory_creates_source() {
765        let registry = ConnectorRegistry::new();
766        register_delta_lake_source(&registry).unwrap();
767
768        let config = crate::config::ConnectorConfig::new("delta-lake");
769        let source = registry.create_source(&config, None);
770        assert!(source.is_ok());
771    }
772
773    #[test]
774    fn test_register_lakehouse_sinks() {
775        let registry = ConnectorRegistry::new();
776        register_lakehouse_sinks(&registry).unwrap();
777
778        assert!(registry.sink_info("delta-lake").is_some());
779        assert!(registry.sink_info("iceberg").is_some());
780    }
781
782    // ── Iceberg registration tests ──
783
784    #[test]
785    fn test_register_iceberg_sink() {
786        let registry = ConnectorRegistry::new();
787        register_iceberg_sink(&registry).unwrap();
788
789        let info = registry.sink_info("iceberg");
790        assert!(info.is_some());
791        let info = info.unwrap();
792        assert_eq!(info.name, "iceberg");
793        assert!(info.is_sink);
794        assert!(!info.is_source);
795        assert!(!info.config_keys.is_empty());
796    }
797
798    #[test]
799    fn test_register_iceberg_source() {
800        let registry = ConnectorRegistry::new();
801        register_iceberg_source(&registry).unwrap();
802
803        let info = registry.source_info("iceberg");
804        assert!(info.is_some());
805        let info = info.unwrap();
806        assert_eq!(info.name, "iceberg");
807        assert!(info.is_source);
808        assert!(!info.is_sink);
809    }
810
811    #[test]
812    fn test_iceberg_sink_config_keys() {
813        let keys = iceberg_sink_config_keys();
814        let required: Vec<&str> = keys
815            .iter()
816            .filter(|k| k.required)
817            .map(|k| k.key.as_str())
818            .collect();
819        assert!(required.contains(&"catalog.uri"));
820        assert!(required.contains(&"warehouse"));
821        assert!(required.contains(&"namespace"));
822        assert!(required.contains(&"table.name"));
823        assert_eq!(required.len(), 4);
824    }
825
826    #[test]
827    fn test_iceberg_source_config_keys() {
828        let keys = iceberg_source_config_keys();
829        let required: Vec<&str> = keys
830            .iter()
831            .filter(|k| k.required)
832            .map(|k| k.key.as_str())
833            .collect();
834        assert!(required.contains(&"catalog.uri"));
835        assert!(required.contains(&"warehouse"));
836        assert!(required.contains(&"namespace"));
837        assert!(required.contains(&"table.name"));
838        assert_eq!(required.len(), 4);
839    }
840
841    #[test]
842    fn test_factory_creates_iceberg_sink() {
843        let registry = ConnectorRegistry::new();
844        register_iceberg_sink(&registry).unwrap();
845
846        let mut config = crate::config::ConnectorConfig::new("iceberg");
847        config.set("catalog.uri", "http://localhost:8181");
848        config.set("warehouse", "s3://bucket/warehouse");
849        config.set("namespace", "default");
850        config.set("table.name", "events");
851        let sink = registry.create_sink(&config, None);
852        assert!(sink.is_ok());
853    }
854
855    #[test]
856    fn test_factory_creates_iceberg_source() {
857        let registry = ConnectorRegistry::new();
858        register_iceberg_source(&registry).unwrap();
859
860        let config = crate::config::ConnectorConfig::new("iceberg");
861        let source = registry.create_source(&config, None);
862        assert!(source.is_ok());
863    }
864}