1use std::fs::{self, OpenOptions};
2use std::io::{Read, Write};
3#[cfg(unix)]
4use std::os::unix::fs::OpenOptionsExt;
5use std::path::PathBuf;
6
7use ahri_tre_secrets::{MAX_SECRET_MATERIAL_BYTES, SecretMaterial};
8use ahri_tre_security::{REDACTED_VALUE, legacy_object_storage_environment_variable};
9use ahri_tre_types::{DataStoreProfile, EncryptionMode};
10use duckdb::OptionalExt;
11
12use crate::LakeCatalogObservations;
13use crate::alias::LAKE_ALIAS;
14use crate::error::LakeError;
15
16#[derive(Debug, Clone)]
18pub struct DuckLakeAdapter {
19 pub lake_root: String,
20}
21
22#[derive(Debug, Clone, PartialEq, Eq)]
24pub struct FilesystemLakeSessionConfig {
25 filesystem_base: String,
26 datastore_prefix: String,
27 catalog_host: String,
28 catalog_port: u16,
29 catalog_database: String,
30 catalog_schema: String,
31 catalog_role: String,
32 tls: FilesystemCatalogTls,
33 encryption_mode: EncryptionMode,
34}
35
36#[derive(Debug, Clone, PartialEq, Eq)]
38pub enum FilesystemCatalogTls {
39 System,
40 CustomCa { certificates: Vec<String> },
41}
42
43#[derive(Debug, Clone, PartialEq, Eq)]
45pub enum ObjectStorageLocation {
46 AwsS3 {
47 region: String,
48 bucket: String,
49 prefix: String,
50 },
51 S3Compatible {
52 endpoint: String,
53 url_style: S3UrlStyle,
54 bucket: String,
55 prefix: String,
56 },
57 AzureBlob {
58 account_endpoint: String,
59 container: String,
60 prefix: String,
61 },
62}
63
64impl ObjectStorageLocation {
65 pub fn aws_s3_for_datastore(
66 region: impl Into<String>,
67 bucket: impl Into<String>,
68 base_prefix: Option<&str>,
69 datastore_prefix: &str,
70 ) -> Result<Self, LakeError> {
71 Ok(Self::AwsS3 {
72 region: region.into(),
73 bucket: bucket.into(),
74 prefix: datastore_object_prefix(base_prefix, datastore_prefix)?,
75 })
76 }
77
78 pub fn s3_compatible_for_datastore(
79 endpoint: impl Into<String>,
80 url_style: S3UrlStyle,
81 bucket: impl Into<String>,
82 base_prefix: Option<&str>,
83 datastore_prefix: &str,
84 ) -> Result<Self, LakeError> {
85 Ok(Self::S3Compatible {
86 endpoint: endpoint.into(),
87 url_style,
88 bucket: bucket.into(),
89 prefix: datastore_object_prefix(base_prefix, datastore_prefix)?,
90 })
91 }
92
93 pub fn azure_blob_for_datastore(
94 account_endpoint: impl Into<String>,
95 container: impl Into<String>,
96 base_prefix: Option<&str>,
97 datastore_prefix: &str,
98 ) -> Result<Self, LakeError> {
99 Ok(Self::AzureBlob {
100 account_endpoint: account_endpoint.into(),
101 container: container.into(),
102 prefix: datastore_object_prefix(base_prefix, datastore_prefix)?,
103 })
104 }
105
106 pub fn canonical_data_path(&self) -> String {
107 match self {
108 Self::AwsS3 { bucket, prefix, .. } | Self::S3Compatible { bucket, prefix, .. } => {
109 format!("s3://{bucket}/{prefix}/")
110 }
111 Self::AzureBlob {
112 account_endpoint,
113 container,
114 prefix,
115 } => {
116 let authority = account_endpoint
117 .strip_prefix("https://")
118 .unwrap_or(account_endpoint);
119 format!("az://{authority}/{container}/{prefix}/")
120 }
121 }
122 }
123
124 pub fn service_authority(&self) -> &str {
125 match self {
126 Self::AwsS3 { region, .. } => region,
127 Self::S3Compatible { endpoint, .. } => endpoint,
128 Self::AzureBlob {
129 account_endpoint, ..
130 } => account_endpoint,
131 }
132 }
133
134 fn validate(&self) -> Result<(), LakeError> {
135 let (namespace, prefix) = match self {
136 Self::AwsS3 {
137 region,
138 bucket,
139 prefix,
140 } => {
141 if region.trim().is_empty() {
142 return Err(LakeError::InvalidObjectStorageIdentity);
143 }
144 (bucket, prefix)
145 }
146 Self::S3Compatible {
147 endpoint,
148 bucket,
149 prefix,
150 ..
151 } => {
152 if !endpoint.starts_with("https://") || endpoint.ends_with('/') {
153 return Err(LakeError::InvalidObjectStorageIdentity);
154 }
155 (bucket, prefix)
156 }
157 Self::AzureBlob {
158 account_endpoint,
159 container,
160 prefix,
161 } => {
162 if !account_endpoint.starts_with("https://") || account_endpoint.ends_with('/') {
163 return Err(LakeError::InvalidObjectStorageIdentity);
164 }
165 (container, prefix)
166 }
167 };
168 if namespace.trim().is_empty() || !canonical_datastore_prefix(prefix) {
169 return Err(LakeError::InvalidObjectStorageIdentity);
170 }
171 Ok(())
172 }
173}
174
175fn datastore_object_prefix(
176 base_prefix: Option<&str>,
177 datastore_prefix: &str,
178) -> Result<String, LakeError> {
179 if !canonical_datastore_prefix(datastore_prefix) {
180 return Err(LakeError::InvalidDatastorePrefix);
181 }
182 let base = base_prefix.map(str::trim).unwrap_or("").trim_matches('/');
183 if base.is_empty() {
184 Ok(datastore_prefix.to_string())
185 } else {
186 Ok(format!("{base}/{datastore_prefix}"))
187 }
188}
189
190#[derive(Debug, Clone, Copy, PartialEq, Eq)]
191pub enum S3UrlStyle {
192 Path,
193 VirtualHosted,
194}
195
196#[derive(Debug, Clone, Copy, PartialEq, Eq)]
198pub enum AwsWorkloadIdentitySource {
199 EcsTask,
200 Ec2Instance,
201}
202
203#[derive(Debug, Clone, PartialEq, Eq)]
205pub enum ObjectStorageTls {
206 System,
207 CustomCa { certificates: Vec<String> },
208}
209
210#[derive(Debug, Clone, PartialEq, Eq)]
212pub struct ObjectLakeSessionConfig {
213 location: ObjectStorageLocation,
214 catalog_host: String,
215 catalog_port: u16,
216 catalog_database: String,
217 catalog_schema: String,
218 catalog_role: String,
219 catalog_tls: FilesystemCatalogTls,
220 storage_tls: ObjectStorageTls,
221 encryption_mode: EncryptionMode,
222}
223
224impl ObjectLakeSessionConfig {
225 #[allow(clippy::too_many_arguments)]
226 pub fn new(
227 location: ObjectStorageLocation,
228 catalog_host: impl Into<String>,
229 catalog_port: u16,
230 catalog_database: impl Into<String>,
231 catalog_schema: impl Into<String>,
232 catalog_role: impl Into<String>,
233 catalog_tls: FilesystemCatalogTls,
234 storage_tls: ObjectStorageTls,
235 encryption_mode: EncryptionMode,
236 ) -> Result<Self, LakeError> {
237 location.validate()?;
238 let config = Self {
239 location,
240 catalog_host: catalog_host.into(),
241 catalog_port,
242 catalog_database: catalog_database.into(),
243 catalog_schema: catalog_schema.into(),
244 catalog_role: catalog_role.into(),
245 catalog_tls,
246 storage_tls,
247 encryption_mode,
248 };
249 if config.catalog_host.trim().is_empty()
250 || config.catalog_port == 0
251 || config.catalog_database.trim().is_empty()
252 || config.catalog_schema.trim().is_empty()
253 || config.catalog_role.trim().is_empty()
254 {
255 return Err(LakeError::IncompleteCatalogIdentity);
256 }
257 Ok(config)
258 }
259
260 pub fn location(&self) -> &ObjectStorageLocation {
261 &self.location
262 }
263
264 fn legacy_profile(&self) -> DataStoreProfile {
265 DataStoreProfile {
266 server: self.catalog_host.clone(),
267 port: self.catalog_port,
268 dbname: self.catalog_database.clone(),
269 sslmode: Some("verify-full".to_string()),
270 lake: ahri_tre_types::LakeAttachProfile {
271 lake_data: self.location.canonical_data_path(),
272 lake_db: self.catalog_database.clone(),
273 catalog_schema: Some(self.catalog_schema.clone()),
274 auth: ahri_tre_types::LakeAuthProfile {
275 lake_user: self.catalog_role.clone(),
276 encryption_mode: self.encryption_mode,
277 },
278 },
279 }
280 }
281}
282
283#[derive(Clone, Copy)]
287pub enum ObjectStorageAuthentication<'a> {
288 S3Static {
289 access_key_id: &'a SecretMaterial,
290 secret_access_key: &'a SecretMaterial,
291 },
292 AwsWorkloadIdentity {
293 source: AwsWorkloadIdentitySource,
294 },
295 AzureManagedIdentity {
296 client_id: Option<uuid::Uuid>,
297 },
298 AzureServicePrincipal {
299 tenant_id: uuid::Uuid,
300 client_id: uuid::Uuid,
301 client_secret: &'a SecretMaterial,
302 },
303}
304
305impl std::fmt::Debug for ObjectStorageAuthentication<'_> {
306 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
307 match self {
308 Self::S3Static { .. } => formatter.write_str("S3Static(<redacted>)"),
309 Self::AwsWorkloadIdentity { source } => formatter
310 .debug_struct("AwsWorkloadIdentity")
311 .field("source", source)
312 .finish(),
313 Self::AzureManagedIdentity { client_id } => formatter
314 .debug_struct("AzureManagedIdentity")
315 .field("client_id", client_id)
316 .finish(),
317 Self::AzureServicePrincipal {
318 tenant_id,
319 client_id,
320 ..
321 } => formatter
322 .debug_struct("AzureServicePrincipal")
323 .field("tenant_id", tenant_id)
324 .field("client_id", client_id)
325 .field("client_secret", &REDACTED_VALUE)
326 .finish(),
327 }
328 }
329}
330
331impl FilesystemLakeSessionConfig {
332 #[allow(clippy::too_many_arguments)]
333 pub fn new(
334 filesystem_base: impl Into<String>,
335 datastore_prefix: impl Into<String>,
336 catalog_host: impl Into<String>,
337 catalog_port: u16,
338 catalog_database: impl Into<String>,
339 catalog_schema: impl Into<String>,
340 catalog_role: impl Into<String>,
341 tls: FilesystemCatalogTls,
342 encryption_mode: EncryptionMode,
343 ) -> Result<Self, LakeError> {
344 let config = Self {
345 filesystem_base: filesystem_base.into(),
346 datastore_prefix: datastore_prefix.into(),
347 catalog_host: catalog_host.into(),
348 catalog_port,
349 catalog_database: catalog_database.into(),
350 catalog_schema: catalog_schema.into(),
351 catalog_role: catalog_role.into(),
352 tls,
353 encryption_mode,
354 };
355 normalized_root(&config.filesystem_base)?;
356 if !canonical_datastore_prefix(&config.datastore_prefix) {
357 return Err(LakeError::InvalidDatastorePrefix);
358 }
359 if config.catalog_host.trim().is_empty()
360 || config.catalog_port == 0
361 || config.catalog_database.trim().is_empty()
362 || config.catalog_schema.trim().is_empty()
363 || config.catalog_role.trim().is_empty()
364 {
365 return Err(LakeError::IncompleteCatalogIdentity);
366 }
367 Ok(config)
368 }
369
370 #[allow(clippy::too_many_arguments)]
372 pub fn from_binding(
373 filesystem_base: impl Into<String>,
374 canonical_data_path: impl Into<String>,
375 catalog_host: impl Into<String>,
376 catalog_port: u16,
377 catalog_database: impl Into<String>,
378 catalog_schema: impl Into<String>,
379 catalog_role: impl Into<String>,
380 tls: FilesystemCatalogTls,
381 encryption_mode: EncryptionMode,
382 ) -> Result<Self, LakeError> {
383 let filesystem_base = filesystem_base.into();
384 let canonical_data_path = canonical_data_path.into();
385 let base = normalized_root(&filesystem_base)?;
386 let data = normalized_root(&canonical_data_path)?;
387 let prefix = if base == "/" {
388 data.strip_prefix('/')
389 } else {
390 data.strip_prefix(base)
391 .and_then(|suffix| suffix.strip_prefix('/'))
392 }
393 .filter(|prefix| canonical_datastore_prefix(prefix))
394 .ok_or(LakeError::InvalidDatastorePrefix)?;
395 Self::new(
396 filesystem_base,
397 prefix,
398 catalog_host,
399 catalog_port,
400 catalog_database,
401 catalog_schema,
402 catalog_role,
403 tls,
404 encryption_mode,
405 )
406 }
407
408 pub fn filesystem_base(&self) -> &str {
409 &self.filesystem_base
410 }
411
412 pub fn datastore_prefix(&self) -> &str {
413 &self.datastore_prefix
414 }
415
416 pub fn canonical_data_path(&self) -> String {
417 if self.filesystem_base == "/" {
418 format!("/{}", self.datastore_prefix)
419 } else {
420 format!(
421 "{}/{}",
422 self.filesystem_base.trim_end_matches('/'),
423 self.datastore_prefix
424 )
425 }
426 }
427
428 fn legacy_profile(&self) -> DataStoreProfile {
429 DataStoreProfile {
430 server: self.catalog_host.clone(),
431 port: self.catalog_port,
432 dbname: self.catalog_database.clone(),
433 sslmode: Some("verify-full".to_string()),
434 lake: ahri_tre_types::LakeAttachProfile {
435 lake_data: self.canonical_data_path(),
436 lake_db: self.catalog_database.clone(),
437 catalog_schema: Some(self.catalog_schema.clone()),
438 auth: ahri_tre_types::LakeAuthProfile {
439 lake_user: self.catalog_role.clone(),
440 encryption_mode: self.encryption_mode,
441 },
442 },
443 }
444 }
445}
446
447#[derive(Debug, Clone, PartialEq, Eq)]
449pub struct DuckLakeLocalLayout {
450 pub data_path: String,
451 pub staging_path: String,
452 pub tmp_path: String,
453 pub test_runs_path: String,
454}
455
456#[derive(Clone, PartialEq, Eq)]
458pub struct DuckLakeAttachPlan {
459 pub(crate) metadata_schema: String,
460 pub alias: String,
461 pub attach_description: String,
462 pub(crate) pre_attach_sql: Vec<String>,
463 pub(crate) attach_sql: String,
464 pub data_path: String,
465 pub credential_kind: DuckLakeCredentialKind,
466 pub create_if_not_exists: bool,
467 pub automatic_migration: bool,
468 pub configured_encryption_mode: EncryptionMode,
469}
470
471impl std::fmt::Debug for DuckLakeAttachPlan {
472 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
473 f.debug_struct("DuckLakeAttachPlan")
474 .field("alias", &self.alias)
475 .field("attach_description", &self.attach_description)
476 .field("pre_attach_sql", &REDACTED_VALUE)
477 .field("attach_sql", &REDACTED_VALUE)
478 .field("data_path", &self.data_path)
479 .field("credential_kind", &self.credential_kind)
480 .field("create_if_not_exists", &self.create_if_not_exists)
481 .field("automatic_migration", &self.automatic_migration)
482 .field(
483 "configured_encryption_mode",
484 &self.configured_encryption_mode,
485 )
486 .finish()
487 }
488}
489
490pub struct DuckLakeOpenedCatalog {
492 pub connection: duckdb::Connection,
493 pub attach_plan: DuckLakeAttachPlan,
494 pub local_layout: DuckLakeLocalLayout,
495 pub health: DuckLakeHealth,
496 pub tls_guard: Option<DuckLakeTlsGuard>,
497}
498
499impl std::fmt::Debug for DuckLakeOpenedCatalog {
500 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
501 formatter
502 .debug_struct("DuckLakeOpenedCatalog")
503 .field("attach_plan", &self.attach_plan)
504 .field("local_layout", &self.local_layout)
505 .field("health", &self.health)
506 .finish_non_exhaustive()
507 }
508}
509
510pub struct DuckLakeTlsGuard {
511 paths: Vec<PathBuf>,
512}
513
514impl std::fmt::Debug for DuckLakeTlsGuard {
515 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
516 formatter.write_str("DuckLakeTlsGuard(<adapter-owned>)")
517 }
518}
519
520impl Drop for DuckLakeTlsGuard {
521 fn drop(&mut self) {
522 for path in &self.paths {
523 let _ = fs::remove_file(path);
524 }
525 }
526}
527
528impl DuckLakeTlsGuard {
529 fn create(certificates: &[String]) -> Result<Self, LakeError> {
530 if certificates.is_empty() {
531 return Err(LakeError::IncompleteCatalogIdentity);
532 }
533 let path = std::env::temp_dir().join(format!(
534 "ahri-tre-ducklake-ca-{}.pem",
535 uuid::Uuid::new_v4().simple()
536 ));
537 let mut options = OpenOptions::new();
538 options.write(true).create_new(true);
539 #[cfg(unix)]
540 options.mode(0o600);
541 let mut file = options.open(&path)?;
542 let projected = Self { paths: vec![path] };
543 for certificate in certificates {
544 file.write_all(certificate.as_bytes())?;
545 if !certificate.ends_with('\n') {
546 file.write_all(b"\n")?;
547 }
548 }
549 file.sync_all()?;
550 Ok(projected)
551 }
552
553 fn path(&self) -> &std::path::Path {
554 &self.paths[0]
555 }
556
557 fn combine(mut self, mut other: Self) -> Self {
558 self.paths.append(&mut other.paths);
559 self
560 }
561}
562
563#[derive(Debug, Clone, PartialEq, Eq)]
565pub struct DuckLakeHealth {
566 pub alias: String,
567 pub catalog_type: String,
568 pub extension_version: String,
569 pub data_path: String,
570 pub detected_encryption_mode: EncryptionMode,
571 pub snapshot_count: i64,
572 pub table_count: i64,
573}
574
575#[derive(Clone, PartialEq, Eq)]
577enum DuckLakeCatalogCredential {
578 SessionSecret {
579 secret_name: String,
580 password: String,
581 },
582}
583
584impl std::fmt::Debug for DuckLakeCatalogCredential {
585 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
586 match self {
587 Self::SessionSecret { secret_name, .. } => f
588 .debug_struct("SessionSecret")
589 .field("secret_name", secret_name)
590 .field("password", &REDACTED_VALUE)
591 .finish(),
592 }
593 }
594}
595
596#[derive(Debug, Clone, Copy, PartialEq, Eq)]
598pub enum DuckLakeCredentialKind {
599 SessionSecret,
600}
601
602#[derive(Debug, Clone, Copy, PartialEq, Eq)]
604pub struct DuckLakeOpenOptions {
605 pub create_if_not_exists: bool,
606 pub automatic_migration: bool,
607 pub apply_configured_encryption: bool,
608}
609
610impl Default for DuckLakeOpenOptions {
611 fn default() -> Self {
612 Self {
613 create_if_not_exists: true,
614 automatic_migration: false,
615 apply_configured_encryption: true,
616 }
617 }
618}
619
620impl DuckLakeOpenOptions {
621 pub const fn initialize() -> Self {
623 Self {
624 create_if_not_exists: true,
625 automatic_migration: false,
626 apply_configured_encryption: true,
627 }
628 }
629
630 pub const fn attach_existing() -> Self {
632 Self {
633 create_if_not_exists: false,
634 automatic_migration: false,
635 apply_configured_encryption: false,
636 }
637 }
638
639 pub const fn with_automatic_migration(mut self, automatic_migration: bool) -> Self {
641 self.automatic_migration = automatic_migration;
642 self
643 }
644}
645
646impl DuckLakeCatalogCredential {
647 fn session_secret(secret_name: impl Into<String>, password: impl Into<String>) -> Self {
649 Self::SessionSecret {
650 secret_name: secret_name.into(),
651 password: password.into(),
652 }
653 }
654
655 fn kind(&self) -> DuckLakeCredentialKind {
656 match self {
657 Self::SessionSecret { .. } => DuckLakeCredentialKind::SessionSecret,
658 }
659 }
660
661 fn validate(&self) -> Result<(), LakeError> {
662 let (kind, value) = match self {
663 Self::SessionSecret {
664 secret_name,
665 password,
666 } => {
667 if secret_name.trim().is_empty() {
668 return Err(LakeError::EmptyCredential {
669 kind: "session_secret",
670 });
671 }
672 ("session_secret_password", password)
673 }
674 };
675
676 if value.trim().is_empty() {
677 return Err(LakeError::EmptyCredential { kind });
678 }
679
680 Ok(())
681 }
682}
683
684impl DuckLakeAdapter {
685 pub fn validate_object_storage_inputs<'a>(
688 config: &ObjectLakeSessionConfig,
689 storage_authentication: ObjectStorageAuthentication<'a>,
690 ambient_environment_names: impl IntoIterator<Item = impl AsRef<str>>,
691 ) -> Result<(), LakeError> {
692 let ambient_environment_names = ambient_environment_names
693 .into_iter()
694 .map(|name| name.as_ref().to_string())
695 .collect::<Vec<_>>();
696 reject_ambient_storage_authority(
697 ambient_environment_names.iter().map(String::as_str),
698 &storage_authentication,
699 )?;
700 validate_storage_authentication(config.location(), &storage_authentication)
701 }
702
703 pub fn open_existing_object_storage_catalog(
706 config: &ObjectLakeSessionConfig,
707 catalog_credential: &SecretMaterial,
708 storage_authentication: ObjectStorageAuthentication<'_>,
709 scratch: &crate::ScratchAttempt,
710 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
711 Self::open_object_storage_catalog(
712 config,
713 catalog_credential,
714 storage_authentication,
715 scratch,
716 DuckLakeOpenOptions::attach_existing(),
717 &mut LakeCatalogObservations::default(),
718 )
719 }
720
721 pub fn open_existing_object_storage_catalog_observed(
723 context: Option<&ahri_tre_observability::CorrelationContext>,
724 config: &ObjectLakeSessionConfig,
725 catalog_credential: &SecretMaterial,
726 storage_authentication: ObjectStorageAuthentication<'_>,
727 scratch: &crate::ScratchAttempt,
728 observations: &mut LakeCatalogObservations,
729 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
730 let span = context.map(|context| context.span(ahri_tre_observability::Stage::Lake));
731 let result = Self::open_object_storage_catalog(
732 config,
733 catalog_credential,
734 storage_authentication,
735 scratch,
736 DuckLakeOpenOptions::attach_existing(),
737 observations,
738 );
739 finish_lake_open(span, &result);
740 result
741 }
742
743 pub fn initialize_object_storage_catalog(
745 config: &ObjectLakeSessionConfig,
746 catalog_credential: &SecretMaterial,
747 storage_authentication: ObjectStorageAuthentication<'_>,
748 scratch: &crate::ScratchAttempt,
749 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
750 Self::open_object_storage_catalog(
751 config,
752 catalog_credential,
753 storage_authentication,
754 scratch,
755 DuckLakeOpenOptions::initialize(),
756 &mut LakeCatalogObservations::default(),
757 )
758 }
759
760 fn open_object_storage_catalog(
761 config: &ObjectLakeSessionConfig,
762 catalog_credential: &SecretMaterial,
763 storage_authentication: ObjectStorageAuthentication<'_>,
764 scratch: &crate::ScratchAttempt,
765 options: DuckLakeOpenOptions,
766 observations: &mut LakeCatalogObservations,
767 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
768 reject_ambient_storage_authority(
769 std::env::vars_os().map(|(name, _)| name.to_string_lossy().into_owned()),
770 &storage_authentication,
771 )?;
772 validate_storage_authentication(config.location(), &storage_authentication)?;
773
774 let catalog_guard = match &config.catalog_tls {
775 FilesystemCatalogTls::System => None,
776 FilesystemCatalogTls::CustomCa { certificates } => {
777 Some(DuckLakeTlsGuard::create(certificates)?)
778 }
779 };
780 let storage_guard = match &config.storage_tls {
781 ObjectStorageTls::System => None,
782 ObjectStorageTls::CustomCa { certificates } => {
783 Some(DuckLakeTlsGuard::create(certificates)?)
784 }
785 };
786 let catalog_root = catalog_guard
787 .as_ref()
788 .map(|guard| guard.path().to_string_lossy().into_owned())
789 .unwrap_or_else(|| "system".to_string());
790
791 let profile = config.legacy_profile();
792 let adapter = Self::new(config.location.canonical_data_path());
793 let mut plan = catalog_credential.expose(|bytes| {
794 let password =
795 std::str::from_utf8(bytes).map_err(|_| LakeError::InvalidCredentialEncoding)?;
796 let mut plan = adapter.build_attach_plan_with_credential_and_root(
797 &profile,
798 DuckLakeCatalogCredential::session_secret("ahri_tre_ducklake_catalog", password),
799 options,
800 Some(&catalog_root),
801 )?;
802 plan.data_path = config.location().canonical_data_path();
803 plan.pre_attach_sql.insert(
804 0,
805 storage_secret_sql(config.location(), storage_authentication)?,
806 );
807 if let Some(guard) = storage_guard.as_ref() {
808 plan.pre_attach_sql.insert(
809 0,
810 format!(
811 "SET ca_cert_file = {}; SET enable_server_cert_verification = true",
812 sql_literal(&guard.path().to_string_lossy())
813 ),
814 );
815 }
816 Ok::<_, LakeError>(plan)
817 })?;
818
819 let connection = duckdb::Connection::open_in_memory()?;
820 connection.execute_batch(&format!(
821 "SET temp_directory = {}",
822 sql_literal(&scratch.path().to_string_lossy())
823 ))?;
824 connection.execute_batch(
825 "INSTALL postgres; LOAD postgres; INSTALL ducklake; LOAD ducklake; INSTALL httpfs; LOAD httpfs;",
826 )?;
827 if matches!(config.location(), ObjectStorageLocation::AzureBlob { .. }) {
828 connection.execute_batch("INSTALL azure; LOAD azure;")?;
829 }
830 for statement in &plan.pre_attach_sql {
831 connection
832 .execute_batch(statement)
833 .map_err(|_| LakeError::CredentialProjectionFailed)?;
834 }
835 if options.create_if_not_exists {
836 let pattern = object_namespace_glob(config.location());
837 let count: i64 = connection.query_row(
838 &format!("SELECT count(*) FROM glob({})", sql_literal(&pattern)),
839 [],
840 |row| row.get(0),
841 )?;
842 if count != 0 {
843 return Err(LakeError::NamespaceNotEmpty);
844 }
845 }
846 observations.attach(&connection, &plan.attach_sql, &config.catalog_role)?;
847 #[cfg(feature = "acceptance-catalog-barrier")]
848 crate::catalog_observations::after_attach_barrier()?;
849 connection.execute_batch(&format!("USE {LAKE_ALIAS};"))?;
850 let health = observations.health_check(&adapter, &connection)?;
851 if health.data_path != config.location().canonical_data_path() {
852 return Err(LakeError::PersistedDataPathMismatch);
853 }
854 let tls_guard = match (catalog_guard, storage_guard) {
855 (Some(first), Some(second)) => Some(first.combine(second)),
856 (Some(guard), None) | (None, Some(guard)) => Some(guard),
857 (None, None) => None,
858 };
859 plan.data_path = config.location().canonical_data_path();
860 Ok(DuckLakeOpenedCatalog {
861 connection,
862 attach_plan: plan,
863 local_layout: DuckLakeLocalLayout {
864 data_path: config.location().canonical_data_path(),
865 staging_path: String::new(),
866 tmp_path: String::new(),
867 test_runs_path: String::new(),
868 },
869 health,
870 tls_guard,
871 })
872 }
873
874 pub fn build_filesystem_attach_plan(
877 config: &FilesystemLakeSessionConfig,
878 credential: &SecretMaterial,
879 ) -> Result<DuckLakeAttachPlan, LakeError> {
880 if matches!(config.tls, FilesystemCatalogTls::CustomCa { .. }) {
881 return Err(LakeError::CustomTlsRequiresOpen);
882 }
883 let profile = config.legacy_profile();
884 let adapter = Self::new(config.canonical_data_path());
885 credential.expose(|bytes| {
886 let password =
887 std::str::from_utf8(bytes).map_err(|_| LakeError::InvalidCredentialEncoding)?;
888 adapter.build_attach_plan_with_credential_and_root(
889 &profile,
890 DuckLakeCatalogCredential::session_secret("ahri_tre_ducklake_catalog", password),
891 DuckLakeOpenOptions::attach_existing(),
892 Some("system"),
893 )
894 })
895 }
896
897 pub fn open_existing_filesystem_catalog(
899 config: &FilesystemLakeSessionConfig,
900 credential: &SecretMaterial,
901 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
902 Self::open_existing_filesystem_catalog_inner(
903 config,
904 credential,
905 None,
906 &mut LakeCatalogObservations::default(),
907 )
908 }
909
910 fn open_existing_filesystem_catalog_inner(
911 config: &FilesystemLakeSessionConfig,
912 credential: &SecretMaterial,
913 scratch: Option<&std::path::Path>,
914 observations: &mut LakeCatalogObservations,
915 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
916 let profile = config.legacy_profile();
917 let adapter = Self::new(config.canonical_data_path());
918 let root_certificate = match &config.tls {
919 FilesystemCatalogTls::System => None,
920 FilesystemCatalogTls::CustomCa { certificates } => {
921 Some(DuckLakeTlsGuard::create(certificates)?)
922 }
923 };
924 let root = match &root_certificate {
925 Some(certificate) => certificate.path().to_string_lossy().into_owned(),
926 None => "system".to_string(),
927 };
928 credential.expose(|bytes| {
929 let password =
930 std::str::from_utf8(bytes).map_err(|_| LakeError::InvalidCredentialEncoding)?;
931 let mut opened = adapter.open_attached_catalog_with_credential_and_root(
932 &profile,
933 DuckLakeCatalogCredential::session_secret("ahri_tre_ducklake_catalog", password),
934 DuckLakeOpenOptions::attach_existing(),
935 Some(&root),
936 scratch,
937 observations,
938 )?;
939 opened.tls_guard = root_certificate;
940 Ok(opened)
941 })
942 }
943
944 pub fn open_existing_filesystem_catalog_with_scratch(
947 config: &FilesystemLakeSessionConfig,
948 credential: &SecretMaterial,
949 scratch: &crate::ScratchAttempt,
950 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
951 let mut opened = Self::open_existing_filesystem_catalog_inner(
952 config,
953 credential,
954 Some(scratch.path()),
955 &mut LakeCatalogObservations::default(),
956 )?;
957 opened.local_layout.staging_path.clear();
958 opened.local_layout.tmp_path.clear();
959 opened.local_layout.test_runs_path.clear();
960 Ok(opened)
961 }
962
963 pub fn open_existing_filesystem_catalog_observed(
965 context: Option<&ahri_tre_observability::CorrelationContext>,
966 config: &FilesystemLakeSessionConfig,
967 credential: &SecretMaterial,
968 scratch: &crate::ScratchAttempt,
969 observations: &mut LakeCatalogObservations,
970 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
971 let span = context.map(|context| context.span(ahri_tre_observability::Stage::Lake));
972 let result = Self::open_existing_filesystem_catalog_inner(
973 config,
974 credential,
975 Some(scratch.path()),
976 observations,
977 )
978 .map(|mut opened| {
979 opened.local_layout.staging_path.clear();
980 opened.local_layout.tmp_path.clear();
981 opened.local_layout.test_runs_path.clear();
982 opened
983 });
984 finish_lake_open(span, &result);
985 result
986 }
987
988 pub fn initialize_filesystem_catalog_with_scratch(
990 config: &FilesystemLakeSessionConfig,
991 credential: &SecretMaterial,
992 scratch: &crate::ScratchAttempt,
993 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
994 match fs::symlink_metadata(config.canonical_data_path()) {
995 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
996 Ok(_) => return Err(LakeError::NamespaceNotEmpty),
997 Err(error) => return Err(error.into()),
998 }
999 let profile = config.legacy_profile();
1000 let adapter = Self::new(config.canonical_data_path());
1001 let root_certificate = match &config.tls {
1002 FilesystemCatalogTls::System => None,
1003 FilesystemCatalogTls::CustomCa { certificates } => {
1004 Some(DuckLakeTlsGuard::create(certificates)?)
1005 }
1006 };
1007 let root = root_certificate
1008 .as_ref()
1009 .map(|certificate| certificate.path().to_string_lossy().into_owned())
1010 .unwrap_or_else(|| "system".to_string());
1011 let mut opened = credential.expose(|bytes| {
1012 let password =
1013 std::str::from_utf8(bytes).map_err(|_| LakeError::InvalidCredentialEncoding)?;
1014 adapter.open_attached_catalog_with_credential_and_root(
1015 &profile,
1016 DuckLakeCatalogCredential::session_secret("ahri_tre_ducklake_catalog", password),
1017 DuckLakeOpenOptions::initialize(),
1018 Some(&root),
1019 Some(scratch.path()),
1020 &mut LakeCatalogObservations::default(),
1021 )
1022 })?;
1023 opened.tls_guard = root_certificate;
1024 opened.local_layout.staging_path.clear();
1025 opened.local_layout.tmp_path.clear();
1026 opened.local_layout.test_runs_path.clear();
1027 Ok(opened)
1028 }
1029
1030 pub fn new(lake_root: impl Into<String>) -> Self {
1032 Self {
1033 lake_root: lake_root.into(),
1034 }
1035 }
1036
1037 pub fn local_layout(&self) -> Result<DuckLakeLocalLayout, LakeError> {
1043 let root = normalized_root(&self.lake_root)?;
1044
1045 Ok(DuckLakeLocalLayout {
1046 data_path: root.to_string(),
1047 staging_path: format!("{root}/__tre_duckdb_stage"),
1048 tmp_path: format!("{root}/__tre_duckdb_stage"),
1049 test_runs_path: format!("{root}/test-runs"),
1050 })
1051 }
1052
1053 fn build_attach_plan_with_credential_and_root(
1054 &self,
1055 profile: &DataStoreProfile,
1056 credential: DuckLakeCatalogCredential,
1057 options: DuckLakeOpenOptions,
1058 root_certificate: Option<&str>,
1059 ) -> Result<DuckLakeAttachPlan, LakeError> {
1060 credential.validate()?;
1061
1062 let layout = self.local_layout()?;
1063 let data_path = normalized_root(&layout.data_path)?.to_string();
1064 let pre_attach_sql = ducklake_pre_attach_sql(profile, &credential, false);
1065 let conninfo = ducklake_postgres_conninfo(profile, &credential, root_certificate, false);
1066 let conninfo_description =
1067 ducklake_postgres_conninfo(profile, &credential, root_certificate, true);
1068 let mut attach_options = self.attach_options(profile, &data_path, options)?;
1069 let DuckLakeCatalogCredential::SessionSecret { secret_name, .. } = &credential;
1070 attach_options.push(format!(
1071 "META_SECRET {}",
1072 sql_literal(&session_secret_name(secret_name))
1073 ));
1074
1075 let attach_sql = format!(
1076 "ATTACH {} AS {LAKE_ALIAS} ({})",
1077 sql_literal(&conninfo),
1078 attach_options.join(", ")
1079 );
1080 let attach_description = format!(
1081 "ATTACH {} AS {LAKE_ALIAS} ({})",
1082 sql_literal(&conninfo_description),
1083 attach_options.join(", ")
1084 );
1085
1086 Ok(DuckLakeAttachPlan {
1087 metadata_schema: profile
1088 .lake
1089 .catalog_schema
1090 .clone()
1091 .unwrap_or_else(|| "main".into()),
1092 alias: LAKE_ALIAS.to_string(),
1093 attach_description,
1094 pre_attach_sql,
1095 attach_sql,
1096 data_path,
1097 credential_kind: credential.kind(),
1098 create_if_not_exists: options.create_if_not_exists,
1099 automatic_migration: options.automatic_migration,
1100 configured_encryption_mode: profile.lake.auth.encryption_mode,
1101 })
1102 }
1103
1104 fn open_attached_catalog_with_credential_and_root(
1105 &self,
1106 profile: &DataStoreProfile,
1107 credential: DuckLakeCatalogCredential,
1108 options: DuckLakeOpenOptions,
1109 root_certificate: Option<&str>,
1110 temp_directory: Option<&std::path::Path>,
1111 observations: &mut LakeCatalogObservations,
1112 ) -> Result<DuckLakeOpenedCatalog, LakeError> {
1113 let local_layout = self.prepare_local_layout(options)?;
1114 let (connection, attach_plan) = self.attach_catalog(
1115 profile,
1116 credential,
1117 options,
1118 root_certificate,
1119 temp_directory,
1120 observations,
1121 )?;
1122 #[cfg(feature = "acceptance-catalog-barrier")]
1123 crate::catalog_observations::after_attach_barrier()?;
1124 let health = observations.health_check(self, &connection)?;
1125 if !options.create_if_not_exists
1126 && normalize_ducklake_data_path(&health.data_path)
1127 != normalize_ducklake_data_path(&local_layout.data_path)
1128 {
1129 return Err(LakeError::PersistedDataPathMismatch);
1130 }
1131
1132 Ok(DuckLakeOpenedCatalog {
1133 connection,
1134 attach_plan,
1135 local_layout,
1136 health,
1137 tls_guard: None,
1138 })
1139 }
1140
1141 pub fn health_check_observed(
1143 &self,
1144 connection: &duckdb::Connection,
1145 observations: &mut LakeCatalogObservations,
1146 ) -> Result<DuckLakeHealth, LakeError> {
1147 observations.health_check(self, connection)
1148 }
1149
1150 pub fn health_check(
1157 &self,
1158 connection: &duckdb::Connection,
1159 ) -> Result<DuckLakeHealth, LakeError> {
1160 let settings = ducklake_settings(connection)?;
1161 let detected_encryption_mode = detected_encryption_mode(connection, &settings.data_path)?;
1162 let snapshot_count = connection.query_row(
1163 &format!(
1164 "SELECT count(*) FROM ducklake_snapshots({})",
1165 sql_literal(LAKE_ALIAS)
1166 ),
1167 [],
1168 |row| row.get(0),
1169 )?;
1170 let table_count = connection.query_row(
1171 &format!(
1172 "SELECT count(*) FROM ducklake_table_info({})",
1173 sql_literal(LAKE_ALIAS)
1174 ),
1175 [],
1176 |row| row.get(0),
1177 )?;
1178
1179 Ok(DuckLakeHealth {
1180 alias: LAKE_ALIAS.to_string(),
1181 catalog_type: settings.catalog_type,
1182 extension_version: settings.extension_version,
1183 data_path: settings.data_path,
1184 detected_encryption_mode,
1185 snapshot_count,
1186 table_count,
1187 })
1188 }
1189}
1190
1191fn object_namespace_glob(location: &ObjectStorageLocation) -> String {
1192 format!("{}**", location.canonical_data_path())
1193}
1194
1195fn reject_ambient_storage_authority(
1196 names: impl IntoIterator<Item = impl AsRef<str>>,
1197 authentication: &ObjectStorageAuthentication<'_>,
1198) -> Result<(), LakeError> {
1199 let names = names
1200 .into_iter()
1201 .map(|name| name.as_ref().to_string())
1202 .collect::<Vec<_>>();
1203 if let Some(name) = legacy_object_storage_environment_variable(|candidate| {
1204 names.iter().any(|name| name == candidate)
1205 }) {
1206 return Err(LakeError::AmbientStorageAuthority {
1207 name: name.to_string(),
1208 });
1209 }
1210 match authentication {
1211 ObjectStorageAuthentication::AwsWorkloadIdentity {
1212 source: AwsWorkloadIdentitySource::EcsTask,
1213 } if !names
1214 .iter()
1215 .any(|name| name == "AWS_CONTAINER_CREDENTIALS_RELATIVE_URI") =>
1216 {
1217 Err(LakeError::MissingWorkloadIdentitySource { kind: "ecs_task" })
1218 }
1219 ObjectStorageAuthentication::AwsWorkloadIdentity {
1220 source: AwsWorkloadIdentitySource::Ec2Instance,
1221 } if names
1222 .iter()
1223 .any(|name| name == "AWS_CONTAINER_CREDENTIALS_RELATIVE_URI") =>
1224 {
1225 Err(LakeError::AmbientStorageAuthority {
1226 name: "AWS_CONTAINER_CREDENTIALS_RELATIVE_URI".into(),
1227 })
1228 }
1229 ObjectStorageAuthentication::S3Static { .. }
1230 | ObjectStorageAuthentication::AzureManagedIdentity { .. }
1231 | ObjectStorageAuthentication::AzureServicePrincipal { .. }
1232 if names
1233 .iter()
1234 .any(|name| name == "AWS_CONTAINER_CREDENTIALS_RELATIVE_URI") =>
1235 {
1236 Err(LakeError::AmbientStorageAuthority {
1237 name: "AWS_CONTAINER_CREDENTIALS_RELATIVE_URI".into(),
1238 })
1239 }
1240 _ => Ok(()),
1241 }
1242}
1243
1244fn validate_storage_authentication(
1245 location: &ObjectStorageLocation,
1246 authentication: &ObjectStorageAuthentication<'_>,
1247) -> Result<(), LakeError> {
1248 let compatible = matches!(
1249 (location, authentication),
1250 (
1251 ObjectStorageLocation::AwsS3 { .. },
1252 ObjectStorageAuthentication::S3Static { .. }
1253 | ObjectStorageAuthentication::AwsWorkloadIdentity { .. }
1254 ) | (
1255 ObjectStorageLocation::S3Compatible { .. },
1256 ObjectStorageAuthentication::S3Static { .. }
1257 ) | (
1258 ObjectStorageLocation::AzureBlob { .. },
1259 ObjectStorageAuthentication::AzureManagedIdentity { .. }
1260 | ObjectStorageAuthentication::AzureServicePrincipal { .. }
1261 )
1262 );
1263 if compatible {
1264 Ok(())
1265 } else {
1266 Err(LakeError::IncompatibleStorageAuthentication)
1267 }
1268}
1269
1270fn storage_secret_sql(
1271 location: &ObjectStorageLocation,
1272 authentication: ObjectStorageAuthentication<'_>,
1273) -> Result<String, LakeError> {
1274 let scope = location.canonical_data_path();
1275 let fields = match authentication {
1276 ObjectStorageAuthentication::S3Static {
1277 access_key_id,
1278 secret_access_key,
1279 } => access_key_id.expose(|access_bytes| {
1280 secret_access_key.expose(|secret_bytes| -> Result<Vec<String>, LakeError> {
1281 let access = std::str::from_utf8(access_bytes)
1282 .map_err(|_| LakeError::InvalidCredentialEncoding)?;
1283 let secret = std::str::from_utf8(secret_bytes)
1284 .map_err(|_| LakeError::InvalidCredentialEncoding)?;
1285 let mut fields = vec![
1286 "TYPE s3".to_string(),
1287 "PROVIDER config".to_string(),
1288 format!("KEY_ID {}", sql_literal(access)),
1289 format!("SECRET {}", sql_literal(secret)),
1290 ];
1291 match location {
1292 ObjectStorageLocation::AwsS3 { region, .. } => {
1293 fields.push(format!("REGION {}", sql_literal(region)));
1294 }
1295 ObjectStorageLocation::S3Compatible {
1296 endpoint,
1297 url_style,
1298 ..
1299 } => {
1300 let endpoint = endpoint.strip_prefix("https://").unwrap_or(endpoint);
1301 fields.push(format!("ENDPOINT {}", sql_literal(endpoint)));
1302 fields.push(format!(
1303 "URL_STYLE {}",
1304 sql_literal(match url_style {
1305 S3UrlStyle::Path => "path",
1306 S3UrlStyle::VirtualHosted => "vhost",
1307 })
1308 ));
1309 fields.push("USE_SSL true".to_string());
1310 fields.push("VERIFY_SSL true".to_string());
1311 }
1312 ObjectStorageLocation::AzureBlob { .. } => {
1313 return Err(LakeError::IncompatibleStorageAuthentication);
1314 }
1315 }
1316 Ok(fields)
1317 })
1318 })?,
1319 ObjectStorageAuthentication::AwsWorkloadIdentity {
1320 source: AwsWorkloadIdentitySource::Ec2Instance,
1321 } => {
1322 let ObjectStorageLocation::AwsS3 { region, .. } = location else {
1323 return Err(LakeError::IncompatibleStorageAuthentication);
1324 };
1325 aws_session_credential_fields(resolve_ec2_instance_credentials()?, region)?
1326 }
1327 ObjectStorageAuthentication::AwsWorkloadIdentity {
1328 source: AwsWorkloadIdentitySource::EcsTask,
1329 } => {
1330 let ObjectStorageLocation::AwsS3 { region, .. } = location else {
1331 return Err(LakeError::IncompatibleStorageAuthentication);
1332 };
1333 aws_session_credential_fields(resolve_ecs_task_credentials()?, region)?
1334 }
1335 ObjectStorageAuthentication::AzureManagedIdentity { client_id } => {
1336 let ObjectStorageLocation::AzureBlob {
1337 account_endpoint, ..
1338 } = location
1339 else {
1340 return Err(LakeError::IncompatibleStorageAuthentication);
1341 };
1342 let account_name = azure_account_name(account_endpoint)?;
1343 let mut fields = vec![
1344 "TYPE azure".to_string(),
1345 "PROVIDER managed_identity".to_string(),
1346 format!("ACCOUNT_NAME {}", sql_literal(account_name)),
1347 ];
1348 if let Some(client_id) = client_id {
1349 fields.push(format!("CLIENT_ID {}", sql_literal(&client_id.to_string())));
1350 }
1351 fields
1352 }
1353 ObjectStorageAuthentication::AzureServicePrincipal {
1354 tenant_id,
1355 client_id,
1356 client_secret,
1357 } => {
1358 let ObjectStorageLocation::AzureBlob {
1359 account_endpoint, ..
1360 } = location
1361 else {
1362 return Err(LakeError::IncompatibleStorageAuthentication);
1363 };
1364 let account_name = azure_account_name(account_endpoint)?;
1365 client_secret.expose(|bytes| -> Result<Vec<String>, LakeError> {
1366 let secret =
1367 std::str::from_utf8(bytes).map_err(|_| LakeError::InvalidCredentialEncoding)?;
1368 Ok(vec![
1369 "TYPE azure".to_string(),
1370 "PROVIDER service_principal".to_string(),
1371 format!("TENANT_ID {}", sql_literal(&tenant_id.to_string())),
1372 format!("CLIENT_ID {}", sql_literal(&client_id.to_string())),
1373 format!("CLIENT_SECRET {}", sql_literal(secret)),
1374 format!("ACCOUNT_NAME {}", sql_literal(account_name)),
1375 ])
1376 })?
1377 }
1378 };
1379 let mut fields = fields;
1380 fields.push(format!("SCOPE {}", sql_literal(&scope)));
1381 Ok(format!(
1382 "CREATE OR REPLACE SECRET ahri_tre_lake_storage ({})",
1383 fields.join(", ")
1384 ))
1385}
1386
1387struct AwsSessionCredentials {
1388 access_key_id: SecretMaterial,
1389 secret_access_key: SecretMaterial,
1390 session_token: SecretMaterial,
1391}
1392
1393#[derive(serde::Deserialize)]
1394#[serde(rename_all = "PascalCase")]
1395struct AwsCredentialResponse {
1396 access_key_id: String,
1397 secret_access_key: String,
1398 token: String,
1399}
1400
1401fn aws_session_credential_fields(
1402 credentials: AwsSessionCredentials,
1403 region: &str,
1404) -> Result<Vec<String>, LakeError> {
1405 credentials.access_key_id.expose(|access_bytes| {
1406 credentials.secret_access_key.expose(|secret_bytes| {
1407 credentials.session_token.expose(|token_bytes| {
1408 let access = std::str::from_utf8(access_bytes)
1409 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?;
1410 let secret = std::str::from_utf8(secret_bytes)
1411 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?;
1412 let token = std::str::from_utf8(token_bytes)
1413 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?;
1414 Ok(vec![
1415 "TYPE s3".to_string(),
1416 "PROVIDER config".to_string(),
1417 format!("KEY_ID {}", sql_literal(access)),
1418 format!("SECRET {}", sql_literal(secret)),
1419 format!("SESSION_TOKEN {}", sql_literal(token)),
1420 format!("REGION {}", sql_literal(region)),
1421 ])
1422 })
1423 })
1424 })
1425}
1426
1427fn workload_identity_client() -> Result<reqwest::blocking::Client, LakeError> {
1428 reqwest::blocking::Client::builder()
1429 .no_proxy()
1430 .redirect(reqwest::redirect::Policy::none())
1431 .timeout(std::time::Duration::from_secs(2))
1432 .build()
1433 .map_err(|_| LakeError::WorkloadIdentityUnavailable)
1434}
1435
1436fn read_bounded_workload_response(
1437 response: reqwest::blocking::Response,
1438 limit: usize,
1439) -> Result<Vec<u8>, LakeError> {
1440 if !response.status().is_success() {
1441 return Err(LakeError::WorkloadIdentityUnavailable);
1442 }
1443 let mut encoded = Vec::new();
1444 response
1445 .take((limit + 1) as u64)
1446 .read_to_end(&mut encoded)
1447 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?;
1448 if encoded.len() > limit {
1449 return Err(LakeError::WorkloadIdentityUnavailable);
1450 }
1451 Ok(encoded)
1452}
1453
1454fn resolve_ec2_instance_credentials() -> Result<AwsSessionCredentials, LakeError> {
1455 const IMDS: &str = "http://169.254.169.254";
1456 const TOKEN_HEADER: &str = "x-aws-ec2-metadata-token";
1457 let client = workload_identity_client()?;
1458 let token = read_bounded_workload_response(
1459 client
1460 .put(format!("{IMDS}/latest/api/token"))
1461 .header("x-aws-ec2-metadata-token-ttl-seconds", "60")
1462 .send()
1463 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?,
1464 4096,
1465 )?;
1466 if token.is_empty() {
1467 return Err(LakeError::WorkloadIdentityUnavailable);
1468 }
1469 let token = reqwest::header::HeaderValue::from_bytes(&token)
1470 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?;
1471 let role = read_bounded_workload_response(
1472 client
1473 .get(format!("{IMDS}/latest/meta-data/iam/security-credentials/"))
1474 .header(TOKEN_HEADER, token.clone())
1475 .send()
1476 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?,
1477 128,
1478 )?;
1479 let role = validated_imds_role_name(&role)?;
1480 let encoded = read_bounded_workload_response(
1481 client
1482 .get(format!(
1483 "{IMDS}/latest/meta-data/iam/security-credentials/{role}"
1484 ))
1485 .header(TOKEN_HEADER, token)
1486 .send()
1487 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?,
1488 MAX_SECRET_MATERIAL_BYTES,
1489 )?;
1490 aws_credentials_from_response(&encoded)
1491}
1492
1493fn validated_imds_role_name(encoded: &[u8]) -> Result<&str, LakeError> {
1494 let role = std::str::from_utf8(encoded)
1495 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?
1496 .trim_end_matches(['\r', '\n']);
1497 if role.is_empty()
1498 || role.len() > 64
1499 || !role.bytes().all(|byte| {
1500 byte.is_ascii_alphanumeric()
1501 || matches!(byte, b'_' | b'+' | b'=' | b',' | b'.' | b'@' | b'-')
1502 })
1503 {
1504 return Err(LakeError::WorkloadIdentityUnavailable);
1505 }
1506 Ok(role)
1507}
1508
1509fn resolve_ecs_task_credentials() -> Result<AwsSessionCredentials, LakeError> {
1510 let relative = std::env::var("AWS_CONTAINER_CREDENTIALS_RELATIVE_URI")
1511 .map_err(|_| LakeError::MissingWorkloadIdentitySource { kind: "ecs_task" })?;
1512 if relative.len() <= "/v2/credentials/".len()
1513 || relative.len() > 2048
1514 || !relative.starts_with("/v2/credentials/")
1515 || relative.contains("//")
1516 || relative.contains('\n')
1517 || relative.contains('\r')
1518 || relative.contains('\\')
1519 || relative.contains('?')
1520 || relative.contains('#')
1521 || relative
1522 .split('/')
1523 .any(|component| component == "." || component == "..")
1524 {
1525 return Err(LakeError::WorkloadIdentityUnavailable);
1526 }
1527 let response = workload_identity_client()?
1528 .get(format!("http://169.254.170.2{relative}"))
1529 .send()
1530 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?;
1531 let encoded = read_bounded_workload_response(response, MAX_SECRET_MATERIAL_BYTES)?;
1532 aws_credentials_from_response(&encoded)
1533}
1534
1535fn aws_credentials_from_response(encoded: &[u8]) -> Result<AwsSessionCredentials, LakeError> {
1536 let response: AwsCredentialResponse =
1537 serde_json::from_slice(encoded).map_err(|_| LakeError::WorkloadIdentityUnavailable)?;
1538 if response.access_key_id.is_empty()
1539 || response.secret_access_key.is_empty()
1540 || response.token.is_empty()
1541 {
1542 return Err(LakeError::WorkloadIdentityUnavailable);
1543 }
1544 Ok(AwsSessionCredentials {
1545 access_key_id: SecretMaterial::try_from(response.access_key_id.into_bytes())
1546 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?,
1547 secret_access_key: SecretMaterial::try_from(response.secret_access_key.into_bytes())
1548 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?,
1549 session_token: SecretMaterial::try_from(response.token.into_bytes())
1550 .map_err(|_| LakeError::WorkloadIdentityUnavailable)?,
1551 })
1552}
1553
1554fn azure_account_name(endpoint: &str) -> Result<&str, LakeError> {
1555 endpoint
1556 .strip_prefix("https://")
1557 .and_then(|authority| authority.split('.').next())
1558 .filter(|value| !value.is_empty())
1559 .ok_or(LakeError::InvalidObjectStorageIdentity)
1560}
1561
1562fn canonical_datastore_prefix(prefix: &str) -> bool {
1563 !prefix.is_empty()
1564 && !prefix.starts_with('/')
1565 && !prefix.ends_with('/')
1566 && !prefix.contains("//")
1567 && !prefix.contains('\\')
1568 && prefix
1569 .split('/')
1570 .all(|component| !component.is_empty() && component != "." && component != "..")
1571 && !prefix.chars().any(char::is_control)
1572}
1573
1574impl DuckLakeAdapter {
1575 fn prepare_local_layout(
1576 &self,
1577 options: DuckLakeOpenOptions,
1578 ) -> Result<DuckLakeLocalLayout, LakeError> {
1579 let layout = self.local_layout()?;
1580 if options.create_if_not_exists {
1581 fs::create_dir_all(&layout.data_path)?;
1582 fs::create_dir_all(&layout.staging_path)?;
1583 fs::create_dir_all(&layout.tmp_path)?;
1584 fs::create_dir_all(&layout.test_runs_path)?;
1585 } else {
1586 fs::metadata(&layout.data_path)?;
1587 }
1588 Ok(layout)
1589 }
1590
1591 fn attach_catalog(
1592 &self,
1593 profile: &DataStoreProfile,
1594 credential: DuckLakeCatalogCredential,
1595 options: DuckLakeOpenOptions,
1596 root_certificate: Option<&str>,
1597 temp_directory: Option<&std::path::Path>,
1598 observations: &mut LakeCatalogObservations,
1599 ) -> Result<(duckdb::Connection, DuckLakeAttachPlan), LakeError> {
1600 let attach_plan = self.build_attach_plan_with_credential_and_root(
1601 profile,
1602 credential,
1603 options,
1604 root_certificate,
1605 )?;
1606 let conn = duckdb::Connection::open_in_memory()?;
1607 if let Some(temp_directory) = temp_directory {
1608 conn.execute_batch(&format!(
1609 "SET temp_directory = {}",
1610 sql_literal(&temp_directory.to_string_lossy())
1611 ))?;
1612 }
1613 conn.execute_batch("INSTALL postgres;\nLOAD postgres;\nINSTALL ducklake;\nLOAD ducklake;")?;
1614 for statement in &attach_plan.pre_attach_sql {
1615 conn.execute_batch(statement)
1616 .map_err(|_| LakeError::CredentialProjectionFailed)?;
1617 }
1618 observations.attach(&conn, &attach_plan.attach_sql, &profile.lake.auth.lake_user)?;
1619 conn.execute_batch(&format!("USE {LAKE_ALIAS};"))?;
1620 Ok((conn, attach_plan))
1621 }
1622
1623 fn attach_options(
1624 &self,
1625 profile: &DataStoreProfile,
1626 data_path: &str,
1627 options: DuckLakeOpenOptions,
1628 ) -> Result<Vec<String>, LakeError> {
1629 let mut attach_options = vec![
1630 format!(
1631 "CREATE_IF_NOT_EXISTS {}",
1632 sql_bool(options.create_if_not_exists)
1633 ),
1634 format!(
1635 "AUTOMATIC_MIGRATION {}",
1636 sql_bool(options.automatic_migration)
1637 ),
1638 ];
1639 if options.create_if_not_exists {
1640 attach_options.insert(0, format!("DATA_PATH {}", sql_literal(data_path)));
1641 }
1642
1643 if let Some(metadata_schema) = profile.lake.catalog_schema.as_deref() {
1644 attach_options.push(format!("METADATA_SCHEMA {}", sql_literal(metadata_schema)));
1645 }
1646
1647 if options.apply_configured_encryption {
1648 match profile.lake.auth.encryption_mode {
1649 EncryptionMode::None => {}
1650 EncryptionMode::ServerManaged => attach_options.push("ENCRYPTED".to_string()),
1651 EncryptionMode::ClientManaged => {
1652 return Err(LakeError::UnsupportedEncryptionMode(
1653 EncryptionMode::ClientManaged,
1654 ));
1655 }
1656 }
1657 }
1658
1659 Ok(attach_options)
1660 }
1661}
1662
1663struct DuckLakeSettings {
1664 catalog_type: String,
1665 extension_version: String,
1666 data_path: String,
1667}
1668
1669#[derive(Debug, Clone, PartialEq, Eq)]
1670struct DuckLakeMetadataRelation {
1671 catalog: String,
1672 schema: String,
1673}
1674
1675impl DuckLakeMetadataRelation {
1676 fn table_sql(&self, table_name: &str) -> String {
1677 format!(
1678 "{}.{}.{}",
1679 quote_identifier(&self.catalog),
1680 quote_identifier(&self.schema),
1681 quote_identifier(table_name)
1682 )
1683 }
1684}
1685
1686fn ducklake_settings(connection: &duckdb::Connection) -> Result<DuckLakeSettings, LakeError> {
1687 connection
1688 .query_row(
1689 &format!(
1690 "SELECT catalog_type, extension_version, data_path FROM ducklake_settings({})",
1691 sql_literal(LAKE_ALIAS)
1692 ),
1693 [],
1694 |row| {
1695 Ok(DuckLakeSettings {
1696 catalog_type: row.get(0)?,
1697 extension_version: row.get(1)?,
1698 data_path: row.get(2)?,
1699 })
1700 },
1701 )
1702 .map_err(LakeError::from)
1703}
1704
1705fn detected_encryption_mode(
1706 connection: &duckdb::Connection,
1707 data_path: &str,
1708) -> Result<EncryptionMode, LakeError> {
1709 if let Some(relation) = ducklake_metadata_relation(connection, data_path)?
1710 && let Some(value) = connection
1711 .query_row(
1712 &format!(
1713 "SELECT value FROM {} WHERE lower(key) = 'encrypted' LIMIT 1",
1714 relation.table_sql("ducklake_metadata")
1715 ),
1716 [],
1717 |row| row.get::<_, String>(0),
1718 )
1719 .optional()?
1720 {
1721 return parse_encryption_metadata_value("ducklake_metadata.encrypted", value);
1722 }
1723
1724 let value = connection
1725 .query_row(
1726 &format!(
1727 "SELECT value FROM {LAKE_ALIAS}.options() WHERE lower(option_name) = 'encrypted' LIMIT 1"
1728 ),
1729 [],
1730 |row| row.get::<_, String>(0),
1731 )
1732 .optional()?
1733 .unwrap_or_else(|| "false".to_string());
1734
1735 parse_encryption_metadata_value("encrypted", value)
1736}
1737
1738fn ducklake_metadata_relation(
1739 connection: &duckdb::Connection,
1740 data_path: &str,
1741) -> Result<Option<DuckLakeMetadataRelation>, LakeError> {
1742 let metadata_catalog = format!("__ducklake_metadata_{LAKE_ALIAS}");
1743 let expected_data_path = normalize_ducklake_data_path(data_path);
1744 let mut statement = connection.prepare(
1745 "SELECT table_schema
1746 FROM information_schema.tables
1747 WHERE table_catalog = ?
1748 AND table_name = 'ducklake_metadata'
1749 ORDER BY table_schema",
1750 )?;
1751 let schemas = statement
1752 .query_map([metadata_catalog.as_str()], |row| row.get::<_, String>(0))?
1753 .collect::<Result<Vec<_>, _>>()?;
1754
1755 for schema in schemas {
1756 let metadata_data_path = connection
1757 .query_row(
1758 &format!(
1759 "SELECT value FROM {}.{}.{} WHERE lower(key) = 'data_path' LIMIT 1",
1760 quote_identifier(&metadata_catalog),
1761 quote_identifier(&schema),
1762 quote_identifier("ducklake_metadata")
1763 ),
1764 [],
1765 |row| row.get::<_, String>(0),
1766 )
1767 .optional()?;
1768
1769 if metadata_data_path
1770 .as_deref()
1771 .map(normalize_ducklake_data_path)
1772 .as_deref()
1773 == Some(expected_data_path.as_str())
1774 {
1775 return Ok(Some(DuckLakeMetadataRelation {
1776 catalog: metadata_catalog,
1777 schema,
1778 }));
1779 }
1780 }
1781
1782 Ok(None)
1783}
1784
1785fn parse_encryption_metadata_value(
1786 key: &'static str,
1787 value: String,
1788) -> Result<EncryptionMode, LakeError> {
1789 match value.trim().to_ascii_lowercase().as_str() {
1790 "false" => Ok(EncryptionMode::None),
1791 "true" => Ok(EncryptionMode::ServerManaged),
1792 _ => Err(LakeError::InvalidMetadataValue { key, value }),
1793 }
1794}
1795
1796fn normalize_ducklake_data_path(value: &str) -> String {
1797 value.trim().trim_end_matches('/').to_string()
1798}
1799
1800fn normalized_root(lake_root: &str) -> Result<&str, LakeError> {
1801 let root = lake_root.trim().trim_end_matches('/');
1802 if root.is_empty() {
1803 return Err(LakeError::EmptyLakeRoot);
1804 }
1805 Ok(root)
1806}
1807
1808fn ducklake_postgres_conninfo(
1809 profile: &DataStoreProfile,
1810 credential: &DuckLakeCatalogCredential,
1811 root_certificate: Option<&str>,
1812 redact_secret: bool,
1813) -> String {
1814 let mut base_params = vec![
1815 libpq_param("host", &profile.server),
1816 libpq_param("port", &profile.port.to_string()),
1817 libpq_param("dbname", &profile.lake.lake_db),
1818 libpq_param("user", &profile.lake.auth.lake_user),
1819 ];
1820 if let Some(sslmode) = profile.sslmode.as_deref().filter(|value| !value.is_empty()) {
1821 base_params.push(libpq_param("sslmode", sslmode));
1822 }
1823 if let Some(root_certificate) = root_certificate {
1824 base_params.push(libpq_param(
1825 "sslrootcert",
1826 if redact_secret && root_certificate != "system" {
1827 "<adapter-projected-ca>"
1828 } else {
1829 root_certificate
1830 },
1831 ));
1832 }
1833
1834 let params = match credential {
1835 DuckLakeCatalogCredential::SessionSecret { .. } => base_params,
1836 };
1837
1838 format!("ducklake:postgres:{}", params.join(" "))
1839}
1840
1841fn ducklake_pre_attach_sql(
1842 profile: &DataStoreProfile,
1843 credential: &DuckLakeCatalogCredential,
1844 redact_secret: bool,
1845) -> Vec<String> {
1846 match credential {
1847 DuckLakeCatalogCredential::SessionSecret {
1848 secret_name,
1849 password,
1850 } => {
1851 let password = if redact_secret {
1852 "<redacted>"
1853 } else {
1854 password.as_str()
1855 };
1856 vec![format!(
1857 "CREATE OR REPLACE SECRET {} (TYPE postgres, HOST {}, PORT {}, DATABASE {}, USER {}, PASSWORD {})",
1858 quote_identifier(&session_secret_name(secret_name)),
1859 sql_literal(&profile.server),
1860 profile.port,
1861 sql_literal(&profile.lake.lake_db),
1862 sql_literal(&profile.lake.auth.lake_user),
1863 sql_literal(password),
1864 )]
1865 }
1866 }
1867}
1868
1869fn session_secret_name(value: &str) -> String {
1870 value
1871 .trim()
1872 .chars()
1873 .map(|ch| {
1874 if ch.is_ascii_alphanumeric() || ch == '_' {
1875 ch
1876 } else {
1877 '_'
1878 }
1879 })
1880 .collect()
1881}
1882
1883fn libpq_param(key: &str, value: &str) -> String {
1884 format!(
1885 "{}='{}'",
1886 key,
1887 value.replace('\\', "\\\\").replace('\'', "\\'")
1888 )
1889}
1890
1891fn sql_literal(value: &str) -> String {
1892 format!("'{}'", value.replace('\'', "''"))
1893}
1894
1895fn quote_identifier(value: &str) -> String {
1896 format!("\"{}\"", value.replace('"', "\"\""))
1897}
1898
1899fn sql_bool(value: bool) -> &'static str {
1900 if value { "true" } else { "false" }
1901}
1902
1903fn finish_lake_open(
1904 span: Option<ahri_tre_observability::OperationSpan>,
1905 result: &Result<DuckLakeOpenedCatalog, LakeError>,
1906) {
1907 use ahri_tre_observability::{FailureCategory, Outcome};
1908 if let Some(span) = span {
1909 span.finish_observed(
1910 if result.is_ok() {
1911 Outcome::Success
1912 } else {
1913 Outcome::Unavailable
1914 },
1915 result.as_ref().err().map(|_| FailureCategory::Adapter),
1916 );
1917 }
1918}
1919
1920#[cfg(test)]
1921mod tests {
1922 use std::{
1923 fs,
1924 os::unix::fs::PermissionsExt,
1925 path::PathBuf,
1926 time::{SystemTime, UNIX_EPOCH},
1927 };
1928
1929 use ahri_tre_secrets::SecretMaterial;
1930 use ahri_tre_types::{DataStoreProfile, EncryptionMode, LakeAttachProfile, LakeAuthProfile};
1931
1932 use super::{
1933 DuckLakeAdapter, DuckLakeCatalogCredential, DuckLakeCredentialKind, DuckLakeOpenOptions,
1934 FilesystemCatalogTls, FilesystemLakeSessionConfig, ObjectStorageLocation,
1935 aws_credentials_from_response, object_namespace_glob, validated_imds_role_name,
1936 };
1937 use crate::{LAKE_ALIAS, LakeError, ScratchAttemptId, TrustedScratch};
1938
1939 #[test]
1940 fn local_layout_expands_canonical_subdirectories() {
1941 let adapter = DuckLakeAdapter::new("/lake/root/");
1942 let layout = adapter.local_layout().expect("layout should build");
1943
1944 assert_eq!(layout.data_path, "/lake/root");
1945 assert_eq!(layout.staging_path, "/lake/root/__tre_duckdb_stage");
1946 assert_eq!(layout.tmp_path, "/lake/root/__tre_duckdb_stage");
1947 assert_eq!(layout.test_runs_path, "/lake/root/test-runs");
1948 }
1949
1950 #[test]
1951 fn local_layout_rejects_empty_root() {
1952 let adapter = DuckLakeAdapter::new(" ");
1953 assert!(matches!(
1954 adapter.local_layout(),
1955 Err(LakeError::EmptyLakeRoot)
1956 ));
1957 }
1958
1959 #[test]
1960 fn object_namespace_glob_preserves_the_exact_canonical_prefix() {
1961 let location = ObjectStorageLocation::aws_s3_for_datastore(
1962 "af-south-1",
1963 "research",
1964 Some("/tenant/"),
1965 "datastore",
1966 )
1967 .unwrap();
1968
1969 assert_eq!(
1970 location.canonical_data_path(),
1971 "s3://research/tenant/datastore/"
1972 );
1973 assert_eq!(
1974 object_namespace_glob(&location),
1975 "s3://research/tenant/datastore/**"
1976 );
1977 }
1978
1979 #[test]
1980 fn ecs_task_response_becomes_bounded_protected_session_credentials() {
1981 let credentials = aws_credentials_from_response(
1982 br#"{
1983 "AccessKeyId":"ecs-access-canary",
1984 "SecretAccessKey":"ecs-secret-canary",
1985 "Token":"ecs-token-canary",
1986 "Expiration":"2030-01-01T00:00:00Z"
1987 }"#,
1988 )
1989 .unwrap();
1990
1991 credentials
1992 .access_key_id
1993 .expose(|bytes| assert_eq!(bytes, b"ecs-access-canary"));
1994 credentials
1995 .secret_access_key
1996 .expose(|bytes| assert_eq!(bytes, b"ecs-secret-canary"));
1997 credentials
1998 .session_token
1999 .expose(|bytes| assert_eq!(bytes, b"ecs-token-canary"));
2000 }
2001
2002 #[test]
2003 fn ec2_role_name_is_a_single_closed_path_component() {
2004 assert_eq!(
2005 validated_imds_role_name(b"research-node-role\n").unwrap(),
2006 "research-node-role"
2007 );
2008 for invalid in [
2009 b"".as_slice(),
2010 b"../role".as_slice(),
2011 b"role/other".as_slice(),
2012 b"role?query".as_slice(),
2013 b"role\nother".as_slice(),
2014 ] {
2015 assert!(matches!(
2016 validated_imds_role_name(invalid),
2017 Err(LakeError::WorkloadIdentityUnavailable)
2018 ));
2019 }
2020 }
2021
2022 #[test]
2023 fn explicit_filesystem_inputs_keep_path_joining_and_credentials_in_adapter() {
2024 let config = FilesystemLakeSessionConfig::new(
2025 "/srv/ahri-tre/lake",
2026 "research/data",
2027 "postgres.example.test",
2028 5432,
2029 "ducklake_catalog",
2030 "catalog",
2031 "catalog_owner",
2032 FilesystemCatalogTls::System,
2033 EncryptionMode::ServerManaged,
2034 )
2035 .unwrap();
2036 let credential = SecretMaterial::try_from(b"catalog-secret".to_vec()).unwrap();
2037
2038 let plan = DuckLakeAdapter::build_filesystem_attach_plan(&config, &credential).unwrap();
2039
2040 assert_eq!(config.filesystem_base(), "/srv/ahri-tre/lake");
2041 assert_eq!(config.datastore_prefix(), "research/data");
2042 assert_eq!(plan.data_path, "/srv/ahri-tre/lake/research/data");
2043 assert!(plan.attach_sql.contains("sslrootcert"));
2044 assert!(!plan.attach_description.contains("catalog-secret"));
2045 assert!(!format!("{plan:?}").contains("catalog-secret"));
2046 }
2047
2048 #[test]
2049 fn datastore_owned_path_is_validated_and_relativized_inside_the_adapter() {
2050 let config = FilesystemLakeSessionConfig::from_binding(
2051 "/srv/ahri-tre/lake",
2052 "/srv/ahri-tre/lake/datastores/abc",
2053 "postgres.example.test",
2054 5432,
2055 "ducklake_catalog",
2056 "catalog",
2057 "catalog_owner",
2058 FilesystemCatalogTls::System,
2059 EncryptionMode::ServerManaged,
2060 )
2061 .unwrap();
2062
2063 assert_eq!(config.datastore_prefix(), "datastores/abc");
2064 assert_eq!(
2065 config.canonical_data_path(),
2066 "/srv/ahri-tre/lake/datastores/abc"
2067 );
2068 }
2069
2070 #[test]
2071 fn filesystem_initialization_rejects_an_existing_namespace_before_opening_duckdb() {
2072 let base = temp_lake_root("occupied_namespace");
2073 let data_path = base.join("datastores/occupied");
2074 fs::create_dir_all(&data_path).unwrap();
2075 let scratch_root = temp_lake_root("occupied_namespace_scratch");
2076 fs::set_permissions(&scratch_root, fs::Permissions::from_mode(0o700)).unwrap();
2077 let scratch = TrustedScratch::open(&scratch_root, [data_path.as_path()])
2078 .unwrap()
2079 .create_attempt(ScratchAttemptId::new("0123456789abcdef0123456789abcdef").unwrap())
2080 .unwrap();
2081 let config = FilesystemLakeSessionConfig::new(
2082 base.display().to_string(),
2083 "datastores/occupied",
2084 "postgres.example.test",
2085 5432,
2086 "ducklake_catalog",
2087 "catalog",
2088 "catalog_owner",
2089 FilesystemCatalogTls::System,
2090 EncryptionMode::None,
2091 )
2092 .unwrap();
2093 let credential = SecretMaterial::try_from(b"catalog-secret".to_vec()).unwrap();
2094
2095 let error = DuckLakeAdapter::initialize_filesystem_catalog_with_scratch(
2096 &config,
2097 &credential,
2098 &scratch,
2099 )
2100 .expect_err("an existing namespace must never be adopted or replaced");
2101
2102 assert!(matches!(error, LakeError::NamespaceNotEmpty));
2103 drop(scratch);
2104 fs::remove_dir_all(base).unwrap();
2105 fs::remove_dir_all(scratch_root).unwrap();
2106 }
2107
2108 #[test]
2109 fn attach_existing_layout_does_not_create_runtime_subdirectories() {
2110 let root = temp_lake_root("attach_existing_layout");
2111 let adapter = DuckLakeAdapter::new(root.display().to_string());
2112
2113 let layout = adapter
2114 .prepare_local_layout(DuckLakeOpenOptions::attach_existing())
2115 .expect("existing data root should be accepted without helper dirs");
2116
2117 assert_eq!(layout.data_path, root.display().to_string());
2118 assert!(!root.join("__tre_duckdb_stage").exists());
2119 assert!(!root.join("test-runs").exists());
2120 fs::remove_dir_all(root).expect("temp lake root should be removed");
2121 }
2122
2123 #[test]
2124 fn initialize_layout_creates_runtime_subdirectories() {
2125 let root = temp_lake_root("initialize_layout");
2126 let adapter = DuckLakeAdapter::new(root.display().to_string());
2127
2128 adapter
2129 .prepare_local_layout(DuckLakeOpenOptions::initialize())
2130 .expect("initialization should prepare helper dirs");
2131
2132 assert!(root.join("__tre_duckdb_stage").is_dir());
2133 assert!(root.join("test-runs").is_dir());
2134 fs::remove_dir_all(root).expect("temp lake root should be removed");
2135 }
2136
2137 fn temp_lake_root(name: &str) -> PathBuf {
2138 let suffix = SystemTime::now()
2139 .duration_since(UNIX_EPOCH)
2140 .expect("system time should be after epoch")
2141 .as_nanos();
2142 let root = std::env::temp_dir().join(format!(
2143 "ahri_tre_lake_{name}_{}_{}",
2144 std::process::id(),
2145 suffix
2146 ));
2147 fs::create_dir_all(&root).expect("temp lake root should be created");
2148 root
2149 }
2150
2151 #[test]
2152 fn attach_plan_uses_canonical_data_path_and_migration_options() {
2153 let adapter = DuckLakeAdapter::new("/lake/root");
2154 let plan = adapter
2155 .build_attach_plan_with_credential_and_root(
2156 &test_profile(EncryptionMode::None),
2157 DuckLakeCatalogCredential::session_secret("tre_catalog", "catalog-secret"),
2158 DuckLakeOpenOptions::attach_existing().with_automatic_migration(true),
2159 None,
2160 )
2161 .expect("attach plan should build");
2162
2163 assert_eq!(plan.alias, LAKE_ALIAS);
2164 assert_eq!(plan.credential_kind, DuckLakeCredentialKind::SessionSecret);
2165 assert_eq!(plan.data_path, "/lake/root");
2166 assert!(!plan.create_if_not_exists);
2167 assert!(plan.automatic_migration);
2168 assert_eq!(plan.attach_sql, plan.attach_description);
2169 assert!(
2170 plan.attach_description
2171 .contains("ATTACH 'ducklake:postgres:host=")
2172 );
2173 assert!(plan.attach_sql.contains("META_SECRET 'tre_catalog'"));
2174 assert!(!plan.attach_description.contains("catalog-secret"));
2175 assert!(plan.attach_description.contains("sslmode=''disable''"));
2176 assert!(!plan.attach_description.contains("DATA_PATH"));
2177 assert!(
2178 plan.attach_description
2179 .contains("CREATE_IF_NOT_EXISTS false")
2180 );
2181 assert!(plan.attach_description.contains("AUTOMATIC_MIGRATION true"));
2182 assert!(
2183 plan.attach_description
2184 .contains("METADATA_SCHEMA 'lake_catalog'")
2185 );
2186 }
2187
2188 #[test]
2189 fn attach_plan_supports_session_secret_and_redacts_setup_sql() {
2190 let adapter = DuckLakeAdapter::new("/lake/root");
2191 let plan = adapter
2192 .build_attach_plan_with_credential_and_root(
2193 &test_profile(EncryptionMode::None),
2194 DuckLakeCatalogCredential::session_secret("tre lake secret", "super-secret"),
2195 DuckLakeOpenOptions::attach_existing(),
2196 None,
2197 )
2198 .expect("attach plan should build");
2199
2200 assert_eq!(plan.credential_kind, DuckLakeCredentialKind::SessionSecret);
2201 assert_eq!(plan.pre_attach_sql.len(), 1);
2202 assert!(plan.pre_attach_sql[0].contains("CREATE OR REPLACE SECRET \"tre_lake_secret\""));
2203 assert!(plan.pre_attach_sql[0].contains("PASSWORD 'super-secret'"));
2204 assert!(
2205 plan.attach_sql
2206 .contains("'ducklake:postgres:host=''tre-postgres''")
2207 );
2208 assert!(plan.attach_sql.contains("dbname=''tre_lake''"));
2209 assert!(plan.attach_sql.contains("META_SECRET 'tre_lake_secret'"));
2210 assert!(
2211 plan.attach_description
2212 .contains("'ducklake:postgres:host=''tre-postgres''")
2213 );
2214 assert!(
2215 plan.attach_description
2216 .contains("META_SECRET 'tre_lake_secret'")
2217 );
2218 assert!(!plan.attach_description.contains("super-secret"));
2219 assert!(!format!("{plan:?}").contains("super-secret"));
2220 assert!(!format!("{plan:?}").contains("CREATE OR REPLACE SECRET"));
2221 }
2222
2223 #[test]
2224 fn attach_plan_applies_server_managed_encryption_for_initialization() {
2225 let adapter = DuckLakeAdapter::new("/lake/root");
2226 let plan = adapter
2227 .build_attach_plan_with_credential_and_root(
2228 &test_profile(EncryptionMode::ServerManaged),
2229 DuckLakeCatalogCredential::session_secret("tre_catalog", "catalog-secret"),
2230 DuckLakeOpenOptions::initialize(),
2231 None,
2232 )
2233 .expect("attach plan should build");
2234
2235 assert_eq!(
2236 plan.configured_encryption_mode,
2237 EncryptionMode::ServerManaged
2238 );
2239 assert!(plan.attach_description.contains("DATA_PATH '/lake/root'"));
2240 assert!(plan.attach_sql.contains("ENCRYPTED"));
2241 }
2242
2243 #[test]
2244 fn attach_plan_discovers_existing_catalog_encryption_without_encrypted_option() {
2245 let adapter = DuckLakeAdapter::new("/lake/root");
2246 let plan = adapter
2247 .build_attach_plan_with_credential_and_root(
2248 &test_profile(EncryptionMode::ServerManaged),
2249 DuckLakeCatalogCredential::session_secret("tre_catalog", "catalog-secret"),
2250 DuckLakeOpenOptions::attach_existing(),
2251 None,
2252 )
2253 .expect("attach plan should build");
2254
2255 assert!(!plan.create_if_not_exists);
2256 assert_eq!(
2257 plan.configured_encryption_mode,
2258 EncryptionMode::ServerManaged
2259 );
2260 assert!(!plan.attach_sql.contains("ENCRYPTED"));
2261 }
2262
2263 #[test]
2264 fn attach_plan_rejects_client_managed_encryption() {
2265 let adapter = DuckLakeAdapter::new("/lake/root");
2266 assert!(matches!(
2267 adapter.build_attach_plan_with_credential_and_root(
2268 &test_profile(EncryptionMode::ClientManaged),
2269 DuckLakeCatalogCredential::session_secret("tre_catalog", "catalog-secret"),
2270 DuckLakeOpenOptions::initialize(),
2271 None,
2272 ),
2273 Err(LakeError::UnsupportedEncryptionMode(
2274 EncryptionMode::ClientManaged
2275 ))
2276 ));
2277 }
2278
2279 #[test]
2280 fn attach_plan_rejects_empty_credentials() {
2281 let adapter = DuckLakeAdapter::new("/lake/root");
2282 assert!(matches!(
2283 adapter.build_attach_plan_with_credential_and_root(
2284 &test_profile(EncryptionMode::None),
2285 DuckLakeCatalogCredential::session_secret("tre_catalog", " "),
2286 DuckLakeOpenOptions::default(),
2287 None,
2288 ),
2289 Err(LakeError::EmptyCredential {
2290 kind: "session_secret_password"
2291 })
2292 ));
2293 }
2294
2295 fn test_profile(encryption_mode: EncryptionMode) -> DataStoreProfile {
2296 DataStoreProfile {
2297 server: "tre-postgres".to_string(),
2298 port: 5432,
2299 dbname: "tre".to_string(),
2300 sslmode: Some("disable".to_string()),
2301 lake: LakeAttachProfile {
2302 lake_data: "/lake/root".to_string(),
2303 lake_db: "tre_lake".to_string(),
2304 catalog_schema: Some("lake_catalog".to_string()),
2305 auth: LakeAuthProfile {
2306 lake_user: "lake_user".to_string(),
2307 encryption_mode,
2308 },
2309 },
2310 }
2311 }
2312}