Skip to main content

ahri_tre_app/
configured_datastore.rs

1use std::sync::{Arc, Mutex};
2
3use ahri_tre_lake::{
4    AwsWorkloadIdentitySource, DuckLakeAdapter, FilesystemCatalogTls, FilesystemLakeSessionConfig,
5    ObjectLakeSessionConfig, ObjectStorageAuthentication, ObjectStorageLocation, ObjectStorageTls,
6    S3UrlStyle, ScratchAttemptId, TrustedScratch,
7};
8use ahri_tre_libpq_oauth::{
9    ExplicitLibpqOAuthConnectConfig, LibpqOAuthConnectFailureKind, LibpqOAuthConnector,
10    LibpqOAuthTls,
11};
12use ahri_tre_orcid::OAuthSession;
13use ahri_tre_pgmeta::{
14    datastore_identity_table_exists_oauth, read_datastore_identity_bindings_oauth,
15};
16use ahri_tre_runtime::{
17    AzureManagedIdentity, ConfiguredAwsWorkloadIdentitySource, ConfiguredS3UrlStyle,
18    ConfiguredSecret, DataStoreRuntime, DatastoreSessionPlan, DuckLakeConnection, LakeStoragePlan,
19    PgStoreConnection, ResolvedOperationSecret, ResolvedSecretVersion, SessionAuthenticationRecord,
20    SessionTls, StorageAuthenticationPlan, StorageTlsPlan,
21};
22use ahri_tre_secrets::{ManagedSecretReference, ResolvedManagedSecret};
23use ahri_tre_session::SessionReferenceIntent;
24use ahri_tre_types::{
25    DataStoreProfile, DatastoreIdentityBinding, DatastoreLakeLocation, DatastoreLifecycleState,
26    EncryptionMode, LakeAttachProfile, LakeAuthProfile,
27};
28
29use crate::{AppError, DataStoreSession, StoreSessionConnection};
30
31/// Already-resolved authority required to open one configured Datastore.
32pub struct OpenConfiguredDatastoreRequest {
33    pub plan: DatastoreSessionPlan,
34    pub authentication: OAuthSession,
35    pub session_reference: ConfiguredSessionReference,
36}
37
38/// Failure opening a configured Datastore for an authenticated Session.
39#[derive(Debug, thiserror::Error)]
40#[non_exhaustive]
41pub enum ConfiguredDatastoreOpenError {
42    /// PostgreSQL rejected the Session's OAuth authentication.
43    #[error("configured PostgreSQL rejected OAuth authentication")]
44    AuthenticationRejected,
45    /// A non-authentication application failure prevented the open.
46    #[error(transparent)]
47    Application(#[from] AppError),
48}
49
50impl ConfiguredDatastoreOpenError {
51    fn into_app_error(self) -> AppError {
52        match self {
53            Self::AuthenticationRejected => AppError::Infrastructure(
54                "configured PostgreSQL open failed: OAuth authentication rejected".to_string(),
55            ),
56            Self::Application(error) => error,
57        }
58    }
59}
60
61pub struct ConfiguredSessionReference {
62    session_id: uuid::Uuid,
63    release_capability: Option<uuid::Uuid>,
64    authentication: SessionAuthenticationRecord,
65}
66
67/// Opaque public identifier paired with one process-ephemeral live Session.
68pub struct LiveSessionIdentifier(uuid::Uuid);
69
70impl LiveSessionIdentifier {
71    pub fn session_id(&self) -> uuid::Uuid {
72        self.0
73    }
74}
75
76impl ConfiguredSessionReference {
77    /// Retains the original live Session identity while refreshing its handles
78    /// for newly authorized work. The caller must still own that live Session.
79    pub fn retained_live(
80        authentication: SessionAuthenticationRecord,
81        session_id: uuid::Uuid,
82    ) -> Self {
83        Self {
84            session_id,
85            release_capability: None,
86            authentication,
87        }
88    }
89
90    pub fn new(
91        intent: SessionReferenceIntent,
92        authentication: SessionAuthenticationRecord,
93    ) -> Result<Self, AppError> {
94        let journal_identity = intent.journal().identity();
95        let authenticated_identity = authentication.identity();
96        if journal_identity.issuer() != authenticated_identity.issuer()
97            || journal_identity.client_id() != authenticated_identity.client_id()
98            || journal_identity.subject() != authenticated_identity.subject()
99        {
100            return Err(AppError::Validation(
101                "Session reference intent belongs to a different authenticated identity".into(),
102            ));
103        }
104        Ok(Self {
105            session_id: intent.journal().reference_id(),
106            release_capability: Some(intent.journal().release_capability()),
107            authentication,
108        })
109    }
110
111    /// Creates a process-ephemeral live Session reference. Unlike a persisted
112    /// client Session journal, it retains no durable Managed-secret holder.
113    pub fn live(authentication: SessionAuthenticationRecord) -> (Self, LiveSessionIdentifier) {
114        let session_id = uuid::Uuid::new_v4();
115        (
116            Self {
117                session_id,
118                release_capability: None,
119                authentication,
120            },
121            LiveSessionIdentifier(session_id),
122        )
123    }
124
125    fn session_id(&self) -> uuid::Uuid {
126        self.session_id
127    }
128
129    fn release_capability(&self) -> uuid::Uuid {
130        self.release_capability
131            .expect("durable Session reference has a release capability")
132    }
133
134    fn retains_snapshot(&self) -> bool {
135        self.release_capability.is_some()
136    }
137}
138
139/// Trusted operation capability that resolves only the binding-owned catalog credential.
140pub trait DatastoreSecretResolver {
141    /// Authorizes one cleanup operation using immutable configuration and the
142    /// existing Secret resolver. The adapter exposes no materialization API.
143    fn dataset_maintenance(
144        &self,
145        _plan: &DatastoreSessionPlan,
146        _binding: &DatastoreIdentityBinding,
147    ) -> Result<ahri_tre_pgmeta::PgDatasetMaintenance, AppError> {
148        Err(AppError::Infrastructure(
149            "Dataset recovery authority is unavailable".into(),
150        ))
151    }
152
153    fn resolve_catalog_credential(
154        &self,
155        binding: &DatastoreIdentityBinding,
156        reference: &ManagedSecretReference,
157    ) -> Result<ResolvedManagedSecret, AppError>;
158
159    fn resolve_storage_credential(
160        &self,
161        _plan: &DatastoreSessionPlan,
162        _configured: &ConfiguredSecret,
163    ) -> Result<ResolvedOperationSecret, AppError> {
164        Err(AppError::Validation(
165            "object-storage credential resolution is unavailable".into(),
166        ))
167    }
168
169    fn retain_session_reference_snapshot(
170        &self,
171        _reference: &ConfiguredSessionReference,
172        _resolved_versions: &[ResolvedSecretVersion],
173    ) -> Result<(), AppError> {
174        Err(AppError::Validation(
175            "configured Datastore Session reference retention is unavailable".into(),
176        ))
177    }
178}
179
180impl DatastoreSecretResolver for ahri_tre_runtime::TrustedRuntimeState {
181    fn dataset_maintenance(
182        &self,
183        plan: &DatastoreSessionPlan,
184        binding: &DatastoreIdentityBinding,
185    ) -> Result<ahri_tre_pgmeta::PgDatasetMaintenance, AppError> {
186        let configured = plan.datastore_id().parse().map_err(|_| {
187            AppError::Infrastructure("Dataset maintenance configuration is unavailable".into())
188        })?;
189        let maintenance = self
190            .datastore_creation_plan_with_id(&configured, binding.datastore_id)
191            .map_err(|e| infrastructure("select Dataset maintenance configuration", e))?;
192        if maintenance.postgresql().database() != plan.postgresql().database()
193            || maintenance.postgresql().host() != plan.postgresql().host()
194            || maintenance.postgresql().port() != plan.postgresql().port()
195        {
196            return Err(AppError::Infrastructure(
197                "Dataset maintenance target does not match the Session".into(),
198            ));
199        }
200        let credential = self
201            .resolve_datastore_administrator(&maintenance)
202            .map_err(|e| infrastructure("authorize Dataset maintenance", e))?;
203        let certificates = match maintenance.postgresql().tls() {
204            SessionTls::System => Vec::new(),
205            SessionTls::CustomCa { certificates } => certificates.clone(),
206        };
207        ahri_tre_pgmeta::PgDatastoreCreationConnection::new(
208            maintenance.postgresql().host(),
209            maintenance.postgresql().port(),
210            maintenance.postgresql().bootstrap_database(),
211            maintenance.administrator().username(),
212            certificates,
213        )
214        .dataset_maintenance(maintenance.postgresql().database(), credential.material())
215        .map_err(|e| infrastructure("open Dataset maintenance capability", e))
216    }
217
218    fn resolve_catalog_credential(
219        &self,
220        binding: &DatastoreIdentityBinding,
221        reference: &ManagedSecretReference,
222    ) -> Result<ResolvedManagedSecret, AppError> {
223        self.resolve_datastore_managed_for_new_operation(binding, reference)
224            .map_err(|error| infrastructure("resolve Datastore catalog credential", error))
225    }
226
227    fn resolve_storage_credential(
228        &self,
229        plan: &DatastoreSessionPlan,
230        configured: &ConfiguredSecret,
231    ) -> Result<ResolvedOperationSecret, AppError> {
232        self.resolve_storage_for_new_operation(plan, configured)
233            .map_err(|error| infrastructure("resolve object-storage credential", error))
234    }
235    fn retain_session_reference_snapshot(
236        &self,
237        reference: &ConfiguredSessionReference,
238        resolved_versions: &[ResolvedSecretVersion],
239    ) -> Result<(), AppError> {
240        self.retain_session_reference_snapshot(
241            reference.session_id(),
242            reference.release_capability(),
243            &reference.authentication,
244            resolved_versions,
245        )
246        .map_err(|error| infrastructure("retain Session Secret references", error))
247    }
248}
249
250impl std::fmt::Debug for OpenConfiguredDatastoreRequest {
251    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
252        formatter
253            .debug_struct("OpenConfiguredDatastoreRequest")
254            .field("plan", &self.plan)
255            .field("authentication", &"OAuthSession(<redacted>)")
256            .finish()
257    }
258}
259
260/// Opens one configured Datastore from resolved runtime authority without
261/// requiring an otherwise unrelated application repository set.
262pub fn open_configured_datastore<R: DatastoreSecretResolver>(
263    request: OpenConfiguredDatastoreRequest,
264    resolver: &R,
265) -> Result<DataStoreSession, AppError> {
266    open_configured_datastore_for_session(request, resolver)
267        .map_err(ConfiguredDatastoreOpenError::into_app_error)
268}
269
270/// Authenticates PostgreSQL and validates its binding without opening Lake.
271fn open_configured_metadata(
272    request: &OpenConfiguredDatastoreRequest,
273) -> Result<
274    (
275        ahri_tre_libpq_oauth::LibpqOAuthConnection,
276        DatastoreIdentityBinding,
277    ),
278    ConfiguredDatastoreOpenError,
279> {
280    open_configured_metadata_observed(request, None)
281}
282
283fn open_configured_metadata_observed(
284    request: &OpenConfiguredDatastoreRequest,
285    context: Option<&ahri_tre_observability::CorrelationContext>,
286) -> Result<
287    (
288        ahri_tre_libpq_oauth::LibpqOAuthConnection,
289        DatastoreIdentityBinding,
290    ),
291    ConfiguredDatastoreOpenError,
292> {
293    validate_authentication_identity(
294        request.session_reference.authentication.identity(),
295        &request.authentication,
296    )?;
297    validate_authentication(&request.plan, &request.authentication)?;
298    let tls = match request.plan.postgresql().tls() {
299        SessionTls::System => LibpqOAuthTls::System,
300        SessionTls::CustomCa { certificates } => LibpqOAuthTls::CustomCa {
301            certificates: certificates.clone(),
302        },
303    };
304    let postgresql = request.plan.postgresql();
305    let authentication = request.plan.authentication();
306    let config = ExplicitLibpqOAuthConnectConfig {
307        host: postgresql.host().to_string(),
308        routing_address: postgresql.routing_address().map(str::to_string),
309        port: postgresql.port(),
310        database: postgresql.database().to_string(),
311        tls,
312        oauth_issuer: authentication.issuer().to_string(),
313        oauth_client_id: authentication.client_id().to_string(),
314        oauth_scope: Some(authentication.scopes().join(" ")),
315    };
316    let mut store = match context {
317        Some(context) => LibpqOAuthConnector.open_explicit_session_observed(
318            context,
319            &request.authentication,
320            &config,
321        ),
322        None => LibpqOAuthConnector.open_explicit_session(&request.authentication, &config),
323    }
324    .map_err(configured_postgresql_open_error)?;
325    store
326        .health_check()
327        .map_err(|error| infrastructure("configured PostgreSQL health check failed", error))?;
328    let binding = resolve_binding(&mut store, &request.plan)?;
329    Ok((store, binding))
330}
331
332/// Opens only PostgreSQL under current user authority; never attaches Lake or resolves its Secrets.
333pub fn open_operation_control(
334    request: OpenConfiguredDatastoreRequest,
335    deployment_id: uuid::Uuid,
336) -> Result<crate::service::OperationControl, ConfiguredDatastoreOpenError> {
337    let (store, binding) = open_configured_metadata(&request)?;
338    let repository = ahri_tre_pgmeta::PgMetadataRepository::operation_control(store)
339        .map_err(|_| AppError::Infrastructure("Operation inspection is unavailable".into()))?;
340    Ok(crate::service::OperationControl::new(
341        repository,
342        deployment_id,
343        binding.datastore_id,
344    ))
345}
346
347/// Opens a Datastore while preserving OAuth authentication rejection.
348pub fn open_configured_datastore_for_session<R: DatastoreSecretResolver>(
349    request: OpenConfiguredDatastoreRequest,
350    resolver: &R,
351) -> Result<DataStoreSession, ConfiguredDatastoreOpenError> {
352    open_configured_datastore_for_session_observed(None, request, resolver)
353}
354
355/// Times the app workflow using its already authorized inputs and resolver.
356pub fn open_configured_datastore_for_session_observed<R: DatastoreSecretResolver>(
357    context: Option<&ahri_tre_observability::CorrelationContext>,
358    request: OpenConfiguredDatastoreRequest,
359    resolver: &R,
360) -> Result<DataStoreSession, ConfiguredDatastoreOpenError> {
361    open_configured_datastore_with_observations(
362        context,
363        request,
364        resolver,
365        crate::SessionObservations::default(),
366        ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::SessionOpen,
367    )
368}
369
370/// Runs the existing open workflow while retaining its actual adapter evidence.
371pub fn open_configured_datastore_with_observations<R: DatastoreSecretResolver>(
372    context: Option<&ahri_tre_observability::CorrelationContext>,
373    request: OpenConfiguredDatastoreRequest,
374    resolver: &R,
375    observations: crate::SessionObservations,
376    origin: ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin,
377) -> Result<DataStoreSession, ConfiguredDatastoreOpenError> {
378    use ahri_tre_observability::{FailureCategory, Outcome, Stage};
379    let span = context.map(|context| context.span(Stage::Application));
380    let child = span.as_ref().map(|span| span.context());
381    let result = open_configured_session(child.as_ref(), request, resolver, observations, origin);
382    if let Some(span) = span {
383        let (outcome, category) = match &result {
384            Ok(_) => (Outcome::Success, None),
385            Err(ConfiguredDatastoreOpenError::AuthenticationRejected) => {
386                (Outcome::Rejected, Some(FailureCategory::Authentication))
387            }
388            Err(ConfiguredDatastoreOpenError::Application(AppError::Validation(_))) => {
389                (Outcome::Rejected, Some(FailureCategory::Compatibility))
390            }
391            Err(_) => (Outcome::Unavailable, Some(FailureCategory::Dependency)),
392        };
393        span.finish(outcome, category);
394    }
395    result
396}
397
398fn open_configured_session<R: DatastoreSecretResolver>(
399    context: Option<&ahri_tre_observability::CorrelationContext>,
400    request: OpenConfiguredDatastoreRequest,
401    resolver: &R,
402    observations: crate::SessionObservations,
403    origin: ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin,
404) -> Result<DataStoreSession, ConfiguredDatastoreOpenError> {
405    let metadata = open_configured_metadata_observed(&request, context);
406    observations.record(
407        ahri_tre_protocol::diagnostics::DiagnosticComponent::Metadata,
408        origin,
409        metadata.is_ok(),
410    );
411    let (store, binding) = metadata?;
412    let postgresql = request.plan.postgresql();
413    let sslmode = "verify-full";
414    let lake_tls = match postgresql.tls() {
415        SessionTls::System => FilesystemCatalogTls::System,
416        SessionTls::CustomCa { certificates } => FilesystemCatalogTls::CustomCa {
417            certificates: certificates.clone(),
418        },
419    };
420
421    let catalog_schema = binding.ducklake_catalog_schema.as_deref().ok_or_else(|| {
422        AppError::Validation("persisted Datastore has no DuckLake catalog schema".to_string())
423    })?;
424    let catalog_role = binding.lake_catalog_role_name.as_deref().ok_or_else(|| {
425        AppError::Validation("persisted Datastore has no DuckLake catalog role".to_string())
426    })?;
427    let managed_reference = binding
428        .managed_secret_ref
429        .as_deref()
430        .ok_or_else(|| AppError::Validation("Datastore has no Managed catalog credential".into()))?
431        .parse::<ManagedSecretReference>()
432        .map_err(|_| {
433            AppError::Validation("Datastore catalog credential reference is invalid".into())
434        })?;
435    let expected_reference = format!(
436        "managed://datastore/{}/ducklake-password",
437        binding.datastore_id
438    );
439    if managed_reference.as_secret_reference().as_str() != expected_reference {
440        return Err(AppError::Validation(
441            "Datastore catalog credential reference is not binding-owned".into(),
442        )
443        .into());
444    }
445    let resolved = resolver.resolve_catalog_credential(&binding, &managed_reference)?;
446    if resolved.metadata().reference() != &managed_reference
447        || u64::try_from(binding.credential_version).ok()
448            != Some(resolved.metadata().version().get())
449    {
450        return Err(AppError::Validation(
451            "resolved DuckLake credential does not match the persisted Datastore identity".into(),
452        )
453        .into());
454    }
455    let lake_credential_version =
456        ResolvedSecretVersion::managed(managed_reference, resolved.metadata().version().get());
457    let mut catalog_observations = ahri_tre_lake::LakeCatalogObservations::new(context);
458    let lake = open_configured_lake(
459        context,
460        resolver,
461        &request.plan,
462        &binding,
463        catalog_schema,
464        catalog_role,
465        lake_tls,
466        resolved.material(),
467        &mut catalog_observations,
468    );
469    observations.record_catalog(&catalog_observations, origin);
470    observations.record(
471        ahri_tre_protocol::diagnostics::DiagnosticComponent::Lake,
472        origin,
473        lake.is_ok(),
474    );
475    let (mut opened, scratch_attempt, mut storage_versions, canonical_data_path) = lake?;
476    storage_versions.insert(0, lake_credential_version);
477
478    let mut dataset_executor = request.plan.dataset_executor_root().and_then(|root| {
479        ahri_tre_lake::DatasetExecutor::acquire(std::path::Path::new(root), binding.datastore_id)
480            .ok()
481    });
482
483    let mut dataset_maintenance = None;
484    if let (Some(executor), Some(scratch)) = (&dataset_executor, &scratch_attempt) {
485        match resolver.dataset_maintenance(&request.plan, &binding) {
486            Ok(maintenance) => {
487                let lake = DuckLakeAdapter::new(&canonical_data_path);
488                if reconcile_dataset_attempts(
489                    &maintenance,
490                    executor,
491                    scratch,
492                    &mut opened.connection,
493                    &lake,
494                )
495                .is_err()
496                {
497                    dataset_executor = None;
498                } else {
499                    dataset_maintenance = Some(maintenance);
500                }
501            }
502            Err(_) => dataset_executor = None,
503        }
504    }
505
506    let profile = DataStoreProfile {
507        server: postgresql.host().to_string(),
508        port: postgresql.port(),
509        dbname: postgresql.database().to_string(),
510        sslmode: Some(sslmode.to_string()),
511        lake: LakeAttachProfile {
512            lake_data: canonical_data_path,
513            lake_db: binding.ducklake_catalog_database.clone(),
514            catalog_schema: Some(catalog_schema.to_string()),
515            auth: LakeAuthProfile {
516                lake_user: catalog_role.to_string(),
517                encryption_mode: if binding.ducklake_encryption {
518                    EncryptionMode::ServerManaged
519                } else {
520                    EncryptionMode::None
521                },
522            },
523        },
524    };
525    let capture = ahri_tre_lake::DisclosureCaptureAuthority::from_opened(&mut opened).ok();
526    let mut runtime = DataStoreRuntime::new(
527        profile,
528        PgStoreConnection {
529            connection_description: store.connection_description().to_string(),
530        },
531        DuckLakeConnection {
532            alias: opened.attach_plan.alias,
533            data_path: opened.local_layout.data_path,
534            staging_path: opened.local_layout.staging_path,
535            tmp_path: opened.local_layout.tmp_path,
536            test_runs_path: opened.local_layout.test_runs_path,
537            catalog_database: binding.ducklake_catalog_database.clone(),
538            catalog_schema: Some(catalog_schema.to_string()),
539            attach_description: opened.attach_plan.attach_description,
540            detected_encryption_mode: opened.health.detected_encryption_mode,
541            automatic_migration: opened.attach_plan.automatic_migration,
542            snapshot_count: opened.health.snapshot_count,
543            table_count: opened.health.table_count,
544            catalog_type: opened.health.catalog_type,
545            extension_version: opened.health.extension_version,
546        },
547    )
548    .with_identity_binding(binding)
549    .with_session_provenance(request.plan.provenance().clone(), storage_versions);
550    runtime.disclosure = request.plan.disclosure().clone();
551    if request.session_reference.retains_snapshot() {
552        resolver.retain_session_reference_snapshot(
553            &request.session_reference,
554            &runtime.resolved_secret_versions,
555        )?;
556    }
557    let session = DataStoreSession::new(
558        runtime,
559        StoreSessionConnection::OAuth(Arc::new(Mutex::new(store))),
560        opened.connection,
561    )
562    .with_authenticated_actor(request.authentication.role.clone())
563    .with_disclosure_context(
564        request.session_reference.session_id,
565        request.session_reference.authentication.clone(),
566    )
567    .with_observations(observations)
568    .with_lake_tls_guard(opened.tls_guard)
569    .with_scratch_attempt(scratch_attempt)
570    .with_dataset_maintenance(dataset_maintenance);
571    let mut session = match dataset_executor {
572        Some(executor) => session.with_dataset_executor(executor),
573        None => session,
574    };
575    session.disclosure_capture = capture;
576    // Older pending evidence must not impose a Datastore-wide admission pause.
577    // A new disclosure still persists and verifies its own required evidence.
578    let _ = session.reconcile_disclosures();
579    Ok(session)
580}
581
582/// Startup maintenance for a configured Datastore, before the listener schedules
583/// user work. Missing/unavailable Datastores are retried on explicit Session open.
584/// No Session, user identity, Current study, or output admission is created here.
585pub fn reconcile_configured_dataset_operations(
586    runtime: &ahri_tre_runtime::TrustedRuntimeState,
587    plan: &DatastoreSessionPlan,
588) -> Result<bool, AppError> {
589    let fail = || AppError::Infrastructure("Configured Dataset recovery is unavailable".into());
590    let configured = plan.datastore_id().parse().map_err(|_| fail())?;
591    let creation = runtime
592        .datastore_creation_plan(&configured)
593        .map_err(|_| fail())?;
594    let credential = runtime
595        .resolve_datastore_administrator(&creation)
596        .map_err(|_| fail())?;
597    let certificates = match creation.postgresql().tls() {
598        SessionTls::System => Vec::new(),
599        SessionTls::CustomCa { certificates } => certificates.clone(),
600    };
601    let connection = ahri_tre_pgmeta::PgDatastoreCreationConnection::new(
602        creation.postgresql().host(),
603        creation.postgresql().port(),
604        creation.postgresql().bootstrap_database(),
605        creation.administrator().username(),
606        certificates,
607    );
608    let Some(binding) = ahri_tre_pgmeta::find_existing_datastore_binding(
609        &connection,
610        credential.material(),
611        creation.postgresql().database(),
612    )
613    .map_err(|_| fail())?
614    else {
615        return Ok(false);
616    };
617    validate_binding(&binding, plan)?;
618    let executor = ahri_tre_lake::DatasetExecutor::acquire(
619        std::path::Path::new(plan.dataset_executor_root().ok_or_else(fail)?),
620        binding.datastore_id,
621    )
622    .map_err(|_| fail())?;
623    let maintenance = connection
624        .dataset_maintenance(creation.postgresql().database(), credential.material())
625        .map_err(|_| fail())?;
626    // Record interruptions even if Lake recovery is temporarily unavailable.
627    maintenance
628        .reconcile_operations(executor.clone(), super::service::interrupt_operation)
629        .map_err(|_| fail())?;
630    let reference = format!(
631        "managed://datastore/{}/ducklake-password",
632        binding.datastore_id
633    )
634    .parse()
635    .map_err(|_| fail())?;
636    let catalog = runtime.resolve_catalog_credential(&binding, &reference)?;
637    if catalog.metadata().version().get()
638        != u64::try_from(binding.credential_version).map_err(|_| fail())?
639    {
640        return Err(fail());
641    }
642    let tls = match plan.postgresql().tls() {
643        SessionTls::System => FilesystemCatalogTls::System,
644        SessionTls::CustomCa { certificates } => FilesystemCatalogTls::CustomCa {
645            certificates: certificates.clone(),
646        },
647    };
648    let (mut lake, scratch, _, path) = open_configured_lake(
649        None,
650        runtime,
651        plan,
652        &binding,
653        binding
654            .ducklake_catalog_schema
655            .as_deref()
656            .ok_or_else(fail)?,
657        binding.lake_catalog_role_name.as_deref().ok_or_else(fail)?,
658        tls,
659        catalog.material(),
660        &mut ahri_tre_lake::LakeCatalogObservations::default(),
661    )?;
662    reconcile_dataset_attempts(
663        &maintenance,
664        &executor,
665        &scratch.ok_or_else(fail)?,
666        &mut lake.connection,
667        &DuckLakeAdapter::new(&path),
668    )
669    .map(|()| true)
670}
671
672/// Reconcile only stopped, durably recorded attempts. One blocked cleanup keeps
673/// its reservation; it does not prevent independent output identities proceeding.
674pub(crate) fn reconcile_dataset_attempts(
675    maintenance: &ahri_tre_pgmeta::PgDatasetMaintenance,
676    executor: &Arc<ahri_tre_lake::DatasetExecutor>,
677    scratch: &ahri_tre_lake::ScratchAttempt,
678    connection: &mut duckdb::Connection,
679    lake: &DuckLakeAdapter,
680) -> Result<(), AppError> {
681    maintenance
682        .reconcile_operations(executor.clone(), super::service::interrupt_operation)
683        .map_err(|e| infrastructure("reconcile interrupted operations", e))?;
684    let attempts = maintenance
685        .recoverable_attempts(executor.clone())
686        .map_err(|e| infrastructure("inspect Dataset recovery", e))?;
687    for attempt in attempts {
688        let cleanup = || -> Result<(), AppError> {
689            let mut recovery = maintenance
690                .claim(executor.clone(), attempt.attempt_id)
691                .map_err(|e| infrastructure("claim Dataset recovery", e))?;
692            scratch
693                .cleanup_dataset_attempt(&mut recovery)
694                .map_err(|e| infrastructure("clean interrupted Dataset source", e))?;
695            if recovery.requires_output_cleanup() {
696                lake.cleanup_reserved_dataset_output(connection, &mut recovery)
697                    .map_err(|e| infrastructure("clean interrupted Dataset output", e))?;
698            }
699            recovery
700                .release_after_cleanup()
701                .map_err(|e| infrastructure("release interrupted Dataset output", e))
702        };
703        // Ownership remains durable after any failed or concurrently claimed cleanup.
704        let recovery = ahri_tre_observability::LifecycleSpan::start_operation(
705            ahri_tre_observability::Phase::RecoveryCleanup,
706            None,
707            attempt.operation_id,
708        );
709        let result = cleanup();
710        recovery.complete(&result, ahri_tre_observability::FailureCategory::Cleanup);
711    }
712    Ok(())
713}
714
715fn validate_authentication(
716    plan: &DatastoreSessionPlan,
717    session: &OAuthSession,
718) -> Result<(), AppError> {
719    if session.issuer != plan.authentication().issuer()
720        || session.client_id != plan.authentication().client_id()
721    {
722        return Err(AppError::Validation(
723            "Datastore Session credential does not match the selected authentication policy"
724                .to_string(),
725        ));
726    }
727    if session.expires_at <= chrono::Utc::now() {
728        return Err(AppError::Validation(
729            "Datastore Session credential has expired".to_string(),
730        ));
731    }
732    Ok(())
733}
734
735fn validate_authentication_identity(
736    identity: &ahri_tre_runtime::AuthenticationIdentityBinding,
737    session: &OAuthSession,
738) -> Result<(), AppError> {
739    if session.issuer != identity.issuer()
740        || session.client_id != identity.client_id()
741        || session.subject != identity.subject()
742    {
743        return Err(AppError::Validation(
744            "Datastore Session credential does not match the persisted Session identity"
745                .to_string(),
746        ));
747    }
748    Ok(())
749}
750
751fn reap_session_scratch(scratch: &TrustedScratch) {
752    if scratch.reap_abandoned_sessions().is_err() {
753        tracing::warn!("Abandoned Session scratch cleanup requires a later retry");
754    }
755}
756
757#[allow(clippy::too_many_arguments)]
758fn open_configured_lake<R: DatastoreSecretResolver>(
759    context: Option<&ahri_tre_observability::CorrelationContext>,
760    resolver: &R,
761    plan: &DatastoreSessionPlan,
762    binding: &DatastoreIdentityBinding,
763    catalog_schema: &str,
764    catalog_role: &str,
765    catalog_tls: FilesystemCatalogTls,
766    catalog_credential: &ahri_tre_secrets::SecretMaterial,
767    catalog_observations: &mut ahri_tre_lake::LakeCatalogObservations,
768) -> Result<
769    (
770        ahri_tre_lake::DuckLakeOpenedCatalog,
771        Option<ahri_tre_lake::ScratchAttempt>,
772        Vec<ResolvedSecretVersion>,
773        String,
774    ),
775    AppError,
776> {
777    let lake_plan = plan.lake();
778    let encryption = if binding.ducklake_encryption {
779        EncryptionMode::ServerManaged
780    } else {
781        EncryptionMode::None
782    };
783    if let LakeStoragePlan::Filesystem { .. } = lake_plan.storage() {
784        let config = FilesystemLakeSessionConfig::from_binding(
785            lake_plan
786                .filesystem_base()
787                .expect("filesystem variant has a base"),
788            &binding.lake_path,
789            plan.postgresql().host(),
790            plan.postgresql().port(),
791            &binding.ducklake_catalog_database,
792            catalog_schema,
793            catalog_role,
794            catalog_tls,
795            encryption,
796        )
797        .map_err(|error| infrastructure("configured filesystem Lake inputs are invalid", error))?;
798        validate_configured_physical_intent(lake_plan, &config, binding)?;
799        let data_path = config.canonical_data_path();
800        let scratch = TrustedScratch::open(
801            std::path::Path::new(plan.configuration_scratch_root()),
802            [std::path::Path::new(&data_path)],
803        )
804        .map_err(|error| infrastructure("trusted scratch validation failed", error))?;
805        reap_session_scratch(&scratch);
806        let attempt = scratch
807            .create_session_attempt(
808                ScratchAttemptId::new(&uuid::Uuid::new_v4().simple().to_string())
809                    .expect("UUID simple form is a canonical scratch attempt identifier"),
810            )
811            .map_err(|error| infrastructure("trusted scratch attempt creation failed", error))?;
812        let opened = DuckLakeAdapter::open_existing_filesystem_catalog_observed(
813            context,
814            &config,
815            catalog_credential,
816            &attempt,
817            catalog_observations,
818        )
819        .map_err(|error| infrastructure("configured filesystem Lake open failed", error))?;
820        return Ok((opened, Some(attempt), Vec::new(), data_path));
821    }
822
823    let location = verified_object_location(plan, binding)?;
824    let storage_tls = match lake_plan.storage() {
825        LakeStoragePlan::S3Compatible { tls, .. } => match tls {
826            StorageTlsPlan::System => ObjectStorageTls::System,
827            StorageTlsPlan::CustomCa { certificates } => ObjectStorageTls::CustomCa {
828                certificates: certificates.clone(),
829            },
830        },
831        _ => ObjectStorageTls::System,
832    };
833    let config = ObjectLakeSessionConfig::new(
834        location,
835        plan.postgresql().host(),
836        plan.postgresql().port(),
837        &binding.ducklake_catalog_database,
838        catalog_schema,
839        catalog_role,
840        catalog_tls,
841        storage_tls,
842        encryption,
843    )
844    .map_err(|error| infrastructure("configured object-storage Lake inputs are invalid", error))?;
845    let scratch = TrustedScratch::open(
846        std::path::Path::new(plan.configuration_scratch_root()),
847        std::iter::empty::<&std::path::Path>(),
848    )
849    .map_err(|error| infrastructure("trusted scratch validation failed", error))?;
850    reap_session_scratch(&scratch);
851    let attempt = scratch
852        .create_session_attempt(
853            ScratchAttemptId::new(&uuid::Uuid::new_v4().simple().to_string())
854                .expect("UUID simple form is a canonical scratch attempt identifier"),
855        )
856        .map_err(|error| infrastructure("trusted scratch attempt creation failed", error))?;
857
858    let mut versions = Vec::new();
859    let opened = match lake_plan.storage().authentication().ok_or_else(|| {
860        AppError::Validation("object-storage policy has no authentication capability".into())
861    })? {
862        StorageAuthenticationPlan::S3Static {
863            access_key_id,
864            secret_access_key,
865        } => {
866            let access = resolver.resolve_storage_credential(plan, access_key_id)?;
867            let secret = resolver.resolve_storage_credential(plan, secret_access_key)?;
868            versions.push(access.version().clone());
869            versions.push(secret.version().clone());
870            DuckLakeAdapter::open_existing_object_storage_catalog_observed(
871                context,
872                &config,
873                catalog_credential,
874                ObjectStorageAuthentication::S3Static {
875                    access_key_id: access.material(),
876                    secret_access_key: secret.material(),
877                },
878                &attempt,
879                catalog_observations,
880            )
881        }
882        StorageAuthenticationPlan::AwsWorkloadIdentity { source } => {
883            let source = match source {
884                ConfiguredAwsWorkloadIdentitySource::EcsTask => AwsWorkloadIdentitySource::EcsTask,
885                ConfiguredAwsWorkloadIdentitySource::Ec2Instance => {
886                    AwsWorkloadIdentitySource::Ec2Instance
887                }
888            };
889            DuckLakeAdapter::open_existing_object_storage_catalog_observed(
890                context,
891                &config,
892                catalog_credential,
893                ObjectStorageAuthentication::AwsWorkloadIdentity { source },
894                &attempt,
895                catalog_observations,
896            )
897        }
898        StorageAuthenticationPlan::AzureManagedIdentity { identity } => {
899            let client_id = match identity {
900                AzureManagedIdentity::SystemAssigned => None,
901                AzureManagedIdentity::UserAssigned { client_id } => Some(*client_id),
902            };
903            DuckLakeAdapter::open_existing_object_storage_catalog_observed(
904                context,
905                &config,
906                catalog_credential,
907                ObjectStorageAuthentication::AzureManagedIdentity { client_id },
908                &attempt,
909                catalog_observations,
910            )
911        }
912        StorageAuthenticationPlan::AzureServicePrincipal {
913            tenant_id,
914            client_id,
915            client_secret,
916        } => {
917            let secret = resolver.resolve_storage_credential(plan, client_secret)?;
918            versions.push(secret.version().clone());
919            DuckLakeAdapter::open_existing_object_storage_catalog_observed(
920                context,
921                &config,
922                catalog_credential,
923                ObjectStorageAuthentication::AzureServicePrincipal {
924                    tenant_id: *tenant_id,
925                    client_id: *client_id,
926                    client_secret: secret.material(),
927                },
928                &attempt,
929                catalog_observations,
930            )
931        }
932    }
933    .map_err(|error| infrastructure("configured object-storage Lake open failed", error))?;
934    let data_path = config.location().canonical_data_path();
935    Ok((opened, Some(attempt), versions, data_path))
936}
937
938fn verified_object_location(
939    plan: &DatastoreSessionPlan,
940    binding: &DatastoreIdentityBinding,
941) -> Result<ObjectStorageLocation, AppError> {
942    if binding.storage_policy_id.as_deref() != Some(plan.storage_policy_id()) {
943        return Err(AppError::Validation(
944            "persisted object-storage policy identity does not match configuration".into(),
945        ));
946    }
947    let persisted = binding.lake_location.as_ref().ok_or_else(|| {
948        AppError::Validation("persisted Datastore has no structured Lake namespace".into())
949    })?;
950    let configured_prefix = plan.lake().datastore_prefix();
951    let (location, base_prefix, actual_prefix) = match (plan.lake().storage(), persisted) {
952        (
953            LakeStoragePlan::AwsS3 {
954                region,
955                bucket,
956                base_prefix,
957                ..
958            },
959            DatastoreLakeLocation::AwsS3 {
960                region: saved_region,
961                bucket: saved_bucket,
962                prefix,
963            },
964        ) if region == saved_region && bucket == saved_bucket => (
965            ObjectStorageLocation::AwsS3 {
966                region: saved_region.clone(),
967                bucket: saved_bucket.clone(),
968                prefix: prefix.clone(),
969            },
970            base_prefix.as_deref(),
971            prefix.as_str(),
972        ),
973        (
974            LakeStoragePlan::S3Compatible {
975                endpoint,
976                url_style,
977                bucket,
978                base_prefix,
979                ..
980            },
981            DatastoreLakeLocation::S3Compatible {
982                endpoint: saved_endpoint,
983                url_style: saved_style,
984                bucket: saved_bucket,
985                prefix,
986            },
987        ) if endpoint == saved_endpoint
988            && bucket == saved_bucket
989            && saved_style
990                == match url_style {
991                    ConfiguredS3UrlStyle::Path => "path",
992                    ConfiguredS3UrlStyle::VirtualHosted => "virtual_hosted",
993                } =>
994        {
995            (
996                ObjectStorageLocation::S3Compatible {
997                    endpoint: saved_endpoint.clone(),
998                    url_style: match url_style {
999                        ConfiguredS3UrlStyle::Path => S3UrlStyle::Path,
1000                        ConfiguredS3UrlStyle::VirtualHosted => S3UrlStyle::VirtualHosted,
1001                    },
1002                    bucket: saved_bucket.clone(),
1003                    prefix: prefix.clone(),
1004                },
1005                base_prefix.as_deref(),
1006                prefix.as_str(),
1007            )
1008        }
1009        (
1010            LakeStoragePlan::AzureBlob {
1011                account_endpoint,
1012                container,
1013                base_prefix,
1014                ..
1015            },
1016            DatastoreLakeLocation::AzureBlob {
1017                account_endpoint: saved_endpoint,
1018                container: saved_container,
1019                prefix,
1020            },
1021        ) if account_endpoint == saved_endpoint && container == saved_container => (
1022            ObjectStorageLocation::AzureBlob {
1023                account_endpoint: saved_endpoint.clone(),
1024                container: saved_container.clone(),
1025                prefix: prefix.clone(),
1026            },
1027            base_prefix.as_deref(),
1028            prefix.as_str(),
1029        ),
1030        _ => {
1031            return Err(AppError::Validation(
1032                "persisted object-storage namespace does not match configuration".into(),
1033            ));
1034        }
1035    };
1036    if !prefix_matches_intent(
1037        base_prefix,
1038        configured_prefix,
1039        binding.datastore_id,
1040        actual_prefix,
1041    ) || binding.lake_path != location.canonical_data_path()
1042    {
1043        return Err(AppError::Validation(
1044            "persisted object-storage namespace or DuckLake data path does not match configuration"
1045                .into(),
1046        ));
1047    }
1048    Ok(location)
1049}
1050
1051fn prefix_matches_intent(
1052    base: Option<&str>,
1053    datastore: Option<&str>,
1054    datastore_id: uuid::Uuid,
1055    actual: &str,
1056) -> bool {
1057    let generated;
1058    let datastore = match datastore {
1059        Some(configured) => configured,
1060        None => {
1061            generated = format!("datastores/{}", datastore_id.simple());
1062            &generated
1063        }
1064    };
1065    match base {
1066        Some(base) => actual == format!("{base}/{datastore}"),
1067        None => actual == datastore,
1068    }
1069}
1070
1071fn resolve_binding(
1072    store: &mut ahri_tre_libpq_oauth::LibpqOAuthConnection,
1073    plan: &DatastoreSessionPlan,
1074) -> Result<DatastoreIdentityBinding, AppError> {
1075    if !datastore_identity_table_exists_oauth(store)
1076        .map_err(|error| infrastructure("check Datastore identity table", error))?
1077    {
1078        return Err(AppError::Validation(
1079            "configured Datastore has no identity binding".to_string(),
1080        ));
1081    }
1082    let bindings = read_datastore_identity_bindings_oauth(store)
1083        .map_err(|error| infrastructure("read Datastore identity binding", error))?;
1084    let [binding] = bindings.as_slice() else {
1085        return Err(AppError::Validation(
1086            "configured Datastore must have exactly one identity binding".to_string(),
1087        ));
1088    };
1089    validate_binding(binding, plan)?;
1090    Ok(binding.clone())
1091}
1092
1093fn validate_binding(
1094    binding: &DatastoreIdentityBinding,
1095    plan: &DatastoreSessionPlan,
1096) -> Result<(), AppError> {
1097    if binding.lifecycle_state != DatastoreLifecycleState::Ready {
1098        return Err(AppError::Validation(
1099            "configured Datastore identity binding is not ready".to_string(),
1100        ));
1101    }
1102    if binding.datastore_name != plan.postgresql().database() {
1103        return Err(AppError::Validation(
1104            "configured Datastore database does not match the persisted Datastore identity".into(),
1105        ));
1106    }
1107    if binding.binding_fingerprint != binding.computed_binding_fingerprint() {
1108        return Err(AppError::Validation(
1109            "persisted Datastore identity fingerprint is invalid".to_string(),
1110        ));
1111    }
1112    Ok(())
1113}
1114
1115fn validate_configured_physical_intent(
1116    configured: &ahri_tre_runtime::FilesystemLakeSessionPlan,
1117    resolved: &FilesystemLakeSessionConfig,
1118    binding: &DatastoreIdentityBinding,
1119) -> Result<(), AppError> {
1120    let mismatched = configured
1121        .datastore_prefix()
1122        .is_some_and(|value| value != resolved.datastore_prefix())
1123        || configured
1124            .catalog_database()
1125            .is_some_and(|value| value != binding.ducklake_catalog_database)
1126        || configured
1127            .catalog_schema()
1128            .is_some_and(|value| binding.ducklake_catalog_schema.as_deref() != Some(value))
1129        || configured
1130            .catalog_role()
1131            .is_some_and(|value| binding.lake_catalog_role_name.as_deref() != Some(value));
1132    if mismatched {
1133        Err(AppError::Validation(
1134            "configured physical Datastore intent does not match its persisted binding".into(),
1135        ))
1136    } else {
1137        Ok(())
1138    }
1139}
1140
1141fn infrastructure(context: &str, error: impl std::fmt::Display) -> AppError {
1142    AppError::Infrastructure(format!("{context}: {error}"))
1143}
1144
1145fn configured_postgresql_open_error(
1146    error: ahri_tre_libpq_oauth::LibpqOAuthError,
1147) -> ConfiguredDatastoreOpenError {
1148    match error.connect_failure_kind() {
1149        LibpqOAuthConnectFailureKind::AuthenticationRejected => {
1150            ConfiguredDatastoreOpenError::AuthenticationRejected
1151        }
1152        _ => ConfiguredDatastoreOpenError::Application(infrastructure(
1153            "configured PostgreSQL open failed",
1154            error,
1155        )),
1156    }
1157}
1158
1159#[cfg(test)]
1160mod tests {
1161    use ahri_tre_config::{ConfigurationAuthority, EffectiveConfigurationTarget};
1162    use ahri_tre_orcid::TokenSource;
1163    use ahri_tre_runtime::DatastoreSessionPlan;
1164    use ahri_tre_types::{
1165        DatastoreIdentityBinding, DatastoreLakeLocation, DatastoreLifecycleState,
1166        LakeCatalogCredentialMode,
1167    };
1168    use chrono::{Duration, Utc};
1169
1170    #[test]
1171    fn configured_postgresql_open_preserves_oauth_rejection_as_authorization() {
1172        let rejected = configured_postgresql_open_error(
1173            ahri_tre_libpq_oauth::LibpqOAuthError::AuthenticationRejected {
1174                token_source: "id_token",
1175            },
1176        );
1177        assert!(matches!(
1178            rejected,
1179            ConfiguredDatastoreOpenError::AuthenticationRejected
1180        ));
1181
1182        let unavailable = configured_postgresql_open_error(
1183            ahri_tre_libpq_oauth::LibpqOAuthError::ConnectFailed {
1184                token_source: "id_token",
1185                message: "connection refused".to_string(),
1186            },
1187        );
1188        assert!(matches!(
1189            unavailable,
1190            ConfiguredDatastoreOpenError::Application(AppError::Infrastructure(_))
1191        ));
1192    }
1193    use secrecy::SecretString;
1194
1195    use super::{
1196        AppError, ConfiguredDatastoreOpenError, ConfiguredSessionReference,
1197        configured_postgresql_open_error, prefix_matches_intent, validate_authentication,
1198        validate_authentication_identity, verified_object_location,
1199    };
1200
1201    #[test]
1202    fn live_server_session_reference_does_not_request_durable_secret_retention() {
1203        let authentication: ahri_tre_runtime::SessionAuthenticationRecord =
1204            serde_json::from_value(serde_json::json!({
1205                "owner": "11111111-1111-4111-8111-111111111111",
1206                "identity": {
1207                    "issuer": "https://orcid.org",
1208                    "client_id": "ahri-tre-example",
1209                    "subject": "0000-0001-2345-6789"
1210                },
1211                "oauth": {
1212                    "artifact_id": "22222222-2222-4222-8222-222222222222",
1213                    "kind": "o_auth",
1214                    "version": 1,
1215                    "expires_at": "2030-01-01T00:00:00Z"
1216                }
1217            }))
1218            .unwrap();
1219
1220        let (reference, identifier) = ConfiguredSessionReference::live(authentication);
1221
1222        assert!(!reference.retains_snapshot());
1223        assert_eq!(reference.session_id(), identifier.session_id());
1224        assert!(reference.release_capability.is_none());
1225    }
1226
1227    #[test]
1228    fn configured_open_rejects_expired_material_at_the_authentication_seam() {
1229        let effective = ConfigurationAuthority::parse_application(include_str!(
1230            "../../../examples/config/application-v1.toml"
1231        ))
1232        .unwrap()
1233        .resolve_effective(EffectiveConfigurationTarget::ExecutionProfile { explicit: None })
1234        .unwrap();
1235        let plan = DatastoreSessionPlan::try_from(&effective).unwrap();
1236        let authentication = ahri_tre_orcid::OAuthSession {
1237            id_token: SecretString::from("postgres-token-canary".to_string()),
1238            issuer: "https://orcid.org".to_string(),
1239            client_id: "ahri-tre-example".to_string(),
1240            subject: "subject".to_string(),
1241            orcid: "0000-0000-0000-0001".to_string(),
1242            role: "orcid_0000000000000001".to_string(),
1243            expires_at: Utc::now() - Duration::seconds(1),
1244            token_source: TokenSource::JupyterHub,
1245        };
1246
1247        let rendered = format!("{authentication:?}");
1248        assert!(!rendered.contains("postgres-token-canary"));
1249        assert!(
1250            validate_authentication(&plan, &authentication)
1251                .expect_err("expired authentication must be rejected")
1252                .to_string()
1253                .contains("has expired")
1254        );
1255    }
1256
1257    #[test]
1258    fn configured_open_rejects_an_oauth_subject_different_from_the_journal_identity() {
1259        let identity: ahri_tre_runtime::AuthenticationIdentityBinding =
1260            serde_json::from_value(serde_json::json!({
1261                "issuer": "https://orcid.org",
1262                "client_id": "ahri-tre-example",
1263                "subject": "subject-a"
1264            }))
1265            .unwrap();
1266        let authentication = ahri_tre_orcid::OAuthSession {
1267            id_token: SecretString::from("postgres-token-canary".to_string()),
1268            issuer: "https://orcid.org".to_string(),
1269            client_id: "ahri-tre-example".to_string(),
1270            subject: "subject-b".to_string(),
1271            orcid: "0000-0000-0000-0001".to_string(),
1272            role: "orcid_0000000000000001".to_string(),
1273            expires_at: Utc::now() + Duration::minutes(5),
1274            token_source: TokenSource::JupyterHub,
1275        };
1276
1277        let error = validate_authentication_identity(&identity, &authentication)
1278            .expect_err("different OAuth subjects must be rejected before opening");
1279        assert!(error.to_string().contains("persisted Session identity"));
1280    }
1281
1282    #[test]
1283    fn configured_object_open_requires_the_exact_persisted_service_authority() {
1284        let document = include_str!("../../../examples/config/application-v1.toml").replace(
1285            "[storage_policies.filesystem]\nkind = \"filesystem\"\nbase_path = \"/srv/ahri-tre/lake\"",
1286            "[storage_authentication_policies.object]\n\
1287             kind = \"s3_static\"\n\
1288             access_key_id = { uri = \"managed://storage/access-key\" }\n\
1289             secret_access_key = { uri = \"managed://storage/secret-key\" }\n\n\
1290             [storage_policies.filesystem]\n\
1291             kind = \"s3_compatible\"\n\
1292             endpoint = \"https://rustfs.example.test\"\n\
1293             url_style = \"path\"\n\
1294             bucket = \"ahri-research\"\n\
1295             base_prefix = \"governed\"\n\
1296             authentication_policy = \"object\"",
1297        );
1298        let effective = ConfigurationAuthority::parse_application(&document)
1299            .unwrap()
1300            .resolve_effective(EffectiveConfigurationTarget::ExecutionProfile { explicit: None })
1301            .unwrap();
1302        let plan = DatastoreSessionPlan::try_from(&effective).unwrap();
1303        let mut binding = DatastoreIdentityBinding {
1304            datastore_id: uuid::Uuid::new_v4(),
1305            datastore_name: "research".into(),
1306            ducklake_catalog_database: "ducklake_catalog".into(),
1307            ducklake_catalog_schema: Some("catalog".into()),
1308            lake_path: "s3://ahri-research/governed/datasets/".into(),
1309            storage_policy_id: Some("filesystem".into()),
1310            lake_location: Some(DatastoreLakeLocation::S3Compatible {
1311                endpoint: "https://rustfs.example.test".into(),
1312                url_style: "path".into(),
1313                bucket: "ahri-research".into(),
1314                prefix: "governed/datasets".into(),
1315            }),
1316            ducklake_encryption: false,
1317            lake_catalog_credential_mode: LakeCatalogCredentialMode::ManagedLocal,
1318            lake_catalog_role_name: Some("catalog_owner".into()),
1319            managed_secret_ref: None,
1320            credential_version: 1,
1321            credential_last_rotated_at: None,
1322            created_at: Utc::now(),
1323            creating_tool_version: None,
1324            binding_fingerprint: String::new(),
1325            lifecycle_state: DatastoreLifecycleState::Ready,
1326            failure_code: None,
1327            failure_summary: None,
1328        };
1329        binding.binding_fingerprint = binding.computed_binding_fingerprint();
1330
1331        let location = verified_object_location(&plan, &binding).unwrap();
1332        assert_eq!(
1333            location.canonical_data_path(),
1334            "s3://ahri-research/governed/datasets/"
1335        );
1336
1337        let Some(DatastoreLakeLocation::S3Compatible { endpoint, .. }) =
1338            binding.lake_location.as_mut()
1339        else {
1340            panic!("test binding must use S3-compatible storage");
1341        };
1342        *endpoint = "https://other.example.test".into();
1343        assert!(
1344            verified_object_location(&plan, &binding)
1345                .unwrap_err()
1346                .to_string()
1347                .contains("namespace does not match")
1348        );
1349    }
1350
1351    #[test]
1352    fn omitted_object_prefix_requires_the_uuid_derived_namespace() {
1353        let datastore_id = uuid::Uuid::parse_str("10000000-0000-4000-8000-000000000001").unwrap();
1354        assert!(prefix_matches_intent(
1355            Some("governed"),
1356            None,
1357            datastore_id,
1358            "governed/datastores/10000000000040008000000000000001"
1359        ));
1360        assert!(!prefix_matches_intent(
1361            Some("governed"),
1362            None,
1363            datastore_id,
1364            "governed/arbitrary"
1365        ));
1366    }
1367}