1#[cfg(feature = "delta-lake")]
5mod commit_descriptor;
6
7pub 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
24pub 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
37pub mod metrics;
39#[cfg(any(test, feature = "delta-lake", feature = "iceberg"))]
40mod snapshot_schema;
41
42pub 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
54pub 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
70pub 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
99pub 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 #[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 #[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 let source = DeltaLookupSource::open(lookup_config).await?;
195
196 Ok(Arc::new(source) as Arc<dyn laminar_core::lookup::source::LookupSourceDyn>)
197 }
198}
199
200pub 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
229pub 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 #[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 #[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
327pub 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
339pub 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 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 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 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 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 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 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(®istry).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 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(®istry).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 #[test]
727 fn test_register_delta_lake_source() {
728 let registry = ConnectorRegistry::new();
729 register_delta_lake_source(®istry).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 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(®istry).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(®istry).unwrap();
777
778 assert!(registry.sink_info("delta-lake").is_some());
779 assert!(registry.sink_info("iceberg").is_some());
780 }
781
782 #[test]
785 fn test_register_iceberg_sink() {
786 let registry = ConnectorRegistry::new();
787 register_iceberg_sink(®istry).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(®istry).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(®istry).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(®istry).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}