laminar_connectors/lakehouse/delta_io/
catalog.rs1#[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#[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 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