Skip to main content

ahri_tre_lake/
adapter.rs

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/// Adapter rooted at a TRE Lake location.
17#[derive(Debug, Clone)]
18pub struct DuckLakeAdapter {
19    pub lake_root: String,
20}
21
22/// Explicit filesystem and PostgreSQL catalog inputs for one existing Lake.
23#[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/// Explicit PostgreSQL trust used by DuckLake's catalog connection.
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub enum FilesystemCatalogTls {
39    System,
40    CustomCa { certificates: Vec<String> },
41}
42
43/// Canonical object-storage namespace owned by one Datastore binding.
44#[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/// The single AWS platform metadata source authorized by configuration.
197#[derive(Debug, Clone, Copy, PartialEq, Eq)]
198pub enum AwsWorkloadIdentitySource {
199    EcsTask,
200    Ec2Instance,
201}
202
203/// Explicit object-storage TLS trust. Verification cannot be disabled.
204#[derive(Debug, Clone, PartialEq, Eq)]
205pub enum ObjectStorageTls {
206    System,
207    CustomCa { certificates: Vec<String> },
208}
209
210/// Explicit catalog and object namespace inputs for one existing Lake.
211#[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/// Connection-scoped storage authentication. Secret material has no formatting
284/// or serialization authority and is consumed only while constructing a DuckDB
285/// temporary secret.
286#[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    /// Validates a datastore-owned canonical path and derives its relative prefix.
371    #[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/// Local directories used by DuckDB and DuckLake for a TRE lake.
448#[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/// Rendered DuckLake `ATTACH` statement and its non-secret description.
457#[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
490/// Open DuckDB connection with the TRE DuckLake catalog attached.
491pub 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/// Basic health and catalog-state information for an attached DuckLake catalog.
564#[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/// Credential material used to connect DuckLake to its PostgreSQL catalog.
576#[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/// Non-secret credential category recorded in attach plans.
597#[derive(Debug, Clone, Copy, PartialEq, Eq)]
598pub enum DuckLakeCredentialKind {
599    SessionSecret,
600}
601
602/// Options controlling whether DuckLake creates or migrates catalog metadata.
603#[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    /// Options for first-time catalog initialization.
622    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    /// Options for opening an existing catalog without creating or migrating it.
631    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    /// Enables or disables DuckLake automatic migration for this open.
640    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    /// Builds a session-local DuckDB secret credential.
648    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    /// Validates object-storage provider and ambient-authority inputs without
686    /// rendering Secret-bearing SQL.
687    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    /// Opens an existing object-backed Lake with explicit catalog and storage
704    /// credentials and a private trusted-scratch attempt.
705    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    /// Observes the existing attach; it never probes or retries independently.
722    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    /// Initializes a new object-backed Lake after proving its configured prefix is empty.
744    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    /// Builds an existing-catalog attach plan from explicit filesystem and
875    /// credential inputs. Lake location composition remains inside this adapter.
876    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    /// Opens an existing filesystem Lake from explicit, non-ambient inputs.
898    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    /// Opens an existing filesystem Lake while directing all DuckDB temporary
945    /// work into one private trusted-scratch attempt rather than the Lake.
946    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    /// Records the time authenticated catalog evidence was actually obtained.
964    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    /// Initializes a new filesystem Lake in a previously absent namespace.
989    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    /// Creates an adapter for the supplied lake root.
1031    pub fn new(lake_root: impl Into<String>) -> Self {
1032        Self {
1033            lake_root: lake_root.into(),
1034        }
1035    }
1036
1037    /// Returns the deterministic local layout below the lake root.
1038    ///
1039    /// # Errors
1040    ///
1041    /// Returns [`LakeError::EmptyLakeRoot`] if the configured root is blank.
1042    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    /// Observe an explicitly authorized health read on the retained connection.
1142    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    /// Reads DuckLake settings and lightweight catalog counts.
1151    ///
1152    /// # Errors
1153    ///
1154    /// Returns an error if DuckDB metadata queries fail or DuckLake reports an
1155    /// unexpected encryption setting.
1156    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}