Skip to main content

laminar_connectors/lakehouse/delta_io/
catalog.rs

1//! Catalog-specific table URI and storage-option resolution.
2
3#[cfg(feature = "delta-lake-glue")]
4use super::info;
5use super::{ConnectorError, HashMap};
6#[cfg(feature = "delta-lake-glue")]
7use laminar_core::storage_location::StorageProvider;
8
9#[cfg(feature = "delta-lake-glue")]
10fn glue_init_error_category(error: &deltalake_catalog_glue::GlueError) -> &'static str {
11    match error {
12        deltalake_catalog_glue::GlueError::MissingMetadata { .. } => "missing-metadata",
13        deltalake_catalog_glue::GlueError::AWSError { source } => aws_glue_error_category(source),
14    }
15}
16
17#[cfg(feature = "delta-lake-glue")]
18fn aws_glue_error_category(error: &aws_sdk_glue::Error) -> &'static str {
19    use aws_sdk_glue::Error;
20
21    match error {
22        Error::AccessDeniedException(_)
23        | Error::KmsKeyNotAccessibleFault(_)
24        | Error::PermissionTypeMismatchException(_) => "authorization",
25        Error::EntityNotFoundException(_)
26        | Error::IntegrationNotFoundFault(_)
27        | Error::ResourceNotFoundException(_)
28        | Error::TargetResourceNotFound(_) => "not-found",
29        Error::ConcurrentRunsExceededException(_)
30        | Error::IntegrationQuotaExceededFault(_)
31        | Error::ResourceNumberLimitExceededException(_)
32        | Error::ThrottlingException(_) => "throttled",
33        Error::OperationTimeoutException(_) => "service-timeout",
34        Error::FederationSourceRetryableException(_)
35        | Error::InternalServerException(_)
36        | Error::InternalServiceException(_)
37        | Error::ResourceNotReadyException(_) => "service-unavailable",
38        Error::InvalidInputException(_) | Error::ValidationException(_) => "invalid-request",
39        Error::AlreadyExistsException(_)
40        | Error::ConcurrentModificationException(_)
41        | Error::ConflictException(_)
42        | Error::VersionMismatchException(_) => "catalog-conflict",
43        _ => aws_glue_pre_service_error_category(error).unwrap_or("unclassified-service-error"),
44    }
45}
46
47#[cfg(feature = "delta-lake-glue")]
48fn aws_glue_pre_service_error_category(error: &aws_sdk_glue::Error) -> Option<&'static str> {
49    type GetTableSdkError =
50        aws_sdk_glue::error::SdkError<aws_sdk_glue::operation::get_table::GetTableError>;
51
52    let mut source = std::error::Error::source(error);
53    while let Some(current) = source {
54        if current
55            .downcast_ref::<aws_sdk_glue::error::BuildError>()
56            .is_some()
57        {
58            return Some("client-input");
59        }
60        if let Some(sdk_error) = current.downcast_ref::<GetTableSdkError>() {
61            return match sdk_error {
62                GetTableSdkError::ConstructionFailure(_) => Some("client-input"),
63                GetTableSdkError::TimeoutError(_) => Some("client-timeout"),
64                GetTableSdkError::DispatchFailure(_) => Some("transport"),
65                GetTableSdkError::ResponseError(_) => Some("invalid-response"),
66                GetTableSdkError::ServiceError(_) => None,
67                _ => Some("unclassified-sdk-error"),
68            };
69        }
70        source = current.source();
71    }
72
73    None
74}
75
76#[cfg(feature = "delta-lake-glue")]
77fn glue_lookup_error_category(error: &deltalake::data_catalog::DataCatalogError) -> &'static str {
78    match error {
79        deltalake::data_catalog::DataCatalogError::Generic { source, .. } => source
80            .downcast_ref::<deltalake_catalog_glue::GlueError>()
81            .map_or("catalog-provider", glue_init_error_category),
82        deltalake::data_catalog::DataCatalogError::InvalidDataCatalog { .. } => "invalid-catalog",
83        deltalake::data_catalog::DataCatalogError::UnknownConfigKey { .. } => {
84            "invalid-catalog-configuration"
85        }
86        deltalake::data_catalog::DataCatalogError::RequestError { .. } => "catalog-request",
87    }
88}
89
90/// Resolves catalog-aware table URI and merges catalog-specific storage options.
91///
92/// - `None`: returns table path and storage options as-is.
93/// - `Glue`: calls AWS Glue API to resolve the table's S3 location.
94/// - `Unity`: injects workspace URL and access token into storage options.
95///
96/// # Errors
97///
98/// Returns `ConnectorError` if catalog resolution fails.
99#[cfg(feature = "delta-lake")]
100#[allow(clippy::implicit_hasher, clippy::unused_async)]
101pub async fn resolve_catalog_options(
102    catalog: &super::super::delta_config::DeltaCatalogType,
103    #[allow(unused_variables)] catalog_database: Option<&str>,
104    #[allow(unused_variables)] catalog_name: Option<&str>,
105    _catalog_schema: Option<&str>,
106    table_path: &str,
107    base_storage_options: &HashMap<String, String>,
108) -> Result<(String, HashMap<String, String>), ConnectorError> {
109    use super::super::delta_config::DeltaCatalogType;
110
111    match catalog {
112        DeltaCatalogType::None => Ok((table_path.to_string(), base_storage_options.clone())),
113        #[cfg(feature = "delta-lake-glue")]
114        DeltaCatalogType::Glue => {
115            use deltalake::DataCatalog;
116            let database = catalog_database.ok_or_else(|| {
117                ConnectorError::ConfigurationError(
118                    "Glue catalog requires 'catalog.database'".into(),
119                )
120            })?;
121            let glue = deltalake_catalog_glue::GlueDataCatalog::from_env()
122                .await
123                .map_err(|error| {
124                    ConnectorError::ConnectionFailed(
125                        format!(
126                            "failed to initialize Glue catalog using the downstream AWS credential chain ({}); verify AWS region and credential-chain setup",
127                            glue_init_error_category(&error)
128                        ),
129                    )
130                })?;
131            let resolved = glue
132                .get_table_storage_location(catalog_name.map(String::from), database, table_path)
133                .await
134                .map_err(|error| {
135                    ConnectorError::ConnectionFailed(
136                        format!(
137                            "Glue catalog lookup failed ({}); verify catalog identifiers, AWS region, credentials, and glue:GetTable permission",
138                            glue_lookup_error_category(&error)
139                        ),
140                    )
141                })?;
142            info!(
143                glue_database = database,
144                storage_provider =
145                    StorageProvider::detect_uri(&resolved).map_or("unknown", StorageProvider::name),
146                "resolved table path via Glue catalog"
147            );
148            Ok((resolved, base_storage_options.clone()))
149        }
150        #[cfg(not(feature = "delta-lake-glue"))]
151        DeltaCatalogType::Glue => Err(ConnectorError::ConfigurationError(
152            "Glue catalog requires the 'delta-lake-glue' feature. \
153             Build with: cargo build --features delta-lake-glue"
154                .into(),
155        )),
156        #[cfg(feature = "delta-lake-unity")]
157        DeltaCatalogType::Unity {
158            workspace_url,
159            access_token,
160        } => {
161            // Resolve the table's actual storage location from Unity Catalog
162            // via REST API, then return that direct path (s3://, az://, gs://)
163            // instead of the uc:// URI. This bypasses delta-rs's built-in
164            // uc:// handling which requires credential vending — a feature
165            // that is denied outside Databricks compute environments.
166            let full_name = table_path.strip_prefix("uc://").unwrap_or(table_path);
167
168            let storage_location = super::super::unity_catalog::get_table_storage_location(
169                workspace_url,
170                access_token,
171                full_name,
172            )
173            .await?;
174
175            Ok((storage_location, base_storage_options.clone()))
176        }
177        #[cfg(not(feature = "delta-lake-unity"))]
178        DeltaCatalogType::Unity { .. } => Err(ConnectorError::ConfigurationError(
179            "Unity catalog requires the 'delta-lake-unity' feature. \
180             Build with: cargo build --features delta-lake-unity"
181                .into(),
182        )),
183    }
184}
185
186#[cfg(all(test, feature = "delta-lake-glue"))]
187mod tests {
188    use super::*;
189
190    #[test]
191    fn glue_categories_are_actionable_without_provider_details() {
192        let error: deltalake::data_catalog::DataCatalogError =
193            deltalake_catalog_glue::GlueError::MissingMetadata {
194                metadata: "secret-table-location".into(),
195            }
196            .into();
197        assert_eq!(glue_lookup_error_category(&error), "missing-metadata");
198        assert!(!glue_lookup_error_category(&error).contains("secret"));
199
200        let error = deltalake_catalog_glue::GlueError::AWSError {
201            source: aws_sdk_glue::Error::AccessDeniedException(
202                aws_sdk_glue::types::error::AccessDeniedException::builder()
203                    .message("private provider detail")
204                    .build(),
205            ),
206        };
207        assert_eq!(glue_init_error_category(&error), "authorization");
208        assert!(!glue_init_error_category(&error).contains("private"));
209
210        let unhandled = aws_sdk_glue::types::TableInput::builder()
211            .build()
212            .expect_err("missing table name must produce a client-side build error");
213        let error = deltalake_catalog_glue::GlueError::AWSError {
214            source: unhandled.into(),
215        };
216        assert_eq!(glue_init_error_category(&error), "client-input");
217
218        let timeout: aws_sdk_glue::error::SdkError<
219            aws_sdk_glue::operation::get_table::GetTableError,
220        > = aws_sdk_glue::error::SdkError::timeout_error(std::io::Error::other(
221            "private transport detail",
222        ));
223        let error = deltalake_catalog_glue::GlueError::AWSError {
224            source: timeout.into(),
225        };
226        assert_eq!(glue_init_error_category(&error), "client-timeout");
227
228        let no_code_service_error = aws_sdk_glue::operation::get_table::GetTableError::unhandled(
229            std::io::Error::other("private provider response"),
230        );
231        let error = deltalake_catalog_glue::GlueError::AWSError {
232            source: no_code_service_error.into(),
233        };
234        assert_eq!(
235            glue_init_error_category(&error),
236            "unclassified-service-error"
237        );
238    }
239}
240
241// ============================================================================
242// Integration tests (require delta-lake feature)
243// ============================================================================