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
31pub struct OpenConfiguredDatastoreRequest {
33 pub plan: DatastoreSessionPlan,
34 pub authentication: OAuthSession,
35 pub session_reference: ConfiguredSessionReference,
36}
37
38#[derive(Debug, thiserror::Error)]
40#[non_exhaustive]
41pub enum ConfiguredDatastoreOpenError {
42 #[error("configured PostgreSQL rejected OAuth authentication")]
44 AuthenticationRejected,
45 #[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
67pub struct LiveSessionIdentifier(uuid::Uuid);
69
70impl LiveSessionIdentifier {
71 pub fn session_id(&self) -> uuid::Uuid {
72 self.0
73 }
74}
75
76impl ConfiguredSessionReference {
77 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 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
139pub trait DatastoreSecretResolver {
141 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
260pub 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
270fn 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
332pub 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
347pub 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
355pub 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
370pub 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 let _ = session.reconcile_disclosures();
579 Ok(session)
580}
581
582pub 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 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
672pub(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 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}